From 943596da6d96d65154c493eeb9cca82fbb4bcf21 Mon Sep 17 00:00:00 2001 From: Florian Egger Date: Tue, 8 Sep 2026 16:54:09 +0200 Subject: [PATCH] feat(app): add PLI feedback, jitter reordering, and hardware decode Loss recovery for the streaming path: - PLI over signaling: the depacketizer now reports damaged frames (DepacketizeResult) and the receiver asks the sender for a keyframe (SessionPli, rate-limited to one per 500 ms). The sender keeps the signaling channel open during the session and honors PLIs through the new thread-safe SenderPipeline::request_keyframe(). Recovery takes one frame time instead of waiting out the GOP. - RtpJitterBuffer: reorders RTP packets by sequence number (16 packets / 60 ms) before the in-order depacketizer, so Wi-Fi reordering is not misread as loss; in-order streams release immediately, and a straggler older than the delivered sequence is discarded. - Hardware H.264 decode probe: DecoderFactory tries h264_v4l2m2m (the VideoCore path on the Pi) with an automatic software fallback and a clear journal line for the chosen path; --swdecode opts out. Validated: PLI end-to-end with a probe that drops a mid-keyframe packet over real UDP (receiver logged the damaged frame and the PLI arrived with the session id); hardware probe fails cleanly and falls back on this desktop; jitter reordering covered by unit tests. meson test 5/5 in both build configurations, valgrind clean. --- include/screencast/app/cli.h | 1 + include/screencast/app/pipeline.h | 4 + include/screencast/codec/decoder.h | 3 + include/screencast/network/h264_packetizer.h | 17 +++- include/screencast/network/rtp_packet.h | 29 ++++++ include/screencast/network/signaling.h | 8 +- src/app/cli.cpp | 4 +- src/app/main.cpp | 13 ++- src/app/pipelines.cpp | 53 ++++++++++- src/codec/ffmpeg_decoder.cpp | 59 ++++++++++--- src/network/h264_packetizer.cpp | 13 ++- src/network/rtp_packet.cpp | 55 ++++++++++++ src/network/signaling.cpp | 18 ++-- tests/app/test_loopback.cpp | 6 +- tests/network/test_rtp.cpp | 93 ++++++++++++++++++-- tests/network/test_signaling.cpp | 21 +++++ 16 files changed, 355 insertions(+), 42 deletions(-) diff --git a/include/screencast/app/cli.h b/include/screencast/app/cli.h index 959cb16..615783d 100644 --- a/include/screencast/app/cli.h +++ b/include/screencast/app/cli.h @@ -18,6 +18,7 @@ struct ReceiveCommand { int local_rtp_port = 5004; int signaling_port = 5005; bool fullscreen = false; + bool software_decode = false; // --swdecode disables the hardware probe }; struct DiscoverCommand { diff --git a/include/screencast/app/pipeline.h b/include/screencast/app/pipeline.h index 1a659ac..7657c3e 100644 --- a/include/screencast/app/pipeline.h +++ b/include/screencast/app/pipeline.h @@ -39,6 +39,10 @@ class SenderPipeline { bool start(); void stop(); + // Ask the sender to encode its next frame as a keyframe. Thread-safe; + // used by the PLI feedback path. + void request_keyframe(); + private: class Impl; std::unique_ptr impl_; diff --git a/include/screencast/codec/decoder.h b/include/screencast/codec/decoder.h index 909f015..c3cb25b 100644 --- a/include/screencast/codec/decoder.h +++ b/include/screencast/codec/decoder.h @@ -22,6 +22,9 @@ struct DecoderConfig { int width = 0; int height = 0; std::vector extradata; // SPS/PPS for H.264 + // Probe a hardware decoder (v4l2 mem2mem) first and fall back to the + // software decoder automatically. Disable with --swdecode. + bool hardware_accel = true; }; class Decoder { diff --git a/include/screencast/network/h264_packetizer.h b/include/screencast/network/h264_packetizer.h index 526d497..dbfd3a7 100644 --- a/include/screencast/network/h264_packetizer.h +++ b/include/screencast/network/h264_packetizer.h @@ -42,14 +42,23 @@ class H264Packetizer { std::uint16_t next_sequence_number_ = 0; }; +// Result of feeding one packet to the depacketizer. +struct DepacketizeResult { + // The completed access unit (Annex-B with 3-byte start codes) when the + // packet closed an undamaged frame. + std::optional> access_unit; + // True when this call discarded a frame as damaged (packet loss or an + // unsupported packetization). Pipelines use it to request a keyframe. + bool frame_dropped = false; +}; + // Reassembles RFC 6184 packet streams (single NAL unit packets and FU-A) // into Annex-B access units. Packets must arrive in order; frames damaged by -// sequence gaps or missing fragments are dropped silently. +// sequence gaps or missing fragments are reported via DepacketizeResult. class H264Depacketizer { public: - // Feed one packet. Returns the completed access unit (Annex-B with 3-byte - // start codes) when the packet closes a frame, nullopt otherwise. - std::optional> depacketize(const RtpPacket& packet); + // Feed one packet. + DepacketizeResult depacketize(const RtpPacket& packet); private: void drop_frame(); diff --git a/include/screencast/network/rtp_packet.h b/include/screencast/network/rtp_packet.h index f18d360..ba4ccc7 100644 --- a/include/screencast/network/rtp_packet.h +++ b/include/screencast/network/rtp_packet.h @@ -1,8 +1,12 @@ #pragma once +#include #include +#include +#include #include #include +#include #include namespace sc { @@ -31,4 +35,29 @@ struct RtpPacket { static std::optional parse(std::span in) noexcept; }; +// Reorders RTP packets by sequence number before depacketization so that a +// reordering link (Wi-Fi) does not read as loss. Delivery stays in order; +// only aged-out or overflowing buffers release out of order, which the +// downstream gap detection still handles for genuine loss. +class RtpJitterBuffer { + public: + explicit RtpJitterBuffer(std::size_t max_depth = 16, + std::chrono::milliseconds max_delay = std::chrono::milliseconds{60}); + + // Insert one packet and return the packets now ready for in-order + // delivery. In-order streams release immediately (zero added latency); + // a straggler older than the next expected sequence is discarded. + std::vector push(RtpPacket packet); + + // Discard everything still buffered. + void clear(); + + private: + std::size_t max_depth_; + std::chrono::milliseconds max_delay_; + std::mutex mutex_; + std::map> buffer_; + std::optional next_expected_; +}; + } // namespace sc diff --git a/include/screencast/network/signaling.h b/include/screencast/network/signaling.h index e7acae7..9c965fd 100644 --- a/include/screencast/network/signaling.h +++ b/include/screencast/network/signaling.h @@ -31,7 +31,13 @@ struct SessionAnswer { Endpoint rtp_endpoint; }; -using SignalingMessage = std::variant; +// Picture Loss Indication: the receiver asks the sender for a keyframe +// after discarding a damaged frame. +struct SessionPli { + std::string session_id; +}; + +using SignalingMessage = std::variant; class SignalingChannel { public: diff --git a/src/app/cli.cpp b/src/app/cli.cpp index 50e03af..cbb2c4f 100644 --- a/src/app/cli.cpp +++ b/src/app/cli.cpp @@ -9,7 +9,7 @@ namespace { void print_usage() { std::fputs("usage: screencast --send [--target monitor|window] [--peer HOST[:PORT]] [--bitrate KBPS]\n" - " screencast --receive [--port PORT] [--signaling-port PORT] [--fullscreen]\n" + " screencast --receive [--port PORT] [--signaling-port PORT] [--fullscreen] [--swdecode]\n" " screencast --discover [--timeout SECONDS]\n" "\n" "--send without --peer discovers a receiver on the LAN and requires\n" @@ -112,6 +112,8 @@ std::optional parse_cli(int argc, const char* const argv[]) { receive.signaling_port = port; } else if (argument == "--fullscreen") { receive.fullscreen = true; + } else if (argument == "--swdecode") { + receive.software_decode = true; } else if (argument == "--timeout") { std::string_view value; if (!next_argument(argc, argv, index, value) || !parse_int(value, discover.timeout_seconds) || diff --git a/src/app/main.cpp b/src/app/main.cpp index 035789f..03af16d 100644 --- a/src/app/main.cpp +++ b/src/app/main.cpp @@ -167,7 +167,8 @@ int negotiate_and_stream(sc::SignalingChannel& channel, const sc::Endpoint& sign return 1; } const sc::SessionAnswer answer = answer_future.get(); - channel.disconnect(); + // The channel stays open: the receiver sends PLI keyframe requests over + // it during the session. if (answer.session_id != offer.session_id) { std::cerr << "screencast: session mismatch in the receiver's answer\n"; @@ -196,12 +197,21 @@ int negotiate_and_stream(sc::SignalingChannel& channel, const sc::Endpoint& sign sc::SenderPipeline pipeline{std::move(config)}; if (!pipeline.start()) { + channel.disconnect(); return 1; } + channel.on_message([&](const sc::SignalingMessage& message) { + const sc::SessionPli* pli = std::get_if(&message); + if (pli != nullptr && pli->session_id == offer.session_id) { + pipeline.request_keyframe(); + std::cerr << "screencast: receiver requested a keyframe\n"; + } + }); while (!g_interrupted.load()) { std::this_thread::sleep_for(std::chrono::milliseconds(100)); } pipeline.stop(); + channel.disconnect(); return 0; } @@ -289,6 +299,7 @@ int run_receiver(const sc::ReceiveCommand& command) { sc::ReceiverPipelineConfig config; config.local_rtp_endpoint = sc::Endpoint{"0.0.0.0", static_cast(command.local_rtp_port)}; config.signaling_port = static_cast(command.signaling_port); + config.decoder.hardware_accel = !command.software_decode; // Fullscreen is automatic under KMSDRM (headless); the flag forces it // on desktop sessions. config.renderer.fullscreen = command.fullscreen; diff --git a/src/app/pipelines.cpp b/src/app/pipelines.cpp index ec345bf..513c6e0 100644 --- a/src/app/pipelines.cpp +++ b/src/app/pipelines.cpp @@ -64,9 +64,17 @@ class SenderPipeline::Impl { transport_->stop(); } + void request_keyframe() { + keyframe_requested_.store(true); + } + private: void run(std::stop_token stop_token) { while (!stop_token.stop_requested()) { + if (keyframe_requested_.exchange(false) && encoder_ != nullptr) { + encoder_->request_keyframe(); + } + const std::optional frame = capture_->next_frame(); if (!frame.has_value()) { break; @@ -109,8 +117,13 @@ class SenderPipeline::Impl { H264Packetizer packetizer_; std::unique_ptr transport_ = RtpTransportFactory::create(); std::jthread run_thread_; + std::atomic keyframe_requested_{false}; }; +void SenderPipeline::request_keyframe() { + impl_->request_keyframe(); +} + SenderPipeline::SenderPipeline(SenderPipelineConfig config) : impl_(std::make_unique(std::move(config))) {} SenderPipeline::~SenderPipeline() = default; @@ -238,6 +251,7 @@ class ReceiverPipeline::Impl { signaling_ = nullptr; discovery_ = nullptr; transport_->stop(); + jitter_.clear(); renderer_ = nullptr; decoder_ = nullptr; { @@ -248,13 +262,24 @@ class ReceiverPipeline::Impl { private: void on_packet(RtpPacket packet) { - std::optional> access_unit = depacketizer_.depacketize(packet); - if (!access_unit.has_value()) { + // Absorb reordering (Wi-Fi) before the in-order depacketizer, so a + // late packet is not misread as loss. + for (RtpPacket& ordered : jitter_.push(std::move(packet))) { + deliver_packet(ordered); + } + } + + void deliver_packet(const RtpPacket& packet) { + const DepacketizeResult result = depacketizer_.depacketize(packet); + if (result.frame_dropped) { + maybe_send_pli(); + } + if (!result.access_unit.has_value()) { return; } EncodedFrame encoded; - encoded.data = std::move(*access_unit); + encoded.data = std::move(*result.access_unit); encoded.rtp_timestamp = packet.header.timestamp; encoded.is_keyframe = false; @@ -282,6 +307,7 @@ class ReceiverPipeline::Impl { if (offer == nullptr) { return; } + session_id_ = offer->session_id; SessionAnswer answer; answer.session_id = offer->session_id; // Empty address: the sender targets the address of its signaling @@ -290,6 +316,23 @@ class ReceiverPipeline::Impl { signaling_->send(answer); } + // Ask the sender for a keyframe after a damaged frame, rate-limited so + // sustained loss cannot flood the signaling channel. + void maybe_send_pli() { + if (signaling_ == nullptr || session_id_.empty()) { + return; + } + const auto now = std::chrono::steady_clock::now(); + if (now - last_pli_time_ < kPliMinInterval) { + return; + } + last_pli_time_ = now; + SessionPli pli; + pli.session_id = session_id_; + signaling_->send(pli); + std::cerr << std::format("screencast: frame damaged; requesting a keyframe\n"); + } + void render_loop(std::stop_token stop_token) { while (!stop_token.stop_requested()) { if (!renderer_->poll_events()) { @@ -315,17 +358,21 @@ class ReceiverPipeline::Impl { } static constexpr std::size_t kMaxQueuedFrames = 3; + static constexpr std::chrono::milliseconds kPliMinInterval{500}; ReceiverPipelineConfig config_; std::unique_ptr renderer_; std::unique_ptr decoder_; H264Depacketizer depacketizer_; + RtpJitterBuffer jitter_; std::unique_ptr transport_ = RtpTransportFactory::create(); std::unique_ptr signaling_; std::unique_ptr discovery_; std::jthread render_thread_; std::mutex queue_mutex_; std::deque queue_; + std::string session_id_; + std::chrono::steady_clock::time_point last_pli_time_{}; bool stream_started_ = false; bool present_error_logged_ = false; }; diff --git a/src/codec/ffmpeg_decoder.cpp b/src/codec/ffmpeg_decoder.cpp index bf3a16e..df8afb9 100644 --- a/src/codec/ffmpeg_decoder.cpp +++ b/src/codec/ffmpeg_decoder.cpp @@ -5,6 +5,8 @@ #include #include #include +#include +#include #include #include #include @@ -199,19 +201,18 @@ class FfmpegDecoder final : public Decoder { mutable AVPixelFormat scaler_input_format_ = AV_PIX_FMT_NONE; }; -CodecResult> DecoderFactory::create(const DecoderConfig& config) { - if (config.codec_name != "h264") { - return CodecError{"only h264 is supported in phase 2"}; - } +namespace { - const AVCodec* codec = avcodec_find_decoder(AV_CODEC_ID_H264); +// Opens the given decoder implementation for the given config. Returns an +// error string on failure (used for the hardware probe + fallback). +std::optional open_decoder(const DecoderConfig& config, const AVCodec* codec, AvCodecContextPtr& ctx) { if (codec == nullptr) { - return CodecError{"h264 decoder not found"}; + return std::string{"decoder not found"}; } - AvCodecContextPtr ctx(avcodec_alloc_context3(codec), AvCodecContextDeleter{}); + ctx.reset(avcodec_alloc_context3(codec)); if (ctx == nullptr) { - return CodecError{"failed to allocate decoder context"}; + return std::string{"failed to allocate decoder context"}; } ctx->codec_type = AVMEDIA_TYPE_VIDEO; @@ -225,24 +226,58 @@ CodecResult> DecoderFactory::create(const DecoderConfig if (!config.extradata.empty()) { if (config.extradata.size() > static_cast(std::numeric_limits::max())) { - return CodecError{"decoder extradata is too large"}; + return std::string{"decoder extradata is too large"}; } // FFmpeg bitstream parsers may read past the end of extradata, so the // buffer must include the padding they require. ctx->extradata = static_cast(av_malloc(config.extradata.size() + AV_INPUT_BUFFER_PADDING_SIZE)); if (ctx->extradata == nullptr) { - return CodecError{"failed to allocate decoder extradata"}; + return std::string{"failed to allocate decoder extradata"}; } std::memcpy(ctx->extradata, config.extradata.data(), config.extradata.size()); std::memset(ctx->extradata + config.extradata.size(), 0, AV_INPUT_BUFFER_PADDING_SIZE); ctx->extradata_size = static_cast(config.extradata.size()); } - int open_ret = avcodec_open2(ctx.get(), codec, nullptr); + const int open_ret = avcodec_open2(ctx.get(), codec, nullptr); if (open_ret < 0) { - return CodecError{std::string{"failed to open h264 decoder: "} + ffmpeg_error(open_ret)}; + return std::string{"failed to open decoder: "} + ffmpeg_error(open_ret); + } + return std::nullopt; +} + +} // namespace + +CodecResult> DecoderFactory::create(const DecoderConfig& config) { + if (config.codec_name != "h264") { + return CodecError{"only h264 is supported"}; } + // Hardware first (v4l2 mem2mem: the Pi's VideoCore H.264 decoder), with + // an automatic software fallback. The kernel's bitstream parser reads + // the stream dimensions from the in-band SPS, so they are not required + // up front. + if (config.hardware_accel) { + if (const AVCodec* hw = avcodec_find_decoder_by_name("h264_v4l2m2m"); hw != nullptr) { + AvCodecContextPtr ctx{nullptr, AvCodecContextDeleter{}}; + if (auto error = open_decoder(config, hw, ctx)) { + std::cerr << std::format("screencast: hardware decode unavailable ({}); falling back to " + "software\n", + *error); + } else { + std::cerr << "screencast: using hardware H.264 decode (h264_v4l2m2m)\n"; + return std::make_unique(std::move(ctx), config); + } + } else { + std::cerr << "screencast: hardware decoder not compiled in; using software decode\n"; + } + } + + const AVCodec* codec = avcodec_find_decoder(AV_CODEC_ID_H264); + AvCodecContextPtr ctx{nullptr, AvCodecContextDeleter{}}; + if (auto error = open_decoder(config, codec, ctx)) { + return CodecError{std::move(*error)}; + } return std::make_unique(std::move(ctx), config); } diff --git a/src/network/h264_packetizer.cpp b/src/network/h264_packetizer.cpp index 6889d23..70a04a4 100644 --- a/src/network/h264_packetizer.cpp +++ b/src/network/h264_packetizer.cpp @@ -139,7 +139,9 @@ void H264Packetizer::append_fu_a_packets(std::span nal, } } -std::optional> H264Depacketizer::depacketize(const RtpPacket& packet) { +DepacketizeResult H264Depacketizer::depacketize(const RtpPacket& packet) { + DepacketizeResult result; + // Track sequence continuity: a gap means packets were lost. if (last_sequence_number_.has_value()) { const std::uint16_t expected = static_cast(*last_sequence_number_ + 1); @@ -157,6 +159,7 @@ std::optional> H264Depacketizer::depacketize(const RtpPac // lost its tail and can no longer be recovered. if (frame_started_ && packet.header.timestamp != frame_timestamp_) { drop_frame(); + result.frame_dropped = true; } if (!frame_started_) { frame_started_ = true; @@ -216,7 +219,7 @@ std::optional> H264Depacketizer::depacketize(const RtpPac } if (!packet.header.marker) { - return std::nullopt; + return result; } if (fu_active_) { @@ -226,9 +229,11 @@ std::optional> H264Depacketizer::depacketize(const RtpPac fu_nal_.clear(); } - std::optional> result; if (!frame_damaged_ && !access_unit_.empty()) { - result = std::move(access_unit_); + result.access_unit = std::move(access_unit_); + } else { + // The frame that just ended is unusable. + result.frame_dropped = true; } drop_frame(); return result; diff --git a/src/network/rtp_packet.cpp b/src/network/rtp_packet.cpp index 4f71ce8..88b37be 100644 --- a/src/network/rtp_packet.cpp +++ b/src/network/rtp_packet.cpp @@ -133,4 +133,59 @@ std::optional RtpPacket::parse(std::span in) noexcep return packet; } +RtpJitterBuffer::RtpJitterBuffer(std::size_t max_depth, std::chrono::milliseconds max_delay) + : max_depth_(max_depth), max_delay_(max_delay) {} + +std::vector RtpJitterBuffer::push(RtpPacket packet) { + std::vector released; + const std::uint16_t sequence = packet.header.sequence_number; + const auto now = std::chrono::steady_clock::now(); + + std::lock_guard lock(mutex_); + if (!next_expected_.has_value()) { + next_expected_ = sequence; + } + + // Serial-number comparison: a difference >= 32768 means the packet is + // older than what we already delivered (a duplicate or a late straggler). + const std::uint16_t distance = static_cast(sequence - *next_expected_); + if (distance >= 32768) { + return released; // discard the straggler + } + + buffer_[sequence] = {now, std::move(packet)}; + + // Release the consecutive run from the expected sequence. + while (true) { + const auto entry = buffer_.find(*next_expected_); + if (entry == buffer_.end()) { + break; + } + released.push_back(std::move(entry->second.second)); + buffer_.erase(entry); + ++(*next_expected_); + } + + // A missing packet stalls the run: age out the backlog (or bound the + // buffer) and release what is there in order, so genuine loss reaches + // the depacketizer's gap detection rather than blocking forever. + if (!buffer_.empty()) { + const auto head_age = now - buffer_.begin()->second.first; + if (head_age > max_delay_ || buffer_.size() > max_depth_) { + for (auto& entry : buffer_) { + released.push_back(std::move(entry.second.second)); + } + next_expected_ = static_cast(buffer_.rbegin()->first + 1); + buffer_.clear(); + } + } + return released; +} + +void RtpJitterBuffer::clear() { + std::lock_guard lock(mutex_); + buffer_.clear(); + next_expected_.reset(); +} + } // namespace sc diff --git a/src/network/signaling.cpp b/src/network/signaling.cpp index 9e4c5b0..c0f0294 100644 --- a/src/network/signaling.cpp +++ b/src/network/signaling.cpp @@ -115,12 +115,15 @@ std::string serialize_message(const SignalingMessage& message) { json["frame_rate_den"] = offer->frame_rate_den; json["rtp_address"] = offer->rtp_endpoint.address; json["rtp_port"] = offer->rtp_endpoint.port; - } else { - const SessionAnswer& answer = std::get(message); + } else if (const SessionAnswer* answer = std::get_if(&message)) { json["type"] = "answer"; - json["session_id"] = answer.session_id; - json["rtp_address"] = answer.rtp_endpoint.address; - json["rtp_port"] = answer.rtp_endpoint.port; + json["session_id"] = answer->session_id; + json["rtp_address"] = answer->rtp_endpoint.address; + json["rtp_port"] = answer->rtp_endpoint.port; + } else { + const SessionPli& pli = std::get(message); + json["type"] = "pli"; + json["session_id"] = pli.session_id; } return json.dump() + "\n"; } @@ -171,6 +174,11 @@ std::optional parse_message(std::string_view line) { answer.rtp_endpoint = rtp_endpoint; return answer; } + if (type == "pli") { + SessionPli pli; + pli.session_id = session_id; + return pli; + } return std::nullopt; } diff --git a/tests/app/test_loopback.cpp b/tests/app/test_loopback.cpp index d83f890..7bfffd9 100644 --- a/tests/app/test_loopback.cpp +++ b/tests/app/test_loopback.cpp @@ -74,13 +74,13 @@ struct ReceiverSink { int last_height = 0; void on_packet(sc::RtpPacket packet) { - std::optional> access_unit = depacketizer.depacketize(packet); - if (!access_unit.has_value()) { + const sc::DepacketizeResult result = depacketizer.depacketize(packet); + if (!result.access_unit.has_value()) { return; } sc::EncodedFrame encoded; - encoded.data = std::move(*access_unit); + encoded.data = std::move(*result.access_unit); encoded.rtp_timestamp = packet.header.timestamp; auto decoded_result = decoder->decode(encoded); diff --git a/tests/network/test_rtp.cpp b/tests/network/test_rtp.cpp index 4b8c35e..e72731f 100644 --- a/tests/network/test_rtp.cpp +++ b/tests/network/test_rtp.cpp @@ -4,6 +4,7 @@ #include #include +#include #include #include #include @@ -288,13 +289,17 @@ void test_depacketize_roundtrip() { sc::H264Depacketizer depacketizer; std::optional> completed; + bool any_dropped = false; for (const sc::RtpPacket& packet : packets) { - if (auto result = depacketizer.depacketize(packet)) { + const sc::DepacketizeResult result = depacketizer.depacketize(packet); + if (result.access_unit.has_value()) { check(!completed.has_value(), "only one completion"); - completed = std::move(result); + completed = std::move(result.access_unit); } + any_dropped = any_dropped || result.frame_dropped; } check(completed.has_value(), "frame completed"); + check(!any_dropped, "no dropped frames in a clean stream"); check(equal_bytes(*completed, access_unit), "access unit round-trip"); } @@ -305,9 +310,11 @@ void test_depacketizer_drops_gapped_frames() { check(packets.size() == 3, "gap test packet count"); sc::H264Depacketizer depacketizer; - check(!depacketizer.depacketize(packets[0]).has_value(), "first fu chunk accepted"); + check(!depacketizer.depacketize(packets[0]).access_unit.has_value(), "first fu chunk accepted"); // packets[1] is lost in transit; the tail cannot complete the frame. - check(!depacketizer.depacketize(packets[2]).has_value(), "tail after gap dropped"); + const sc::DepacketizeResult tail = depacketizer.depacketize(packets[2]); + check(!tail.access_unit.has_value(), "tail after gap dropped"); + check(tail.frame_dropped, "drop reported after gap"); } void test_depacketizer_separate_frames() { @@ -319,12 +326,79 @@ void test_depacketizer_separate_frames() { const std::vector second = packetizer.packetize(make_frame(f2, 90000)); sc::H264Depacketizer depacketizer; - const std::optional> au1 = depacketizer.depacketize(first[0]); - check(au1.has_value() && equal_bytes(*au1, f1), "first frame"); + const sc::DepacketizeResult au1 = depacketizer.depacketize(first[0]); + check(au1.access_unit.has_value() && equal_bytes(*au1.access_unit, f1), "first frame"); + check(!au1.frame_dropped, "first frame not dropped"); // Same RTP timestamp on purpose: the marker alone separates frames. - const std::optional> au2 = depacketizer.depacketize(second[0]); - check(au2.has_value() && equal_bytes(*au2, f2), "second frame with same timestamp"); + const sc::DepacketizeResult au2 = depacketizer.depacketize(second[0]); + check(au2.access_unit.has_value() && equal_bytes(*au2.access_unit, f2), "second frame with same timestamp"); + check(!au2.frame_dropped, "second frame not dropped"); +} + +void test_jitter_buffer_in_order() { + sc::RtpJitterBuffer jitter; + sc::RtpPacket packet; + packet.header.sequence_number = 100; + + std::vector released = jitter.push(packet); + check(released.size() == 1 && released[0].header.sequence_number == 100, "in-order releases immediately"); + + packet.header.sequence_number = 101; + released = jitter.push(std::move(packet)); + check(released.size() == 1 && released[0].header.sequence_number == 101, "next packet releases too"); +} + +void test_jitter_buffer_reorders() { + sc::RtpJitterBuffer jitter; + + sc::RtpPacket first; + first.header.sequence_number = 1; + std::vector released = jitter.push(std::move(first)); + check(released.size() == 1 && released[0].header.sequence_number == 1, "first packet releases"); + + // Arrives ahead of its predecessor: held, not delivered. + sc::RtpPacket third; + third.header.sequence_number = 3; + released = jitter.push(std::move(third)); + check(released.empty(), "gap holds packets"); + + sc::RtpPacket second; + second.header.sequence_number = 2; + released = jitter.push(std::move(second)); + check(released.size() == 2, "held packets release in order"); + check(released[0].header.sequence_number == 2 && released[1].header.sequence_number == 3, + "released in sequence order"); +} + +void test_jitter_buffer_overflow_and_stragglers() { + // Small depth: a persistent gap overflows the buffer and flushes what + // is there, so genuine loss reaches the depacketizer instead of + // stalling delivery. + sc::RtpJitterBuffer jitter(4, std::chrono::milliseconds{500}); + + sc::RtpPacket packet; + packet.header.sequence_number = 10; + check(jitter.push(std::move(packet)).size() == 1, "first releases"); + + std::vector released; + for (std::uint16_t sequence = 12; sequence < 17; ++sequence) { + sc::RtpPacket missing; + missing.header.sequence_number = sequence; + for (sc::RtpPacket out : jitter.push(std::move(missing))) { + released.push_back(std::move(out)); + } + } + check(released.size() == 5, "overflow flushes the backlog"); + for (std::size_t i = 0; i < released.size(); ++i) { + check(released[i].header.sequence_number == 12 + i, "flushed in order"); + } + + // A straggler older than the delivered sequence is discarded, not + // re-inserted out of order. + sc::RtpPacket straggler; + straggler.header.sequence_number = 11; + check(jitter.push(std::move(straggler)).empty(), "straggler discarded"); } void test_sequence_wrap() { @@ -370,6 +444,9 @@ int main() { test_depacketize_roundtrip(); test_depacketizer_drops_gapped_frames(); test_depacketizer_separate_frames(); + test_jitter_buffer_in_order(); + test_jitter_buffer_reorders(); + test_jitter_buffer_overflow_and_stragglers(); test_sequence_wrap(); test_default_config_randomizes(); test_empty_inputs(); diff --git a/tests/network/test_signaling.cpp b/tests/network/test_signaling.cpp index dbb4757..3c94324 100644 --- a/tests/network/test_signaling.cpp +++ b/tests/network/test_signaling.cpp @@ -93,6 +93,27 @@ int main() { check(answer.session_id == "test-session-0001", "session id echoes"); check(answer.rtp_endpoint.port == advertised_rtp_port, "rtp port in answer"); + // PLI: the receiver side asks for a keyframe mid-session; the client + // must receive it with the session id intact. + std::promise pli_promise; + auto pli_future = pli_promise.get_future(); + std::atomic pli_seen{false}; + client->on_message([&](const sc::SignalingMessage& message) { + if (const sc::SessionPli* pli = std::get_if(&message)) { + if (!pli_seen.exchange(true)) { + pli_promise.set_value(*pli); + } + } + }); + + sc::SessionPli pli; + pli.session_id = "test-session-0001"; + server->send(pli); + + check(pli_future.wait_for(std::chrono::seconds(5)) == std::future_status::ready, "pli received"); + const sc::SessionPli received_pli = pli_future.get(); + check(received_pli.session_id == "test-session-0001", "pli session id echoes"); + client->disconnect(); server->disconnect();