FERS 0.1.0
The Flexible Extensible Radar Simulator
Loading...
Searching...
No Matches
vita49_output_sink.h
Go to the documentation of this file.
1// SPDX-License-Identifier: GPL-2.0-only
2//
3// Copyright (c) 2026-present FERS Contributors (see AUTHORS.md).
4//
5// See the GNU GPLv2 LICENSE file in the FERS project root for more information.
6
7#pragma once
8
9#include <chrono>
10#include <memory>
11#include <mutex>
12#include <optional>
13#include <unordered_map>
14#include <unordered_set>
15#include <vector>
16
22
23namespace serial::vita49
24{
26 {
27 public:
28 explicit Vita49OutputSink(std::unique_ptr<DatagramSender> sender = nullptr,
30 ~Vita49OutputSink() override;
35
36 void initializeRun(const core::OutputConfig& config, std::string simulation_name) override;
37 std::uint32_t registerStream(const core::ReceiverStreamDescriptor& stream) override;
38 void openStream(std::uint32_t stream_id, RealType first_sample_time) override;
39 void submitBlock(const core::ReceiverSampleBlock& block) override;
40 void submitBlocks(std::span<const core::ReceiverSampleBlock> blocks) override;
41 void emitContextHeartbeat(RealType simulation_time) override;
42 void closeStream(std::uint32_t stream_id) override;
43 core::OutputStats finalize() override;
44 [[nodiscard]] core::OutputStats snapshotStats() const override;
45
46 private:
47 struct StreamState
48 {
51 PacketCountSequencer packet_counts;
52 bool opened = false;
53 bool closed = false;
54 bool sample_loss_pending = false;
55 bool over_range_pending = false;
56 RealType last_context_time = -1.0e300;
57 };
58
59 struct PendingContext
60 {
61 std::uint32_t stream_id = 0;
62 RealType simulation_time = 0.0;
63 bool stream_open = false;
64 bool stream_close = false;
65 };
66
67 [[nodiscard]] StreamState& stateFor(std::uint32_t stream_id);
68 [[nodiscard]] const StreamState& stateFor(std::uint32_t stream_id) const;
69 [[nodiscard]] core::OutputStats snapshotStatsLocked() const;
70 void ensurePacketizer();
71 void startPacing();
72 void appendPendingContexts(std::vector<SerializedPacket>& packets,
73 std::vector<std::uint32_t>& context_stream_ids);
74 [[nodiscard]] bool enqueuePacket(SerializedPacket&& packet);
75 [[nodiscard]] bool enqueuePackets(std::vector<SerializedPacket> packets);
76 void emitTelemetry(std::vector<core::ReceiverOutputPacketTrace> packets = {}, bool force_stats = false);
77 [[nodiscard]] std::vector<core::ReceiverOutputPacketTrace> consumeSenderDropsLocked();
78 [[nodiscard]] core::ReceiverOutputPacketTrace makeTrace(const SerializedPacket& packet,
79 std::string event) const;
80 [[nodiscard]] core::ReceiverOutputPacketTrace makeDropTrace(const DroppedDatagram& dropped) const;
81 [[nodiscard]] SerializedPacket buildContextPacket(std::uint32_t stream_id, RealType simulation_time,
82 bool stream_open, bool stream_close);
83 void emitContext(std::uint32_t stream_id, RealType simulation_time, bool stream_open, bool stream_close);
84 void applyDropped(const DroppedDatagram& dropped);
85
86 core::OutputConfig _config;
87 std::string _simulation_name;
88 core::ReceiverOutputTelemetryCallback _telemetry_callback;
89 std::unique_ptr<DatagramSender> _provided_sender;
90 StreamRegistry _registry;
91 std::unique_ptr<Vita49Packetizer> _packetizer;
92 std::unique_ptr<PacedSender> _sender;
93 std::unordered_map<std::uint32_t, StreamState> _streams;
94 std::vector<PendingContext> _pending_contexts;
95 mutable std::recursive_mutex _mutex;
96 std::chrono::steady_clock::time_point _last_stats_emit = std::chrono::steady_clock::time_point::min();
97 std::chrono::steady_clock::time_point _last_packet_trace_emit = std::chrono::steady_clock::time_point::min();
98 std::vector<core::ReceiverOutputPacketTrace> _pending_packet_traces;
99 std::uint64_t _trace_sequence = 0;
100 bool _initialized = false;
101 bool _finalized = false;
102 bool _pacing_started = false;
103 };
104
105 [[nodiscard]] std::unique_ptr<core::ReceiverOutputSink>
107}
Vita49OutputSink & operator=(Vita49OutputSink &&)=delete
void closeStream(std::uint32_t stream_id) override
core::OutputStats snapshotStats() const override
Vita49OutputSink & operator=(const Vita49OutputSink &)=delete
void emitContextHeartbeat(RealType simulation_time) override
Vita49OutputSink(Vita49OutputSink &&)=delete
void submitBlocks(std::span< const core::ReceiverSampleBlock > blocks) override
core::OutputStats finalize() override
void openStream(std::uint32_t stream_id, RealType first_sample_time) override
std::uint32_t registerStream(const core::ReceiverStreamDescriptor &stream) override
void submitBlock(const core::ReceiverSampleBlock &block) override
void initializeRun(const core::OutputConfig &config, std::string simulation_name) override
Vita49OutputSink(const Vita49OutputSink &)=delete
double RealType
Type for real numbers.
Definition config.h:27
std::function< void(const std::optional< OutputStats > &, std::span< const ReceiverOutputPacketTrace >)> ReceiverOutputTelemetryCallback
std::unique_ptr< core::ReceiverOutputSink > makeVita49OutputSink(core::ReceiverOutputTelemetryCallback telemetry_callback)
math::Vec3 max