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.
This commit is contained in:
@@ -18,6 +18,7 @@ struct ReceiveCommand {
|
|||||||
int local_rtp_port = 5004;
|
int local_rtp_port = 5004;
|
||||||
int signaling_port = 5005;
|
int signaling_port = 5005;
|
||||||
bool fullscreen = false;
|
bool fullscreen = false;
|
||||||
|
bool software_decode = false; // --swdecode disables the hardware probe
|
||||||
};
|
};
|
||||||
|
|
||||||
struct DiscoverCommand {
|
struct DiscoverCommand {
|
||||||
|
|||||||
@@ -39,6 +39,10 @@ class SenderPipeline {
|
|||||||
bool start();
|
bool start();
|
||||||
void stop();
|
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:
|
private:
|
||||||
class Impl;
|
class Impl;
|
||||||
std::unique_ptr<Impl> impl_;
|
std::unique_ptr<Impl> impl_;
|
||||||
|
|||||||
@@ -22,6 +22,9 @@ struct DecoderConfig {
|
|||||||
int width = 0;
|
int width = 0;
|
||||||
int height = 0;
|
int height = 0;
|
||||||
std::vector<std::byte> extradata; // SPS/PPS for H.264
|
std::vector<std::byte> 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 {
|
class Decoder {
|
||||||
|
|||||||
@@ -42,14 +42,23 @@ class H264Packetizer {
|
|||||||
std::uint16_t next_sequence_number_ = 0;
|
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<std::vector<std::byte>> 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)
|
// 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
|
// 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 {
|
class H264Depacketizer {
|
||||||
public:
|
public:
|
||||||
// Feed one packet. Returns the completed access unit (Annex-B with 3-byte
|
// Feed one packet.
|
||||||
// start codes) when the packet closes a frame, nullopt otherwise.
|
DepacketizeResult depacketize(const RtpPacket& packet);
|
||||||
std::optional<std::vector<std::byte>> depacketize(const RtpPacket& packet);
|
|
||||||
|
|
||||||
private:
|
private:
|
||||||
void drop_frame();
|
void drop_frame();
|
||||||
|
|||||||
@@ -1,8 +1,12 @@
|
|||||||
#pragma once
|
#pragma once
|
||||||
|
|
||||||
|
#include <chrono>
|
||||||
#include <cstdint>
|
#include <cstdint>
|
||||||
|
#include <map>
|
||||||
|
#include <mutex>
|
||||||
#include <optional>
|
#include <optional>
|
||||||
#include <span>
|
#include <span>
|
||||||
|
#include <utility>
|
||||||
#include <vector>
|
#include <vector>
|
||||||
|
|
||||||
namespace sc {
|
namespace sc {
|
||||||
@@ -31,4 +35,29 @@ struct RtpPacket {
|
|||||||
static std::optional<RtpPacket> parse(std::span<const std::byte> in) noexcept;
|
static std::optional<RtpPacket> parse(std::span<const std::byte> 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<RtpPacket> 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<std::uint16_t, std::pair<std::chrono::steady_clock::time_point, RtpPacket>> buffer_;
|
||||||
|
std::optional<std::uint16_t> next_expected_;
|
||||||
|
};
|
||||||
|
|
||||||
} // namespace sc
|
} // namespace sc
|
||||||
|
|||||||
@@ -31,7 +31,13 @@ struct SessionAnswer {
|
|||||||
Endpoint rtp_endpoint;
|
Endpoint rtp_endpoint;
|
||||||
};
|
};
|
||||||
|
|
||||||
using SignalingMessage = std::variant<SessionOffer, SessionAnswer>;
|
// 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<SessionOffer, SessionAnswer, SessionPli>;
|
||||||
|
|
||||||
class SignalingChannel {
|
class SignalingChannel {
|
||||||
public:
|
public:
|
||||||
|
|||||||
+3
-1
@@ -9,7 +9,7 @@ namespace {
|
|||||||
|
|
||||||
void print_usage() {
|
void print_usage() {
|
||||||
std::fputs("usage: screencast --send [--target monitor|window] [--peer HOST[:PORT]] [--bitrate KBPS]\n"
|
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"
|
" screencast --discover [--timeout SECONDS]\n"
|
||||||
"\n"
|
"\n"
|
||||||
"--send without --peer discovers a receiver on the LAN and requires\n"
|
"--send without --peer discovers a receiver on the LAN and requires\n"
|
||||||
@@ -112,6 +112,8 @@ std::optional<Command> parse_cli(int argc, const char* const argv[]) {
|
|||||||
receive.signaling_port = port;
|
receive.signaling_port = port;
|
||||||
} else if (argument == "--fullscreen") {
|
} else if (argument == "--fullscreen") {
|
||||||
receive.fullscreen = true;
|
receive.fullscreen = true;
|
||||||
|
} else if (argument == "--swdecode") {
|
||||||
|
receive.software_decode = true;
|
||||||
} else if (argument == "--timeout") {
|
} else if (argument == "--timeout") {
|
||||||
std::string_view value;
|
std::string_view value;
|
||||||
if (!next_argument(argc, argv, index, value) || !parse_int(value, discover.timeout_seconds) ||
|
if (!next_argument(argc, argv, index, value) || !parse_int(value, discover.timeout_seconds) ||
|
||||||
|
|||||||
+12
-1
@@ -167,7 +167,8 @@ int negotiate_and_stream(sc::SignalingChannel& channel, const sc::Endpoint& sign
|
|||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
const sc::SessionAnswer answer = answer_future.get();
|
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) {
|
if (answer.session_id != offer.session_id) {
|
||||||
std::cerr << "screencast: session mismatch in the receiver's answer\n";
|
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)};
|
sc::SenderPipeline pipeline{std::move(config)};
|
||||||
if (!pipeline.start()) {
|
if (!pipeline.start()) {
|
||||||
|
channel.disconnect();
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
|
channel.on_message([&](const sc::SignalingMessage& message) {
|
||||||
|
const sc::SessionPli* pli = std::get_if<sc::SessionPli>(&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()) {
|
while (!g_interrupted.load()) {
|
||||||
std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
||||||
}
|
}
|
||||||
pipeline.stop();
|
pipeline.stop();
|
||||||
|
channel.disconnect();
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -289,6 +299,7 @@ int run_receiver(const sc::ReceiveCommand& command) {
|
|||||||
sc::ReceiverPipelineConfig config;
|
sc::ReceiverPipelineConfig config;
|
||||||
config.local_rtp_endpoint = sc::Endpoint{"0.0.0.0", static_cast<std::uint16_t>(command.local_rtp_port)};
|
config.local_rtp_endpoint = sc::Endpoint{"0.0.0.0", static_cast<std::uint16_t>(command.local_rtp_port)};
|
||||||
config.signaling_port = static_cast<std::uint16_t>(command.signaling_port);
|
config.signaling_port = static_cast<std::uint16_t>(command.signaling_port);
|
||||||
|
config.decoder.hardware_accel = !command.software_decode;
|
||||||
// Fullscreen is automatic under KMSDRM (headless); the flag forces it
|
// Fullscreen is automatic under KMSDRM (headless); the flag forces it
|
||||||
// on desktop sessions.
|
// on desktop sessions.
|
||||||
config.renderer.fullscreen = command.fullscreen;
|
config.renderer.fullscreen = command.fullscreen;
|
||||||
|
|||||||
+50
-3
@@ -64,9 +64,17 @@ class SenderPipeline::Impl {
|
|||||||
transport_->stop();
|
transport_->stop();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void request_keyframe() {
|
||||||
|
keyframe_requested_.store(true);
|
||||||
|
}
|
||||||
|
|
||||||
private:
|
private:
|
||||||
void run(std::stop_token stop_token) {
|
void run(std::stop_token stop_token) {
|
||||||
while (!stop_token.stop_requested()) {
|
while (!stop_token.stop_requested()) {
|
||||||
|
if (keyframe_requested_.exchange(false) && encoder_ != nullptr) {
|
||||||
|
encoder_->request_keyframe();
|
||||||
|
}
|
||||||
|
|
||||||
const std::optional<CapturedFrame> frame = capture_->next_frame();
|
const std::optional<CapturedFrame> frame = capture_->next_frame();
|
||||||
if (!frame.has_value()) {
|
if (!frame.has_value()) {
|
||||||
break;
|
break;
|
||||||
@@ -109,8 +117,13 @@ class SenderPipeline::Impl {
|
|||||||
H264Packetizer packetizer_;
|
H264Packetizer packetizer_;
|
||||||
std::unique_ptr<RtpTransport> transport_ = RtpTransportFactory::create();
|
std::unique_ptr<RtpTransport> transport_ = RtpTransportFactory::create();
|
||||||
std::jthread run_thread_;
|
std::jthread run_thread_;
|
||||||
|
std::atomic<bool> keyframe_requested_{false};
|
||||||
};
|
};
|
||||||
|
|
||||||
|
void SenderPipeline::request_keyframe() {
|
||||||
|
impl_->request_keyframe();
|
||||||
|
}
|
||||||
|
|
||||||
SenderPipeline::SenderPipeline(SenderPipelineConfig config) : impl_(std::make_unique<Impl>(std::move(config))) {}
|
SenderPipeline::SenderPipeline(SenderPipelineConfig config) : impl_(std::make_unique<Impl>(std::move(config))) {}
|
||||||
|
|
||||||
SenderPipeline::~SenderPipeline() = default;
|
SenderPipeline::~SenderPipeline() = default;
|
||||||
@@ -238,6 +251,7 @@ class ReceiverPipeline::Impl {
|
|||||||
signaling_ = nullptr;
|
signaling_ = nullptr;
|
||||||
discovery_ = nullptr;
|
discovery_ = nullptr;
|
||||||
transport_->stop();
|
transport_->stop();
|
||||||
|
jitter_.clear();
|
||||||
renderer_ = nullptr;
|
renderer_ = nullptr;
|
||||||
decoder_ = nullptr;
|
decoder_ = nullptr;
|
||||||
{
|
{
|
||||||
@@ -248,13 +262,24 @@ class ReceiverPipeline::Impl {
|
|||||||
|
|
||||||
private:
|
private:
|
||||||
void on_packet(RtpPacket packet) {
|
void on_packet(RtpPacket packet) {
|
||||||
std::optional<std::vector<std::byte>> access_unit = depacketizer_.depacketize(packet);
|
// Absorb reordering (Wi-Fi) before the in-order depacketizer, so a
|
||||||
if (!access_unit.has_value()) {
|
// 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;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
EncodedFrame encoded;
|
EncodedFrame encoded;
|
||||||
encoded.data = std::move(*access_unit);
|
encoded.data = std::move(*result.access_unit);
|
||||||
encoded.rtp_timestamp = packet.header.timestamp;
|
encoded.rtp_timestamp = packet.header.timestamp;
|
||||||
encoded.is_keyframe = false;
|
encoded.is_keyframe = false;
|
||||||
|
|
||||||
@@ -282,6 +307,7 @@ class ReceiverPipeline::Impl {
|
|||||||
if (offer == nullptr) {
|
if (offer == nullptr) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
session_id_ = offer->session_id;
|
||||||
SessionAnswer answer;
|
SessionAnswer answer;
|
||||||
answer.session_id = offer->session_id;
|
answer.session_id = offer->session_id;
|
||||||
// Empty address: the sender targets the address of its signaling
|
// Empty address: the sender targets the address of its signaling
|
||||||
@@ -290,6 +316,23 @@ class ReceiverPipeline::Impl {
|
|||||||
signaling_->send(answer);
|
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) {
|
void render_loop(std::stop_token stop_token) {
|
||||||
while (!stop_token.stop_requested()) {
|
while (!stop_token.stop_requested()) {
|
||||||
if (!renderer_->poll_events()) {
|
if (!renderer_->poll_events()) {
|
||||||
@@ -315,17 +358,21 @@ class ReceiverPipeline::Impl {
|
|||||||
}
|
}
|
||||||
|
|
||||||
static constexpr std::size_t kMaxQueuedFrames = 3;
|
static constexpr std::size_t kMaxQueuedFrames = 3;
|
||||||
|
static constexpr std::chrono::milliseconds kPliMinInterval{500};
|
||||||
|
|
||||||
ReceiverPipelineConfig config_;
|
ReceiverPipelineConfig config_;
|
||||||
std::unique_ptr<Renderer> renderer_;
|
std::unique_ptr<Renderer> renderer_;
|
||||||
std::unique_ptr<Decoder> decoder_;
|
std::unique_ptr<Decoder> decoder_;
|
||||||
H264Depacketizer depacketizer_;
|
H264Depacketizer depacketizer_;
|
||||||
|
RtpJitterBuffer jitter_;
|
||||||
std::unique_ptr<RtpTransport> transport_ = RtpTransportFactory::create();
|
std::unique_ptr<RtpTransport> transport_ = RtpTransportFactory::create();
|
||||||
std::unique_ptr<SignalingChannel> signaling_;
|
std::unique_ptr<SignalingChannel> signaling_;
|
||||||
std::unique_ptr<DiscoveryService> discovery_;
|
std::unique_ptr<DiscoveryService> discovery_;
|
||||||
std::jthread render_thread_;
|
std::jthread render_thread_;
|
||||||
std::mutex queue_mutex_;
|
std::mutex queue_mutex_;
|
||||||
std::deque<DecodedFrame> queue_;
|
std::deque<DecodedFrame> queue_;
|
||||||
|
std::string session_id_;
|
||||||
|
std::chrono::steady_clock::time_point last_pli_time_{};
|
||||||
bool stream_started_ = false;
|
bool stream_started_ = false;
|
||||||
bool present_error_logged_ = false;
|
bool present_error_logged_ = false;
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -5,6 +5,8 @@
|
|||||||
#include <array>
|
#include <array>
|
||||||
#include <cstdint>
|
#include <cstdint>
|
||||||
#include <cstring>
|
#include <cstring>
|
||||||
|
#include <format>
|
||||||
|
#include <iostream>
|
||||||
#include <limits>
|
#include <limits>
|
||||||
#include <optional>
|
#include <optional>
|
||||||
#include <string>
|
#include <string>
|
||||||
@@ -199,19 +201,18 @@ class FfmpegDecoder final : public Decoder {
|
|||||||
mutable AVPixelFormat scaler_input_format_ = AV_PIX_FMT_NONE;
|
mutable AVPixelFormat scaler_input_format_ = AV_PIX_FMT_NONE;
|
||||||
};
|
};
|
||||||
|
|
||||||
CodecResult<std::unique_ptr<Decoder>> DecoderFactory::create(const DecoderConfig& config) {
|
namespace {
|
||||||
if (config.codec_name != "h264") {
|
|
||||||
return CodecError{"only h264 is supported in phase 2"};
|
|
||||||
}
|
|
||||||
|
|
||||||
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<std::string> open_decoder(const DecoderConfig& config, const AVCodec* codec, AvCodecContextPtr& ctx) {
|
||||||
if (codec == nullptr) {
|
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) {
|
if (ctx == nullptr) {
|
||||||
return CodecError{"failed to allocate decoder context"};
|
return std::string{"failed to allocate decoder context"};
|
||||||
}
|
}
|
||||||
|
|
||||||
ctx->codec_type = AVMEDIA_TYPE_VIDEO;
|
ctx->codec_type = AVMEDIA_TYPE_VIDEO;
|
||||||
@@ -225,24 +226,58 @@ CodecResult<std::unique_ptr<Decoder>> DecoderFactory::create(const DecoderConfig
|
|||||||
|
|
||||||
if (!config.extradata.empty()) {
|
if (!config.extradata.empty()) {
|
||||||
if (config.extradata.size() > static_cast<std::size_t>(std::numeric_limits<int>::max())) {
|
if (config.extradata.size() > static_cast<std::size_t>(std::numeric_limits<int>::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
|
// FFmpeg bitstream parsers may read past the end of extradata, so the
|
||||||
// buffer must include the padding they require.
|
// buffer must include the padding they require.
|
||||||
ctx->extradata = static_cast<uint8_t*>(av_malloc(config.extradata.size() + AV_INPUT_BUFFER_PADDING_SIZE));
|
ctx->extradata = static_cast<uint8_t*>(av_malloc(config.extradata.size() + AV_INPUT_BUFFER_PADDING_SIZE));
|
||||||
if (ctx->extradata == nullptr) {
|
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::memcpy(ctx->extradata, config.extradata.data(), config.extradata.size());
|
||||||
std::memset(ctx->extradata + config.extradata.size(), 0, AV_INPUT_BUFFER_PADDING_SIZE);
|
std::memset(ctx->extradata + config.extradata.size(), 0, AV_INPUT_BUFFER_PADDING_SIZE);
|
||||||
ctx->extradata_size = static_cast<int>(config.extradata.size());
|
ctx->extradata_size = static_cast<int>(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) {
|
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<std::unique_ptr<Decoder>> 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<FfmpegDecoder>(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<FfmpegDecoder>(std::move(ctx), config);
|
return std::make_unique<FfmpegDecoder>(std::move(ctx), config);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -139,7 +139,9 @@ void H264Packetizer::append_fu_a_packets(std::span<const std::byte> nal,
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
std::optional<std::vector<std::byte>> H264Depacketizer::depacketize(const RtpPacket& packet) {
|
DepacketizeResult H264Depacketizer::depacketize(const RtpPacket& packet) {
|
||||||
|
DepacketizeResult result;
|
||||||
|
|
||||||
// Track sequence continuity: a gap means packets were lost.
|
// Track sequence continuity: a gap means packets were lost.
|
||||||
if (last_sequence_number_.has_value()) {
|
if (last_sequence_number_.has_value()) {
|
||||||
const std::uint16_t expected = static_cast<std::uint16_t>(*last_sequence_number_ + 1);
|
const std::uint16_t expected = static_cast<std::uint16_t>(*last_sequence_number_ + 1);
|
||||||
@@ -157,6 +159,7 @@ std::optional<std::vector<std::byte>> H264Depacketizer::depacketize(const RtpPac
|
|||||||
// lost its tail and can no longer be recovered.
|
// lost its tail and can no longer be recovered.
|
||||||
if (frame_started_ && packet.header.timestamp != frame_timestamp_) {
|
if (frame_started_ && packet.header.timestamp != frame_timestamp_) {
|
||||||
drop_frame();
|
drop_frame();
|
||||||
|
result.frame_dropped = true;
|
||||||
}
|
}
|
||||||
if (!frame_started_) {
|
if (!frame_started_) {
|
||||||
frame_started_ = true;
|
frame_started_ = true;
|
||||||
@@ -216,7 +219,7 @@ std::optional<std::vector<std::byte>> H264Depacketizer::depacketize(const RtpPac
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (!packet.header.marker) {
|
if (!packet.header.marker) {
|
||||||
return std::nullopt;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
if (fu_active_) {
|
if (fu_active_) {
|
||||||
@@ -226,9 +229,11 @@ std::optional<std::vector<std::byte>> H264Depacketizer::depacketize(const RtpPac
|
|||||||
fu_nal_.clear();
|
fu_nal_.clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
std::optional<std::vector<std::byte>> result;
|
|
||||||
if (!frame_damaged_ && !access_unit_.empty()) {
|
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();
|
drop_frame();
|
||||||
return result;
|
return result;
|
||||||
|
|||||||
@@ -133,4 +133,59 @@ std::optional<RtpPacket> RtpPacket::parse(std::span<const std::byte> in) noexcep
|
|||||||
return packet;
|
return packet;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
RtpJitterBuffer::RtpJitterBuffer(std::size_t max_depth, std::chrono::milliseconds max_delay)
|
||||||
|
: max_depth_(max_depth), max_delay_(max_delay) {}
|
||||||
|
|
||||||
|
std::vector<RtpPacket> RtpJitterBuffer::push(RtpPacket packet) {
|
||||||
|
std::vector<RtpPacket> 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<std::uint16_t>(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<std::uint16_t>(buffer_.rbegin()->first + 1);
|
||||||
|
buffer_.clear();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return released;
|
||||||
|
}
|
||||||
|
|
||||||
|
void RtpJitterBuffer::clear() {
|
||||||
|
std::lock_guard lock(mutex_);
|
||||||
|
buffer_.clear();
|
||||||
|
next_expected_.reset();
|
||||||
|
}
|
||||||
|
|
||||||
} // namespace sc
|
} // namespace sc
|
||||||
|
|||||||
@@ -115,12 +115,15 @@ std::string serialize_message(const SignalingMessage& message) {
|
|||||||
json["frame_rate_den"] = offer->frame_rate_den;
|
json["frame_rate_den"] = offer->frame_rate_den;
|
||||||
json["rtp_address"] = offer->rtp_endpoint.address;
|
json["rtp_address"] = offer->rtp_endpoint.address;
|
||||||
json["rtp_port"] = offer->rtp_endpoint.port;
|
json["rtp_port"] = offer->rtp_endpoint.port;
|
||||||
} else {
|
} else if (const SessionAnswer* answer = std::get_if<SessionAnswer>(&message)) {
|
||||||
const SessionAnswer& answer = std::get<SessionAnswer>(message);
|
|
||||||
json["type"] = "answer";
|
json["type"] = "answer";
|
||||||
json["session_id"] = answer.session_id;
|
json["session_id"] = answer->session_id;
|
||||||
json["rtp_address"] = answer.rtp_endpoint.address;
|
json["rtp_address"] = answer->rtp_endpoint.address;
|
||||||
json["rtp_port"] = answer.rtp_endpoint.port;
|
json["rtp_port"] = answer->rtp_endpoint.port;
|
||||||
|
} else {
|
||||||
|
const SessionPli& pli = std::get<SessionPli>(message);
|
||||||
|
json["type"] = "pli";
|
||||||
|
json["session_id"] = pli.session_id;
|
||||||
}
|
}
|
||||||
return json.dump() + "\n";
|
return json.dump() + "\n";
|
||||||
}
|
}
|
||||||
@@ -171,6 +174,11 @@ std::optional<SignalingMessage> parse_message(std::string_view line) {
|
|||||||
answer.rtp_endpoint = rtp_endpoint;
|
answer.rtp_endpoint = rtp_endpoint;
|
||||||
return answer;
|
return answer;
|
||||||
}
|
}
|
||||||
|
if (type == "pli") {
|
||||||
|
SessionPli pli;
|
||||||
|
pli.session_id = session_id;
|
||||||
|
return pli;
|
||||||
|
}
|
||||||
return std::nullopt;
|
return std::nullopt;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -74,13 +74,13 @@ struct ReceiverSink {
|
|||||||
int last_height = 0;
|
int last_height = 0;
|
||||||
|
|
||||||
void on_packet(sc::RtpPacket packet) {
|
void on_packet(sc::RtpPacket packet) {
|
||||||
std::optional<std::vector<std::byte>> access_unit = depacketizer.depacketize(packet);
|
const sc::DepacketizeResult result = depacketizer.depacketize(packet);
|
||||||
if (!access_unit.has_value()) {
|
if (!result.access_unit.has_value()) {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
sc::EncodedFrame encoded;
|
sc::EncodedFrame encoded;
|
||||||
encoded.data = std::move(*access_unit);
|
encoded.data = std::move(*result.access_unit);
|
||||||
encoded.rtp_timestamp = packet.header.timestamp;
|
encoded.rtp_timestamp = packet.header.timestamp;
|
||||||
|
|
||||||
auto decoded_result = decoder->decode(encoded);
|
auto decoded_result = decoder->decode(encoded);
|
||||||
|
|||||||
@@ -4,6 +4,7 @@
|
|||||||
|
|
||||||
#include <algorithm>
|
#include <algorithm>
|
||||||
#include <array>
|
#include <array>
|
||||||
|
#include <chrono>
|
||||||
#include <cstdint>
|
#include <cstdint>
|
||||||
#include <cstdio>
|
#include <cstdio>
|
||||||
#include <cstdlib>
|
#include <cstdlib>
|
||||||
@@ -288,13 +289,17 @@ void test_depacketize_roundtrip() {
|
|||||||
|
|
||||||
sc::H264Depacketizer depacketizer;
|
sc::H264Depacketizer depacketizer;
|
||||||
std::optional<std::vector<std::byte>> completed;
|
std::optional<std::vector<std::byte>> completed;
|
||||||
|
bool any_dropped = false;
|
||||||
for (const sc::RtpPacket& packet : packets) {
|
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");
|
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(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");
|
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");
|
check(packets.size() == 3, "gap test packet count");
|
||||||
|
|
||||||
sc::H264Depacketizer depacketizer;
|
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.
|
// 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() {
|
void test_depacketizer_separate_frames() {
|
||||||
@@ -319,12 +326,79 @@ void test_depacketizer_separate_frames() {
|
|||||||
const std::vector<sc::RtpPacket> second = packetizer.packetize(make_frame(f2, 90000));
|
const std::vector<sc::RtpPacket> second = packetizer.packetize(make_frame(f2, 90000));
|
||||||
|
|
||||||
sc::H264Depacketizer depacketizer;
|
sc::H264Depacketizer depacketizer;
|
||||||
const std::optional<std::vector<std::byte>> au1 = depacketizer.depacketize(first[0]);
|
const sc::DepacketizeResult au1 = depacketizer.depacketize(first[0]);
|
||||||
check(au1.has_value() && equal_bytes(*au1, f1), "first frame");
|
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.
|
// Same RTP timestamp on purpose: the marker alone separates frames.
|
||||||
const std::optional<std::vector<std::byte>> au2 = depacketizer.depacketize(second[0]);
|
const sc::DepacketizeResult au2 = depacketizer.depacketize(second[0]);
|
||||||
check(au2.has_value() && equal_bytes(*au2, f2), "second frame with same timestamp");
|
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<sc::RtpPacket> 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<sc::RtpPacket> 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<sc::RtpPacket> 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() {
|
void test_sequence_wrap() {
|
||||||
@@ -370,6 +444,9 @@ int main() {
|
|||||||
test_depacketize_roundtrip();
|
test_depacketize_roundtrip();
|
||||||
test_depacketizer_drops_gapped_frames();
|
test_depacketizer_drops_gapped_frames();
|
||||||
test_depacketizer_separate_frames();
|
test_depacketizer_separate_frames();
|
||||||
|
test_jitter_buffer_in_order();
|
||||||
|
test_jitter_buffer_reorders();
|
||||||
|
test_jitter_buffer_overflow_and_stragglers();
|
||||||
test_sequence_wrap();
|
test_sequence_wrap();
|
||||||
test_default_config_randomizes();
|
test_default_config_randomizes();
|
||||||
test_empty_inputs();
|
test_empty_inputs();
|
||||||
|
|||||||
@@ -93,6 +93,27 @@ int main() {
|
|||||||
check(answer.session_id == "test-session-0001", "session id echoes");
|
check(answer.session_id == "test-session-0001", "session id echoes");
|
||||||
check(answer.rtp_endpoint.port == advertised_rtp_port, "rtp port in answer");
|
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<sc::SessionPli> pli_promise;
|
||||||
|
auto pli_future = pli_promise.get_future();
|
||||||
|
std::atomic<bool> pli_seen{false};
|
||||||
|
client->on_message([&](const sc::SignalingMessage& message) {
|
||||||
|
if (const sc::SessionPli* pli = std::get_if<sc::SessionPli>(&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();
|
client->disconnect();
|
||||||
server->disconnect();
|
server->disconnect();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user