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); } }