TCP stream: tolerate writer backpressure via BUSY heartbeats
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 10m3s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 11m14s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 11m34s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 11m32s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 11m42s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 12m30s
Build Packages / XDS test (durin plugin) (push) Successful in 6m49s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 8m52s
Build Packages / build:rpm (rocky8) (push) Successful in 11m24s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 10m35s
Build Packages / Create release (push) Skipped
Build Packages / Generate python client (push) Successful in 25s
Build Packages / Build documentation (push) Successful in 49s
Build Packages / build:rpm (rocky9) (push) Successful in 12m7s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 11m0s
Build Packages / DIALS test (push) Successful in 13m11s
Build Packages / XDS test (neggia plugin) (push) Successful in 6m20s
Build Packages / Unit tests (push) Successful in 1h7m39s

A slow filesystem could stall the writer's receive pipeline, propagating
TCP backpressure to the pusher. The pusher's 3s per-frame send deadline
then treated that backpressure as a dead peer and force-closed the
connection mid-run. The writer reconnected as a brand-new connection that
was not part of the active session, so the rest of the run was silently
dropped and the half-written HDF5 file was later finalized with holes.

Replace the throughput-based send timeout with a peer-liveness timeout:

- Add TCPFrameType::BUSY (wire version 2 -> 3).
- TCPImagePuller runs a heartbeat thread that sends BUSY every 250ms on a
  thread independent of the stalled write path, so liveness keeps flowing
  during deep stalls. A send_mutex serializes ACK/pong/heartbeat writes.
- TCPStreamPusher refreshes last_peer_activity on every inbound frame and
  only declares a connection dead after peer_liveness_timeout (15s) of
  complete silence, tolerating arbitrarily long backpressure while still
  catching a genuinely dead peer (and immediate EPIPE/ECONNRESET).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-06-23 16:30:16 +02:00
co-authored by Claude Opus 4.8
parent 185b7bd377
commit 9b1eeb1342
5 changed files with 94 additions and 5 deletions
+42 -1
View File
@@ -54,6 +54,7 @@ TCPImagePuller::TCPImagePuller(const std::string &tcp_addr,
receiver_thread = std::thread(&TCPImagePuller::ReceiverThread, this);
cbor_thread = std::thread(&TCPImagePuller::CBORThread, this);
heartbeat_thread = std::thread(&TCPImagePuller::HeartbeatThread, this);
if (!repub_address.empty()) {
repub_socket = std::make_unique<ZMQSocket>(ZMQSocketType::Push);
@@ -91,6 +92,8 @@ bool TCPImagePuller::SendAll(const void *buf, size_t len) {
}
bool TCPImagePuller::SendAck(const PullerAckMessage &ack) {
std::lock_guard lg(send_mutex);
TcpFrameHeader h{};
h.type = static_cast<uint16_t>(TCPFrameType::ACK);
h.run_number = ack.run_number;
@@ -319,7 +322,12 @@ void TCPImagePuller::ReceiverThread() {
TcpFrameHeader pong{};
pong.type = static_cast<uint16_t>(TCPFrameType::KEEPALIVE);
pong.payload_size = 0;
if (!SendAll(&pong, sizeof(pong))) {
bool pong_ok;
{
std::lock_guard lg(send_mutex);
pong_ok = SendAll(&pong, sizeof(pong));
}
if (!pong_ok) {
logger.Info("Keepalive pong send failed, reconnecting to " + addr);
CloseSocket();
std::this_thread::sleep_for(std::chrono::milliseconds(20));
@@ -365,6 +373,37 @@ void TCPImagePuller::ReceiverThread() {
cbor_fifo.PutBlocking(ImagePullerOutput{});
}
void TCPImagePuller::HeartbeatThread() {
// While connected, periodically tell the pusher we are alive even if the
// consuming pipeline is stalled (e.g. blocked on a slow filesystem). This lets
// the pusher distinguish a busy-but-healthy writer from a dead one and keep
// waiting through arbitrarily long backpressure instead of dropping the run.
while (!disconnect) {
// Sleep in small slices so shutdown stays prompt.
for (int i = 0; i < 5 && !disconnect; i++)
std::this_thread::sleep_for(HeartbeatInterval / 5);
if (disconnect)
break;
{
std::unique_lock ul(fd_mutex);
if (fd < 0)
continue; // Not connected; ReceiverThread is (re)establishing the link.
}
TcpFrameHeader h{};
h.type = static_cast<uint16_t>(TCPFrameType::BUSY);
h.payload_size = 0;
h.ack_fifo_occupancy = cbor_fifo.GetCurrentUtilization();
h.ack_fifo_max_occupancy = cbor_fifo.Size();
// Best effort: a failure here just means the socket is gone, which
// ReceiverThread will detect and reconnect on its own.
std::lock_guard lg(send_mutex);
SendAll(&h, sizeof(h));
}
}
void TCPImagePuller::Disconnect() {
if (disconnect.exchange(true))
return;
@@ -377,4 +416,6 @@ void TCPImagePuller::Disconnect() {
cbor_thread.join();
if (repub_thread.joinable())
repub_thread.join();
if (heartbeat_thread.joinable())
heartbeat_thread.join();
}