Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7723082795 | ||
|
|
22d50ab7be | ||
|
|
352a843d93 | ||
|
|
564461a16d | ||
|
|
969a500dfd | ||
|
|
3f2933ae3f | ||
|
|
48a74e8c67 |
+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
|
||||
)
|
||||
|
||||
@@ -11,6 +11,7 @@ You may obtain a copy of the License at
|
||||
|
||||
#pragma once
|
||||
|
||||
#include <chrono>
|
||||
#include <memory>
|
||||
#include <string>
|
||||
|
||||
@@ -43,17 +44,31 @@ 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.
|
||||
int connect_timeout_ms = 15000;
|
||||
};
|
||||
|
||||
/// How long the stall-recovery watchdog waits for a decoded video frame
|
||||
/// before treating the subscription as stalled and forcing a fresh
|
||||
/// keyframe (see StallWatchdog's comment in session_types.h for why this
|
||||
/// exists -- unrecoverable packet loss with no PLI/keyframe-request API in
|
||||
/// the pinned SDK). Long enough that ordinary jitter never trips it (a
|
||||
/// healthy 30fps subscription delivers a frame at least every ~33ms);
|
||||
/// short enough a director barely has time to notice before it recovers.
|
||||
constexpr std::chrono::milliseconds kStallRecoveryTimeout{2000};
|
||||
|
||||
/// Ceiling for the backoff between repeated recovery attempts against the
|
||||
/// SAME stall. Starts at kStallRecoveryTimeout and doubles each attempt, so
|
||||
/// a genuinely gone publisher (crashed encoder, dead upstream network) is
|
||||
/// retried every 2s, 4s, 8s, ... 30s rather than hammered every 2 seconds
|
||||
/// for the rest of the show.
|
||||
constexpr std::chrono::milliseconds kStallRecoveryMaxBackoff{30000};
|
||||
|
||||
/// Wraps livekit::Room for exactly one subscribed slot.
|
||||
///
|
||||
/// Threading contract, which the OBS adapter depends on:
|
||||
@@ -67,6 +82,11 @@ struct SessionConfig {
|
||||
/// call.
|
||||
/// - The state handler is invoked from whichever thread observed the
|
||||
/// change. It must not block and must not call back into this object.
|
||||
/// - The diagnostic handler (currently just the stall-recovery watchdog,
|
||||
/// see kStallRecoveryTimeout below) may be invoked from the internal
|
||||
/// command-queue worker thread or a video reader thread. Same rules as
|
||||
/// the state handler: must not block, must not call back into this
|
||||
/// object.
|
||||
/// - All handlers must be installed before connect(); they are not
|
||||
/// synchronised against a running session.
|
||||
class LiveKitSession {
|
||||
@@ -80,6 +100,7 @@ public:
|
||||
void setVideoHandler(VideoFrameHandler handler);
|
||||
void setAudioHandler(AudioFrameHandler handler);
|
||||
void setStateHandler(SessionStateHandler handler);
|
||||
void setDiagnosticHandler(DiagnosticHandler handler);
|
||||
|
||||
/// Connect and start subscribing. Returns true once the room is up; the
|
||||
/// selected slot's tracks may still arrive later (or not at all, if the
|
||||
|
||||
@@ -19,6 +19,7 @@ You may obtain a copy of the License at
|
||||
// self-consistent -- all live here, and LiveKitSession is the (much thinner)
|
||||
// piece that wires real SDK callbacks into them.
|
||||
|
||||
#include <chrono>
|
||||
#include <cstddef>
|
||||
#include <cstdint>
|
||||
#include <functional>
|
||||
@@ -69,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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -137,6 +146,82 @@ private:
|
||||
bool has_audio_ = false;
|
||||
};
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Stall-recovery watchdog (pure timing/decision logic)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Decides WHEN to force a fresh keyframe on an already-subscribed video
|
||||
/// track. It does not touch LiveKit or OBS at all -- LiveKitSession::Impl
|
||||
/// (session.cpp) is what actually carries the decision out
|
||||
/// (RemoteTrackPublication::setEnabled(false) then setEnabled(true)), and
|
||||
/// only from its SDK command-queue thread, the same rule every other SDK
|
||||
/// interaction in that file already follows.
|
||||
///
|
||||
/// Why this exists: measured on the live server, comparing subscribers in
|
||||
/// the same LiveKit room over the same 30-minute window, every OBS plugin
|
||||
/// connection racked up ~862-1099 nackMisses and ~2200-3100 nackRepeated --
|
||||
/// nackMisses means the subscriber asked the SFU to retransmit a packet
|
||||
/// that had already aged out of its send buffer, i.e. unrecoverable loss --
|
||||
/// while a browser subscriber in the same room saw 0 and 0. A decoder that
|
||||
/// loses a frame that way cannot resync without a fresh keyframe. Every OBS
|
||||
/// plugin connection also sat at `plis` == 2 for a multi-hour session (a
|
||||
/// browser adapts and asks for keyframes normally), and grepping the pinned
|
||||
/// client-sdk-cpp (1.10.1) headers turns up no PLI/keyframe-request API at
|
||||
/// all. So today, once that happens, the source just stays broken for the
|
||||
/// rest of the show -- "drops at random and never recovers". Toggling the
|
||||
/// subscription off and back on is the one lever this SDK exposes that
|
||||
/// forces the SFU to stop and restart delivery of the track, and a restart
|
||||
/// always begins with a keyframe. This class is the "have we gone too long
|
||||
/// without a frame, and is it still worth trying again" clock behind that
|
||||
/// lever; see LiveKitSession::Impl::recoverVideo() for where it is pulled.
|
||||
///
|
||||
/// Not thread-safe on its own, deliberately -- same contract as
|
||||
/// SessionStateMachine above: LiveKitSession::Impl owns the lock (in
|
||||
/// practice the same state_mutex that guards `machine`).
|
||||
class StallWatchdog {
|
||||
public:
|
||||
StallWatchdog(std::chrono::milliseconds timeout, std::chrono::milliseconds max_backoff);
|
||||
|
||||
/// Call whenever whether a frame could legitimately arrive right now
|
||||
/// changes: true once a video track is subscribed (and unmuted), false
|
||||
/// on detach/unsubscribe/mute/disconnect. Flipping to false always
|
||||
/// clears all timing state; flipping back to true always starts a
|
||||
/// brand-new grace period rather than measuring from a stale timestamp.
|
||||
/// That is what stops an unmute -- or an ordinary publisher swap --
|
||||
/// from firing the INSTANT it resumes, off a "last frame" that might
|
||||
/// actually be minutes old: a muted, disabled or unsubscribed track, an
|
||||
/// audio-only source, or a disconnected session must never trip this.
|
||||
void setExpectingFrames(bool expecting, std::chrono::steady_clock::time_point now);
|
||||
|
||||
/// Call every time a decoded video frame is actually delivered.
|
||||
void onFrameDelivered(std::chrono::steady_clock::time_point now);
|
||||
|
||||
/// Call periodically (finer-grained than the configured timeout).
|
||||
/// Returns true exactly when a recovery attempt should be made right
|
||||
/// now; each true also arms the backoff before the next one is even
|
||||
/// considered, so a caller polling in a tight loop still cannot fire
|
||||
/// back-to-back attempts against a publisher that never comes back --
|
||||
/// see the class comment: hammering every 2 seconds forever against a
|
||||
/// genuinely gone publisher is worse than a frozen source.
|
||||
bool poll(std::chrono::steady_clock::time_point now);
|
||||
|
||||
/// Attempts made since the current stall started (since the last frame,
|
||||
/// or since expecting-frames most recently became true). Reset by
|
||||
/// onFrameDelivered and by setExpectingFrames.
|
||||
int attemptsThisStall() const { return attempts_; }
|
||||
|
||||
private:
|
||||
std::chrono::milliseconds timeout_;
|
||||
std::chrono::milliseconds max_backoff_;
|
||||
|
||||
bool expecting_ = false;
|
||||
bool have_baseline_ = false;
|
||||
std::chrono::steady_clock::time_point baseline_{};
|
||||
std::chrono::milliseconds backoff_{};
|
||||
std::chrono::steady_clock::time_point next_attempt_allowed_{};
|
||||
int attempts_ = 0;
|
||||
};
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Frames handed to the OBS adapter
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -177,4 +262,13 @@ using VideoFrameHandler = std::function<void(const VideoFrameData &)>;
|
||||
using AudioFrameHandler = std::function<void(const AudioFrameData &)>;
|
||||
using SessionStateHandler = std::function<void(SessionState state, const std::string &detail)>;
|
||||
|
||||
/// Severity for LiveKitSession's own diagnostic log lines (currently just
|
||||
/// the stall-recovery watchdog). Kept separate from OBS's LOG_* levels and
|
||||
/// from the SDK's own livekit::LogLevel so core/ stays free of any OBS
|
||||
/// dependency -- the adapter maps this onto obs_log the same way it already
|
||||
/// maps livekit::LogLevel (see plugin-main.cpp's livekit log bridge).
|
||||
enum class DiagnosticLevel { Info, Warning };
|
||||
|
||||
using DiagnosticHandler = std::function<void(DiagnosticLevel level, const std::string &message)>;
|
||||
|
||||
} // namespace stplugin
|
||||
|
||||
+359
-58
@@ -17,6 +17,7 @@ You may obtain a copy of the License at
|
||||
#include <deque>
|
||||
#include <exception>
|
||||
#include <mutex>
|
||||
#include <string>
|
||||
#include <thread>
|
||||
#include <utility>
|
||||
#include <vector>
|
||||
@@ -30,6 +31,7 @@ You may obtain a copy of the License at
|
||||
#include <livekit/room_delegate.h>
|
||||
#include <livekit/room_event_types.h>
|
||||
#include <livekit/track.h>
|
||||
#include <livekit/track_publication.h>
|
||||
#include <livekit/video_frame.h>
|
||||
#include <livekit/video_stream.h>
|
||||
|
||||
@@ -138,6 +140,13 @@ int &globalRefCount()
|
||||
return n;
|
||||
}
|
||||
|
||||
// --- Stall-recovery watchdog tuning -----------------------------------------
|
||||
|
||||
/// How often the watchdog thread checks StallWatchdog's clock. Deliberately
|
||||
/// finer than kStallRecoveryTimeout (session.h) so detection latency tracks
|
||||
/// the threshold itself, not the threshold plus a whole polling period.
|
||||
constexpr std::chrono::milliseconds kWatchdogPollInterval{250};
|
||||
|
||||
} // namespace
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -145,11 +154,12 @@ int &globalRefCount()
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
enum class CommandType { AttachVideo, DetachVideo, AttachAudio, DetachAudio, Stop };
|
||||
enum class CommandType { AttachVideo, DetachVideo, AttachAudio, DetachAudio, RecoverVideo, Stop };
|
||||
|
||||
struct Command {
|
||||
CommandType type;
|
||||
std::shared_ptr<livekit::Track> track;
|
||||
std::shared_ptr<livekit::RemoteTrackPublication> publication;
|
||||
};
|
||||
|
||||
livekit::Room room;
|
||||
@@ -157,10 +167,14 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
|
||||
mutable std::mutex state_mutex;
|
||||
SessionStateMachine machine;
|
||||
// Guarded by state_mutex too, same contract as `machine` -- see
|
||||
// StallWatchdog's own comment for what it decides and why it exists.
|
||||
StallWatchdog stall_watchdog{kStallRecoveryTimeout, kStallRecoveryMaxBackoff};
|
||||
|
||||
VideoFrameHandler on_video;
|
||||
AudioFrameHandler on_audio;
|
||||
SessionStateHandler on_state;
|
||||
DiagnosticHandler on_diagnostic;
|
||||
|
||||
std::atomic<std::uint64_t> video_frames{0};
|
||||
std::atomic<std::uint64_t> audio_frames{0};
|
||||
@@ -170,7 +184,10 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
// livekit::AudioStream happens on `worker`, never on a room event thread:
|
||||
// the SDK's room callbacks run on its own event thread and blocking or
|
||||
// re-entering there stalls every other event (and Room::disconnect() from
|
||||
// inside one is documented to deadlock outright).
|
||||
// inside one is documented to deadlock outright). The stall-recovery
|
||||
// watchdog thread follows the same rule: it never touches
|
||||
// RemoteTrackPublication itself, only posts CommandType::RecoverVideo
|
||||
// and lets the worker thread do it (see recoverVideo() below).
|
||||
std::mutex queue_mutex;
|
||||
std::condition_variable queue_cv;
|
||||
std::deque<Command> queue;
|
||||
@@ -183,6 +200,26 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
std::shared_ptr<livekit::AudioStream> audio_stream;
|
||||
std::thread audio_thread;
|
||||
|
||||
// The publication/track backing the CURRENT video subscription, kept
|
||||
// around purely so recoverVideo() and the mute-change handlers have
|
||||
// something to act on without reaching into `video_stream` (which is
|
||||
// worker-thread-exclusive, per the comment above). Set together in
|
||||
// attachVideo(), cleared together in detachVideo(), both on the worker
|
||||
// thread; read from the room event thread (handleMuteChange) and the
|
||||
// worker thread (recoverVideo()) under this mutex.
|
||||
std::mutex video_track_mutex;
|
||||
std::shared_ptr<livekit::Track> current_video_track;
|
||||
std::shared_ptr<livekit::RemoteTrackPublication> current_video_publication;
|
||||
|
||||
// The watchdog's own timer thread. It owns no SDK state and calls no SDK
|
||||
// method directly -- see the queue comment above. `watchdog_mutex` only
|
||||
// ever guards the shutdown flag/condvar pair, never `stall_watchdog`
|
||||
// (that is guarded by `state_mutex`, alongside `machine`).
|
||||
std::mutex watchdog_mutex;
|
||||
std::condition_variable watchdog_cv;
|
||||
bool watchdog_running = false;
|
||||
std::thread watchdog_thread;
|
||||
|
||||
bool connected = false;
|
||||
|
||||
~Impl() override = default;
|
||||
@@ -207,52 +244,46 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
handler(state, detail);
|
||||
}
|
||||
|
||||
void post(CommandType type, std::shared_ptr<livekit::Track> track = nullptr)
|
||||
void post(CommandType type, std::shared_ptr<livekit::Track> track = nullptr,
|
||||
std::shared_ptr<livekit::RemoteTrackPublication> publication = nullptr)
|
||||
{
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(queue_mutex);
|
||||
if (!worker_running)
|
||||
return;
|
||||
queue.push_back(Command{type, std::move(track)});
|
||||
queue.push_back(Command{type, std::move(track), std::move(publication)});
|
||||
}
|
||||
queue_cv.notify_one();
|
||||
}
|
||||
|
||||
void logDiagnostic(DiagnosticLevel level, const std::string &message)
|
||||
{
|
||||
DiagnosticHandler handler;
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
handler = on_diagnostic;
|
||||
}
|
||||
if (handler)
|
||||
handler(level, message);
|
||||
}
|
||||
|
||||
// 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);
|
||||
@@ -261,7 +292,42 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
// it already had, which is the pre-existing behaviour.
|
||||
}
|
||||
}
|
||||
post(CommandType::AttachVideo, track);
|
||||
// `publication` rides along so attachVideo() can remember it: it is
|
||||
// the handle the stall-recovery watchdog later toggles
|
||||
// (setEnabled(false)/(true)) to force a fresh keyframe. See
|
||||
// recoverVideo() and StallWatchdog's comment in session_types.h.
|
||||
post(CommandType::AttachVideo, track, publication);
|
||||
}
|
||||
|
||||
// Fired for ANY track (any participant, any kind) muting or unmuting.
|
||||
// Filtered down to "is this the video publication we are currently
|
||||
// watching" by SID, which also naturally excludes every audio mute and
|
||||
// every other participant's tracks without a separate identity/kind
|
||||
// check.
|
||||
//
|
||||
// Why this exists: the stall-recovery watchdog (see StallWatchdog's
|
||||
// comment) treats "no decoded frame for kStallRecoveryTimeout" as a
|
||||
// stall worth toggling the subscription over. A publisher who
|
||||
// legitimately turned their camera off produces exactly that symptom on
|
||||
// purpose, and toggling their subscription every couple of seconds for
|
||||
// as long as they stay off would be an endless, pointless loop against
|
||||
// healthy behaviour. Muting suspends the watchdog's clock entirely;
|
||||
// unmuting starts a brand-new grace period rather than reading "muted
|
||||
// for twenty minutes" as "stalled for twenty minutes".
|
||||
void handleMuteChange(const std::shared_ptr<livekit::TrackPublication> &publication, bool unmuted)
|
||||
{
|
||||
if (!publication || publication->kind() != livekit::TrackKind::KIND_VIDEO)
|
||||
return;
|
||||
std::string current_sid;
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(video_track_mutex);
|
||||
if (current_video_publication)
|
||||
current_sid = current_video_publication->sid();
|
||||
}
|
||||
if (current_sid.empty() || publication->sid() != current_sid)
|
||||
return;
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
stall_watchdog.setExpectingFrames(unmuted, std::chrono::steady_clock::now());
|
||||
}
|
||||
|
||||
// --- RoomDelegate ------------------------------------------------------
|
||||
@@ -275,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)
|
||||
@@ -304,6 +404,16 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
post(CommandType::DetachAudio);
|
||||
}
|
||||
|
||||
void onTrackMuted(livekit::Room &, const livekit::TrackMutedEvent &event) override
|
||||
{
|
||||
handleMuteChange(event.publication, false);
|
||||
}
|
||||
|
||||
void onTrackUnmuted(livekit::Room &, const livekit::TrackUnmutedEvent &event) override
|
||||
{
|
||||
handleMuteChange(event.publication, true);
|
||||
}
|
||||
|
||||
void onReconnecting(livekit::Room &, const livekit::ReconnectingEvent &) override
|
||||
{
|
||||
mutateState([](SessionStateMachine &m) { m.onReconnecting(); });
|
||||
@@ -312,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
|
||||
@@ -359,7 +474,7 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
void workerLoop()
|
||||
{
|
||||
for (;;) {
|
||||
Command command{CommandType::Stop, nullptr};
|
||||
Command command{CommandType::Stop, nullptr, nullptr};
|
||||
{
|
||||
std::unique_lock<std::mutex> lock(queue_mutex);
|
||||
queue_cv.wait(lock, [this] { return !queue.empty(); });
|
||||
@@ -369,7 +484,7 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
|
||||
switch (command.type) {
|
||||
case CommandType::AttachVideo:
|
||||
attachVideo(command.track);
|
||||
attachVideo(command.track, command.publication);
|
||||
break;
|
||||
case CommandType::DetachVideo:
|
||||
detachVideo();
|
||||
@@ -380,6 +495,9 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
case CommandType::DetachAudio:
|
||||
detachAudio();
|
||||
break;
|
||||
case CommandType::RecoverVideo:
|
||||
recoverVideo();
|
||||
break;
|
||||
case CommandType::Stop:
|
||||
detachVideo();
|
||||
detachAudio();
|
||||
@@ -388,7 +506,103 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
}
|
||||
}
|
||||
|
||||
void attachVideo(const std::shared_ptr<livekit::Track> &track)
|
||||
// --- stall-recovery watchdog --------------------------------------------
|
||||
//
|
||||
// See StallWatchdog's comment (session_types.h) for the measured
|
||||
// evidence and why toggling the publication is the only lever
|
||||
// available. This thread does exactly one thing: tick StallWatchdog's
|
||||
// clock and, when it says to, post CommandType::RecoverVideo so the
|
||||
// worker thread does the actual SDK call. It never touches
|
||||
// livekit::VideoStream, livekit::Room or RemoteTrackPublication itself.
|
||||
|
||||
void startWatchdog()
|
||||
{
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(watchdog_mutex);
|
||||
watchdog_running = true;
|
||||
}
|
||||
watchdog_thread = std::thread([this] { watchdogLoop(); });
|
||||
}
|
||||
|
||||
void stopWatchdog()
|
||||
{
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(watchdog_mutex);
|
||||
if (!watchdog_running)
|
||||
return;
|
||||
watchdog_running = false;
|
||||
}
|
||||
watchdog_cv.notify_all();
|
||||
if (watchdog_thread.joinable())
|
||||
watchdog_thread.join();
|
||||
}
|
||||
|
||||
void watchdogLoop()
|
||||
{
|
||||
std::unique_lock<std::mutex> lock(watchdog_mutex);
|
||||
while (watchdog_running) {
|
||||
// kWatchdogPollInterval is finer than kStallRecoveryTimeout so
|
||||
// detection latency tracks the threshold itself rather than the
|
||||
// threshold plus a whole polling period; woken early on
|
||||
// shutdown by stopWatchdog()'s notify_all().
|
||||
watchdog_cv.wait_for(lock, kWatchdogPollInterval);
|
||||
if (!watchdog_running)
|
||||
break;
|
||||
lock.unlock();
|
||||
|
||||
bool should_fire;
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
should_fire = stall_watchdog.poll(std::chrono::steady_clock::now());
|
||||
}
|
||||
if (should_fire)
|
||||
post(CommandType::RecoverVideo);
|
||||
|
||||
lock.lock();
|
||||
}
|
||||
}
|
||||
|
||||
// Actually pulls the lever: toggles the video publication off and back
|
||||
// on, which makes the SFU stop and restart delivery of that track --
|
||||
// and a restart always begins with a keyframe. Only ever called from
|
||||
// the worker thread (via CommandType::RecoverVideo), same as every
|
||||
// other RemoteTrackPublication/VideoStream call in this file.
|
||||
void recoverVideo()
|
||||
{
|
||||
std::shared_ptr<livekit::RemoteTrackPublication> publication;
|
||||
std::shared_ptr<livekit::Track> track;
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(video_track_mutex);
|
||||
publication = current_video_publication;
|
||||
track = current_video_track;
|
||||
}
|
||||
// Detached, replaced, or muted between the watchdog deciding to
|
||||
// fire and the worker getting to this command -- nothing to do, and
|
||||
// silently: this is the expected shape of the race, not a failure
|
||||
// worth logging.
|
||||
if (!publication || !track || track->muted())
|
||||
return;
|
||||
|
||||
int attempt = 0;
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
attempt = stall_watchdog.attemptsThisStall();
|
||||
}
|
||||
|
||||
logDiagnostic(DiagnosticLevel::Warning,
|
||||
"no decoded video frame for >= " + std::to_string(kStallRecoveryTimeout.count()) +
|
||||
"ms; toggling the subscription to force a fresh keyframe (attempt " +
|
||||
std::to_string(attempt) + ")");
|
||||
try {
|
||||
publication->setEnabled(false);
|
||||
publication->setEnabled(true);
|
||||
} catch (const std::exception &e) {
|
||||
logDiagnostic(DiagnosticLevel::Warning, std::string("stall-recovery toggle failed: ") + e.what());
|
||||
}
|
||||
}
|
||||
|
||||
void attachVideo(const std::shared_ptr<livekit::Track> &track,
|
||||
const std::shared_ptr<livekit::RemoteTrackPublication> &publication)
|
||||
{
|
||||
if (!track)
|
||||
return;
|
||||
@@ -413,6 +627,18 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
if (!stream)
|
||||
return;
|
||||
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(video_track_mutex);
|
||||
current_video_track = track;
|
||||
current_video_publication = publication;
|
||||
}
|
||||
{
|
||||
// A fresh subscription (or a publisher swap) starts a brand-new
|
||||
// grace period -- see StallWatchdog::setExpectingFrames.
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
stall_watchdog.setExpectingFrames(true, std::chrono::steady_clock::now());
|
||||
}
|
||||
|
||||
video_stream = stream;
|
||||
video_thread = std::thread([this, stream] { videoReaderLoop(stream); });
|
||||
mutateState([](SessionStateMachine &m) { m.onVideoAttached(); });
|
||||
@@ -426,6 +652,20 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
video_thread.join();
|
||||
const bool had = static_cast<bool>(video_stream);
|
||||
video_stream.reset();
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(video_track_mutex);
|
||||
current_video_track.reset();
|
||||
current_video_publication.reset();
|
||||
}
|
||||
{
|
||||
// Nothing subscribed means nothing expected -- see
|
||||
// StallWatchdog::setExpectingFrames. Unconditional, not gated on
|
||||
// `had`: this also covers the detachVideo() at the top of
|
||||
// attachVideo() above, which is exactly the publisher-swap
|
||||
// moment the grace period needs to restart from.
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
stall_watchdog.setExpectingFrames(false, std::chrono::steady_clock::now());
|
||||
}
|
||||
if (had)
|
||||
mutateState([](SessionStateMachine &m) { m.onVideoDetached(); });
|
||||
}
|
||||
@@ -548,6 +788,23 @@ struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
||||
out.plane_count = count;
|
||||
|
||||
video_frames.fetch_add(1);
|
||||
|
||||
// Tell the stall-recovery watchdog a frame actually made it all the
|
||||
// way to "about to hand to OBS" -- not merely that the SDK's queue
|
||||
// produced an event, which the drop paths above also see. See
|
||||
// StallWatchdog's comment for why this exists.
|
||||
int attempts_before_this_frame = 0;
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(state_mutex);
|
||||
attempts_before_this_frame = stall_watchdog.attemptsThisStall();
|
||||
stall_watchdog.onFrameDelivered(std::chrono::steady_clock::now());
|
||||
}
|
||||
if (attempts_before_this_frame > 0) {
|
||||
logDiagnostic(DiagnosticLevel::Info, "video resumed after " +
|
||||
std::to_string(attempts_before_this_frame) +
|
||||
" stall-recovery attempt(s)");
|
||||
}
|
||||
|
||||
handler(out);
|
||||
}
|
||||
|
||||
@@ -577,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();
|
||||
@@ -592,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);
|
||||
@@ -633,6 +921,12 @@ void LiveKitSession::setStateHandler(SessionStateHandler handler)
|
||||
impl_->on_state = std::move(handler);
|
||||
}
|
||||
|
||||
void LiveKitSession::setDiagnosticHandler(DiagnosticHandler handler)
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
||||
impl_->on_diagnostic = std::move(handler);
|
||||
}
|
||||
|
||||
bool LiveKitSession::connect(const SessionConfig &config)
|
||||
{
|
||||
if (impl_->connected)
|
||||
@@ -652,20 +946,23 @@ bool LiveKitSession::connect(const SessionConfig &config)
|
||||
}
|
||||
|
||||
impl_->startWorker();
|
||||
// Runs for the lifetime of the worker: connected-but-nothing-subscribed
|
||||
// is a no-op for StallWatchdog (see setExpectingFrames), so there is no
|
||||
// reason to start/stop it separately from the worker it posts to.
|
||||
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.
|
||||
@@ -680,6 +977,7 @@ bool LiveKitSession::connect(const SessionConfig &config)
|
||||
} catch (const std::exception &e) {
|
||||
ok = false;
|
||||
impl_->mutateState([&](SessionStateMachine &m) { m.onConnectFailed(e.what()); });
|
||||
impl_->stopWatchdog();
|
||||
impl_->stopWorker();
|
||||
impl_->room.setDelegate(nullptr);
|
||||
return false;
|
||||
@@ -689,6 +987,7 @@ bool LiveKitSession::connect(const SessionConfig &config)
|
||||
impl_->mutateState([](SessionStateMachine &m) {
|
||||
m.onConnectFailed("could not connect to LiveKit (check the server URL, or the token may have expired)");
|
||||
});
|
||||
impl_->stopWatchdog();
|
||||
impl_->stopWorker();
|
||||
impl_->room.setDelegate(nullptr);
|
||||
return false;
|
||||
@@ -705,9 +1004,11 @@ void LiveKitSession::disconnect()
|
||||
if (!impl_)
|
||||
return;
|
||||
|
||||
// Order matters: stop the readers first so nothing is mid-read on a
|
||||
// stream the room is about to tear down, then disconnect the room, then
|
||||
// drop the delegate so no event can arrive at a half-destroyed object.
|
||||
// Order matters: stop the watchdog and the readers first so nothing is
|
||||
// mid-read (or about to post a recovery command) on a stream the room
|
||||
// is about to tear down, then disconnect the room, then drop the
|
||||
// delegate so no event can arrive at a half-destroyed object.
|
||||
impl_->stopWatchdog();
|
||||
impl_->stopWorker();
|
||||
|
||||
if (impl_->connected) {
|
||||
|
||||
@@ -70,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) {
|
||||
@@ -175,4 +183,89 @@ void SessionStateMachine::onAudioDetached()
|
||||
has_audio_ = false;
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// StallWatchdog
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
StallWatchdog::StallWatchdog(std::chrono::milliseconds timeout, std::chrono::milliseconds max_backoff)
|
||||
: timeout_(timeout), max_backoff_(max_backoff), backoff_(timeout)
|
||||
{
|
||||
}
|
||||
|
||||
void StallWatchdog::setExpectingFrames(bool expecting, std::chrono::steady_clock::time_point now)
|
||||
{
|
||||
expecting_ = expecting;
|
||||
// Always re-baseline from `now`, whichever direction this flips.
|
||||
// Losing the baseline (rather than, say, keeping the old one around for
|
||||
// when expecting_ next becomes true) is what stops a track that was
|
||||
// muted for the last twenty minutes from reading as "twenty minutes
|
||||
// stalled" the instant it unmutes.
|
||||
have_baseline_ = expecting;
|
||||
baseline_ = now;
|
||||
backoff_ = timeout_;
|
||||
next_attempt_allowed_ = now;
|
||||
attempts_ = 0;
|
||||
}
|
||||
|
||||
void StallWatchdog::onFrameDelivered(std::chrono::steady_clock::time_point now)
|
||||
{
|
||||
have_baseline_ = true;
|
||||
baseline_ = now;
|
||||
backoff_ = timeout_;
|
||||
// A recovered stream must be able to fire again the moment a FRESH
|
||||
// stall clears the (now-reset) timeout, not sit throttled by whatever
|
||||
// backoff a previous, unrelated stall had climbed to -- next_attempt_
|
||||
// allowed_ belongs to that old stall and is meaningless once frames are
|
||||
// flowing again.
|
||||
next_attempt_allowed_ = now;
|
||||
attempts_ = 0;
|
||||
}
|
||||
|
||||
bool StallWatchdog::poll(std::chrono::steady_clock::time_point now)
|
||||
{
|
||||
if (!expecting_ || !have_baseline_)
|
||||
return false;
|
||||
if (now - baseline_ < timeout_)
|
||||
return false;
|
||||
if (now < next_attempt_allowed_)
|
||||
return false;
|
||||
|
||||
++attempts_;
|
||||
// Next attempt against this SAME stall is not allowed until the backoff
|
||||
// elapses, and the backoff itself doubles (capped) each time -- 2s, 4s,
|
||||
// 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.
|
||||
//
|
||||
// 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;
|
||||
}
|
||||
|
||||
} // namespace stplugin
|
||||
|
||||
@@ -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(),
|
||||
|
||||
+233
-2
@@ -14,8 +14,10 @@ You may obtain a copy of the License at
|
||||
// COVERED headlessly -- the session's own decision-making: the state
|
||||
// machine's transitions (including the publisher-swap and reconnect paths
|
||||
// that motivated this plugin), track selection, frame geometry validation,
|
||||
// and the real connect() failure paths against the real SDK (bad URL,
|
||||
// unreachable host, garbage token).
|
||||
// the stall-recovery watchdog's timing/backoff decisions (StallWatchdog,
|
||||
// driven with a fake clock -- see its own section below), and the real
|
||||
// connect() failure paths against the real SDK (bad URL, unreachable
|
||||
// host, garbage token).
|
||||
//
|
||||
// NOT COVERED here -- anything that needs a LiveKit server to answer:
|
||||
// a successful connect, actual subscription, and actual decoded frames
|
||||
@@ -90,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;
|
||||
@@ -209,6 +239,201 @@ void testFailureAndRecovery()
|
||||
ST_ASSERT(idle.state() == SessionState::Idle);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// StallWatchdog -- the stall-recovery watchdog's pure timing/decision logic.
|
||||
//
|
||||
// This is deliberately driven with an explicit, fake clock (arbitrary
|
||||
// steady_clock::time_points built by hand, never std::chrono::...::now())
|
||||
// rather than real sleeps: every case below needs to be exact about
|
||||
// "1999ms in" vs "2001ms in" and about backoff boundaries, and a test that
|
||||
// actually slept for 30+ seconds to exercise the backoff ceiling would be
|
||||
// exactly the kind of slow, flaky test this project's whole headless-test
|
||||
// philosophy exists to avoid. See StallWatchdog's own comment
|
||||
// (session_types.h) for the measured server evidence this exists to fix.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
void testStallWatchdogFiresAfterThreshold()
|
||||
{
|
||||
const auto t0 = std::chrono::steady_clock::now();
|
||||
StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000));
|
||||
|
||||
// Nothing subscribed yet: polling is a no-op, no matter how much time
|
||||
// has "passed" -- an audio-only source or a disconnected session must
|
||||
// never fire.
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(10000)));
|
||||
|
||||
// A track becomes subscribed. Still well under the threshold: quiet.
|
||||
w.setExpectingFrames(true, t0);
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(500)));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(1999)));
|
||||
|
||||
// No frame ever arrived, and the threshold has now elapsed: fires
|
||||
// exactly once when asked right at/after the boundary.
|
||||
ST_ASSERT(w.poll(t0 + std::chrono::milliseconds(2000)));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 1);
|
||||
|
||||
// A frame arriving resets the clock -- the far more common case in a
|
||||
// healthy stream, where onFrameDelivered() is called every ~33ms and
|
||||
// poll() (every kWatchdogPollInterval) never sees 2000ms of silence.
|
||||
StallWatchdog healthy(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000));
|
||||
healthy.setExpectingFrames(true, t0);
|
||||
for (int ms = 0; ms <= 5000; ms += 33)
|
||||
healthy.onFrameDelivered(t0 + std::chrono::milliseconds(ms));
|
||||
ST_ASSERT(!healthy.poll(t0 + std::chrono::milliseconds(5010)));
|
||||
ST_ASSERT_EQ(healthy.attemptsThisStall(), 0);
|
||||
}
|
||||
|
||||
void testStallWatchdogDoesNotFireWhenNotExpectingFrames()
|
||||
{
|
||||
const auto t0 = std::chrono::steady_clock::now();
|
||||
StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000));
|
||||
|
||||
// A muted (or disabled/unsubscribed) track is expected silence, not a
|
||||
// stall -- setExpectingFrames(false, ...) is exactly what
|
||||
// LiveKitSession::Impl::handleMuteChange (and detachVideo()) call in
|
||||
// that case. It must not fire no matter how long it stays that way.
|
||||
w.setExpectingFrames(false, t0);
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(2000)));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(60000)));
|
||||
ST_ASSERT(!w.poll(t0 + std::chrono::milliseconds(600000)));
|
||||
|
||||
// Un-muting (setExpectingFrames(true, ...)) starts a BRAND NEW grace
|
||||
// period from that moment -- it must not read "was silent for ten
|
||||
// minutes" as "stalled for ten minutes" and fire immediately.
|
||||
const auto unmuted_at = t0 + std::chrono::milliseconds(600000);
|
||||
w.setExpectingFrames(true, unmuted_at);
|
||||
ST_ASSERT(!w.poll(unmuted_at + std::chrono::milliseconds(1999)));
|
||||
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.
|
||||
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(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.
|
||||
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(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(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(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
|
||||
// never exceeds max_backoff (30000ms) no matter how many attempts have
|
||||
// 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(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(t0 + std::chrono::milliseconds(at_ms));
|
||||
ST_ASSERT_EQ(w.attemptsThisStall(), 0);
|
||||
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)
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -330,11 +555,17 @@ int main()
|
||||
{
|
||||
testFrameGeometry();
|
||||
testTrackSelection();
|
||||
testSubscriptionSelection();
|
||||
testStateMachineHappyPath();
|
||||
testPublisherSwapIsNotAnError();
|
||||
testReconnect();
|
||||
testFailureAndRecovery();
|
||||
|
||||
testStallWatchdogFiresAfterThreshold();
|
||||
testStallWatchdogDoesNotFireWhenNotExpectingFrames();
|
||||
testStallWatchdogBacksOffRatherThanLooping();
|
||||
testStallWatchdogNeverWaitsLongerThanTheCeiling();
|
||||
|
||||
LiveKitSession::globalInitialize();
|
||||
testConnectRejectsIncompleteConfig();
|
||||
testConnectToUnreachableServerFailsCleanly();
|
||||
|
||||
@@ -375,6 +375,13 @@ void *sourceCreate(obs_data_t *settings, obs_source_t *source)
|
||||
|
||||
self->session->setVideoHandler([self](const VideoFrameData &frame) { outputVideoFrame(self, frame); });
|
||||
self->session->setAudioHandler([self](const AudioFrameData &frame) { outputAudioFrame(self, frame); });
|
||||
// The stall-recovery watchdog (core/src/session.cpp) is the only thing
|
||||
// that currently uses this: it logs each toggle-the-subscription
|
||||
// recovery attempt, and its eventual success, so a stalled-and-fixed
|
||||
// camera is diagnosable from an OBS log afterward instead of invisible.
|
||||
self->session->setDiagnosticHandler([](DiagnosticLevel level, const std::string &message) {
|
||||
obs_log(level == DiagnosticLevel::Warning ? LOG_WARNING : LOG_INFO, "%s", message.c_str());
|
||||
});
|
||||
self->session->setStateHandler([self](SessionState state, const std::string &detail) {
|
||||
self->setStatus(detail.empty() ? describeSessionState(state) : detail);
|
||||
self->status_is_error.store(state == SessionState::Failed);
|
||||
|
||||
@@ -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