From f970972163f55bb36713dc198b6d13a5107ea8bc Mon Sep 17 00:00:00 2001 From: Filip Leonarski Date: Thu, 8 Oct 2026 11:35:10 +0200 Subject: [PATCH] Writer and broker: the master file appears only after every data file it links to The final master marks a finished collection for anyone watching the directory (rsync, autoprocessing, a reader following the collection), so it must not appear before a data file it links to. - FileWriter: at END every data file is closed (renamed into place) before the master is finalized and renamed. A data file that fails to close does not cost the master; the failure is reported after it. For NXmxIntegrated, file 0 is the master's own image dataset, so there a failure still stops the master. - With several writers, the master was written by the writer of the first socket as soon as it got END, while the others could still be closing their data files. END now goes to every other writer first: - TCP: the connections are taken in reverse order, and each END waits for its acknowledgement, which a writer sends only after closing its files; - ZeroMQ, with a writer notification socket: END to every other socket, then their notifications (the same ones Finalize collected before, now collected here), then END to the first socket. Without a notification socket there is nothing to wait on and the order is all it changes. Tests: TCPImageCommTest_MasterWriterGetsEndLast (three mock writers, the master's END after the other two acknowledged), StreamWriterTest_ZMQ_MasterAfterAllDataFiles (two real writers, the second started only after the collection was sent; the master's ctime is after every data file's). Both fail without this change. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi --- image_pusher/TCPStreamPusher.cpp | 7 +++- image_pusher/ZMQStream2Pusher.cpp | 23 ++++++++++-- image_pusher/ZMQStream2Pusher.h | 2 + tests/StreamWriterTest.cpp | 60 ++++++++++++++++++++++++++++++ tests/TCPImagePusherTest.cpp | 62 +++++++++++++++++++++++++++++++ writer/FileWriter.cpp | 18 +++++++-- 6 files changed, 163 insertions(+), 9 deletions(-) diff --git a/image_pusher/TCPStreamPusher.cpp b/image_pusher/TCPStreamPusher.cpp index 2fd529fff..539868189 100644 --- a/image_pusher/TCPStreamPusher.cpp +++ b/image_pusher/TCPStreamPusher.cpp @@ -1080,8 +1080,11 @@ bool TCPStreamPusher::EndDataCollection(const EndMessage& message) { local_connections = (!session_connections.empty() ? session_connections : connections); } - for (auto& cptr : local_connections) { - auto& c = *cptr; + // The writer on the first connection writes the master file (write_master_file in START), so it + // gets END last: every other writer acknowledges END only once it has closed - renamed into + // place - its data files, and the master then appears only after every data file it links to. + for (auto it = local_connections.rbegin(); it != local_connections.rend(); ++it) { + auto& c = **it; // Flush all queued DATA before sending END: stopping the writer thread drains the // send queue and joins it, so every image is on the wire before END goes out. This diff --git a/image_pusher/ZMQStream2Pusher.cpp b/image_pusher/ZMQStream2Pusher.cpp index a29ca5fde..659d66fb8 100644 --- a/image_pusher/ZMQStream2Pusher.cpp +++ b/image_pusher/ZMQStream2Pusher.cpp @@ -46,6 +46,7 @@ void ZMQStream2Pusher::StartDataCollection(StartMessage& message) { throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "Images per file cannot be zero or negative"); images_per_file = message.images_per_file; + early_notifications.clear(); run_number = message.run_number; run_name = message.run_name; transmission_error = false; @@ -78,9 +79,21 @@ bool ZMQStream2Pusher::EndDataCollection(const EndMessage& message) { bool ret = true; - for (auto &s: socket) { - s->StopWriterThread(); - if (!s->Send(serialization_buffer.data(), serializer.GetBufferSize())) + // The writer on the first socket writes the master file, so it gets END last, and - where the + // writers report back - only once every other writer has reported that it closed (renamed into + // place) its data files: the master then appears only after every data file it links to. + early_notifications.clear(); + for (size_t i = 1; i < socket.size(); i++) { + socket[i]->StopWriterThread(); + if (!socket[i]->Send(serialization_buffer.data(), serializer.GetBufferSize())) + ret = false; + } + if (writer_notification_socket) + for (size_t i = 1; i < socket.size(); i++) + early_notifications.push_back(writer_notification_socket->Receive(run_number, run_name)); + if (!socket.empty()) { + socket[0]->StopWriterThread(); + if (!socket[0]->Send(serialization_buffer.data(), serializer.GetBufferSize())) ret = false; } transmission_error = !ret; @@ -102,7 +115,9 @@ std::string ZMQStream2Pusher::Finalize() { bool images_saved = false; if (writer_notification_socket) { for (int i = 0; i < socket.size(); i++) { - auto n = writer_notification_socket->Receive(run_number, run_name); + // The notifications EndDataCollection already waited for, then the rest. + auto n = (i < early_notifications.size()) ? early_notifications[i] + : writer_notification_socket->Receive(run_number, run_name); if (!n) ret += "Writer " + std::to_string(i) + ": no end notification received within 1 minute from collection end"; else if (n->socket_number >= socket.size()) diff --git a/image_pusher/ZMQStream2Pusher.h b/image_pusher/ZMQStream2Pusher.h index 3fc17de44..d52643512 100644 --- a/image_pusher/ZMQStream2Pusher.h +++ b/image_pusher/ZMQStream2Pusher.h @@ -13,6 +13,8 @@ class ZMQStream2Pusher : public ImagePusher { std::vector> socket; std::unique_ptr writer_notification_socket; + // Received in EndDataCollection, before END goes to the writer of the master file. + std::vector> early_notifications; int64_t images_per_file = 1; uint64_t run_number = 0; diff --git a/tests/StreamWriterTest.cpp b/tests/StreamWriterTest.cpp index bde595404..662477ccf 100644 --- a/tests/StreamWriterTest.cpp +++ b/tests/StreamWriterTest.cpp @@ -3,6 +3,9 @@ #include #include +#include +#include +#include #include "../writer/StreamWriter.h" #include "../image_pusher/ZMQStream2Pusher.h" @@ -186,3 +189,60 @@ TEST_CASE("StreamWriterTest_ZMQ_Update_NoNotification", "[StreamWriter]") { REQUIRE(std::filesystem::remove("subdir/StreamWriterTest3_data_000001.h5")); REQUIRE(std::filesystem::remove("subdir")); } + +// Two writers, the second one writing every other data file: the master, written by the first, must +// appear only after the second has closed its data files - here the second one starts only after the +// whole collection has been sent, so it is still writing when END goes out. +TEST_CASE("StreamWriterTest_ZMQ_MasterAfterAllDataFiles", "[StreamWriter]") { + RegisterHDF5Filter(); + Logger logger("StreamWriterTest_ZMQ_MasterAfterAllDataFiles"); + + DatasetSettings d; + d.FilePrefix("subdir4/StreamWriterTest4").NumTriggers(1).ImagesPerTrigger(4).ImagesPerFile(2) + .RunName("run1").RunNumber(257); + DiffractionExperiment x(DetJF(2)); + x.UseInternalPacketGenerator(true).Raw().PedestalG0Frames(0) + .ImportDatasetSettings(d).OverwriteExistingFiles(true); + PixelMask pixel_mask(x); + + AcquisitionDeviceGroup aq_devices; + for (int i = 0; i < x.GetDataStreamsNum(); i++) + aq_devices.AddHLSDevice(64); + + ZMQStream2Pusher pusher({"ipc://*", "ipc://*"}); + pusher.WriterNotificationSocket("ipc://*"); + JFJochReceiverService fpga_receiver_service(aq_devices, logger, pusher); + + auto pusher_addr = pusher.GetAddress(); + REQUIRE(pusher_addr.size() == 2); + ZMQImagePuller puller0(pusher_addr[0]), puller1(pusher_addr[1]); + StreamWriter writer0(logger, puller0), writer1(logger, puller1); + + REQUIRE_NOTHROW(fpga_receiver_service.Start(x, pixel_mask, nullptr)); + auto w0 = std::async(std::launch::async, [&] { return writer0.Run(); }); + for (int i = 0; i < 1200; i++) { + const auto status = fpga_receiver_service.GetStatus(); + if (status && status->images_sent == 4) + break; + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + } + auto stop = std::async(std::launch::async, [&] { return fpga_receiver_service.Stop(); }); + std::this_thread::sleep_for(std::chrono::seconds(2)); + auto w1 = std::async(std::launch::async, [&] { return writer1.Run(); }); + JFJochReceiverOutput r; + REQUIRE_NOTHROW(r = stop.get()); + w0.get(); + w1.get(); + CHECK(r.writer_err.empty()); + + auto changed = [](const std::string &path) { + struct stat st{}; + REQUIRE(stat(path.c_str(), &st) == 0); + return std::make_pair(st.st_ctim.tv_sec, st.st_ctim.tv_nsec); + }; + const auto master = changed("subdir4/StreamWriterTest4_master.h5"); + for (int f = 1; f <= 2; f++) + CHECK(changed("subdir4/StreamWriterTest4_data_00000" + std::to_string(f) + ".h5") <= master); + + std::filesystem::remove_all("subdir4"); +} diff --git a/tests/TCPImagePusherTest.cpp b/tests/TCPImagePusherTest.cpp index fcc2706ea..a01179def 100644 --- a/tests/TCPImagePusherTest.cpp +++ b/tests/TCPImagePusherTest.cpp @@ -167,6 +167,68 @@ TEST_CASE("TCPImageCommTest_2Writers_WithAck", "[TCP]") { p->Disconnect(); } +// The writer that writes the master file gets END only once every other writer has acknowledged its +// END - which a writer does only after closing (renaming into place) its data files - so the master +// never appears before a data file it links to. +TEST_CASE("TCPImageCommTest_MasterWriterGetsEndLast", "[TCP]") { + const int64_t npullers = 3; + TCPStreamPusher pusher("tcp://127.0.0.1:*", npullers); + std::vector > puller; + for (int i = 0; i < npullers; i++) + puller.push_back(std::make_unique(pusher.GetAddress()[0], 8 * 1024 * 1024)); + for (int attempt = 0; attempt < 100 && pusher.GetConnectedWriters() < static_cast(npullers); ++attempt) + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + REQUIRE(pusher.GetConnectedWriters() == static_cast(npullers)); + + std::atomic others_done{0}; + std::atomic others_done_at_master_end{-1}; + + std::thread sender([&] { + StartMessage start{.images_per_file = 16, .run_number = 12, .write_master_file = true}; + EndMessage end{}; + pusher.StartDataCollection(start); + CHECK(pusher.EndDataCollection(end)); + }); + + std::vector receivers; + for (int w = 0; w < npullers; w++) { + receivers.emplace_back([&, w] { + bool master = false; + for (int polls = 0; polls < 100; polls++) { + auto out = puller[w]->PollImage(std::chrono::milliseconds(100)); + if (!out.has_value() || !out->cbor) + continue; + const auto &h = out->tcp_msg->header; + PullerAckMessage ack; + ack.ok = true; + ack.run_number = h.run_number; + ack.socket_number = h.socket_number; + ack.error_code = TCPAckCode::None; + if (out->cbor->start_message) { + master = out->cbor->start_message->write_master_file.value_or(false); + ack.ack_for = TCPFrameType::START; + puller[w]->SendAck(ack); + } else if (out->cbor->end_message) { + if (master) + others_done_at_master_end = others_done.load(); + else { + std::this_thread::sleep_for(std::chrono::milliseconds(200)); // closing files + others_done++; + } + ack.ack_for = TCPFrameType::END; + puller[w]->SendAck(ack); + break; + } + } + }); + } + sender.join(); + for (auto &t: receivers) t.join(); + CHECK(others_done_at_master_end == npullers - 1); + for (auto &p: puller) + p->Disconnect(); +} + // One writer rejects START (as it would on an overwrite conflict) while its sibling // accepts. The broker must abort the whole collection and cleanly cancel the sibling // that already started - no half-armed collection, no stuck writer. diff --git a/writer/FileWriter.cpp b/writer/FileWriter.cpp index a9fcee718..0b5c32e1a 100644 --- a/writer/FileWriter.cpp +++ b/writer/FileWriter.cpp @@ -306,14 +306,26 @@ void FileWriter::WriteHDF5(const EndMessage &msg) { if (master_file) { std::lock_guard lock(hdf5_mutex); - if (format == FileWriterFormat::NXmxIntegrated) { + // Every data file is closed - and so in place under its final name - before the master + // appears under its own: a reader that finds the final master can then count on the data + // files it links to. A data file that fails to close does not cost the master; the failure + // is reported after it. For NXmxIntegrated file 0 is the master's own image dataset, so + // there a failure does stop the master. + std::exception_ptr first_exception; + for (uint64_t f = 0; f < files.size(); ++f) { try { - CloseFile(0); + CloseFile(f); } catch (...) { - throw; + if (!first_exception) + first_exception = std::current_exception(); } } + if (first_exception && format == FileWriterFormat::NXmxIntegrated) + std::rethrow_exception(first_exception); master_file->Finalize(msg); + + if (first_exception) + std::rethrow_exception(first_exception); } }