mirror of
https://github.com/alexhopeoconnor/firmware.git
synced 2026-10-04 03:18:10 +10:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d4dc73d0f6 | ||
|
|
5e950d065b | ||
|
|
df6c814c92 |
+33
-25
@@ -52,7 +52,8 @@ void StreamAPI::writeStream()
|
||||
do {
|
||||
// Send every packet we can
|
||||
len = getFromRadio(txBuf + HEADER_LEN);
|
||||
emitTxBuffer(len);
|
||||
if (len != 0 && !emitTxBuffer(len))
|
||||
break;
|
||||
} while (len);
|
||||
}
|
||||
}
|
||||
@@ -169,21 +170,41 @@ int32_t StreamAPI::readStream()
|
||||
/**
|
||||
* Send the current txBuffer over our stream
|
||||
*/
|
||||
void StreamAPI::emitTxBuffer(size_t len)
|
||||
bool StreamAPI::writeFrame(uint8_t *buf, size_t len)
|
||||
{
|
||||
if (len != 0) {
|
||||
txBuf[0] = START1;
|
||||
txBuf[1] = START2;
|
||||
txBuf[2] = (len >> 8) & 0xff;
|
||||
txBuf[3] = len & 0xff;
|
||||
if (len == 0 || !canWrite)
|
||||
return false;
|
||||
|
||||
auto totalLen = len + HEADER_LEN;
|
||||
buf[0] = START1;
|
||||
buf[1] = START2;
|
||||
buf[2] = (len >> 8) & 0xff;
|
||||
buf[3] = len & 0xff;
|
||||
|
||||
auto totalLen = len + HEADER_LEN;
|
||||
if (!canWriteFrame(totalLen))
|
||||
return false;
|
||||
|
||||
size_t written;
|
||||
{
|
||||
// Serialize stream writes against `emitLogRecord` so a LOG_ firing
|
||||
// mid-packet-emission can't interleave bytes on the wire.
|
||||
concurrency::LockGuard guard(&streamLock);
|
||||
stream->write(txBuf, totalLen);
|
||||
stream->flush();
|
||||
written = stream->write(buf, totalLen);
|
||||
if (written == totalLen)
|
||||
stream->flush();
|
||||
}
|
||||
|
||||
if (written != totalLen) {
|
||||
onFrameWriteFailed(totalLen, written);
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
bool StreamAPI::emitTxBuffer(size_t len)
|
||||
{
|
||||
return writeFrame(txBuf, len);
|
||||
}
|
||||
|
||||
void StreamAPI::emitRebooted()
|
||||
@@ -221,20 +242,7 @@ void StreamAPI::emitLogRecord(meshtastic_LogRecord_Level level, const char *src,
|
||||
|
||||
size_t len =
|
||||
pb_encode_to_bytes(txBufLog + HEADER_LEN, meshtastic_FromRadio_size, &meshtastic_FromRadio_msg, &fromRadioScratchLog);
|
||||
if (len != 0) {
|
||||
txBufLog[0] = START1;
|
||||
txBufLog[1] = START2;
|
||||
txBufLog[2] = (len >> 8) & 0xff;
|
||||
txBufLog[3] = len & 0xff;
|
||||
|
||||
auto totalLen = len + HEADER_LEN;
|
||||
// Serialize stream writes against `emitTxBuffer` so a packet
|
||||
// emission in flight on another task doesn't interleave bytes
|
||||
// with this log record.
|
||||
concurrency::LockGuard guard(&streamLock);
|
||||
stream->write(txBufLog, totalLen);
|
||||
stream->flush();
|
||||
}
|
||||
writeFrame(txBufLog, len);
|
||||
}
|
||||
|
||||
/// Hookable to find out when connection changes
|
||||
@@ -249,4 +257,4 @@ void StreamAPI::onConnectionChanged(bool connected)
|
||||
// received a packet in a while
|
||||
powerFSM.trigger(EVENT_SERIAL_DISCONNECTED);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -80,7 +80,7 @@ class StreamAPI : public PhoneAPI
|
||||
/**
|
||||
* Send the current txBuffer over our stream
|
||||
*/
|
||||
void emitTxBuffer(size_t len);
|
||||
bool emitTxBuffer(size_t len);
|
||||
|
||||
/// Are we allowed to write packets to our output stream (subclasses can turn this off - i.e. SerialConsole)
|
||||
bool canWrite = true;
|
||||
@@ -91,7 +91,12 @@ class StreamAPI : public PhoneAPI
|
||||
/// Low level function to emit a protobuf encapsulated log record
|
||||
void emitLogRecord(meshtastic_LogRecord_Level level, const char *src, const char *format, va_list arg);
|
||||
|
||||
virtual bool canWriteFrame(size_t frameLen) { return true; }
|
||||
virtual void onFrameWriteFailed(size_t frameLen, size_t writtenLen) {}
|
||||
|
||||
private:
|
||||
bool writeFrame(uint8_t *buf, size_t len);
|
||||
|
||||
/// Dedicated scratch + tx buffer for LogRecord emission.
|
||||
///
|
||||
/// The main packet emission path (`writeStream` -> `getFromRadio` ->
|
||||
@@ -113,4 +118,4 @@ class StreamAPI : public PhoneAPI
|
||||
meshtastic_FromRadio fromRadioScratchLog = {};
|
||||
uint8_t txBufLog[MAX_STREAM_BUF_SIZE] = {0};
|
||||
concurrency::Lock streamLock;
|
||||
};
|
||||
};
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
#pragma once
|
||||
|
||||
#include <type_traits>
|
||||
#include <utility>
|
||||
|
||||
namespace stream_api
|
||||
{
|
||||
template <typename Client, typename = void> struct HasAvailableForWrite : std::false_type
|
||||
{};
|
||||
|
||||
template <typename Client>
|
||||
struct HasAvailableForWrite<Client, std::void_t<decltype(std::declval<Client &>().availableForWrite())>> : std::true_type
|
||||
{};
|
||||
|
||||
template <typename Client> int availableForWriteOrUnknown(Client &client)
|
||||
{
|
||||
if constexpr (HasAvailableForWrite<Client>::value)
|
||||
return client.availableForWrite();
|
||||
|
||||
(void)client;
|
||||
return -1;
|
||||
}
|
||||
|
||||
template <typename Client> bool clientReadyForWrite(Client &client)
|
||||
{
|
||||
if (!client.connected())
|
||||
return false;
|
||||
|
||||
int writable = availableForWriteOrUnknown(client);
|
||||
return writable == -1 || writable > 0;
|
||||
}
|
||||
} // namespace stream_api
|
||||
@@ -28,6 +28,33 @@ template <typename T> bool ServerAPI<T>::checkIsConnected()
|
||||
return client.connected();
|
||||
}
|
||||
|
||||
template <typename T> bool ServerAPI<T>::canWriteFrame(size_t frameLen)
|
||||
{
|
||||
int writable = stream_api::availableForWriteOrUnknown(client);
|
||||
if (!stream_api::clientReadyForWrite(client)) {
|
||||
canWrite = false;
|
||||
enabled = false;
|
||||
if (writable == -1)
|
||||
LOG_WARN("TCP client disconnected before write, closing API service");
|
||||
else
|
||||
LOG_WARN("TCP client not writable before write (%d/%lu bytes), closing API service", writable,
|
||||
(unsigned long)frameLen);
|
||||
close();
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
template <typename T> void ServerAPI<T>::onFrameWriteFailed(size_t frameLen, size_t writtenLen)
|
||||
{
|
||||
canWrite = false;
|
||||
enabled = false;
|
||||
LOG_WARN("TCP client write short (%lu/%lu bytes), closing API service", (unsigned long)writtenLen,
|
||||
(unsigned long)frameLen);
|
||||
close();
|
||||
}
|
||||
|
||||
template <class T> int32_t ServerAPI<T>::runOnce()
|
||||
{
|
||||
if (client.connected()) {
|
||||
@@ -95,4 +122,4 @@ template <class T, class U> int32_t APIServerPort<T, U>::runOnce()
|
||||
waitTime = 100;
|
||||
#endif
|
||||
return 100; // only check occasionally for incoming connections
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
#pragma once
|
||||
|
||||
#include "ClientWriteChecks.h"
|
||||
#include "StreamAPI.h"
|
||||
#include <memory>
|
||||
|
||||
@@ -29,6 +30,8 @@ template <class T> class ServerAPI : public StreamAPI, private concurrency::OSTh
|
||||
/// We override this method to prevent publishing EVENT_SERIAL_CONNECTED/DISCONNECTED for wifi links (we want the board to
|
||||
/// stay in the POWERED state to prevent disabling wifi)
|
||||
virtual void onConnectionChanged(bool connected) override {}
|
||||
virtual bool canWriteFrame(size_t frameLen) override;
|
||||
virtual void onFrameWriteFailed(size_t frameLen, size_t writtenLen) override;
|
||||
|
||||
virtual int32_t runOnce() override; // Check for dropped client connections
|
||||
};
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
#include "TestUtil.h"
|
||||
#include "mesh/api/ClientWriteChecks.h"
|
||||
#include <unity.h>
|
||||
|
||||
namespace
|
||||
{
|
||||
class ClientWithAvailable
|
||||
{
|
||||
public:
|
||||
bool connected() const { return connected_; }
|
||||
int availableForWrite() const { return writable_; }
|
||||
|
||||
bool connected_ = true;
|
||||
int writable_ = 64;
|
||||
};
|
||||
|
||||
class ClientWithoutAvailable
|
||||
{
|
||||
public:
|
||||
bool connected() const { return connected_; }
|
||||
|
||||
bool connected_ = true;
|
||||
};
|
||||
} // namespace
|
||||
|
||||
void setUp(void) {}
|
||||
void tearDown(void) {}
|
||||
|
||||
void test_availableForWrite_is_detected_when_present()
|
||||
{
|
||||
TEST_ASSERT_TRUE((stream_api::HasAvailableForWrite<ClientWithAvailable>::value));
|
||||
TEST_ASSERT_FALSE((stream_api::HasAvailableForWrite<ClientWithoutAvailable>::value));
|
||||
}
|
||||
|
||||
void test_availableForWriteOrUnknown_returns_client_capacity()
|
||||
{
|
||||
ClientWithAvailable client;
|
||||
TEST_ASSERT_EQUAL_INT(64, stream_api::availableForWriteOrUnknown(client));
|
||||
}
|
||||
|
||||
void test_availableForWriteOrUnknown_returns_unknown_when_absent()
|
||||
{
|
||||
ClientWithoutAvailable client;
|
||||
TEST_ASSERT_EQUAL_INT(-1, stream_api::availableForWriteOrUnknown(client));
|
||||
}
|
||||
|
||||
void test_clientReadyForWrite_rejects_disconnected_clients()
|
||||
{
|
||||
ClientWithAvailable client;
|
||||
client.connected_ = false;
|
||||
TEST_ASSERT_FALSE(stream_api::clientReadyForWrite(client));
|
||||
}
|
||||
|
||||
void test_clientReadyForWrite_rejects_non_writable_clients()
|
||||
{
|
||||
ClientWithAvailable client;
|
||||
client.writable_ = 0;
|
||||
TEST_ASSERT_FALSE(stream_api::clientReadyForWrite(client));
|
||||
}
|
||||
|
||||
void test_clientReadyForWrite_allows_clients_without_availableForWrite()
|
||||
{
|
||||
ClientWithoutAvailable client;
|
||||
TEST_ASSERT_TRUE(stream_api::clientReadyForWrite(client));
|
||||
}
|
||||
|
||||
void setup()
|
||||
{
|
||||
initializeTestEnvironment();
|
||||
UNITY_BEGIN();
|
||||
RUN_TEST(test_availableForWrite_is_detected_when_present);
|
||||
RUN_TEST(test_availableForWriteOrUnknown_returns_client_capacity);
|
||||
RUN_TEST(test_availableForWriteOrUnknown_returns_unknown_when_absent);
|
||||
RUN_TEST(test_clientReadyForWrite_rejects_disconnected_clients);
|
||||
RUN_TEST(test_clientReadyForWrite_rejects_non_writable_clients);
|
||||
RUN_TEST(test_clientReadyForWrite_allows_clients_without_availableForWrite);
|
||||
exit(UNITY_END());
|
||||
}
|
||||
|
||||
void loop() {}
|
||||
Reference in New Issue
Block a user