27 const auto now = std::chrono::system_clock::now().time_since_epoch();
28 return static_cast<std::uint64_t
>(std::chrono::duration_cast<std::chrono::nanoseconds>(
now).count());
34 .fractional_picoseconds = timestamp.fractional_picoseconds};
37 [[
nodiscard]] std::optional<core::Vita49Timestamp>
74 std::scoped_lock
const lock(_mutex);
77 throw std::invalid_argument(
"Vita49OutputSink requires VITA output mode");
79 if (
config.vita49.host.empty() ||
config.vita49.port == 0)
81 throw std::invalid_argument(
"VITA output requires destination host and port");
83 if (
config.vita49.queue_depth == 0)
85 throw std::invalid_argument(
"VITA output queue depth must be positive");
89 _simulation_name = std::move(simulation_name);
91 _sender = std::make_unique<PacedSender>(
92 _provided_sender ? std::move(_provided_sender) : std::make_unique<UdpSender>(),
config.vita49.queue_depth);
94 _pending_contexts.clear();
95 _pacing_started =
false;
96 _last_stats_emit = std::chrono::steady_clock::time_point::min();
97 _last_packet_trace_emit = std::chrono::steady_clock::now();
98 _pending_packet_traces.clear();
102 emitTelemetry({},
true);
107 std::scoped_lock
const lock(_mutex);
109 if (!_streams.contains(stream_id))
112 state.descriptor = stream;
115 state.stats.stream_id = stream_id;
116 state.stats.mode = stream.
mode.empty() ?
"unknown" : stream.
mode;
119 _streams.emplace(stream_id, std::move(
state));
120 emitTelemetry({},
true);
127 std::scoped_lock
const lock(_mutex);
128 emitContext(stream_id, first_sample_time,
true,
false);
129 stateFor(stream_id).opened =
true;
130 emitTelemetry({},
true);
137 std::scoped_lock
const lock(_mutex);
138 if (!_initialized || !_sender)
140 throw std::logic_error(
"VITA output sink has not been initialized");
142 emitTelemetry(consumeSenderDropsLocked(),
false);
147 std::uint32_t stream_id = 0;
148 std::uint64_t sample_count = 0;
152 std::optional<core::Vita49Timestamp> end_timestamp;
153 bool over_range =
false;
156 std::vector<SerializedPacket> packets;
162 auto&
state = stateFor(stream_id);
165 emitContext(stream_id,
block.first_sample_time,
true,
false);
174 auto&
state = stateFor(stream_id);
175 auto result = _packetizer->packetize(
block, stream_id,
state.packet_counts,
state.sample_loss_pending);
176 state.sample_loss_pending =
false;
180 const auto end_sample_time = sample_rate > 0.0
182 :
packet.first_sample_time;
184 .stream_id = stream_id,
185 .sample_count =
packet.sample_count,
186 .first_sample_time =
packet.first_sample_time,
187 .end_sample_time = end_sample_time,
190 .over_range =
packet.over_range});
192 state.stats.over_range_count +=
result.over_range_count;
193 packets.insert(packets.end(), std::make_move_iterator(
result.packets.begin()),
194 std::make_move_iterator(
result.packets.end()));
198 { return lhs.first_sample_time < rhs.first_sample_time; });
200 if (!enqueuePackets(std::move(packets)))
207 ++stateFor(stream_id).stats.context_packets;
212 state.stats.samples_emitted +=
packet.sample_count;
213 ++
state.stats.packets_emitted;
214 if (!
state.stats.first_sample_time.has_value())
216 state.stats.first_sample_time =
packet.first_sample_time;
219 state.stats.end_sample_time =
packet.end_sample_time;
223 state.over_range_pending =
true;
230 std::scoped_lock
const lock(_mutex);
231 emitTelemetry(consumeSenderDropsLocked(),
false);
232 for (
auto& [stream_id,
state] : _streams)
234 if (!
state.closed && simulation_time -
state.last_context_time >= 1.0)
236 emitContext(stream_id, simulation_time,
false,
false);
243 std::scoped_lock
const lock(_mutex);
244 emitTelemetry(consumeSenderDropsLocked(),
false);
245 auto&
state = stateFor(stream_id);
248 emitContext(stream_id,
state.last_context_time,
false,
true);
250 emitTelemetry({},
true);
256 std::scoped_lock
const lock(_mutex);
259 auto stats = snapshotStatsLocked();
260 emitTelemetry({},
true);
267 emitTelemetry(consumeSenderDropsLocked(),
false);
270 for (
auto& [stream_id,
state] : _streams)
274 emitContext(stream_id,
state.last_context_time,
false,
true);
279 if (!_pacing_started && !_pending_contexts.empty())
282 std::vector<SerializedPacket> packets;
285 std::stable_sort(packets.begin(), packets.end(),
287 { return lhs.first_sample_time < rhs.first_sample_time; });
289 if (enqueuePackets(std::move(packets)))
293 ++stateFor(stream_id).stats.context_packets;
301 emitTelemetry(consumeSenderDropsLocked(),
false);
305 .epoch_unix_nanoseconds = _packetizer
306 ? std::optional<std::uint64_t>(_packetizer->epochUnixNanoseconds())
309 for (
auto& [stream_id,
state] : _streams)
313 state.stats.late_data_packet_count = _sender->lateDataPacketCount(stream_id);
314 state.stats.late_context_packet_count = _sender->lateContextPacketCount(stream_id);
316 stats.streams.push_back(
state.stats);
319 emitTelemetry({},
true);
325 std::scoped_lock
const lock(_mutex);
326 return snapshotStatsLocked();
329 Vita49OutputSink::StreamState& Vita49OutputSink::stateFor(
const std::uint32_t stream_id)
331 const auto found = _streams.find(stream_id);
332 if (
found == _streams.end())
334 throw std::out_of_range(
"Unknown VITA stream ID");
336 return found->second;
339 const Vita49OutputSink::StreamState& Vita49OutputSink::stateFor(
const std::uint32_t stream_id)
const
341 const auto found = _streams.find(stream_id);
342 if (
found == _streams.end())
344 throw std::out_of_range(
"Unknown VITA stream ID");
346 return found->second;
349 void Vita49OutputSink::ensurePacketizer()
360 void Vita49OutputSink::startPacing()
368 throw std::logic_error(
"VITA paced sender is unavailable");
371 _pacing_started =
true;
374 void Vita49OutputSink::appendPendingContexts(std::vector<SerializedPacket>& packets,
377 packets.reserve(packets.size() + _pending_contexts.size());
379 for (
const auto&
pending : _pending_contexts)
381 packets.push_back(buildContextPacket(
pending.stream_id,
pending.simulation_time,
pending.stream_open,
385 _pending_contexts.clear();
391 .epoch_unix_nanoseconds = _packetizer
392 ? std::optional<std::uint64_t>(_packetizer->epochUnixNanoseconds())
393 : _config.vita49.epoch_unix_nanoseconds,
395 for (
const auto& [stream_id,
state] : _streams)
400 stream_stats.late_data_packet_count = _sender->lateDataPacketCount(stream_id);
401 stream_stats.late_context_packet_count = _sender->lateContextPacketCount(stream_id);
408 std::vector<core::ReceiverOutputPacketTrace> Vita49OutputSink::consumeSenderDropsLocked()
410 std::vector<core::ReceiverOutputPacketTrace>
traces;
416 for (
const auto& dropped : _sender->consumeDroppedDatagrams())
418 applyDropped(dropped);
419 if (_config.vita49.packet_trace_enabled)
421 traces.push_back(makeDropTrace(dropped));
427 bool Vita49OutputSink::enqueuePacket(SerializedPacket&&
packet)
431 throw std::logic_error(
"VITA paced sender is unavailable");
434 std::vector<core::ReceiverOutputPacketTrace>
traces;
435 std::optional<core::ReceiverOutputPacketTrace>
sent_trace;
436 if (_config.vita49.packet_trace_enabled)
440 const auto result = _sender->enqueue(std::move(
packet));
441 if (_config.vita49.packet_trace_enabled &&
result.dropped)
447 applyDropped(*
result.dropped);
453 emitTelemetry(std::move(
traces),
false);
457 bool Vita49OutputSink::enqueuePackets(std::vector<SerializedPacket> packets)
461 throw std::logic_error(
"VITA paced sender is unavailable");
464 std::vector<core::ReceiverOutputPacketTrace>
traces;
465 if (_config.vita49.packet_trace_enabled)
467 traces.reserve(packets.size());
468 for (
const auto&
packet : packets)
473 const bool enqueued = _sender->enqueueBatch(std::move(packets));
478 emitTelemetry(std::move(
traces),
false);
482 void Vita49OutputSink::emitTelemetry(std::vector<core::ReceiverOutputPacketTrace> packets,
const bool force_stats)
484 if (!_telemetry_callback)
489 std::optional<core::OutputStats> stats;
490 const auto now = std::chrono::steady_clock::now();
491 if (
force_stats || _last_stats_emit == std::chrono::steady_clock::time_point::min() ||
494 stats = snapshotStatsLocked();
495 _last_stats_emit =
now;
498 for (
auto&
packet : packets)
500 packet.sequence = ++_trace_sequence;
501 _pending_packet_traces.push_back(std::move(
packet));
504 std::vector<core::ReceiverOutputPacketTrace>
packet_batch;
505 const bool trace_interval_elapsed = _last_packet_trace_emit == std::chrono::steady_clock::time_point::min() ||
507 if (!_pending_packet_traces.empty() &&
511 _pending_packet_traces.clear();
512 _last_packet_trace_emit =
now;
525 .event = std::move(event),
526 .stream_id =
packet.stream_id,
527 .byte_count =
packet.bytes.size(),
528 .sample_count =
packet.sample_count,
529 .first_sample_time =
packet.first_sample_time,
531 .data_packet =
packet.data_packet,
532 .context_packet =
packet.context_packet,
534 .over_range =
packet.over_range,
535 .sample_loss =
packet.sample_loss};
542 .stream_id = dropped.stream_id,
544 .sample_count = dropped.sample_count,
545 .first_sample_time = 0.0,
546 .timestamp = std::nullopt,
547 .data_packet = dropped.data_packet,
548 .context_packet = dropped.context_packet,
551 .sample_loss =
true};
554 SerializedPacket Vita49OutputSink::buildContextPacket(
const std::uint32_t stream_id,
const RealType simulation_time,
555 const bool stream_open,
const bool stream_close)
559 throw std::logic_error(
"VITA packetizer is unavailable");
561 auto&
state = stateFor(stream_id);
564 const ContextBuildRequest
request{.stream =
state.descriptor,
565 .stream_id = stream_id,
566 .simulation_name = _simulation_name,
567 .adc_fullscale = _packetizer->adcFullscale(),
568 .timestamp = timestamp,
569 .packet_count =
state.packet_counts.next(),
571 .calibrated_time =
true,
572 .reference_lock =
true,
573 .over_range =
state.over_range_pending,
574 .sample_loss =
state.sample_loss_pending,
575 .stream_open = stream_open,
576 .stream_close = stream_close};
577 const auto context = Vita49ContextBuilder::build(
request);
578 auto packet = _packetizer->makeContextPacket(context);
581 state.sample_loss_pending =
false;
582 state.over_range_pending =
false;
586 void Vita49OutputSink::emitContext(
const std::uint32_t stream_id,
const RealType simulation_time,
587 const bool stream_open,
const bool stream_close)
589 if (!_pacing_started)
592 _pending_contexts.push_back(PendingContext{.stream_id = stream_id,
594 .stream_open = stream_open,
595 .stream_close = stream_close});
599 auto packet = buildContextPacket(stream_id, simulation_time, stream_open, stream_close);
600 if (enqueuePacket(std::move(
packet)))
602 ++stateFor(stream_id).stats.context_packets;
606 void Vita49OutputSink::applyDropped(
const DroppedDatagram& dropped)
608 if (dropped.stream_id == 0 || !_streams.contains(dropped.stream_id))
612 auto&
state = stateFor(dropped.stream_id);
613 ++
state.stats.packets_dropped;
614 state.stats.samples_dropped += dropped.sample_count;
615 if (dropped.data_packet || (!dropped.context_packet && dropped.sample_count > 0))
617 state.stats.packets_emitted -= std::min<std::uint64_t>(
state.stats.packets_emitted, 1u);
618 state.stats.samples_emitted -= std::min(
state.stats.samples_emitted, dropped.sample_count);
620 else if (dropped.context_packet)
622 state.stats.context_packets -= std::min<std::uint64_t>(
state.stats.context_packets, 1u);
624 state.sample_loss_pending =
true;
627 std::unique_ptr<core::ReceiverOutputSink>
std::uint32_t registerStream(const core::ReceiverStreamDescriptor &stream)
void closeStream(std::uint32_t stream_id) override
core::OutputStats snapshotStats() const override
void emitContextHeartbeat(RealType simulation_time) override
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
~Vita49OutputSink() override
void initializeRun(const core::OutputConfig &config, std::string simulation_name) override
Vita49OutputSink(std::unique_ptr< DatagramSender > sender=nullptr, core::ReceiverOutputTelemetryCallback telemetry_callback=nullptr)
double RealType
Type for real numbers.
std::function< void(const std::optional< OutputStats > &, std::span< const ReceiverOutputPacketTrace >)> ReceiverOutputTelemetryCallback
RealType startTime() noexcept
Get the start time for the simulation.
Timestamp timestampFromEpoch(const std::uint64_t epoch_unix_nanoseconds, const RealType sample_time_seconds)
std::unique_ptr< core::ReceiverOutputSink > makeVita49OutputSink(core::ReceiverOutputTelemetryCallback telemetry_callback)
Defines the Parameters struct and provides methods for managing simulation parameters.
Vita49OutputConfig vita49
std::string receiver_name
RealType reference_frequency
std::uint16_t max_udp_payload
std::optional< std::uint64_t > epoch_unix_nanoseconds
std::uint32_t integer_seconds