Files
Jungfraujoch/image_pusher/TCPStreamPusher.h
T
leonarski_f 511be0c366
Build Packages / build:rpm (rocky8) (push) Successful in 24m0s
Build Packages / Unit tests (push) Skipped
Build Packages / build:windows:nocuda (push) Successful in 16m54s
Build Packages / build:windows:cuda (push) Successful in 19m25s
Build Packages / build:viewer-tgz:cpu (push) Successful in 14m44s
Build Packages / build:viewer-tgz:cuda (push) Successful in 16m3s
Build Packages / build:rugnux-tgz (x86_64) (push) Successful in 13m15s
Build Packages / build:rugnux:windows (push) Successful in 10m45s
Build Packages / build:rugnux:aarch64 (cross) (push) Successful in 9m34s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 19m7s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 18m9s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 24m48s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 18m13s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 24m51s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 22m58s
Build Packages / build:rpm (rocky9) (push) Successful in 21m23s
Build Packages / Generate python client (push) Successful in 1m2s
Build Packages / Build documentation (push) Successful in 1m23s
Build Packages / Create release (push) Skipped
Build Packages / XDS test (durin plugin) (push) Successful in 9m45s
Build Packages / XDS test (neggia plugin) (push) Successful in 10m19s
Build Packages / XDS test (JFJoch plugin) (push) Successful in 11m10s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 22m15s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 17m37s
Build Packages / DIALS test (push) Successful in 17m16s
v1.0.0-rc.165 (#75)
* `rugnux --model` adopts the model's space group as a label where the data were merged in its enantiomorph, instead of reindexing the reflections - which swapped I(+) with I(-).
* `rugnux --model` warns, naming the atom, when the anomalous density at the model's atoms comes out inverted, which means the data and the model are in opposite hands.
* `rugnux --model` writes an anomalous difference map (`<prefix>_anom.ccp4`) when the merge kept the Bijvoet split, and names the ten model atoms it peaks highest on as `ANOMALOUS_SITE_01`..`_10`.
* `MEAN_ATOM_DENSITY_SIGMA` is read from the map by cubic rather than linear interpolation and comes out around a tenth higher; it is no longer comparable with the figure earlier versions printed.
* `rugnux --model` reads an mmCIF coordinate file as well as a PDB one, gzipped or not, taking the format from the file's content rather than its name.
* A model `rugnux --model` cannot use is reported as a `WARNING:` line in the results report instead of only in the log.
* The rugnux results report has a `10. MODEL VALIDATION` section when `--model` was given; `REPORT_VERSION` is 4, `WARNINGS` moves to section 11 and no existing key changed.
* The rugnux results report records how the run was invoked, what it cost and what it ran on: `COMMAND_LINE=`, `WALL_TIME=` and `GPU_COUNT=` / `GPU=`.
* rugnux says which GPUs it can see before it starts processing.
* `rugnux --export-unmerged` also writes `<prefix>_unmerged.mtz` on a `--no-merge` run, and is ignored on a run with no output prefix instead of writing a file called `_unmerged.mtz`.
* `/start` asks the writer whether the run can be written before the detector is armed, so a run whose master file already exists, or whose output directory cannot be created, is refused up front with the writer's own message. This needs the TCP image stream or the built-in HDF5 writer; the ZeroMQ stream is unchanged.
* A calibration that fails goes to `Error` carrying the reason instead of `Inactive`, so `/wait_till_done` and `/wait_until_running` report it; a cancelled calibration still goes to `Inactive`.
* `/wait_till_done` answers 500 with the message when a collection ended in an error. A cancelled collection and a collection that only triggered a warning still answer 200.
* A pending start failure is discarded by `/cancel` and `/deactivate`, as it already was by `/start` and `/initialize`.
* `/scan_result` no longer reports the previous run's images after a collection that failed to start, or after `/deactivate`.
* The TCP image stream protocol version is 4. `jfjoch_writer` and `jfjoch_broker` have to be of the same release, as before.

Reviewed-on: #75
2026-08-27 22:16:54 +02:00

209 lines
9.6 KiB
C++

// SPDX-FileCopyrightText: 2025 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#pragma once
#include <atomic>
#include <future>
#include <mutex>
#include <optional>
#include <string>
#include <vector>
#include <condition_variable>
#include "ImagePusher.h"
#include "ZMQWriterNotificationPuller.h"
#include "../common/ThreadSafeFIFO.h"
#include "../common/Logger.h"
#include "../common/JfjochTCP.h"
#include "../frame_serialize/CBORStream2Serializer.h"
/// TCP-based image stream pusher with persistent connection pool.
///
/// Threading model:
/// - AcceptorThread: accepts new TCP connections, holds connections_mutex briefly
/// - KeepaliveThread: sends periodic keepalive frames when idle (skipped during data collection)
/// - Per-connection WriterThread: drains the connection's queue, sends DATA frames
/// - Per-connection PersistentAckThread: reads ACKs and keepalive pongs from the peer
///
/// Lock ordering: connections_mutex → send_mutex → ack_mutex
/// IMPORTANT: Never call blocking queue operations while holding connections_mutex.
///
/// Concurrency contract:
/// - StartDataCollection, EndDataCollection, SendCalibration, and Finalize
/// are called from a single control thread in a serialized manner.
/// - SendImage may be called concurrently from multiple threads between
/// StartDataCollection and EndDataCollection.
/// - SendCalibration is called between StartDataCollection and SendImage calls.
/// - GetConnectedWriters, GetImagesWritten, and PrintSetup are safe to call at any time.
class TCPStreamPusher : public ImagePusher {
struct Connection {
explicit Connection(size_t queue_size) : queue(queue_size) {}
std::atomic<int> fd{-1};
uint32_t socket_number = 0;
std::atomic<bool> active{false}; // data-collection threads running
std::atomic<bool> broken{false};
std::atomic<bool> connected{false}; // persistent connection is alive
ThreadSafeFIFO<ImagePusherQueueElement> queue;
std::future<void> writer_future;
// Persistent ack/keepalive reader (runs as long as the connection is alive)
std::future<void> persistent_ack_future;
// Serialises tearing this connection down. Both futures below are joined from more than one
// path - the acceptor reaping a dead connection, and the control plane starting or ending a
// run - and calling get() on one future from two threads is undefined and invalidates it.
// Held only around the joins, never around a send, and no thread it joins takes it.
std::mutex teardown_mutex;
std::mutex send_mutex;
std::mutex ack_mutex;
std::condition_variable ack_cv;
bool preflight_ack_received = false;
bool preflight_ack_ok = false;
bool start_ack_received = false;
bool start_ack_ok = false;
bool end_ack_received = false;
bool end_ack_ok = false;
bool cancel_ack_received = false;
bool cancel_ack_ok = false;
std::string last_ack_error;
std::atomic<TCPAckCode> last_ack_code{TCPAckCode::None};
// Soft writer failure reported via DATA ACK (do not break stream on this alone)
std::atomic<bool> data_ack_error_reported{false};
std::string data_ack_error_text;
std::atomic<uint64_t> data_acked_ok{0};
std::atomic<uint64_t> data_acked_bad{0};
std::atomic<uint64_t> data_acked_total{0};
std::atomic<uint64_t> last_ack_fifo_occupancy{0};
std::chrono::steady_clock::time_point last_keepalive_sent{};
std::chrono::steady_clock::time_point last_keepalive_recv{};
// Last time ANY frame (ACK / keepalive pong / busy heartbeat) was received from
// the peer, as steady_clock nanoseconds. Used to keep a healthy-but-busy writer
// alive through long backpressure while still detecting a genuinely dead peer.
std::atomic<int64_t> last_peer_activity_ns{0};
};
std::string endpoint;
size_t max_connections;
std::optional<int32_t> send_buffer_size;
size_t send_queue_size = 128;
// Persistent connection pool, guarded by connections_mutex.
// IMPORTANT: never call PutBlocking/GetBlocking on a queue while holding this mutex.
mutable std::mutex connections_mutex;
std::vector<std::shared_ptr<Connection>> connections;
std::vector<std::shared_ptr<Connection>> session_connections;
std::shared_ptr<Connection> calibration_connection;
// Acceptor thread state
std::atomic<int> listen_fd{-1};
std::atomic<bool> acceptor_running{false};
std::future<void> acceptor_future;
std::future<void> keepalive_future;
std::chrono::milliseconds send_poll_timeout{250};
// Maximum time a send (or the post-END ACK wait) may block with NO sign of life from
// the peer before the connection is declared dead. A busy writer refreshes its liveness
// every ~250 ms via BUSY heartbeats (and via DATA ACKs), so genuine backpressure of any
// duration is tolerated; only a truly silent (frozen/dead) peer trips this.
std::chrono::milliseconds peer_liveness_timeout{15000};
// Hard upper bound on backpressure: if the socket accepts no bytes for this long the
// writer is wedged and is declared dead even if it keeps heartbeating, so a
// misbehaving writer cannot block the run (or its finalization) forever. Generous
// relative to peer_liveness_timeout, since a heartbeating peer is given more grace
// than a silent one — but still finite.
std::chrono::milliseconds max_backpressure_timeout{60000};
int64_t images_per_file = 1;
// Written by the control thread (Preflight/StartDataCollection), read by every
// PersistentAckThread to discard ACKs belonging to an earlier run.
std::atomic<uint64_t> run_number{0};
std::string run_name;
std::atomic<bool> transmission_error = false;
std::atomic<bool> data_collection_active{false};
std::atomic<uint64_t> total_data_acked_ok{0};
std::atomic<uint64_t> total_data_acked_bad{0};
std::atomic<uint64_t> total_data_acked_total{0};
Logger logger{"TCPStreamPusher"};
static std::pair<std::string, std::optional<uint16_t>> ParseTcpAddress(const std::string& addr);
static std::pair<int, std::string> OpenListenSocket(const std::string& addr);
static int AcceptOne(int listen_fd, std::chrono::milliseconds timeout);
static void CloseFd(std::atomic<int>& fd);
bool IsConnectionAlive(const Connection& c) const;
bool SendAll(Connection& c, const void* buf, size_t len);
bool ReadExact(Connection& c, void* buf, size_t len);
bool ReadExactPersistent(Connection& c, void* buf, size_t len);
bool SendFrame(Connection& c, const uint8_t* data, size_t size, TCPFrameType type, int64_t image_number);
void WriterThread(Connection* c);
void PersistentAckThread(Connection* c);
void AcceptorThread();
void KeepaliveThread();
void SetupNewConnection(int new_fd, uint32_t socket_number);
// Unlink dead connections from the pool (connections_mutex held) and close them (mutex released).
// Split because closing joins a writer thread that can be blocked in a send.
std::vector<std::shared_ptr<Connection>> DetachDeadConnections();
void CloseDeadConnections(const std::vector<std::shared_ptr<Connection>> &dead);
void TearDownConnection(Connection& c);
void StartDataCollectionThreads(Connection& c);
void StopDataCollectionThreads(Connection& c);
void JoinPersistentAck(Connection& c);
bool WaitForAck(Connection& c, TCPFrameType ack_for, std::chrono::milliseconds timeout, std::string* error_text);
bool WaitForEndAck(Connection& c, std::chrono::milliseconds liveness_timeout, std::string* error_text);
public:
explicit TCPStreamPusher(const std::string& addr,
size_t in_max_connections,
std::optional<int32_t> in_send_buffer_size = {});
~TCPStreamPusher() override;
/// Max time a send may block on backpressure with no sign of life from the peer
/// before the connection is declared dead. A busy-but-alive writer keeps it fresh
/// via BUSY heartbeats, so this only catches a genuinely silent peer. Must be set
/// before data collection starts.
void SetPeerLivenessTimeout(std::chrono::milliseconds t) { peer_liveness_timeout = t; }
/// Hard upper bound on backpressure. Even while the peer keeps heartbeating, if no
/// bytes can be sent for this long the writer is declared dead so a wedged writer
/// cannot block the run or its finalization forever. Must be set before data
/// collection starts.
void SetMaxBackpressureTimeout(std::chrono::milliseconds t) { max_backpressure_timeout = t; }
std::vector<std::string> GetAddress() const override { return {endpoint}; }
/// Returns the number of currently connected writers (can be called at any time)
size_t GetConnectedWriters() const override;
void Preflight(StartMessage& message) override;
void StartDataCollection(StartMessage& message) override;
bool EndDataCollection(const EndMessage& message) override;
bool SendImage(const uint8_t *image_data, size_t image_size, int64_t image_number) override;
bool SendImage(ZeroCopyReturnValue &z) override;
bool SendCalibration(const CompressedImage& message) override;
std::string Finalize() override;
std::string PrintSetup() const override;
std::optional<uint64_t> GetImagesWritten() const override;
std::optional<uint64_t> GetImagesWriteError() const override;
std::vector<int64_t> GetWriterFifoUtilization() const override;
ImagePusherType GetType() const override { return ImagePusherType::TCP; }
};