Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7723082795 | ||
|
|
22d50ab7be | ||
|
|
352a843d93 | ||
|
|
564461a16d | ||
|
|
969a500dfd |
+1
-1
@@ -1,7 +1,7 @@
|
||||
cmake_minimum_required(VERSION 3.19)
|
||||
|
||||
project(obs-streamer-tools-plugin
|
||||
VERSION 0.1.0
|
||||
VERSION 0.1.1
|
||||
DESCRIPTION "OBS Studio source plugin for streamer-tools camera feeds"
|
||||
LANGUAGES C CXX
|
||||
)
|
||||
|
||||
@@ -44,11 +44,9 @@ struct SessionConfig {
|
||||
|
||||
bool subscribe_audio = true;
|
||||
|
||||
/// False for an audio-only source (the soundboard, say): the wanted
|
||||
/// video track is never attached (no AttachVideo command posted), and
|
||||
/// its publication is explicitly disabled server-side (RemoteTrack-
|
||||
/// Publication::setEnabled(false)) so the SFU stops sending it at all --
|
||||
/// not just "decoded and discarded here", genuinely not delivered.
|
||||
/// False for an audio-only source (the soundboard, say): the slot's
|
||||
/// video track is never subscribed to (see shouldSubscribe), so the SFU
|
||||
/// never sends it at all -- not just "decoded and discarded here".
|
||||
bool subscribe_video = true;
|
||||
|
||||
/// How long connect() waits for the room to come up before giving up.
|
||||
|
||||
@@ -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,
|
||||
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
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
+98
-47
@@ -269,39 +269,21 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
|
||||
// Handles the wanted video track once matched, shared by onTrackSubscribed
|
||||
// (a fresh subscription) and attachExistingTracks (one already up when
|
||||
// this session started watching). Two responsibilities that only make
|
||||
// sense together, both keyed off the SAME publication:
|
||||
// this session started watching). An audio-only source (the soundboard)
|
||||
// never gets here for video: shouldSubscribe() never asks for it.
|
||||
//
|
||||
// - subscribe_video: an audio-only source (the soundboard) never wants
|
||||
// this video at all. Rather than attach it and let the OBS adapter
|
||||
// discard every decoded frame, disable the publication itself
|
||||
// (RemoteTrackPublication::setEnabled(false)) so the SFU stops
|
||||
// sending it -- real bandwidth saved, not just wasted decode.
|
||||
// - Fixed video quality: LiveKit's default subscriber behaviour lets
|
||||
// the SFU switch simulcast layers per its own adaptive/bandwidth
|
||||
// logic, which for a source with no rendered-size hint (this is a
|
||||
// native C++ subscriber, not a sized <video> element) means the
|
||||
// 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).
|
||||
// Fixed video quality: LiveKit's default subscriber behaviour lets the
|
||||
// SFU switch simulcast layers per its own adaptive/bandwidth logic, which
|
||||
// for a source with no rendered-size hint (this is a native C++
|
||||
// subscriber, not a sized <video> element) means the 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,
|
||||
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) {
|
||||
try {
|
||||
publication->setVideoQuality(livekit::VideoQuality::HIGH);
|
||||
@@ -359,12 +341,46 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
const MediaSource source =
|
||||
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))
|
||||
handleWantedVideoTrack(event.track, event.publication);
|
||||
else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source))
|
||||
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
|
||||
{
|
||||
if (!event.participant || !event.track)
|
||||
@@ -406,6 +422,11 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
void onReconnected(livekit::Room &, const livekit::ReconnectedEvent &) override
|
||||
{
|
||||
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
|
||||
@@ -813,10 +834,39 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
}
|
||||
}
|
||||
|
||||
/// After connect(), the target slot may already be in the room with its
|
||||
/// tracks subscribed, in which case no onTrackSubscribed event is coming.
|
||||
/// Sweep what is already there so a source added mid-show shows video
|
||||
/// immediately instead of waiting for the publisher to republish.
|
||||
/// Subscribes to `publication` if, and only if, this source shows it.
|
||||
/// Idempotent: an already-subscribed publication is left alone. The
|
||||
/// resulting onTrackSubscribed does the attaching.
|
||||
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()
|
||||
{
|
||||
auto participant = room.remoteParticipant(config.participant_identity).lock();
|
||||
@@ -828,11 +878,13 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
if (!publication)
|
||||
continue;
|
||||
const std::shared_ptr<livekit::Track> track = publication->track();
|
||||
if (!track)
|
||||
continue; // published but not subscribed yet
|
||||
if (!track) {
|
||||
requestSubscription(identity, publication);
|
||||
continue;
|
||||
}
|
||||
const MediaKind kind = toMediaKind(track->kind());
|
||||
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);
|
||||
else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source))
|
||||
post(CommandType::AttachAudio, track);
|
||||
@@ -900,18 +952,17 @@ bool LiveKitSession::connect(const SessionConfig &config)
|
||||
impl_->startWatchdog();
|
||||
|
||||
livekit::RoomOptions options;
|
||||
// auto_subscribe is what makes track_subscribed events (and therefore any
|
||||
// media at all) happen; the SDK is emphatic about this.
|
||||
// Selective subscription. The SDK's headers warn that without
|
||||
// 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
|
||||
// participant's published track, not just the one camera this session
|
||||
// actually wants, and this client discards the unwanted ones
|
||||
// client-side. In a multi-camera room that is real, wasted bandwidth
|
||||
// and decode CPU that scales with room size, not with what this source
|
||||
// displays. Selectively unsubscribing from unwanted publications (the
|
||||
// SDK exposes per-publication subscribe/unsubscribe) is a real
|
||||
// follow-up optimization, deliberately out of scope here.
|
||||
options.auto_subscribe = true;
|
||||
// With auto_subscribe on, every source pulled every camera in the room
|
||||
// and threw away all but one. Measured live 2026-10-04 (eight sources,
|
||||
// seven cameras): each source carried 4-7 Mbps, and all eight hit
|
||||
// congestion in the same instant because they shared the director's one
|
||||
// downlink -- starving the cameras actually on screen.
|
||||
options.auto_subscribe = false;
|
||||
options.dynacast = false;
|
||||
// This client never publishes, so a single peer connection is all it
|
||||
// needs.
|
||||
|
||||
@@ -11,8 +11,6 @@ You may obtain a copy of the License at
|
||||
|
||||
#include "stplugin/session_types.h"
|
||||
|
||||
#include <algorithm>
|
||||
|
||||
namespace stplugin {
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
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)
|
||||
{
|
||||
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
|
||||
// gets progressively less frequent toggles instead of one every 2
|
||||
// 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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(),
|
||||
|
||||
+118
-21
@@ -92,6 +92,34 @@ void testTrackSelection()
|
||||
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()
|
||||
{
|
||||
SessionStateMachine m;
|
||||
@@ -278,45 +306,81 @@ void testStallWatchdogDoesNotFireWhenNotExpectingFrames()
|
||||
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()
|
||||
{
|
||||
const auto t0 = std::chrono::steady_clock::now();
|
||||
StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000));
|
||||
w.setExpectingFrames(true, t0);
|
||||
|
||||
// Milliseconds since t0. A plain integer cursor, advanced explicitly.
|
||||
long long at_ms = 2000;
|
||||
|
||||
// First attempt at the threshold.
|
||||
auto now = t0 + std::chrono::milliseconds(2000);
|
||||
ST_ASSERT(w.poll(now));
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 1);
|
||||
|
||||
// A genuinely gone publisher: no frame ever comes back. Immediately
|
||||
// asking again (the naive "retry every poll interval forever" a
|
||||
// watchdog without backoff would do) must NOT fire -- that is precisely
|
||||
// the "hammered every 2 seconds forever" this backoff exists to avoid.
|
||||
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(250)));
|
||||
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(1999)));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 250)));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 1999)));
|
||||
|
||||
// The backoff after attempt 1 is the base timeout (2000ms): the second
|
||||
// attempt is allowed at +2000ms from the first, not before.
|
||||
now += std::chrono::milliseconds(2000);
|
||||
ST_ASSERT(w.poll(now));
|
||||
at_ms += 2000;
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 2);
|
||||
|
||||
// Backoff doubles: the third attempt needs a 4000ms gap, not 2000ms.
|
||||
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(3999)));
|
||||
now += std::chrono::milliseconds(4000);
|
||||
ST_ASSERT(w.poll(now));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 3999)));
|
||||
at_ms += 4000;
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 3);
|
||||
|
||||
// ... and again to 8000ms, and again to 16000ms.
|
||||
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(7999)));
|
||||
now += std::chrono::milliseconds(8000);
|
||||
ST_ASSERT(w.poll(now));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 7999)));
|
||||
at_ms += 8000;
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 4);
|
||||
|
||||
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(15999)));
|
||||
now += std::chrono::milliseconds(16000);
|
||||
ST_ASSERT(w.poll(now));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 15999)));
|
||||
at_ms += 16000;
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 5);
|
||||
|
||||
// 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
|
||||
// retried at a bounded cadence, not abandoned.
|
||||
for (int i = 0; i < 6; ++i) {
|
||||
ST_ASSERT(!w.poll(now + std::chrono::milliseconds(29999)));
|
||||
now += std::chrono::milliseconds(30000);
|
||||
ST_ASSERT(w.poll(now));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 29999)));
|
||||
at_ms += 30000;
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms)));
|
||||
}
|
||||
|
||||
// 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
|
||||
// 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(!w.poll(now + std::chrono::milliseconds(1999)));
|
||||
ST_ASSERT(w.poll(now + std::chrono::milliseconds(2000)));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(at_ms + 1999)));
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(at_ms + 2000)));
|
||||
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)
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -460,6 +555,7 @@ int main()
|
||||
{
|
||||
testFrameGeometry();
|
||||
testTrackSelection();
|
||||
testSubscriptionSelection();
|
||||
testStateMachineHappyPath();
|
||||
testPublisherSwapIsNotAnError();
|
||||
testReconnect();
|
||||
@@ -468,6 +564,7 @@ int main()
|
||||
testStallWatchdogFiresAfterThreshold();
|
||||
testStallWatchdogDoesNotFireWhenNotExpectingFrames();
|
||||
testStallWatchdogBacksOffRatherThanLooping();
|
||||
testStallWatchdogNeverWaitsLongerThanTheCeiling();
|
||||
|
||||
LiveKitSession::globalInitialize();
|
||||
testConnectRejectsIncompleteConfig();
|
||||
|
||||
@@ -6,7 +6,7 @@ Usage:
|
||||
eval "$(python3 scripts/livekit-dev-room.py)"
|
||||
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
|
||||
`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("--publisher-identity", default="cam-test")
|
||||
parser.add_argument("--subscriber-identity", default="obs:obs-plugin-it:test")
|
||||
parser.add_argument("--bystander-identity", default="cam-bystander")
|
||||
args = parser.parse_args()
|
||||
|
||||
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,
|
||||
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_PUBLISH_TOKEN="{publish_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_BYSTANDER_TOKEN="{bystander_token}"')
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
Reference in New Issue
Block a user