1#include <gemc/gstreamer/sro/gSROImplementation.h>
2#include <gemc/gdynamicDigitization/gdynamicdigitization.h>
11 Data(GSROEventId
id, std::size_t thread) : event(id), worker(thread) {}
12 std::size_t size_bytes() const noexcept
override {
return sizeof(*this); }
16 std::optional<GSROTime> earliest_remaining_time(GSROEventId
id)
const override {
17 return GSROTime{
static_cast<std::int64_t
>(id) * 10};
23 const std::thread::id owner = std::this_thread::get_id();
24 explicit Sink(
const std::string& filename) {
25 file.exceptions(std::ios::failbit | std::ios::badbit);
28 void write_frame(GSROFrame frame)
override {
29 if (owner != std::this_thread::get_id()) {
throw std::runtime_error(
"Wrong SRO output thread"); }
30 const auto& data =
dynamic_cast<const Data&
>(*
frame.data);
31 file <<
"F " << data.event <<
' ' <<
frame.frame_id <<
' ' <<
frame.begin.count()
32 <<
' ' << data.worker <<
'\n';
34 void finish_output()
override { file.close(); }
39 std::vector<GSROPayload> pending;
40 Crate(GSROFrameSink& sink, GSROCrateId crate) : GSROCratePlugin(sink), id(crate) {}
41 void consume_payload(GSROPayload payload)
override { pending.push_back(std::move(payload)); }
42 void advance_time(GSROTime boundary)
override {
43 for (
auto& payload : pending) {
44 if (payload.
time >= boundary) {
throw std::runtime_error(
"Unproven SRO payload released"); }
46 std::move(payload.
data)});
50 void finish_run(
const GSROEndContext& context)
override {
52 auto& sink =
dynamic_cast<Sink&
>(output);
53 sink.file <<
"E " << (context.
reason == GSROEndReason::completed ?
"completed" :
"interrupted")
61 GSROConfiguration configure_run(
const GSRORunContext&)
override {
62 GSROConfiguration config;
63 config.
timing = std::make_shared<Timing>();
67 GSROCrateResources create_crate(GSROCrateId
id,
const GSRORunContext& run)
const override {
68 auto sink = std::make_unique<Sink>(
run.output_basename +
"_r" + std::to_string(
run.run_id) +
69 "_c" + std::to_string(
id) +
".txt");
70 auto crate = std::make_unique<Crate>(*sink,
id);
71 return {std::move(sink), std::move(crate)};
78 bool defineReadoutSpecsImpl()
override {
return true; }
79 std::unique_ptr<GTrueInfoData> collectTrueInformationImpl(GHit*, std::size_t)
override {
return nullptr; }
80 void stream_hit(GHit*, std::size_t,
const GSROEventContext& event,
const GSROEmit& emit)
const override {
81 if (gopts->getSwitch(
"sro_test_worker_failure")) {
82 throw std::runtime_error(
"synthetic SRO worker failure");
84 const auto worker = std::hash<std::thread::id>{}(std::this_thread::get_id());
85 for (GSROCrateId crate : {1u, 2u}) {
86 emit(crate, GSROTime{
static_cast<std::int64_t
>(
event.event_id) * 10},
87 std::make_unique<Data>(
event.event_id, worker));
94 return new Implementation(options);
97 return new Digitizer(options);
GDynamicDigitization(const std::shared_ptr< GOptions > &g)
GSROImplementation(const std::shared_ptr< GOptions > &opts)
std::chrono::duration< std::int64_t, std::nano > GSROTime
std::uint32_t GSROCrateId
std::uint64_t GSROEventId
GDynamicDigitization * GDynamicDigitizationFactory(const std::shared_ptr< GOptions > &options)
GSROImplementation * GSROImplementationFactory(const std::shared_ptr< GOptions > &options)
GSROCrateLimits crate_limits
std::shared_ptr< const GSROTiming > timing
std::size_t queue_messages
std::optional< GSROTime > safe_time
std::unique_ptr< const GSROData > data