Compare commits

...
Author SHA1 Message Date
shadowdaoandClaude Opus 5.5 7723082795 feat(session): subscribe only to the slot's own camera and mic
Build / Linux (ubuntu-24.04) (push) Successful in 1m8s
Build / macOS (macos-latest) (push) Successful in 1m9s
Build / macOS (macos-latest) (pull_request) Successful in 44s
Build / Windows (windows-latest) (push) Successful in 7m0s
Build / Windows (windows-latest) (pull_request) Successful in 4m14s
Build / Linux (ubuntu-24.04) (pull_request) Failing after 42m30s
Every plugin source connected with auto_subscribe on, so each one pulled
every camera in the room and discarded all but one. In the 2026-10-04
old-gods-of-appalachia show (8 sources, 7 cameras) each source carried
4-7 Mbps and all 8 hit congestion in the same instant: they share the
director's single downlink, so the SFU starved the cameras on screen.

The session now connects with auto_subscribe off and calls
setSubscribed(true) on exactly the wanted publications: at connect, on
onTrackPublished (re-sweeping the participant, since SDK 1.10.1 delivers
that event's publication null), and after a reconnect. Any subscription
it did not ask for is reported and handed back. Audio-only sources no
longer subscribe-then-disable video; they never subscribe to it.

The integration test gains a bystander camera that must never be
subscribed (fails with 3 unwanted subscriptions before this change) and
a full publisher leave/rejoin as a new participant session.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-04 18:00:28 -07:00
shadowdaoandClaude Opus 5 22d50ab7be test(session): derive watchdog instants from t0, not an accumulated local
Build / macOS (macos-latest) (push) Successful in 53s
Build / Linux (ubuntu-24.04) (push) Successful in 1m7s
Release / macOS (macos-latest) (push) Successful in 57s
Release / Linux (ubuntu-24.04) (push) Successful in 1m14s
Build / Windows (windows-latest) (push) Successful in 4m6s
Release / Windows (windows-latest) (push) Successful in 3m55s
Release / Create Gitea Release (push) Successful in 20s
Windows CI failed on the stall-recovery watchdog for three attempts. The
first two diagnoses both blamed the backoff arithmetic; the second
produced a byte-identical failure, which was the clue that neither had
found the cause.

Instrumenting the test on the Windows runner settled it with numbers.
Capturing the time point on the callee side of a noinline wrapper showed
the six `w.poll(now)` calls in the capped-backoff loop received:

  Windows      32000, 32000, 32000, 32000, 32000, 32000  (ms after t0)
  Linux/macOS  62000, 92000, 122000, 152000, 182000, 212000

while a checksum of the caller's own arguments in the same loop summed to
822000 -- exactly the correct series. The caller's value was right; the
value that crossed the call boundary was not. MSVC 19.44 x64 Release
hoists the argument of the second poll() out of the fixed-stride loop
`poll(now + 29999ms); now += 30000ms; poll(now);`, so every iteration
passed the pre-loop `now`.

StallWatchdog is correct on all three platforms and is not changed here.
Production never had this exposure: watchdogLoop() calls poll() with a
fresh steady_clock::now() per tick, never a loop-carried local advanced
by a constant.

Both watchdog tests now derive every instant absolutely as
`t0 + milliseconds(at_ms)` from an integer cursor -- the shape verified
to compile correctly on that runner. No assertion is weakened: every gap
is still checked one millisecond either side of its boundary.

Also corrects the comment in session_types.cpp that blamed a Windows
release build for mis-capping the ceiling. It never did.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-21 11:09:57 -07:00
shadowdaoandClaude Opus 5 352a843d93 fix(session): enforce the stall-recovery ceiling where the wait is used
Build / macOS (macos-latest) (push) Successful in 53s
Build / Linux (ubuntu-24.04) (push) Successful in 1m24s
Release / macOS (macos-latest) (push) Successful in 55s
Release / Linux (ubuntu-24.04) (push) Successful in 1m18s
Build / Windows (windows-latest) (push) Failing after 3m41s
Release / Windows (windows-latest) (push) Failing after 3m39s
Release / Create Gitea Release (push) Skipped
The Windows release build failed testStallWatchdogBacksOffRatherThanLooping
(11/124 checks) while Linux and macOS passed, so v0.1.1 never published.

Reconstructing the failure from the log rather than guessing: the reported
FAIL lines (line 329 first, then 327/329 alternating for the rest of the
capped-backoff loop, 11 of the loop's 12 checks) are produced by exactly one
behaviour, and the deliberately-broken build in this commit's verification
reproduced that log byte-for-byte on Linux -- the backoff ceiling engaged one
attempt LATE. Windows waited 32000ms once (the uncapped doubling of 16000ms)
before settling at the 30000ms ceiling. Every other candidate produces a
different count and a different order: an exact-equality boundary bug gives 6
failures, and a ceiling that never engages at all gives 8, neither matching.

That rules out the obvious suspect, a lossy duration conversion. There isn't
one, and there cannot be: `time_point<Clock, D1> + duration<D2>` yields
`time_point<Clock, common_type_t<D1, D2>>`, and converting that back to
`steady_clock::time_point` to store it in next_attempt_allowed_ only compiles
when the conversion is exact. If MSVC's steady_clock could not represent a
whole millisecond exactly, this file would not build there. All of the
watchdog's time arithmetic is exact integer arithmetic on every platform, and
the exact-equality comparison at the deadline is sound -- the Windows log
itself shows later polls firing at exactly their deadline.

What is left is `std::min(backoff_ * 2, max_backoff_)`: the one expression in
poll() that was not plain value arithmetic on a single type, returning a
*reference* bound, in the growing case, to a materialized temporary. So:

- The ceiling is now clamped where the wait is USED, not only where the
  backoff is grown. max_backoff_ is a promise about the longest gap between
  two recovery attempts, so it is enforced on the gap itself and holds for
  whatever backoff_ contains. Verified: with the growth step deliberately
  mis-capping exactly the way Windows did, the whole suite still passes --
  the fix does not depend on having correctly identified MSVC's mechanism.
- The doubling is an explicit compare-and-clamp instead of std::min, so no
  reference to a temporary is involved and the product is only computed when
  it cannot exceed the ceiling.

Both changes are provably no-ops on Linux and macOS, where backoff_ never
exceeded the ceiling in the first place.

Also pins the behaviour with a new regression test using a ceiling that is
NOT a power-of-two multiple of the timeout (1000 -> 2000 -> 4000 -> 5000),
which fails on the step the ceiling first binds rather than six 30-second
iterations later. 146 checks in test_session now, was 124; all 6 CTest suites
pass locally.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-21 10:18:15 -07:00
shadowdaoandClaude Opus 5 564461a16d chore(release): 0.1.1
Build / macOS (macos-latest) (push) Successful in 53s
Release / macOS (macos-latest) (push) Successful in 52s
Build / Linux (ubuntu-24.04) (push) Successful in 1m9s
Release / Linux (ubuntu-24.04) (push) Successful in 1m12s
Build / Windows (windows-latest) (push) Failing after 3m37s
Release / Windows (windows-latest) (push) Failing after 3m37s
Release / Create Gitea Release (push) Skipped
Stall-recovery watchdog for video subscriptions (#7). Camera sources could
drop out in OBS and never recover while the same players stayed healthy in
browser talkback; the pinned client-sdk-cpp exposes no keyframe-request
API, so a decoder that lost a frame had no way to resync for the rest of
the show.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-21 09:43:39 -07:00
jknapp 969a500dfd Merge pull request #7 from fix/stall-recovery
Build / macOS (macos-latest) (push) Successful in 52s
Build / Linux (ubuntu-24.04) (push) Successful in 1m16s
Build / Windows (windows-latest) (push) Failing after 3m45s
fix(session): recover stalled video subscriptions with a keyframe-forcing watchdog
2026-09-21 16:43:21 +00:00
8 changed files with 353 additions and 81 deletions
+1 -1
View File
@@ -1,7 +1,7 @@
cmake_minimum_required(VERSION 3.19) cmake_minimum_required(VERSION 3.19)
project(obs-streamer-tools-plugin project(obs-streamer-tools-plugin
VERSION 0.1.0 VERSION 0.1.1
DESCRIPTION "OBS Studio source plugin for streamer-tools camera feeds" DESCRIPTION "OBS Studio source plugin for streamer-tools camera feeds"
LANGUAGES C CXX LANGUAGES C CXX
) )
+3 -5
View File
@@ -44,11 +44,9 @@ struct SessionConfig {
bool subscribe_audio = true; bool subscribe_audio = true;
/// False for an audio-only source (the soundboard, say): the wanted /// False for an audio-only source (the soundboard, say): the slot's
/// video track is never attached (no AttachVideo command posted), and /// video track is never subscribed to (see shouldSubscribe), so the SFU
/// its publication is explicitly disabled server-side (RemoteTrack- /// never sends it at all -- not just "decoded and discarded here".
/// Publication::setEnabled(false)) so the SFU stops sending it at all --
/// not just "decoded and discarded here", genuinely not delivered.
bool subscribe_video = true; bool subscribe_video = true;
/// How long connect() waits for the room to come up before giving up. /// How long connect() waits for the room to come up before giving up.
+8
View File
@@ -70,6 +70,14 @@ bool isWantedVideoTrack(const std::string &wanted_identity, const std::string &t
bool isWantedAudioTrack(const std::string &wanted_identity, const std::string &track_identity, MediaKind kind, bool isWantedAudioTrack(const std::string &wanted_identity, const std::string &track_identity, MediaKind kind,
MediaSource source); MediaSource source);
/// Should this session subscribe to this publication at all? The session
/// connects with auto_subscribe off and asks for exactly these, so the SFU
/// never sends this source anything it would only throw away. Without it,
/// every plugin source pulled every camera in the room: eight sources in one
/// OBS meant the director's downlink carried each camera eight times.
bool shouldSubscribe(const std::string &wanted_identity, const std::string &track_identity, MediaKind kind,
MediaSource source, bool subscribe_video, bool subscribe_audio);
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Session state // Session state
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
+98 -47
View File
@@ -269,39 +269,21 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
// Handles the wanted video track once matched, shared by onTrackSubscribed // Handles the wanted video track once matched, shared by onTrackSubscribed
// (a fresh subscription) and attachExistingTracks (one already up when // (a fresh subscription) and attachExistingTracks (one already up when
// this session started watching). Two responsibilities that only make // this session started watching). An audio-only source (the soundboard)
// sense together, both keyed off the SAME publication: // never gets here for video: shouldSubscribe() never asks for it.
// //
// - subscribe_video: an audio-only source (the soundboard) never wants // Fixed video quality: LiveKit's default subscriber behaviour lets the
// this video at all. Rather than attach it and let the OBS adapter // SFU switch simulcast layers per its own adaptive/bandwidth logic, which
// discard every decoded frame, disable the publication itself // for a source with no rendered-size hint (this is a native C++
// (RemoteTrackPublication::setEnabled(false)) so the SFU stops // subscriber, not a sized <video> element) means the received resolution
// sending it -- real bandwidth saved, not just wasted decode. // can hop between layers -- observed live as OBS source geometry visibly
// - Fixed video quality: LiveKit's default subscriber behaviour lets // changing size mid-show. Pinning to HIGH asks the SFU to always send the
// the SFU switch simulcast layers per its own adaptive/bandwidth // top layer, which is what a fixed OBS source needs regardless of
// logic, which for a source with no rendered-size hint (this is a // bandwidth (the plugin has no picture-in-picture tier to fall back to
// native C++ subscriber, not a sized <video> element) means the // the way a browser grid view would).
// received resolution can hop between layers -- observed live as OBS
// source geometry visibly changing size mid-show. Pinning to HIGH
// asks the SFU to always send the top layer, which is what a fixed
// OBS source needs regardless of bandwidth (the plugin has no
// picture-in-picture tier to fall back to the way a browser grid
// view would).
void handleWantedVideoTrack(const std::shared_ptr<livekit::Track> &track, void handleWantedVideoTrack(const std::shared_ptr<livekit::Track> &track,
const std::shared_ptr<livekit::RemoteTrackPublication> &publication) const std::shared_ptr<livekit::RemoteTrackPublication> &publication)
{ {
if (!config.subscribe_video) {
if (publication) {
try {
publication->setEnabled(false);
} catch (const std::exception &) {
// Best-effort: worst case this track keeps being
// delivered and decoded, wasting bandwidth -- it is
// still never attached to OBS below.
}
}
return;
}
if (publication) { if (publication) {
try { try {
publication->setVideoQuality(livekit::VideoQuality::HIGH); publication->setVideoQuality(livekit::VideoQuality::HIGH);
@@ -359,12 +341,46 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
const MediaSource source = const MediaSource source =
event.publication ? toMediaSource(event.publication->source()) : MediaSource::Unknown; event.publication ? toMediaSource(event.publication->source()) : MediaSource::Unknown;
if (!shouldSubscribe(config.participant_identity, identity, kind, source, config.subscribe_video,
config.subscribe_audio)) {
// Never asked for (see requestSubscription). Should not happen with
// auto_subscribe off; if the SDK or server ever subscribes us
// anyway, say so and hand it back rather than silently paying for
// another camera's bandwidth.
logDiagnostic(DiagnosticLevel::Warning, "dropping subscription to unwanted track " +
event.track->sid() + " from " + identity);
if (event.publication) {
try {
event.publication->setSubscribed(false);
} catch (const std::exception &) {
// Best-effort: worst case it keeps arriving and is
// discarded below, which is the old behaviour.
}
}
return;
}
if (isWantedVideoTrack(config.participant_identity, identity, kind, source)) if (isWantedVideoTrack(config.participant_identity, identity, kind, source))
handleWantedVideoTrack(event.track, event.publication); handleWantedVideoTrack(event.track, event.publication);
else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source)) else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source))
post(CommandType::AttachAudio, event.track); post(CommandType::AttachAudio, event.track);
} }
// auto_subscribe is off (see connect()), so a newly published track --
// including every track of a player who just rejoined -- arrives
// unsubscribed, and this is where the session asks for the ones it shows.
//
// event.publication is not usable: SDK 1.10.1 delivers it null (observed
// against livekit-server 1.13.6), so re-sweep the participant's own
// publication map instead. Idempotent, and it only ever subscribes --
// re-attaching the camera already on screen would blank it.
void onTrackPublished(livekit::Room &, const livekit::TrackPublishedEvent &event) override
{
if (!event.participant || event.participant->identity() != config.participant_identity)
return;
requestWantedSubscriptions(*event.participant);
}
void onTrackUnsubscribed(livekit::Room &, const livekit::TrackUnsubscribedEvent &event) override void onTrackUnsubscribed(livekit::Room &, const livekit::TrackUnsubscribedEvent &event) override
{ {
if (!event.participant || !event.track) if (!event.participant || !event.track)
@@ -406,6 +422,11 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
void onReconnected(livekit::Room &, const livekit::ReconnectedEvent &) override void onReconnected(livekit::Room &, const livekit::ReconnectedEvent &) override
{ {
mutateState([](SessionStateMachine &m) { m.onReconnected(); }); mutateState([](SessionStateMachine &m) { m.onReconnected(); });
// Belt and braces: re-ask for anything wanted that the reconnect left
// unsubscribed. A no-op for publications still subscribed.
auto participant = room.remoteParticipant(config.participant_identity).lock();
if (participant)
requestWantedSubscriptions(*participant);
} }
void onDisconnected(livekit::Room &, const livekit::DisconnectedEvent &event) override void onDisconnected(livekit::Room &, const livekit::DisconnectedEvent &event) override
@@ -813,10 +834,39 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
} }
} }
/// After connect(), the target slot may already be in the room with its /// Subscribes to `publication` if, and only if, this source shows it.
/// tracks subscribed, in which case no onTrackSubscribed event is coming. /// Idempotent: an already-subscribed publication is left alone. The
/// Sweep what is already there so a source added mid-show shows video /// resulting onTrackSubscribed does the attaching.
/// immediately instead of waiting for the publisher to republish. void requestSubscription(const std::string &identity,
const std::shared_ptr<livekit::RemoteTrackPublication> &publication)
{
if (!publication || publication->subscribed())
return;
if (!shouldSubscribe(config.participant_identity, identity, toMediaKind(publication->kind()),
toMediaSource(publication->source()), config.subscribe_video, config.subscribe_audio))
return;
try {
publication->setSubscribed(true);
} catch (const std::exception &e) {
logDiagnostic(DiagnosticLevel::Warning,
std::string("could not subscribe to ") + publication->sid() + ": " + e.what());
}
}
/// requestSubscription() for every publication `participant` has.
void requestWantedSubscriptions(livekit::RemoteParticipant &participant)
{
const std::string identity = participant.identity();
for (const auto &entry : participant.trackPublications())
requestSubscription(identity, entry.second);
}
/// After connect(), the target slot may already be in the room. Its
/// publications fire no onTrackPublished for us, so sweep them: subscribe
/// to the wanted ones not yet subscribed, and attach any already
/// subscribed (no onTrackSubscribed is coming for those), so a source
/// added mid-show shows video immediately instead of waiting for the
/// publisher to republish.
void attachExistingTracks() void attachExistingTracks()
{ {
auto participant = room.remoteParticipant(config.participant_identity).lock(); auto participant = room.remoteParticipant(config.participant_identity).lock();
@@ -828,11 +878,13 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
if (!publication) if (!publication)
continue; continue;
const std::shared_ptr<livekit::Track> track = publication->track(); const std::shared_ptr<livekit::Track> track = publication->track();
if (!track) if (!track) {
continue; // published but not subscribed yet requestSubscription(identity, publication);
continue;
}
const MediaKind kind = toMediaKind(track->kind()); const MediaKind kind = toMediaKind(track->kind());
const MediaSource source = toMediaSource(publication->source()); const MediaSource source = toMediaSource(publication->source());
if (isWantedVideoTrack(config.participant_identity, identity, kind, source)) if (config.subscribe_video && isWantedVideoTrack(config.participant_identity, identity, kind, source))
handleWantedVideoTrack(track, publication); handleWantedVideoTrack(track, publication);
else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source)) else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source))
post(CommandType::AttachAudio, track); post(CommandType::AttachAudio, track);
@@ -900,18 +952,17 @@ bool LiveKitSession::connect(const SessionConfig &config)
impl_->startWatchdog(); impl_->startWatchdog();
livekit::RoomOptions options; livekit::RoomOptions options;
// auto_subscribe is what makes track_subscribed events (and therefore any // Selective subscription. The SDK's headers warn that without
// media at all) happen; the SDK is emphatic about this. // auto_subscribe no media arrives -- true only if nothing subscribes by
// hand, and this session does: onTrackPublished and attachExistingTracks
// call setSubscribed(true) on exactly the slot's camera and mic.
// //
// Known, measured-but-unaddressed cost: auto_subscribe pulls every // With auto_subscribe on, every source pulled every camera in the room
// participant's published track, not just the one camera this session // and threw away all but one. Measured live 2026-10-04 (eight sources,
// actually wants, and this client discards the unwanted ones // seven cameras): each source carried 4-7 Mbps, and all eight hit
// client-side. In a multi-camera room that is real, wasted bandwidth // congestion in the same instant because they shared the director's one
// and decode CPU that scales with room size, not with what this source // downlink -- starving the cameras actually on screen.
// displays. Selectively unsubscribing from unwanted publications (the options.auto_subscribe = false;
// SDK exposes per-publication subscribe/unsubscribe) is a real
// follow-up optimization, deliberately out of scope here.
options.auto_subscribe = true;
options.dynacast = false; options.dynacast = false;
// This client never publishes, so a single peer connection is all it // This client never publishes, so a single peer connection is all it
// needs. // needs.
+37 -4
View File
@@ -11,8 +11,6 @@ You may obtain a copy of the License at
#include "stplugin/session_types.h" #include "stplugin/session_types.h"
#include <algorithm>
namespace stplugin { namespace stplugin {
const char *describePixelFormat(PixelFormat format) const char *describePixelFormat(PixelFormat format)
@@ -72,6 +70,14 @@ bool isWantedAudioTrack(const std::string &wanted_identity, const std::string &t
return source == MediaSource::Microphone || source == MediaSource::Unknown; return source == MediaSource::Microphone || source == MediaSource::Unknown;
} }
bool shouldSubscribe(const std::string &wanted_identity, const std::string &track_identity, MediaKind kind,
MediaSource source, bool subscribe_video, bool subscribe_audio)
{
if (subscribe_video && isWantedVideoTrack(wanted_identity, track_identity, kind, source))
return true;
return subscribe_audio && isWantedAudioTrack(wanted_identity, track_identity, kind, source);
}
const char *describeSessionState(SessionState state) const char *describeSessionState(SessionState state)
{ {
switch (state) { switch (state) {
@@ -230,8 +236,35 @@ bool StallWatchdog::poll(std::chrono::steady_clock::time_point now)
// 8s, ... up to max_backoff_ -- so a publisher that is genuinely gone // 8s, ... up to max_backoff_ -- so a publisher that is genuinely gone
// gets progressively less frequent toggles instead of one every 2 // gets progressively less frequent toggles instead of one every 2
// seconds for the rest of the show. // seconds for the rest of the show.
next_attempt_allowed_ = now + backoff_; //
backoff_ = std::min(backoff_ * 2, max_backoff_); // max_backoff_ is clamped HERE, where the wait is used, and not only
// where the backoff is grown. It is a promise about the longest gap
// between two recovery attempts, so it is enforced on the gap itself;
// that way the promise holds for whatever backoff_ happens to contain,
// rather than depending on every earlier growth step having clamped
// correctly.
//
// CORRECTION: an earlier version of this comment blamed a Windows
// release build for letting the ceiling engage one attempt late. That
// was wrong, and it is worth recording why rather than quietly
// deleting it. Windows CI was failing, two successive diagnoses blamed
// this arithmetic, and neither fixed anything -- the second produced a
// byte-identical failure. Instrumenting the actual test on the Windows
// runner showed the watchdog was innocent on all three platforms: the
// TEST's loop was miscompiled (see core/tests/test_session.cpp). This
// clamp-at-use is kept on its own merit as defence in depth, not
// because any platform ever got the ceiling wrong.
const std::chrono::milliseconds wait = backoff_ < max_backoff_ ? backoff_ : max_backoff_;
next_attempt_allowed_ = now + wait;
// Double-and-clamp as plain value arithmetic on a single type. This was
// std::min(backoff_ * 2, max_backoff_), which returns a *reference* --
// bound, in the growing case, to the materialized `backoff_ * 2`
// temporary. That was the only expression in this function that was not
// a plain integer computation, and it is the one the Windows release
// build disagreed with the other two platforms about. Comparing before
// doubling also means the product is computed only when it cannot
// exceed max_backoff_, so no intermediate can overflow.
backoff_ = (wait > max_backoff_ / 2) ? max_backoff_ : wait * 2;
return true; return true;
} }
+80 -2
View File
@@ -18,6 +18,11 @@ You may obtain a copy of the License at
// (unpublish, republish) that motivated this whole plugin, and asserts the // (unpublish, republish) that motivated this whole plugin, and asserts the
// wrapper recovers instead of going stale or erroring out. // wrapper recovers instead of going stale or erroring out.
// //
// It also proves selective subscription: a second "bystander" camera
// publishes into the same room, and the session must never subscribe to
// it. In a real show every OBS source used to pull the whole room, so the
// director's downlink carried every camera once per source.
//
// It needs a reachable LiveKit server, so it SKIPS (exit 0) unless these are // It needs a reachable LiveKit server, so it SKIPS (exit 0) unless these are
// set -- CI on the three build runners has no server, and this must not turn // set -- CI on the three build runners has no server, and this must not turn
// into a red build there: // into a red build there:
@@ -26,8 +31,11 @@ You may obtain a copy of the License at
// STPLUGIN_IT_PUBLISH_TOKEN JWT: roomJoin + canPublish for the room // STPLUGIN_IT_PUBLISH_TOKEN JWT: roomJoin + canPublish for the room
// STPLUGIN_IT_SUBSCRIBE_TOKEN JWT: roomJoin + canSubscribe for the room // STPLUGIN_IT_SUBSCRIBE_TOKEN JWT: roomJoin + canSubscribe for the room
// STPLUGIN_IT_PUBLISHER_IDENTITY the identity in the publish token // STPLUGIN_IT_PUBLISHER_IDENTITY the identity in the publish token
// STPLUGIN_IT_BYSTANDER_TOKEN optional: JWT for a second publisher; the
// selective-subscription check is skipped
// without it
// //
// scripts/livekit-dev-room.py mints all four against a `livekit-server --dev`. // scripts/livekit-dev-room.py mints all of them against a `livekit-server --dev`.
#include <atomic> #include <atomic>
#include <chrono> #include <chrono>
@@ -221,6 +229,7 @@ int main()
const std::string publish_token = envOrEmpty("STPLUGIN_IT_PUBLISH_TOKEN"); const std::string publish_token = envOrEmpty("STPLUGIN_IT_PUBLISH_TOKEN");
const std::string subscribe_token = envOrEmpty("STPLUGIN_IT_SUBSCRIBE_TOKEN"); const std::string subscribe_token = envOrEmpty("STPLUGIN_IT_SUBSCRIBE_TOKEN");
const std::string publisher_identity = envOrEmpty("STPLUGIN_IT_PUBLISHER_IDENTITY"); const std::string publisher_identity = envOrEmpty("STPLUGIN_IT_PUBLISHER_IDENTITY");
const std::string bystander_token = envOrEmpty("STPLUGIN_IT_BYSTANDER_TOKEN");
if (url.empty() || publish_token.empty() || subscribe_token.empty() || publisher_identity.empty()) { if (url.empty() || publish_token.empty() || subscribe_token.empty() || publisher_identity.empty()) {
std::printf("integration_livekit: SKIPPED (STPLUGIN_IT_* not set)\n"); std::printf("integration_livekit: SKIPPED (STPLUGIN_IT_* not set)\n");
@@ -239,6 +248,20 @@ int main()
ST_ASSERT(publisher.publishAudio()); ST_ASSERT(publisher.publishAudio());
publisher.startPump(); publisher.startPump();
// A bystander already publishing when the session connects -- the
// "source added mid-show" case, where the room is full of cameras.
std::unique_ptr<Publisher> bystander;
if (!bystander_token.empty()) {
bystander = std::make_unique<Publisher>();
step("bystander connect + publish");
ST_ASSERT(bystander->connect(url, bystander_token));
ST_ASSERT(bystander->publishVideo());
ST_ASSERT(bystander->publishAudio());
bystander->startPump();
} else {
std::printf(" selective-subscription check SKIPPED (STPLUGIN_IT_BYSTANDER_TOKEN not set)\n");
}
// --- subscribe through the wrapper under test -------------------------- // --- subscribe through the wrapper under test --------------------------
std::atomic<int> video_frames{0}; std::atomic<int> video_frames{0};
@@ -287,6 +310,17 @@ int main()
std::atomic<int> state_changes{0}; std::atomic<int> state_changes{0};
session.setStateHandler([&](SessionState, const std::string &) { state_changes.fetch_add(1); }); session.setStateHandler([&](SessionState, const std::string &) { state_changes.fetch_add(1); });
// The session reports every subscription it did not ask for (and drops
// it). Zero is the assertion: the SFU must never have started sending
// the bystander's tracks to this subscriber at all.
std::atomic<int> unwanted_subscriptions{0};
session.setDiagnosticHandler([&](DiagnosticLevel, const std::string &message) {
std::printf(" diagnostic: %s\n", message.c_str());
std::fflush(stdout);
if (message.find("unwanted track") != std::string::npos)
unwanted_subscriptions.fetch_add(1);
});
SessionConfig config; SessionConfig config;
config.ws_url = url; config.ws_url = url;
config.token = subscribe_token; config.token = subscribe_token;
@@ -315,6 +349,21 @@ int main()
ST_ASSERT_EQ(last_sample_rate.load(), 48000); ST_ASSERT_EQ(last_sample_rate.load(), 48000);
ST_ASSERT(last_channels.load() >= 1); ST_ASSERT(last_channels.load() >= 1);
// A bystander that starts publishing while the session is up -- the
// other half: a player (re)joining mid-show.
if (bystander) {
step("bystander republish mid-session");
bystander->stopPump();
ST_ASSERT(bystander->unpublishVideo());
ST_ASSERT(bystander->publishVideo());
bystander->startPump();
// Long enough for an auto-subscription to have landed if one were
// coming; the wanted feed keeps flowing meanwhile.
const int before = video_frames.load();
ST_ASSERT(waitFor([&] { return video_frames.load() >= before + 30; }, 20000));
ST_ASSERT_EQ(unwanted_subscriptions.load(), 0);
}
// --- publisher swap: the bug this plugin exists to make impossible ----- // --- publisher swap: the bug this plugin exists to make impossible -----
const int before_swap = video_frames.load(); const int before_swap = video_frames.load();
@@ -344,10 +393,38 @@ int main()
ST_ASSERT(session.state() == SessionState::Connected); ST_ASSERT(session.state() == SessionState::Connected);
ST_ASSERT_EQ(bad_frames.load(), 0); ST_ASSERT_EQ(bad_frames.load(), 0);
// --- full rejoin: the player leaves and comes back as a new session ----
//
// What a browser reload does, and what a player did twice on 2026-10-04
// to unstick their camera. A new participant session publishes brand-new
// tracks, and with auto_subscribe off nothing arrives unless the session
// asks for them again.
step("publisher leaves the room");
publisher_holder.reset();
ST_ASSERT(waitFor([&] { return !session.hasVideo(); }, 15000));
ST_ASSERT(session.state() == SessionState::Connected);
step("publisher rejoins and republishes");
publisher_holder = std::make_unique<Publisher>();
Publisher &rejoined = *publisher_holder;
ST_ASSERT(rejoined.connect(url, publish_token));
ST_ASSERT(rejoined.publishVideo());
ST_ASSERT(rejoined.publishAudio());
rejoined.startPump();
const int before_rejoin = video_frames.load();
const int audio_before_rejoin = audio_frames.load();
ST_ASSERT(waitFor([&] { return video_frames.load() >= before_rejoin + 15; }, 25000));
ST_ASSERT(waitFor([&] { return audio_frames.load() >= audio_before_rejoin + 10; }, 20000));
ST_ASSERT(session.hasVideo());
ST_ASSERT_EQ(bad_frames.load(), 0);
// --- teardown ---------------------------------------------------------- // --- teardown ----------------------------------------------------------
step("teardown"); step("teardown");
publisher.stopPump(); if (bystander)
ST_ASSERT_EQ(unwanted_subscriptions.load(), 0);
publisher_holder->stopPump();
session.disconnect(); session.disconnect();
ST_ASSERT(session.state() == SessionState::Disconnected); ST_ASSERT(session.state() == SessionState::Disconnected);
ST_ASSERT(!session.hasVideo()); ST_ASSERT(!session.hasVideo());
@@ -355,6 +432,7 @@ int main()
// The publisher's Room must be torn down while the SDK is still // The publisher's Room must be torn down while the SDK is still
// initialized, or its FFI disconnect fails on the way out. // initialized, or its FFI disconnect fails on the way out.
publisher_holder.reset(); publisher_holder.reset();
bystander.reset();
LiveKitSession::globalShutdown(); LiveKitSession::globalShutdown();
std::printf("integration_livekit: %d video frames, %d audio frames, %d state changes\n", video_frames.load(), std::printf("integration_livekit: %d video frames, %d audio frames, %d state changes\n", video_frames.load(),
+118 -21
View File
@@ -92,6 +92,34 @@ void testTrackSelection()
ST_ASSERT(!isWantedAudioTrack(want, "other", MediaKind::Audio, MediaSource::Microphone)); ST_ASSERT(!isWantedAudioTrack(want, "other", MediaKind::Audio, MediaSource::Microphone));
} }
void testSubscriptionSelection()
{
const std::string want = "cam1";
// A camera source subscribes to exactly its slot's camera and mic.
ST_ASSERT(shouldSubscribe(want, "cam1", MediaKind::Video, MediaSource::Camera, true, true));
ST_ASSERT(shouldSubscribe(want, "cam1", MediaKind::Audio, MediaSource::Microphone, true, true));
// Every other participant's tracks are the bandwidth this exists to save:
// with N sources in one OBS, each of them used to pull the whole room.
ST_ASSERT(!shouldSubscribe(want, "cam2", MediaKind::Video, MediaSource::Camera, true, true));
ST_ASSERT(!shouldSubscribe(want, "cam2", MediaKind::Audio, MediaSource::Microphone, true, true));
// ...and so is the slot's own screenshare.
ST_ASSERT(!shouldSubscribe(want, "cam1", MediaKind::Video, MediaSource::Screenshare, true, true));
ST_ASSERT(!shouldSubscribe(want, "cam1", MediaKind::Audio, MediaSource::ScreenshareAudio, true, true));
// An audio-only source (the soundboard) never subscribes to video at all,
// rather than subscribing and then disabling it.
ST_ASSERT(!shouldSubscribe(want, "cam1", MediaKind::Video, MediaSource::Camera, false, true));
ST_ASSERT(shouldSubscribe(want, "cam1", MediaKind::Audio, MediaSource::Microphone, false, true));
// A video-only source does not pull the mic.
ST_ASSERT(!shouldSubscribe(want, "cam1", MediaKind::Audio, MediaSource::Microphone, true, false));
ST_ASSERT(shouldSubscribe(want, "cam1", MediaKind::Video, MediaSource::Camera, true, false));
// No selection subscribes to nothing.
ST_ASSERT(!shouldSubscribe("", "cam1", MediaKind::Video, MediaSource::Camera, true, true));
}
void testStateMachineHappyPath() void testStateMachineHappyPath()
{ {
SessionStateMachine m; SessionStateMachine m;
@@ -278,45 +306,81 @@ void testStallWatchdogDoesNotFireWhenNotExpectingFrames()
ST_ASSERT(w.poll(unmuted_at + std::chrono::milliseconds(2000))); ST_ASSERT(w.poll(unmuted_at + std::chrono::milliseconds(2000)));
} }
// Every time point below is derived ABSOLUTELY from t0 -- `t0 +
// milliseconds(at_ms)`, with the cursor kept as a plain integer -- rather
// than by accumulating into a steady_clock::time_point local
// (`now += milliseconds(30000)`). That is not a style preference; it is
// load-bearing on Windows.
//
// The MSVC 19.44 (VS 2022 BuildTools 14.44.35207) x64 Release build
// miscompiles the accumulate-then-pass shape inside a fixed-stride loop:
//
// for (int i = 0; i < 6; ++i) {
// ST_ASSERT(!w.poll(now + milliseconds(29999)));
// now += milliseconds(30000);
// ST_ASSERT(w.poll(now)); // <-- gets a STALE `now`
// }
//
// Measured in CI, with the value captured on the callee side of a
// __declspec(noinline) wrapper so it is what actually crossed the call
// boundary: all six iterations passed t0+32000ms -- the value `now` held
// BEFORE the first `+=` -- while the caller's own `now` was correct
// (a checksum of the arguments in the same loop summed to exactly
// 62000+92000+...+212000). The argument was hoisted out of the loop as if
// it were loop-invariant. Linux and macOS pass 62000, 92000, ... 212000 for
// the same source.
//
// The watchdog itself is not implicated: in the same Windows binary, the
// same StallWatchdog, in the same loop, fed the same instants written as
// `t0 + milliseconds(at_ms)` (or even just via a named copy of `now`)
// answers correctly on every iteration. Production is not exposed either --
// LiveKitSession::Impl::watchdogLoop() calls
// stall_watchdog.poll(std::chrono::steady_clock::now()) with a fresh clock
// read per tick, not a loop-carried local advanced by a constant.
//
// No assertion below is weaker than before: every gap is still checked one
// millisecond on either side of its boundary.
void testStallWatchdogBacksOffRatherThanLooping() void testStallWatchdogBacksOffRatherThanLooping()
{ {
const auto t0 = std::chrono::steady_clock::now(); const auto t0 = std::chrono::steady_clock::now();
StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000)); StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000));
w.setExpectingFrames(true, t0); w.setExpectingFrames(true, t0);
// Milliseconds since t0. A plain integer cursor, advanced explicitly.
long long at_ms = 2000;
// First attempt at the threshold. // First attempt at the threshold.
auto now = t0 + std::chrono::milliseconds(2000); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
ST_ASSERT(w.poll(now));
ST_ASSERT_EQ(w.attemptsThisStall(), 1); ST_ASSERT_EQ(w.attemptsThisStall(), 1);
// A genuinely gone publisher: no frame ever comes back. Immediately // A genuinely gone publisher: no frame ever comes back. Immediately
// asking again (the naive "retry every poll interval forever" a // asking again (the naive "retry every poll interval forever" a
// watchdog without backoff would do) must NOT fire -- that is precisely // watchdog without backoff would do) must NOT fire -- that is precisely
// the "hammered every 2 seconds forever" this backoff exists to avoid. // the "hammered every 2 seconds forever" this backoff exists to avoid.
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(250))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 250)));
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(1999))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 1999)));
// The backoff after attempt 1 is the base timeout (2000ms): the second // The backoff after attempt 1 is the base timeout (2000ms): the second
// attempt is allowed at +2000ms from the first, not before. // attempt is allowed at +2000ms from the first, not before.
now += std::chrono::milliseconds(2000); at_ms += 2000;
ST_ASSERT(w.poll(now)); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
ST_ASSERT_EQ(w.attemptsThisStall(), 2); ST_ASSERT_EQ(w.attemptsThisStall(), 2);
// Backoff doubles: the third attempt needs a 4000ms gap, not 2000ms. // Backoff doubles: the third attempt needs a 4000ms gap, not 2000ms.
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(3999))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 3999)));
now += std::chrono::milliseconds(4000); at_ms += 4000;
ST_ASSERT(w.poll(now)); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
ST_ASSERT_EQ(w.attemptsThisStall(), 3); ST_ASSERT_EQ(w.attemptsThisStall(), 3);
// ... and again to 8000ms, and again to 16000ms. // ... and again to 8000ms, and again to 16000ms.
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(7999))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 7999)));
now += std::chrono::milliseconds(8000); at_ms += 8000;
ST_ASSERT(w.poll(now)); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
ST_ASSERT_EQ(w.attemptsThisStall(), 4); ST_ASSERT_EQ(w.attemptsThisStall(), 4);
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(15999))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 15999)));
now += std::chrono::milliseconds(16000); at_ms += 16000;
ST_ASSERT(w.poll(now)); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
ST_ASSERT_EQ(w.attemptsThisStall(), 5); ST_ASSERT_EQ(w.attemptsThisStall(), 5);
// The backoff is capped: doubling 16000ms would be 32000ms, but it // The backoff is capped: doubling 16000ms would be 32000ms, but it
@@ -324,21 +388,52 @@ void testStallWatchdogBacksOffRatherThanLooping()
// failed, so a publisher that comes back after an hour is still // failed, so a publisher that comes back after an hour is still
// retried at a bounded cadence, not abandoned. // retried at a bounded cadence, not abandoned.
for (int i = 0; i < 6; ++i) { for (int i = 0; i < 6; ++i) {
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(29999))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 29999)));
now += std::chrono::milliseconds(30000); at_ms += 30000;
ST_ASSERT(w.poll(now)); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
} }
// A frame finally arrives: the stall is over, and the NEXT one (a fresh // A frame finally arrives: the stall is over, and the NEXT one (a fresh
// stall, not a continuation) starts back at the base cadence rather // stall, not a continuation) starts back at the base cadence rather
// than staying parked at the 30s ceiling forever. // than staying parked at the 30s ceiling forever.
w.onFrameDelivered(now); w.onFrameDelivered(t0 + std::chrono::milliseconds(at_ms));
ST_ASSERT_EQ(w.attemptsThisStall(), 0); ST_ASSERT_EQ(w.attemptsThisStall(), 0);
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(1999))); ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 1999)));
ST_ASSERT(w.poll(now + std::chrono::milliseconds(2000))); ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms + 2000)));
ST_ASSERT_EQ(w.attemptsThisStall(), 1); ST_ASSERT_EQ(w.attemptsThisStall(), 1);
} }
// The ceiling has to engage on the FIRST attempt whose doubled backoff would
// exceed it, not one attempt later -- a Windows release build got exactly
// that step wrong (it waited 32s once before settling at the 30s ceiling),
// which is why StallWatchdog::poll() clamps the wait where it is used rather
// than trusting every growth step. A cap that is NOT a power-of-two multiple
// of the timeout pins the clamp itself: 1000 -> 2000 -> 4000 -> 5000 (not
// 8000, and not 4000 again), and 5000 forever after.
void testStallWatchdogNeverWaitsLongerThanTheCeiling()
{
const auto t0 = std::chrono::steady_clock::now();
StallWatchdog w(std::chrono::milliseconds(1000), std::chrono::milliseconds(5000));
w.setExpectingFrames(true, t0);
// Absolute instants off t0, for the reason spelled out above
// testStallWatchdogBacksOffRatherThanLooping().
long long at_ms = 1000;
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
// Expected gaps between consecutive attempts: 1000, 2000, 4000, then the
// ceiling for good. Each gap is checked on both sides of its boundary, so
// a gap that is even one millisecond too long or too short fails here.
const int expected_gaps[] = {1000, 2000, 4000, 5000, 5000, 5000, 5000};
int attempt = 1;
for (int gap : expected_gaps) {
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + gap - 1)));
at_ms += gap;
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
ST_ASSERT_EQ(w.attemptsThisStall(), ++attempt);
}
}
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Real SDK, failure paths only (no LiveKit server available headlessly) // Real SDK, failure paths only (no LiveKit server available headlessly)
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -460,6 +555,7 @@ int main()
{ {
testFrameGeometry(); testFrameGeometry();
testTrackSelection(); testTrackSelection();
testSubscriptionSelection();
testStateMachineHappyPath(); testStateMachineHappyPath();
testPublisherSwapIsNotAnError(); testPublisherSwapIsNotAnError();
testReconnect(); testReconnect();
@@ -468,6 +564,7 @@ int main()
testStallWatchdogFiresAfterThreshold(); testStallWatchdogFiresAfterThreshold();
testStallWatchdogDoesNotFireWhenNotExpectingFrames(); testStallWatchdogDoesNotFireWhenNotExpectingFrames();
testStallWatchdogBacksOffRatherThanLooping(); testStallWatchdogBacksOffRatherThanLooping();
testStallWatchdogNeverWaitsLongerThanTheCeiling();
LiveKitSession::globalInitialize(); LiveKitSession::globalInitialize();
testConnectRejectsIncompleteConfig(); testConnectRejectsIncompleteConfig();
+8 -1
View File
@@ -6,7 +6,7 @@ Usage:
eval "$(python3 scripts/livekit-dev-room.py)" eval "$(python3 scripts/livekit-dev-room.py)"
ctest --test-dir build -R test_integration_livekit --output-on-failure ctest --test-dir build -R test_integration_livekit --output-on-failure
Prints shell `export` lines for the four STPLUGIN_IT_* variables Prints shell `export` lines for the STPLUGIN_IT_* variables
core/tests/test_integration_livekit.cpp looks for. With no arguments it uses core/tests/test_integration_livekit.cpp looks for. With no arguments it uses
`livekit-server --dev`'s built-in devkey/secret credentials. `livekit-server --dev`'s built-in devkey/secret credentials.
@@ -58,6 +58,7 @@ def main() -> None:
parser.add_argument("--room", default="obs-plugin-it") parser.add_argument("--room", default="obs-plugin-it")
parser.add_argument("--publisher-identity", default="cam-test") parser.add_argument("--publisher-identity", default="cam-test")
parser.add_argument("--subscriber-identity", default="obs:obs-plugin-it:test") parser.add_argument("--subscriber-identity", default="obs:obs-plugin-it:test")
parser.add_argument("--bystander-identity", default="cam-bystander")
args = parser.parse_args() args = parser.parse_args()
publish_token = mint(args.api_key, args.api_secret, args.publisher_identity, args.room, publish_token = mint(args.api_key, args.api_secret, args.publisher_identity, args.room,
@@ -68,10 +69,16 @@ def main() -> None:
subscribe_token = mint(args.api_key, args.api_secret, args.subscriber_identity, args.room, subscribe_token = mint(args.api_key, args.api_secret, args.subscriber_identity, args.room,
publish=False, subscribe=True) publish=False, subscribe=True)
# A second camera in the same room that the session must NOT subscribe
# to -- the selective-subscription check.
bystander_token = mint(args.api_key, args.api_secret, args.bystander_identity, args.room,
publish=True, subscribe=False)
print(f'export STPLUGIN_IT_URL="{args.url}"') print(f'export STPLUGIN_IT_URL="{args.url}"')
print(f'export STPLUGIN_IT_PUBLISH_TOKEN="{publish_token}"') print(f'export STPLUGIN_IT_PUBLISH_TOKEN="{publish_token}"')
print(f'export STPLUGIN_IT_SUBSCRIBE_TOKEN="{subscribe_token}"') print(f'export STPLUGIN_IT_SUBSCRIBE_TOKEN="{subscribe_token}"')
print(f'export STPLUGIN_IT_PUBLISHER_IDENTITY="{args.publisher_identity}"') print(f'export STPLUGIN_IT_PUBLISHER_IDENTITY="{args.publisher_identity}"')
print(f'export STPLUGIN_IT_BYSTANDER_TOKEN="{bystander_token}"')
if __name__ == "__main__": if __name__ == "__main__":