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) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
This commit is contained in:
@@ -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<std::unique_ptr<TCPImagePuller> > puller;
|
||||
for (int i = 0; i < npullers; i++)
|
||||
puller.push_back(std::make_unique<TCPImagePuller>(pusher.GetAddress()[0], 8 * 1024 * 1024));
|
||||
for (int attempt = 0; attempt < 100 && pusher.GetConnectedWriters() < static_cast<size_t>(npullers); ++attempt)
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(50));
|
||||
REQUIRE(pusher.GetConnectedWriters() == static_cast<size_t>(npullers));
|
||||
|
||||
std::atomic<int> others_done{0};
|
||||
std::atomic<int> 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<std::thread> 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.
|
||||
|
||||
Reference in New Issue
Block a user