From 49dc3cefe79cf69f2aa2eac2022d6d94db54fba2 Mon Sep 17 00:00:00 2001 From: Andrey Bychkov Date: Mon, 24 Aug 2026 15:14:45 +0300 Subject: [PATCH] fix(picloud-dispatcher): drop data to a connection whose write is stalled While the peer does not drain its receive buffer, every forwarded packet would re-stall the reader thread for up to writeTimeout(), so a queued Disconnect frame could wait behind the whole backlog. Remember a stalled target (send slower than 2 s) and drop data for a short cooldown instead of stalling; keep dispatcher writes bounded at 3 s. Control frames (Connected/Disconnected) are unaffected. --- utils/cloud_dispatcher/cloudserver.cpp | 17 ++++++++++++++++- utils/cloud_dispatcher/dispatcherclient.cpp | 12 ++++++++++++ utils/cloud_dispatcher/dispatcherclient.h | 5 +++++ 3 files changed, 33 insertions(+), 1 deletion(-) diff --git a/utils/cloud_dispatcher/cloudserver.cpp b/utils/cloud_dispatcher/cloudserver.cpp index 087f708f..62aa74f9 100644 --- a/utils/cloud_dispatcher/cloudserver.cpp +++ b/utils/cloud_dispatcher/cloudserver.cpp @@ -9,7 +9,17 @@ CloudServer::CloudServer(DispatcherClient * c, const PIByteArray & sname): serve mutex_clients.lock(); last_ping.reset(); DispatcherClient * cl = index_clients.value(id, nullptr); - if (cl) cl->sendData(ba); + if (cl) { + // a healthy forward is fast: a slow one means + // the peer is not draining, remember it and + // drop further data instead of stalling + PITimeMeasurer tm; + cl->sendData(ba); + if (tm.elapsed_s() > 2.0) { + cl->markWriteStalled(); + piCoutObj << "write stalled to client" << id; + } + } mutex_clients.unlock(); })); connects << CONNECTL(c, pingReceived, [this]() { @@ -43,7 +53,12 @@ void CloudServer::addClient(DispatcherClient * c) { server->sendConnected(cid); CONNECTL(c, dataReaded, ([this, cid](PIByteArray & ba) { // piCoutObj << c->clientId() << "dataReaded"; + PITimeMeasurer tm; server->sendDataToClient(ba, cid); + if (tm.elapsed_s() > 2.0) { + server->markWriteStalled(); + piCoutObj << "write stalled to server" << cid; + } })); } diff --git a/utils/cloud_dispatcher/dispatcherclient.cpp b/utils/cloud_dispatcher/dispatcherclient.cpp index 5d600de6..f7aeb854 100644 --- a/utils/cloud_dispatcher/dispatcherclient.cpp +++ b/utils/cloud_dispatcher/dispatcherclient.cpp @@ -10,6 +10,10 @@ DispatcherClient::DispatcherClient(PIEthernet * eth_, int id) , tcp(&streampacker) , client_id(id) { eth->setName(PIString::fromNumber(id)); + // bound the EAGAIN retry loop in writeDevice: a write to a peer that + // does not drain its receive buffer must not block the reader thread + // for longer than this + eth->setWriteTimeout(PISystemTime::fromSeconds(3)); CONNECTU(&streampacker, packetReceiveEvent, this, readed); CONNECTU(eth, disconnected, this, disconnected); piCoutObj << "client connected" << client_id << eth->sendAddress(); @@ -58,6 +62,10 @@ void DispatcherClient::sendDisconnected(uint client_id) { void DispatcherClient::sendData(const PIByteArray & data) { if (tcp.role() == PICloud::TCP::Client) { + if (isWriteStalled()) { + piCoutObj << "sendData dropped, write stalled"; + return; + } PIMutexLocker send_lock(mutex_send); tcp.sendData(data); } else @@ -67,6 +75,10 @@ void DispatcherClient::sendData(const PIByteArray & data) { void DispatcherClient::sendDataToClient(const PIByteArray & data, uint client_id) { if (tcp.role() == PICloud::TCP::Server) { + if (isWriteStalled()) { + piCoutObj << "sendDataToClient dropped, write stalled"; + return; + } PIMutexLocker send_lock(mutex_send); tcp.sendData(data, client_id); } else diff --git a/utils/cloud_dispatcher/dispatcherclient.h b/utils/cloud_dispatcher/dispatcherclient.h index d8b01e73..62b0d6c8 100644 --- a/utils/cloud_dispatcher/dispatcherclient.h +++ b/utils/cloud_dispatcher/dispatcherclient.h @@ -21,6 +21,10 @@ public: void sendDataToClient(const PIByteArray & data, uint client_id); void authorise(bool ok); bool isAuthorised() const { return authorised; } + // forwarding to this connection stalled (peer not draining the send + // buffer) recently: drop data instead of stalling the reader thread + bool isWriteStalled() const { return (PISystemTime::current().toMilliseconds() - write_stalled_at) < 5000.0; } + void markWriteStalled() { write_stalled_at = PISystemTime::current().toMilliseconds(); } PICloud::TCP::Role role() const { return tcp.role(); } PIString address(); uint clientId() const { return client_id; } @@ -38,6 +42,7 @@ private: EVENT_HANDLER1(void, disconnected, bool, withError); std::atomic_bool authorised; + std::atomic write_stalled_at{0.0}; PIEthernet * eth; PIStreamPacker streampacker; PICloud::TCP tcp;