Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d097ea4300 | ||
|
|
a3344c249b | ||
|
|
3e8affea03 | ||
|
|
dc5ffc9938 | ||
|
|
7723082795 |
@@ -3,5 +3,13 @@
|
|||||||
# both need the exact same Linux build dependencies. Edit once, here.
|
# both need the exact same Linux build dependencies. Edit once, here.
|
||||||
set -euo pipefail
|
set -euo pipefail
|
||||||
|
|
||||||
sudo apt-get update -qq
|
# Non-interactive, or a self-hosted Debian runner hangs forever: after an
|
||||||
sudo apt-get install -y -qq cmake ninja-build libobs-dev libcurl4-openssl-dev
|
# install, needrestart asks which daemons to restart ("1. postfix.service
|
||||||
|
# 2. none of the above") and waits on a terminal CI does not have. That
|
||||||
|
# stalled the v0.1.2 Linux release job on 2026-10-05. NEEDRESTART_MODE=a
|
||||||
|
# answers it; DEBIAN_FRONTEND covers debconf prompts. Passed through sudo
|
||||||
|
# explicitly, since sudo drops the caller's environment by default.
|
||||||
|
APT=(sudo DEBIAN_FRONTEND=noninteractive NEEDRESTART_MODE=a apt-get)
|
||||||
|
|
||||||
|
"${APT[@]}" update -qq
|
||||||
|
"${APT[@]}" install -y -qq cmake ninja-build libobs-dev libcurl4-openssl-dev
|
||||||
|
|||||||
@@ -41,7 +41,8 @@ jobs:
|
|||||||
# zip is not guaranteed present on a minimal self-hosted runner
|
# zip is not guaranteed present on a minimal self-hosted runner
|
||||||
# image (unlike GitHub-hosted ubuntu-24.04, which build.yml's
|
# image (unlike GitHub-hosted ubuntu-24.04, which build.yml's
|
||||||
# deps script doesn't need to care about).
|
# deps script doesn't need to care about).
|
||||||
command -v zip >/dev/null || sudo apt-get install -y -qq zip
|
# Non-interactive for the same reason as .gitea/scripts/linux-deps.sh.
|
||||||
|
command -v zip >/dev/null || sudo DEBIAN_FRONTEND=noninteractive NEEDRESTART_MODE=a apt-get install -y -qq zip
|
||||||
out="streamer-tools-camera-${GITEA_REF_NAME}-linux-x64.zip"
|
out="streamer-tools-camera-${GITEA_REF_NAME}-linux-x64.zip"
|
||||||
root="$(pwd)"
|
root="$(pwd)"
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -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.1
|
VERSION 0.1.2
|
||||||
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
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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
@@ -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.
|
||||||
|
|||||||
@@ -70,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) {
|
||||||
|
|||||||
@@ -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(),
|
||||||
|
|||||||
@@ -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;
|
||||||
@@ -527,6 +555,7 @@ int main()
|
|||||||
{
|
{
|
||||||
testFrameGeometry();
|
testFrameGeometry();
|
||||||
testTrackSelection();
|
testTrackSelection();
|
||||||
|
testSubscriptionSelection();
|
||||||
testStateMachineHappyPath();
|
testStateMachineHappyPath();
|
||||||
testPublisherSwapIsNotAnError();
|
testPublisherSwapIsNotAnError();
|
||||||
testReconnect();
|
testReconnect();
|
||||||
|
|||||||
@@ -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__":
|
||||||
|
|||||||
Reference in New Issue
Block a user