From c984becd7f172e11dab3df4f75cc886f7a266335 Mon Sep 17 00:00:00 2001 From: vroad <396351+vroad@users.noreply.github.com> Date: Wed, 5 Aug 2026 13:55:29 +0000 Subject: [PATCH 1/2] fix(record-service): correct the help text for the `stop` message `SendStopMessage()` emits `"type": "stop"`, but the help text documented the message as `end`. --- src/main.cc | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/src/main.cc b/src/main.cc index 9bf54ec..987d65d 100644 --- a/src/main.cc +++ b/src/main.cc @@ -690,12 +690,12 @@ JSON Messages: "type": "start" }} - end - The `end` message is sent when `record-service` ends. The message structure - is like below: + stop + The `stop` message is sent when `record-service` ends. The message + structure is like below: {{ - "type": "end", + "type": "stop", "data": {{ "reset": false, }} From 4cbee39321151cb46da22e7a31004392743eda54 Mon Sep 17 00:00:00 2001 From: vroad <396351+vroad@users.noreply.github.com> Date: Sat, 8 Aug 2026 05:12:10 +0000 Subject: [PATCH 2/2] feat: initial packet-stats implementation --- CMakeLists.txt | 2 + src/main.cc | 59 +++- src/packet_stats_collector.hh | 219 +++++++++++++ src/service_recorder.hh | 74 ++++- test/cli_tests.sh | 1 + test/packet_stats_collector_test.cc | 482 ++++++++++++++++++++++++++++ test/service_recorder_test.cc | 170 +++++++++- 7 files changed, 993 insertions(+), 14 deletions(-) create mode 100644 src/packet_stats_collector.hh create mode 100644 test/packet_stats_collector_test.cc diff --git a/CMakeLists.txt b/CMakeLists.txt index d078cdd..6e53e1d 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -807,6 +807,7 @@ add_executable(mirakc-arib src/main.cc src/packet_sink.hh src/packet_source.hh + src/packet_stats_collector.hh src/pcr_synchronizer.hh src/pes_printer.hh src/program_filter.hh @@ -880,6 +881,7 @@ if(MIRAKC_ARIB_TEST) test/eitpf_collector_test.cc test/logo_collector_test.cc test/packet_source_test.cc + test/packet_stats_collector_test.cc test/pcr_synchronizer_test.cc test/program_filter_test.cc test/ring_file_sink_test.cc diff --git a/src/main.cc b/src/main.cc index 987d65d..1ae0bbf 100644 --- a/src/main.cc +++ b/src/main.cc @@ -94,7 +94,7 @@ Tools to process ARIB TS streams. [--pre-streaming] [] mirakc-arib filter-program-metadata [--sid=] [] mirakc-arib record-service --sid= --file= - --chunk-size= --num-chunks= [--start-pos=] [] + --chunk-size= --num-chunks= [--start-pos=] [--packet-stats] [] mirakc-arib track-airtime --sid= --eid= [] mirakc-arib seek-start --sid= [--max-duration=] [--max-packets=] [] @@ -651,7 +651,7 @@ Record a service stream into a ring buffer file Usage: mirakc-arib record-service --sid= --file= - --chunk-size= --num-chunks= [--start-pos=] [] + --chunk-size= --num-chunks= [--start-pos=] [--packet-stats] [] Options: -h --help @@ -674,6 +674,9 @@ Record a service stream into a ring buffer file A file position to start recoring. The value must be a multiple of the chunk size. + --packet-stats + Collect statistics on TS packets and send `packet-stats` messages. + Arguments: Path to a TS file. @@ -761,6 +764,51 @@ JSON Messages: The `event-end` message is sent when ended recoring a program. The message structure is the same as the `event-start` message. + packet-stats + The `packet-stats` message is sent when `--packet-stats` is specified. + The message is sent immediately before each `chunk`, `event-end`, and `stop` + message, except for the first `chunk` message which is sent when recording + starts. Each message reports statistics for the packets recorded since the + previous `packet-stats` message, or since recording started in the case of + the first message. The counters are reset each time the message is sent. + The message structure is like below: + + {{ + "type": "packet-stats", + "data": {{ + "errorPackets": 0, + "scrambledPackets": 0, + "droppedPackets": {{ + "video": 0, + "audio": 0, + "subtitle": 0, + "pmt": 0 + }} + }} + }} + + where: + errorPackets + The number of TS packets with TEI (transport_error_indicator) set to 1. + Only packets retained by service filtering are counted. + + scrambledPackets + The number of TS packets with TSC (transport_scrambling_control) set to + a nonzero value. Only packets retained by service filtering are + counted. Null packets and packets with TEI set are excluded because + their TSC value is not meaningful. + + droppedPackets + The estimated number of missing TS packets, calculated from continuity + counter discontinuities and broken down by category. Only packets + belonging to the selected service are counted. + + video PES packets carrying video. + audio PES packets carrying audio. + subtitle PES packets carrying subtitles, including ARIB subtitles and + superimposed text. + pmt Packets carrying the selected service's PMT. + Environment Variables: MIRAKC_ARIB_KEEP_UNICODE_SYMBOLS Set `1` if you like to keep Unicode symbols like enclosed ideographic @@ -1243,6 +1291,7 @@ void LoadOption(const Args& args, ServiceRecorderOption* opt) { static const std::string kChunkSize = "--chunk-size"; static const std::string kNumChunks = "--num-chunks"; static const std::string kStartPos = "--start-pos"; + static const std::string kPacketStats = "--packet-stats"; opt->sid = static_cast(args.at(kSid).asLong()); opt->file = args.at(kFile).asString(); @@ -1280,9 +1329,11 @@ void LoadOption(const Args& args, ServiceRecorderOption* opt) { std::abort(); } } + opt->packet_stats = args.at(kPacketStats).asBool(); MIRAKC_ARIB_INFO( - "ServiceRecorderOptions: sid={:04X} file={} chunk-size={} num-chunks={} start-pos={}", - opt->sid, opt->file, opt->chunk_size, opt->num_chunks, opt->start_pos); + "ServiceRecorderOptions: sid={:04X} file={} chunk-size={} num-chunks={} start-pos={} " + "packet-stats={}", + opt->sid, opt->file, opt->chunk_size, opt->num_chunks, opt->start_pos, opt->packet_stats); } void LoadOption(const Args& args, AirtimeTrackerOption* opt) { diff --git a/src/packet_stats_collector.hh b/src/packet_stats_collector.hh new file mode 100644 index 0000000..925992e --- /dev/null +++ b/src/packet_stats_collector.hh @@ -0,0 +1,219 @@ +// SPDX-License-Identifier: GPL-2.0-or-later + +// mirakc-arib +// Copyright (C) 2019 masnagam +// +// This program is free software; you can redistribute it and/or modify it under the terms of the +// GNU General Public License as published by the Free Software Foundation; either version 2 of the +// License, or (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; +// without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See +// the GNU General Public License for more details. +// +// You should have received a copy of the GNU General Public License along with this program; if +// not, write to the Free Software Foundation, 51 Franklin Street, Fifth Floor, Boston, MA +// 02110-1301, USA. + +#pragma once + +#include +#include +#include +#include +#include + +#include + +#include "base.hh" +#include "tsduck_helper.hh" + +namespace { + +enum class PacketCategory : uint8_t { + kVideo, + kAudio, + kSubtitle, + kPmt, + kNumCategories, +}; + +constexpr size_t kNumPacketCategories = static_cast(PacketCategory::kNumCategories); + +constexpr std::array kPacketCategoryNames = { + "video", + "audio", + "subtitle", + "pmt", +}; + +class PacketStatsCollector final { + public: + PacketStatsCollector() = default; + + // Changes the PMT PID tracked for packet statistics after a PAT update. + void SetPmtPid(ts::PID pid) { + if (pmt_pid_ == pid) { + return; + } + + const auto old_pid = pmt_pid_; + pmt_pid_ = pid; + + if (old_pid.has_value()) { + RemoveCategory(*old_pid); + } + + SetCategory(pid, PacketCategory::kPmt); + } + + // Classifies each PID carried by `pmt` into a PacketCategory. + void UpdatePidCategories(const ts::PMT& pmt) { + RebuildPidCategories(pmt); + } + + void CollectPacketStats(const ts::TSPacket& packet) { + // Packets with TEI set may have invalid PID or CC values. + // Count them as errors and return early. + if (packet.getTEI()) { + ++error_packets_; + return; + } + + auto pid = packet.getPID(); + if (pid == ts::PID_NULL) { + return; + } + + if (packet.getScrambling() != 0) { + ++scrambled_packets_; + } + + auto& stat = stats_[pid]; + if (!stat.category.has_value()) { + return; + } + + auto cc = packet.getCC(); + + // Although MPEG-2 TS permits a packet to be transmitted twice, we intentionally do not retain + // the previous packet to detect duplicates. + // + // We are not aware of any reports of duplicate packets in Japanese digital broadcasts, and + // we see little value in handling them. If a duplicate payload packet is ever transmitted, + // its repeated CC is counted as 15 dropped packets. + if (stat.last_cc != ts::INVALID_CC && !packet.getDiscontinuityIndicator()) { + // The continuity counter increments only when the packet has a payload. + // Keep the expected CC unchanged for packets without payload. + uint8_t expected_cc = + packet.hasPayload() ? ((stat.last_cc + 1) & ts::CC_MASK) : stat.last_cc; + + uint8_t missed = (cc - expected_cc) & ts::CC_MASK; + if (missed != 0) { + dropped_packets_[static_cast(*stat.category)] += missed; + } + } + + stat.last_cc = cc; + } + + void ResetPacketStats() { + error_packets_ = 0; + scrambled_packets_ = 0; + dropped_packets_.fill(0); + } + + uint64_t GetErrorPackets() const { + return error_packets_; + } + + uint64_t GetScrambledPackets() const { + return scrambled_packets_; + } + + uint64_t GetDroppedPackets(PacketCategory category) const { + MIRAKC_ARIB_ASSERT(static_cast(category) < kNumPacketCategories); + return dropped_packets_[static_cast(category)]; + } + + private: + struct PacketStat { + uint8_t last_cc = ts::INVALID_CC; + std::optional category; + }; + + static void ClearCategory(PacketStat& stat) { + stat.category.reset(); + stat.last_cc = ts::INVALID_CC; + } + + void SetCategory(ts::PID pid, PacketCategory category) { + auto& stat = stats_[pid]; + if (!stat.category.has_value()) { + categorized_pids_.push_back(pid); + } + stat.category = category; + } + + void RemoveCategory(ts::PID pid) { + ClearCategory(stats_[pid]); + const auto it = std::find(categorized_pids_.begin(), categorized_pids_.end(), pid); + if (it != categorized_pids_.end()) { + categorized_pids_.erase(it); + } + } + + // Rebuilds packet-statistics PID categories from the tracked PMT PID and PMT. + void RebuildPidCategories(const ts::PMT& pmt) { + std::vector previous_pids; + previous_pids.swap(categorized_pids_); + + // Clear the previous category assignments before rebuilding them from the current PMT. + // Do not clear the last_cc of existing stats here, as it may still be used by the current PMT. + for (auto pid : previous_pids) { + stats_[pid].category.reset(); + } + + if (pmt_pid_.has_value()) { + SetCategory(*pmt_pid_, PacketCategory::kPmt); + } + + RebuildStreamPidCategories(pmt); + + // Clear both the category and saved CC for PIDs that are no longer tracked. + for (auto pid : previous_pids) { + if (!stats_[pid].category.has_value()) { + ClearCategory(stats_[pid]); + } + } + } + + // Assigns categories to the PIDs referenced by `pmt`. + void RebuildStreamPidCategories(const ts::PMT& pmt) { + for (const auto& [pid, stream] : pmt.streams) { + // Do not classify the PMT PID as other categories. + if (pmt_pid_.has_value() && pid == *pmt_pid_) { + continue; + } + + if (stream.isVideo()) { + SetCategory(pid, PacketCategory::kVideo); + } else if (stream.isAudio()) { + SetCategory(pid, PacketCategory::kAudio); + } else if (stream.isSubtitles() || IsAribSubtitle(stream) || + IsAribSuperimposedText(stream)) { + SetCategory(pid, PacketCategory::kSubtitle); + } + } + } + + MIRAKC_ARIB_NON_COPYABLE(PacketStatsCollector); + std::array stats_; + std::vector categorized_pids_; + std::optional pmt_pid_; + uint64_t error_packets_ = 0; + uint64_t scrambled_packets_ = 0; + std::array dropped_packets_{}; +}; + +} // namespace diff --git a/src/service_recorder.hh b/src/service_recorder.hh index d7b64f5..be99f8a 100644 --- a/src/service_recorder.hh +++ b/src/service_recorder.hh @@ -29,6 +29,7 @@ #include "jsonl_source.hh" #include "logging.hh" #include "packet_sink.hh" +#include "packet_stats_collector.hh" #include "tsduck_helper.hh" #define MIRAKC_ARIB_SERVICE_RECORDER_TRACE(...) MIRAKC_ARIB_TRACE("service-recorder: " __VA_ARGS__) @@ -45,6 +46,7 @@ struct ServiceRecorderOption final { size_t chunk_size = 0; size_t num_chunks = 0; uint64_t start_pos = 0; + bool packet_stats = false; }; class ServiceRecorderTestAccessor; @@ -56,6 +58,9 @@ class ServiceRecorder final : public PacketSink, public: explicit ServiceRecorder(const ServiceRecorderOption& option) : option_(option), demux_(context_) { + if (option_.packet_stats) { + packet_stats_collector_ = std::make_unique(); + } demux_.setTableHandler(this); demux_.addPID(ts::PID_PAT); MIRAKC_ARIB_SERVICE_RECORDER_DEBUG("Demux PAT"); @@ -88,6 +93,7 @@ class ServiceRecorder final : public PacketSink, void End() override { MIRAKC_ARIB_ASSERT(sink_ != nullptr); + SendPacketStatsMessage(); SendStopMessage(sink_->IsBroken()); sink_->End(); } @@ -140,6 +146,7 @@ class ServiceRecorder final : public PacketSink, // internal consistency even if no TV program starts. In this case, the end // time of the record for the last TV program (held by `eit_`) is extended. SendEventUpdateMessage(eit_, now, pos); + SendPacketStatsMessage(); SendChunkMessage(now, pos); } @@ -202,6 +209,10 @@ class ServiceRecorder final : public PacketSink, pmt_pid_ = new_pmt_pid; demux_.addPID(pmt_pid_); MIRAKC_ARIB_SERVICE_RECORDER_DEBUG("PAT: Demux += PMT#{:04X}", pmt_pid_); + + if (packet_stats_collector_) { + packet_stats_collector_->SetPmtPid(pmt_pid_); + } } void HandlePmt(const ts::BinaryTable& table) { @@ -217,6 +228,10 @@ class ServiceRecorder final : public PacketSink, return; } + if (packet_stats_collector_) { + packet_stats_collector_->UpdatePidCategories(pmt); + } + auto pcr_pid = pmt.pcr_pid; if (!clock_.HasPid()) { MIRAKC_ARIB_SERVICE_RECORDER_DEBUG("PMT: PCR#{:04X}", pcr_pid); @@ -326,8 +341,7 @@ class ServiceRecorder final : public PacketSink, if (event_changed) { MIRAKC_ARIB_SERVICE_RECORDER_WARN("Event#{:04X} has started before Event#{:04X} ends", GetEvent(new_eit).event_id, GetEvent(eit).event_id); - UpdateEventBoundary(now, sink_->pos()); - SendEventEndMessage(eit); + HandleEventEnd(now, eit); SendEventStartMessage(new_eit); } else { if (IsUnspecifiedEventEndTime(GetEvent(eit))) { @@ -335,8 +349,7 @@ class ServiceRecorder final : public PacketSink, } else { auto end_time = GetEventEndTime(GetEvent(eit)); if (now >= end_time) { - UpdateEventBoundary(end_time, sink_->pos()); - SendEventEndMessage(eit); + HandleEventEnd(end_time, eit); event_started_ = false; // wait for new event } } @@ -347,6 +360,14 @@ class ServiceRecorder final : public PacketSink, event_started_ = true; } } + + // When option_.packet_stats is true, HandleEventEnd() sends a `packet-stats` message before + // the `event-end` message. Collect the current TS packet after HandleEventEnd() so that the + // `packet-stats` message sent by HandleEventEnd() excludes the TS packet not yet written to + // the ring buffer. + if (packet_stats_collector_) { + packet_stats_collector_->CollectPacketStats(packet); + } return sink_->HandlePacket(packet); } @@ -356,6 +377,12 @@ class ServiceRecorder final : public PacketSink, event_boundary_pos_ = pos; } + void HandleEventEnd(const ts::Time& end_time, const std::shared_ptr& eit) { + UpdateEventBoundary(end_time, sink_->pos()); + SendPacketStatsMessage(); + SendEventEndMessage(eit); + } + void SendStartMessage() { MIRAKC_ARIB_SERVICE_RECORDER_INFO("Started recording SID#{:04X}", option_.sid); @@ -426,6 +453,44 @@ class ServiceRecorder final : public PacketSink, SendEventMessage("event-end", eit, event_boundary_time_, event_boundary_pos_); } + void SendPacketStatsMessage() { + if (!packet_stats_collector_) { + return; + } + + auto error_packets = packet_stats_collector_->GetErrorPackets(); + auto scrambled_packets = packet_stats_collector_->GetScrambledPackets(); + + rapidjson::Document doc(rapidjson::kObjectType); + auto& allocator = doc.GetAllocator(); + + rapidjson::Value dropped_packets(rapidjson::kObjectType); + std::string dropped_packets_log; + for (size_t i = 0; i < kNumPacketCategories; ++i) { + auto category = static_cast(i); + auto count = packet_stats_collector_->GetDroppedPackets(category); + dropped_packets.AddMember(rapidjson::StringRef(kPacketCategoryNames[i]), count, allocator); + if (i != 0) { + dropped_packets_log += ' '; + } + dropped_packets_log += fmt::format("{}={}", kPacketCategoryNames[i], count); + } + + MIRAKC_ARIB_SERVICE_RECORDER_INFO("PacketStats: Error: {}, Scrambled: {}, Dropped: {}", + error_packets, scrambled_packets, dropped_packets_log); + + rapidjson::Value data(rapidjson::kObjectType); + data.AddMember("errorPackets", error_packets, allocator); + data.AddMember("scrambledPackets", scrambled_packets, allocator); + data.AddMember("droppedPackets", dropped_packets, allocator); + + doc.AddMember("type", "packet-stats", allocator); + doc.AddMember("data", data, allocator); + + FeedDocument(doc); + packet_stats_collector_->ResetPacketStats(); + } + void SendEventMessage(const std::string& type, const std::shared_ptr& eit, const ts::Time& time, uint64_t pos) { MIRAKC_ARIB_ASSERT(eit); @@ -480,6 +545,7 @@ class ServiceRecorder final : public PacketSink, ts::PID pmt_pid_ = ts::PID_NULL; State state_ = State::kPreparing; bool event_started_ = false; + std::unique_ptr packet_stats_collector_; friend class ServiceRecorderTestAccessor; diff --git a/test/cli_tests.sh b/test/cli_tests.sh index b2ad10a..8946d9e 100644 --- a/test/cli_tests.sh +++ b/test/cli_tests.sh @@ -89,6 +89,7 @@ assert 0 "$MIRAKC_ARIB filter-program-metadata --sid=0xFFFF" assert 0 "$MIRAKC_ARIB record-service --sid=1 --file=$TMPFILE --chunk-size=8192 --num-chunks=1" assert 0 "$MIRAKC_ARIB record-service --sid=1 --file=$TMPFILE --chunk-size=8192 --num-chunks=1 --start-pos=0" assert 0 "$MIRAKC_ARIB record-service --sid=1 --file=$TMPFILE --chunk-size=8192 --num-chunks=2 --start-pos=8192" +assert 0 "$MIRAKC_ARIB record-service --sid=1 --file=$TMPFILE --chunk-size=8192 --num-chunks=2 --packet-stats" if [ -z "$CI" ] then # This test fails in GitHub Actions. diff --git a/test/packet_stats_collector_test.cc b/test/packet_stats_collector_test.cc new file mode 100644 index 0000000..f332d56 --- /dev/null +++ b/test/packet_stats_collector_test.cc @@ -0,0 +1,482 @@ +// SPDX-License-Identifier: GPL-2.0-or-later + +// mirakc-arib +// Copyright (C) 2019 masnagam +// +// This program is free software; you can redistribute it and/or modify it under the terms of the +// GNU General Public License as published by the Free Software Foundation; either version 2 of the +// License, or (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; +// without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See +// the GNU General Public License for more details. +// +// You should have received a copy of the GNU General Public License along with this program; if +// not, write to the Free Software Foundation, 51 Franklin Street, Fifth Floor, Boston, MA +// 02110-1301, USA. + +#include + +#include +#include +#include + +#include "packet_stats_collector.hh" + +namespace { + +constexpr ts::PID kPid = 0x0100; + +ts::TSPacket MakePacket( + ts::PID pid, uint8_t cc, uint8_t data = 0xFF, bool tei = false, bool discontinuity = false) { + ts::TSPacket packet; + packet.init(pid, cc, data); + packet.setTEI(tei); + if (discontinuity) { + packet.setDiscontinuityIndicator(/*shift_payload=*/true); + } + return packet; +} + +// Builds a packet containing only an adaptation field, i.e. hasPayload() == false. +ts::TSPacket MakeNoPayloadPacket(ts::PID pid, uint8_t cc) { + ts::TSPacket packet; + packet.b[0] = 0x47; + packet.b[1] = static_cast((pid >> 8) & 0x1F); + packet.b[2] = static_cast(pid); + packet.b[3] = 0x20 | (cc & ts::CC_MASK); // 0x20: adaptation field only, no payload + packet.b[4] = ts::PKT_SIZE - 5; // Number of bytes after packet.b[4] + packet.b[5] = 0x00; // 0x00: adaptation field flags (no discontinuity, no PCR, etc.) + ::memset(packet.b + 6, 0xFF, ts::PKT_SIZE - 6); + return packet; +} + +// Configures kPid via the selected service's PMT so that its dropped packets are counted in the +// kVideo category. +void ConfigureVideoPid(PacketStatsCollector& collector) { + ts::PMT pmt; + pmt.streams.try_emplace(kPid, &pmt, ts::ST_MPEG2_VIDEO); + collector.UpdatePidCategories(pmt); +} + +void ExpectNoDroppedPackets(const PacketStatsCollector& collector) { + for (size_t i = 0; i < kNumPacketCategories; ++i) { + EXPECT_EQ(0, collector.GetDroppedPackets(static_cast(i))); + } +} + +} // namespace + +// ----------------------------------------------------------------------------------------------- +// TEI / error-packet handling +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, TeiPacketIsCountedAsErrorAndExcludedFromContinuity) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 0)); + // A packet with TEI set may have a corrupted PID/CC, so it must not become the baseline for + // the next continuity check. + ts::TSPacket tei_packet = MakePacket(kPid, 5, 0xFF, /*tei=*/true); + tei_packet.setScrambling(1); + collector.CollectPacketStats(tei_packet); + + EXPECT_EQ(1, collector.GetErrorPackets()); + // The collector should not count a packet with TEI set as scrambled even if the scrambling bits + // happen to be set. + EXPECT_EQ(0, collector.GetScrambledPackets()); + + // The next regular packet must still be checked against the first packet's CC 0, not against + // the TEI packet's CC 5. + collector.CollectPacketStats(MakePacket(kPid, 1)); + EXPECT_EQ(0, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, TeiWithDiscontinuityIndicatorIsStillCountedAsError) { + PacketStatsCollector collector; + collector.CollectPacketStats(MakePacket(kPid, 0, 0xFF, /*tei=*/true, /*discontinuity=*/true)); + + EXPECT_EQ(1, collector.GetErrorPackets()); +} + +TEST(PacketStatsCollectorTest, TeiPacketIsCountedEvenWhenPidIsIgnored) { + PacketStatsCollector collector; + ts::TSPacket packet = MakePacket(kPid, 0, 0xFF, /*tei=*/true); + packet.setScrambling(1); + collector.CollectPacketStats(packet); + + EXPECT_EQ(1, collector.GetErrorPackets()); + EXPECT_EQ(0, collector.GetScrambledPackets()); + ExpectNoDroppedPackets(collector); +} + +// ----------------------------------------------------------------------------------------------- +// Scrambled-packet handling +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, ScrambledPacketIsCounted) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + ts::TSPacket packet = MakePacket(kPid, 0); + packet.setScrambling(1); + collector.CollectPacketStats(packet); + + EXPECT_EQ(1, collector.GetScrambledPackets()); +} + +TEST(PacketStatsCollectorTest, ScrambledPacketIsCountedEvenWhenPidIsUnclassified) { + PacketStatsCollector collector; + ts::TSPacket packet = MakePacket(kPid, 0); + packet.setScrambling(1); + collector.CollectPacketStats(packet); + + EXPECT_EQ(1, collector.GetScrambledPackets()); +} + +TEST(PacketStatsCollectorTest, ScrambledNullPacketIsNotCounted) { + PacketStatsCollector collector; + ts::TSPacket packet = MakePacket(ts::PID_NULL, 0); + packet.setScrambling(1); + collector.CollectPacketStats(packet); + + EXPECT_EQ(0, collector.GetScrambledPackets()); +} + +// ----------------------------------------------------------------------------------------------- +// PID classification via PMT +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, PmtPidIsClassifiedAsPmt) { + constexpr ts::PID kPmtPid = 0x0030; + PacketStatsCollector collector; + collector.SetPmtPid(kPmtPid); + collector.CollectPacketStats(MakePacket(kPmtPid, 0)); + collector.CollectPacketStats(MakePacket(kPmtPid, 5)); + + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kPmt)); +} + +TEST(PacketStatsCollectorTest, UnclassifiedPidIsIgnored) { + PacketStatsCollector collector; + collector.CollectPacketStats(MakePacket(kPid, 0)); + collector.CollectPacketStats(MakePacket(kPid, 5)); + + ExpectNoDroppedPackets(collector); +} + +TEST(PacketStatsCollectorTest, PmtClassifiesStreamsByCategory) { + constexpr ts::PID kVideoPid = 0x0101; + constexpr ts::PID kAudioPid = 0x0102; + constexpr ts::PID kSubtitlePid = 0x0103; + + ts::DuckContext context; + ts::PMT pmt; + pmt.streams.try_emplace(kVideoPid, &pmt, ts::ST_MPEG2_VIDEO); + pmt.streams.try_emplace(kAudioPid, &pmt, ts::ST_MPEG1_AUDIO); + pmt.streams.try_emplace(kSubtitlePid, &pmt, ts::ST_PES_PRIV); + pmt.streams[kSubtitlePid].descs.add(context, ts::StreamIdentifierDescriptor(0x30)); + + PacketStatsCollector collector; + collector.UpdatePidCategories(pmt); + + collector.CollectPacketStats(MakePacket(kVideoPid, 0)); + collector.CollectPacketStats(MakePacket(kVideoPid, 5)); + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kVideo)); + collector.CollectPacketStats(MakePacket(kAudioPid, 0)); + collector.CollectPacketStats(MakePacket(kAudioPid, 3)); + EXPECT_EQ(2, collector.GetDroppedPackets(PacketCategory::kAudio)); + collector.CollectPacketStats(MakePacket(kSubtitlePid, 0)); + collector.CollectPacketStats(MakePacket(kSubtitlePid, 2)); + EXPECT_EQ(1, collector.GetDroppedPackets(PacketCategory::kSubtitle)); +} + +TEST(PacketStatsCollectorTest, PmtPidCategoryIsNotUpdatedByInvalidPmt) { + constexpr ts::PID kPmtPid = 0x0100; + + ts::PMT pmt; + // This invalid PMT references the selected service's PMT PID as its PCR and video stream. + pmt.pcr_pid = kPmtPid; + pmt.streams.try_emplace(kPmtPid, &pmt, ts::ST_MPEG2_VIDEO); + + PacketStatsCollector collector; + collector.SetPmtPid(kPmtPid); + collector.UpdatePidCategories(pmt); + + collector.CollectPacketStats(MakePacket(kPmtPid, 0)); + collector.CollectPacketStats(MakePacket(kPmtPid, 5)); + + // The PMT PID must remain classified as PMT, not video, even when the invalid PMT specifies it + // as a video stream. + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kPmt)); + EXPECT_EQ(0, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, PmtUpdateClassifiesPreviouslyIgnoredPid) { + PacketStatsCollector collector; + + // This packet is not part of the selected service yet. + collector.CollectPacketStats(MakePacket(kPid, 0)); + + ts::PMT pmt; + pmt.streams.try_emplace(kPid, &pmt, ts::ST_MPEG2_VIDEO); + collector.UpdatePidCategories(pmt); + + // The collector must not record CC 0 from the kPid packet received before kPid + // was classified as a video PID. + // Therefore, the next kPid packet with CC 2 must be treated as the first collected video + // packet, and must not add any dropped video packets. + collector.CollectPacketStats(MakePacket(kPid, 2)); + EXPECT_EQ(0, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, PmtUpdateResetsCategoryOfPidNoLongerReferenced) { + ts::PMT pmt; + pmt.streams.try_emplace(kPid, &pmt, ts::ST_MPEG2_VIDEO); + + PacketStatsCollector collector; + collector.UpdatePidCategories(pmt); + collector.CollectPacketStats(MakePacket(kPid, 0)); + collector.CollectPacketStats(MakePacket(kPid, 5)); + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kVideo)); + + // Because the new PMT no longer references kPid, subsequent kPid packets must be ignored. + ts::PMT next_pmt; + collector.UpdatePidCategories(next_pmt); + collector.CollectPacketStats(MakePacket(kPid, 6)); + collector.CollectPacketStats(MakePacket(kPid, 9)); + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +// ----------------------------------------------------------------------------------------------- +// PMT PID switching (SetPmtPid) +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, SetPmtPidResetsCategoryOfPreviousPmtPid) { + constexpr ts::PID kOldPmtPid = 0x0100; + constexpr ts::PID kNewPmtPid = 0x0200; + + PacketStatsCollector collector; + collector.SetPmtPid(kOldPmtPid); + collector.SetPmtPid(kNewPmtPid); + + // The old PMT PID is no longer part of the selected service and must be ignored. + collector.CollectPacketStats(MakePacket(kOldPmtPid, 0)); + collector.CollectPacketStats(MakePacket(kOldPmtPid, 5)); + ExpectNoDroppedPackets(collector); + + collector.CollectPacketStats(MakePacket(kNewPmtPid, 0)); + collector.CollectPacketStats(MakePacket(kNewPmtPid, 5)); + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kPmt)); +} + +TEST(PacketStatsCollectorTest, ChangingPmtPidKeepsStreamClassificationUntilNewPmt) { + constexpr ts::PID kNewPmtPid = 0x0200; + + ts::PMT pmt; + pmt.streams.try_emplace(kPid, &pmt, ts::ST_MPEG2_VIDEO); + + PacketStatsCollector collector; + collector.UpdatePidCategories(pmt); + + // Simulate a PAT update that advertises kNewPmtPid as the PMT PID for the selected service. + // No PMT has arrived on kNewPmtPid yet. + collector.SetPmtPid(kNewPmtPid); + + collector.CollectPacketStats(MakePacket(kPid, 0)); + collector.CollectPacketStats(MakePacket(kPid, 5)); + + // Until a PMT arrives on kNewPmtPid, kPid must remain classified as video using the last PMT. + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +// ----------------------------------------------------------------------------------------------- +// PCR PID handling +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, PcrPidIsIgnored) { + constexpr ts::PID kVideoPid = 0x0101; + constexpr ts::PID kPcrPid = 0x0102; + + ts::PMT pmt; + pmt.pcr_pid = kPcrPid; + pmt.streams.try_emplace(kVideoPid, &pmt, ts::ST_MPEG2_VIDEO); + + PacketStatsCollector collector; + collector.UpdatePidCategories(pmt); + + collector.CollectPacketStats(MakePacket(kPcrPid, 0)); + collector.CollectPacketStats(MakePacket(kPcrPid, 5)); + + // A PCR-only PID, i.e. a PCR PID that is not referenced by any stream in the PMT, must not be + // classified into any category. + ExpectNoDroppedPackets(collector); +} + +TEST(PacketStatsCollectorTest, PcrPidSharedWithVideoStreamIsClassifiedAsVideo) { + constexpr ts::PID kVideoPid = 0x0101; + + ts::PMT pmt; + pmt.pcr_pid = kVideoPid; + pmt.streams.try_emplace(kVideoPid, &pmt, ts::ST_MPEG2_VIDEO); + + PacketStatsCollector collector; + collector.UpdatePidCategories(pmt); + + collector.CollectPacketStats(MakePacket(kVideoPid, 0)); + collector.CollectPacketStats(MakePacket(kVideoPid, 5)); + EXPECT_EQ(4, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +// ----------------------------------------------------------------------------------------------- +// NULL PID handling +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, NullPidInPmtIsIgnored) { + ts::PMT pmt; + // The collector must not track the NULL PID specified in any PMT field, because the NULL PID's + // CC is meaningless. + pmt.pcr_pid = ts::PID_NULL; + pmt.streams.try_emplace(ts::PID_NULL, &pmt, ts::ST_MPEG2_VIDEO); + + PacketStatsCollector collector; + collector.UpdatePidCategories(pmt); + + collector.CollectPacketStats(MakePacket(ts::PID_NULL, 0)); + collector.CollectPacketStats(MakePacket(ts::PID_NULL, 5)); + + ExpectNoDroppedPackets(collector); +} + +// ----------------------------------------------------------------------------------------------- +// Continuity counter / dropped-packet counting +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, RegularPacketsProduceZeroErrorStatistics) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 0)); + collector.CollectPacketStats(MakePacket(kPid, 1)); + collector.CollectPacketStats(MakePacket(kPid, 2)); + + EXPECT_EQ(0, collector.GetErrorPackets()); + EXPECT_EQ(0, collector.GetScrambledPackets()); + ExpectNoDroppedPackets(collector); +} + +TEST(PacketStatsCollectorTest, SameCcRecordsDroppedPackets) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 0, 0xAA)); + collector.CollectPacketStats(MakePacket(kPid, 0, 0xBB)); + + // A 0-to-0 CC transition must be counted as 15 dropped packets. + EXPECT_EQ(15, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, MultipleMissingPacketsIncreasesDropped) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 5)); + collector.CollectPacketStats(MakePacket(kPid, 9)); + + EXPECT_EQ(3, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, WrapAroundMissingPacketsIncreasesDropped) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 14)); + collector.CollectPacketStats(MakePacket(kPid, 2)); + + EXPECT_EQ(3, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, NoPayloadSameCcDoesNotIncreaseDropped) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakeNoPayloadPacket(kPid, 3)); + collector.CollectPacketStats(MakeNoPayloadPacket(kPid, 3)); + + EXPECT_EQ(0, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, NoPayloadCcChangedIncreasesDropped) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakeNoPayloadPacket(kPid, 3)); + collector.CollectPacketStats(MakeNoPayloadPacket(kPid, 4)); + + EXPECT_EQ(1, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +TEST(PacketStatsCollectorTest, DiscontinuityIndicatorDoesNotIncreaseDropped) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 0)); + collector.CollectPacketStats(MakePacket(kPid, 8, 0xFF, /*tei=*/false, /*discontinuity=*/true)); + + EXPECT_EQ(0, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +// ----------------------------------------------------------------------------------------------- +// CC state across PMT re-categorization +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, RecategorizedPidDropsItsCcState) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 0)); + + // This PMT update leaves kPid uncategorized. Packets received on kPid are ignored and do not + // update its CC state. + collector.UpdatePidCategories(ts::PMT()); + collector.CollectPacketStats(MakePacket(kPid, 1)); + collector.CollectPacketStats(MakePacket(kPid, 2)); + + // A later PMT update categorizes kPid again. CC=3 must not be compared with the earlier CC=0. + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 3)); + collector.CollectPacketStats(MakePacket(kPid, 4)); + + ExpectNoDroppedPackets(collector); +} + +TEST(PacketStatsCollectorTest, RepeatedPmtKeepsCcState) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 0)); + + // Processing the same PMT again must not reset CC state. + // The CC gap (0 to 2) must still be counted. + ConfigureVideoPid(collector); + collector.CollectPacketStats(MakePacket(kPid, 2)); + + EXPECT_EQ(1, collector.GetDroppedPackets(PacketCategory::kVideo)); +} + +// ----------------------------------------------------------------------------------------------- +// Reset behavior +// ----------------------------------------------------------------------------------------------- + +TEST(PacketStatsCollectorTest, ResetPacketStatsClearsAllCounters) { + PacketStatsCollector collector; + ConfigureVideoPid(collector); + + collector.CollectPacketStats(MakePacket(kPid, 0)); + collector.CollectPacketStats(MakePacket(kPid, 5)); + ts::TSPacket tei_packet = MakePacket(kPid, 8, 0xFF, /*tei=*/true); + collector.CollectPacketStats(tei_packet); + ts::TSPacket scrambled_packet = MakePacket(kPid, 9); + scrambled_packet.setScrambling(1); + collector.CollectPacketStats(scrambled_packet); + + EXPECT_EQ(1, collector.GetErrorPackets()); + EXPECT_EQ(1, collector.GetScrambledPackets()); + EXPECT_EQ(7, collector.GetDroppedPackets(PacketCategory::kVideo)); + + collector.ResetPacketStats(); + + EXPECT_EQ(0, collector.GetErrorPackets()); + EXPECT_EQ(0, collector.GetScrambledPackets()); + ExpectNoDroppedPackets(collector); +} diff --git a/test/service_recorder_test.cc b/test/service_recorder_test.cc index 0f30218..e4655c2 100644 --- a/test/service_recorder_test.cc +++ b/test/service_recorder_test.cc @@ -16,6 +16,8 @@ // 02110-1301, USA. #include +#include +#include #include #include @@ -41,6 +43,8 @@ class ServiceRecorderTestAccessor final { return recorder.clock_.IsReady(); } }; + +struct ServiceRecorderTest : testing::TestWithParam {}; } // namespace TEST(ServiceRecorderTest, NoPacket) { @@ -191,8 +195,8 @@ TEST(ServiceRecorderTest, EventStart) { EXPECT_TRUE(src.IsEmpty()); } -TEST(ServiceRecorderTest, EventProgress) { - ServiceRecorderOption option = kOption; +TEST_P(ServiceRecorderTest, EventProgress) { + ServiceRecorderOption option = GetParam(); TableSource src; auto ring_sink = std::make_unique(option.chunk_size, option.num_chunks); @@ -281,6 +285,20 @@ TEST(ServiceRecorderTest, EventProgress) { MockJsonlSink::Stringify(doc)); return true; }); + if (option.packet_stats) { + EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { + EXPECT_EQ(R"({"type":"packet-stats","data":{)" + R"("errorPackets":0,)" + R"("scrambledPackets":0,)" + R"("droppedPackets":{)" + R"("video":0,"audio":0,"subtitle":0,"pmt":0)" + R"(})" + R"(})" + R"(})", + MockJsonlSink::Stringify(doc)); + return true; + }); + } EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { EXPECT_EQ(R"({"type":"chunk","data":{"chunk":{)" R"("timestamp":1609426800000,"pos":16384)" @@ -308,6 +326,20 @@ TEST(ServiceRecorderTest, EventProgress) { MockJsonlSink::Stringify(doc)); return true; }); + if (option.packet_stats) { + EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { + EXPECT_EQ(R"({"type":"packet-stats","data":{)" + R"("errorPackets":0,)" + R"("scrambledPackets":0,)" + R"("droppedPackets":{)" + R"("video":0,"audio":0,"subtitle":0,"pmt":0)" + R"(})" + R"(})" + R"(})", + MockJsonlSink::Stringify(doc)); + return true; + }); + } EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { EXPECT_EQ(R"({"type":"chunk","data":{"chunk":{)" R"("timestamp":1609426800000,"pos":0)" @@ -315,6 +347,20 @@ TEST(ServiceRecorderTest, EventProgress) { MockJsonlSink::Stringify(doc)); return true; }); + if (option.packet_stats) { + EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { + EXPECT_EQ(R"({"type":"packet-stats","data":{)" + R"("errorPackets":0,)" + R"("scrambledPackets":0,)" + R"("droppedPackets":{)" + R"("video":0,"audio":0,"subtitle":0,"pmt":0)" + R"(})" + R"(})" + R"(})", + MockJsonlSink::Stringify(doc)); + return true; + }); + } EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { EXPECT_EQ(R"({"type":"stop","data":{"reset":false}})", MockJsonlSink::Stringify(doc)); return true; @@ -322,7 +368,7 @@ TEST(ServiceRecorderTest, EventProgress) { EXPECT_CALL(*ring_sink, End).WillOnce(testing::Return()); } - auto recorder = std::make_unique(kOption); + auto recorder = std::make_unique(option); recorder->ServiceRecorder::Connect(std::move(ring_sink)); recorder->JsonlSource::Connect(std::move(json_sink)); src.Connect(std::move(recorder)); @@ -330,8 +376,8 @@ TEST(ServiceRecorderTest, EventProgress) { EXPECT_TRUE(src.IsEmpty()); } -TEST(ServiceRecorderTest, EventEnd) { - ServiceRecorderOption option = kOption; +TEST_P(ServiceRecorderTest, EventEnd) { + ServiceRecorderOption option = GetParam(); TableSource src; auto ring_sink = std::make_unique(option.chunk_size, option.num_chunks); @@ -409,6 +455,21 @@ TEST(ServiceRecorderTest, EventEnd) { MockJsonlSink::Stringify(doc)); return true; }); + if (option.packet_stats) { + // Sent from HandleEventEnd(). + EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { + EXPECT_EQ(R"({"type":"packet-stats","data":{)" + R"("errorPackets":0,)" + R"("scrambledPackets":0,)" + R"("droppedPackets":{)" + R"("video":0,"audio":0,"subtitle":0,"pmt":0)" + R"(})" + R"(})" + R"(})", + MockJsonlSink::Stringify(doc)); + return true; + }); + } EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { EXPECT_EQ(R"({"type":"event-end","data":{)" R"("originalNetworkId":1,)" @@ -449,6 +510,21 @@ TEST(ServiceRecorderTest, EventEnd) { MockJsonlSink::Stringify(doc)); return true; }); + if (option.packet_stats) { + // Sent from End(). + EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { + EXPECT_EQ(R"({"type":"packet-stats","data":{)" + R"("errorPackets":0,)" + R"("scrambledPackets":0,)" + R"("droppedPackets":{)" + R"("video":0,"audio":0,"subtitle":0,"pmt":0)" + R"(})" + R"(})" + R"(})", + MockJsonlSink::Stringify(doc)); + return true; + }); + } EXPECT_CALL(*json_sink, HandleDocument).WillOnce([](const rapidjson::Document& doc) { EXPECT_EQ(R"({"type":"stop","data":{"reset":false}})", MockJsonlSink::Stringify(doc)); return true; @@ -456,7 +532,7 @@ TEST(ServiceRecorderTest, EventEnd) { EXPECT_CALL(*ring_sink, End).WillOnce(testing::Return()); } - auto recorder = std::make_unique(kOption); + auto recorder = std::make_unique(option); recorder->ServiceRecorder::Connect(std::move(ring_sink)); recorder->JsonlSource::Connect(std::move(json_sink)); src.Connect(std::move(recorder)); @@ -464,6 +540,84 @@ TEST(ServiceRecorderTest, EventEnd) { EXPECT_TRUE(src.IsEmpty()); } +// The `packet-stats` message sent before the `event-end` message must exclude the current TS +// packet because `sink_->HandlePacket()` has not yet written the packet to the ring buffer. +TEST(ServiceRecorderTest, EventEndDoesNotIncludeItsTriggeringPacketInPacketStats) { + ServiceRecorderOption option{"/dev/null", 3, kChunkSize, kNumChunks, 0, true}; + + TableSource src; + auto ring_sink = std::make_unique(option.chunk_size, option.num_chunks); + auto json_sink = std::make_unique(); + + src.LoadXml(R"( + + + + + + + + + + + + + + + + + + + + )"); + + std::vector message_types; + std::vector dropped_video_packet_counts; + EXPECT_CALL(*json_sink, HandleDocument) + .WillRepeatedly( + [&message_types, &dropped_video_packet_counts](const rapidjson::Document& doc) { + const std::string type = doc["type"].GetString(); + + // Accept every message, but record only `packet-stats` and `event-end` messages + // to verify their order and dropped-packet counts later with EXPECT_THAT calls. + if (type == "packet-stats") { + message_types.emplace_back(type); + dropped_video_packet_counts.emplace_back( + doc["data"]["droppedPackets"]["video"].GetUint64()); + } else if (type == "event-end") { + message_types.emplace_back(type); + } + + return true; + }); + + { + testing::InSequence seq; + EXPECT_CALL(*ring_sink, Start).WillOnce(testing::Return(true)); + EXPECT_CALL(*ring_sink, End).WillOnce(testing::Return()); + } + + auto recorder = std::make_unique(option); + recorder->ServiceRecorder::Connect(std::move(ring_sink)); + recorder->JsonlSource::Connect(std::move(json_sink)); + src.Connect(std::move(recorder)); + + EXPECT_EQ(EXIT_SUCCESS, src.FeedPackets()); + EXPECT_TRUE(src.IsEmpty()); + EXPECT_THAT(message_types, testing::ElementsAre("packet-stats", "event-end", "packet-stats")); + + // The first `packet-stats` message must report zero dropped `video` packets because the final + // PCR packet has not been written to the ring buffer yet. + EXPECT_THAT(dropped_video_packet_counts, testing::ElementsAre(0, 4)); +} + TEST(ServiceRecorderTest, EventStartBeforeEventEnd) { ServiceRecorderOption option = kOption; @@ -991,3 +1145,7 @@ TEST(ServiceRecorderTest, EndOfChunkBeforeNextEvent) { EXPECT_EQ(EXIT_SUCCESS, src.FeedPackets()); EXPECT_TRUE(src.IsEmpty()); } + +INSTANTIATE_TEST_SUITE_P(EventTestWithPacketStats, ServiceRecorderTest, + testing::Values(ServiceRecorderOption{"/dev/null", 3, kChunkSize, kNumChunks, 0, false}, + ServiceRecorderOption{"/dev/null", 3, kChunkSize, kNumChunks, 0, true}));