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>
This commit is contained in:
2026-10-04 18:00:28 -07:00
co-authored by Claude Opus 5.5
parent 22d50ab7be
commit 7723082795
7 changed files with 234 additions and 55 deletions
+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
// 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
// set -- CI on the three build runners has no server, and this must not turn
// 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_SUBSCRIBE_TOKEN JWT: roomJoin + canSubscribe for the room
// 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 <chrono>
@@ -221,6 +229,7 @@ int main()
const std::string publish_token = envOrEmpty("STPLUGIN_IT_PUBLISH_TOKEN");
const std::string subscribe_token = envOrEmpty("STPLUGIN_IT_SUBSCRIBE_TOKEN");
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()) {
std::printf("integration_livekit: SKIPPED (STPLUGIN_IT_* not set)\n");
@@ -239,6 +248,20 @@ int main()
ST_ASSERT(publisher.publishAudio());
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 --------------------------
std::atomic<int> video_frames{0};
@@ -287,6 +310,17 @@ int main()
std::atomic<int> state_changes{0};
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;
config.ws_url = url;
config.token = subscribe_token;
@@ -315,6 +349,21 @@ int main()
ST_ASSERT_EQ(last_sample_rate.load(), 48000);
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 -----
const int before_swap = video_frames.load();
@@ -344,10 +393,38 @@ int main()
ST_ASSERT(session.state() == SessionState::Connected);
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 ----------------------------------------------------------
step("teardown");
publisher.stopPump();
if (bystander)
ST_ASSERT_EQ(unwanted_subscriptions.load(), 0);
publisher_holder->stopPump();
session.disconnect();
ST_ASSERT(session.state() == SessionState::Disconnected);
ST_ASSERT(!session.hasVideo());
@@ -355,6 +432,7 @@ int main()
// The publisher's Room must be torn down while the SDK is still
// initialized, or its FFI disconnect fails on the way out.
publisher_holder.reset();
bystander.reset();
LiveKitSession::globalShutdown();
std::printf("integration_livekit: %d video frames, %d audio frames, %d state changes\n", video_frames.load(),