/* streamer-tools OBS Camera Plugin - session wrapper tests Copyright (C) 2026 CyberCoveLLC Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 */ // What is and is not covered here, stated plainly because it matters: // // 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, // 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 // reaching the handlers. Those can only be verified against a real room, // and the design doc's Testing section puts that in the integration-test / // manual-sign-off bucket. #include #include #include #include #include #include #include "stplugin/session.h" #include "stplugin/session_types.h" #include "test_util.h" using namespace stplugin; namespace { // --------------------------------------------------------------------------- // Pure logic // --------------------------------------------------------------------------- void testFrameGeometry() { ST_ASSERT_EQ(planeCount(PixelFormat::I420), 3); ST_ASSERT_EQ(planeCount(PixelFormat::NV12), 2); ST_ASSERT_EQ(planeCount(PixelFormat::BGRA), 1); // 1280x720 I420: 921600 luma + 2 * 230400 chroma. ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::I420, 1280, 720), std::size_t(1382400)); ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::NV12, 1280, 720), std::size_t(1382400)); ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::BGRA, 1280, 720), std::size_t(3686400)); // Odd dimensions round the chroma planes up, the way libyuv does. ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::I420, 3, 3), std::size_t(9 + 2 * 4)); ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::I420, 1, 1), std::size_t(1 + 2)); // Degenerate geometry is 0, which the reader treats as "drop the frame". ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::I420, 0, 720), std::size_t(0)); ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::I420, 1280, 0), std::size_t(0)); ST_ASSERT_EQ(expectedFrameBytes(PixelFormat::I420, -1, -1), std::size_t(0)); ST_ASSERT_EQ(std::string(describePixelFormat(PixelFormat::I420)), std::string("I420")); } void testTrackSelection() { const std::string want = "cam1"; // The camera we asked for. ST_ASSERT(isWantedVideoTrack(want, "cam1", MediaKind::Video, MediaSource::Camera)); // A video track with no declared source is taken on kind alone. ST_ASSERT(isWantedVideoTrack(want, "cam1", MediaKind::Video, MediaSource::Unknown)); // Someone else's camera. ST_ASSERT(!isWantedVideoTrack(want, "cam2", MediaKind::Video, MediaSource::Camera)); // The right participant's screenshare is explicitly NOT the camera -- // streamer-tools publishes those as separate sources. ST_ASSERT(!isWantedVideoTrack(want, "cam1", MediaKind::Video, MediaSource::Screenshare)); // Their microphone is not a video track. ST_ASSERT(!isWantedVideoTrack(want, "cam1", MediaKind::Audio, MediaSource::Microphone)); // No selection means nothing matches -- never "the first thing we see". ST_ASSERT(!isWantedVideoTrack("", "cam1", MediaKind::Video, MediaSource::Camera)); ST_ASSERT(!isWantedVideoTrack("", "", MediaKind::Video, MediaSource::Camera)); ST_ASSERT(isWantedAudioTrack(want, "cam1", MediaKind::Audio, MediaSource::Microphone)); ST_ASSERT(isWantedAudioTrack(want, "cam1", MediaKind::Audio, MediaSource::Unknown)); ST_ASSERT(!isWantedAudioTrack(want, "cam1", MediaKind::Audio, MediaSource::ScreenshareAudio)); ST_ASSERT(!isWantedAudioTrack(want, "cam1", MediaKind::Video, MediaSource::Camera)); ST_ASSERT(!isWantedAudioTrack(want, "other", MediaKind::Audio, MediaSource::Microphone)); } void testStateMachineHappyPath() { SessionStateMachine m; ST_ASSERT(m.state() == SessionState::Idle); ST_ASSERT(!m.hasVideo()); ST_ASSERT(!m.waitingForCamera()); m.onConnectRequested(); ST_ASSERT(m.state() == SessionState::Connecting); // Connecting is not "waiting for camera": the placeholder belongs to a // live connection with a dark slot, not to a connection in progress. ST_ASSERT(!m.waitingForCamera()); m.onConnectSucceeded(); ST_ASSERT(m.state() == SessionState::Connected); ST_ASSERT(m.waitingForCamera()); m.onVideoAttached(); ST_ASSERT(m.hasVideo()); ST_ASSERT(!m.waitingForCamera()); m.onAudioAttached(); ST_ASSERT(m.hasAudio()); m.onLocalDisconnect(); ST_ASSERT(m.state() == SessionState::Disconnected); ST_ASSERT(!m.hasVideo()); ST_ASSERT(!m.hasAudio()); } void testPublisherSwapIsNotAnError() { // The motivating bug: a slot's publisher restarts mid-show. That must // read as "waiting for camera", never as a failure, and the connection // state must not move at all. SessionStateMachine m; m.onConnectRequested(); m.onConnectSucceeded(); m.onVideoAttached(); m.onVideoDetached(); ST_ASSERT(m.state() == SessionState::Connected); ST_ASSERT(!m.hasVideo()); ST_ASSERT(m.waitingForCamera()); ST_ASSERT(m.detail().empty()); m.onVideoAttached(); ST_ASSERT(m.state() == SessionState::Connected); ST_ASSERT(m.hasVideo()); ST_ASSERT(!m.waitingForCamera()); } void testReconnect() { SessionStateMachine m; m.onConnectRequested(); m.onConnectSucceeded(); m.onVideoAttached(); m.onReconnecting(); ST_ASSERT(m.state() == SessionState::Reconnecting); // Tracks are re-subscribed on the far side, so video is not live yet. ST_ASSERT(!m.hasVideo()); ST_ASSERT(m.waitingForCamera()); ST_ASSERT_EQ(m.detail(), std::string("reconnecting")); m.onReconnected(); ST_ASSERT(m.state() == SessionState::Connected); ST_ASSERT(m.detail().empty()); // A stray reconnect notification after a hard failure must not resurrect // the session. SessionStateMachine dead; dead.onConnectRequested(); dead.onConnectFailed("token rejected"); dead.onReconnecting(); ST_ASSERT(dead.state() == SessionState::Failed); dead.onReconnected(); ST_ASSERT(dead.state() == SessionState::Failed); } void testFailureAndRecovery() { SessionStateMachine m; m.onConnectRequested(); m.onConnectFailed("token rejected"); ST_ASSERT(m.state() == SessionState::Failed); ST_ASSERT_EQ(m.detail(), std::string("token rejected")); ST_ASSERT(!m.waitingForCamera()); // A fresh attempt clears the stale reason, so a healthy connection can // never be shown next to the previous failure's message. m.onConnectRequested(); ST_ASSERT(m.detail().empty()); m.onConnectSucceeded(); ST_ASSERT(m.state() == SessionState::Connected); ST_ASSERT(m.detail().empty()); // A fatal room end (duplicate identity, token rejected) is Failed; an // ordinary drop is Disconnected. SessionStateMachine fatal; fatal.onConnectRequested(); fatal.onConnectSucceeded(); fatal.onRoomEnded("another client joined with the same identity", true); ST_ASSERT(fatal.state() == SessionState::Failed); SessionStateMachine dropped; dropped.onConnectRequested(); dropped.onConnectSucceeded(); dropped.onRoomEnded("the signalling connection closed", false); ST_ASSERT(dropped.state() == SessionState::Disconnected); // Room-ended events after we are already down are ignored, so a late // event cannot overwrite the reason the operator needs to see. SessionStateMachine idle; idle.onRoomEnded("stray", true); 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))); } void testStallWatchdogBacksOffRatherThanLooping() { const auto t0 = std::chrono::steady_clock::now(); StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000)); w.setExpectingFrames(true, t0); // TEMPORARY DEBUG -- remove before merge. auto rel = [t0](std::chrono::steady_clock::time_point tp) { return (long long)std::chrono::duration_cast(tp - t0).count(); }; auto dump = [&](const char *tag, std::chrono::steady_clock::time_point n) { std::fprintf(stderr, "WDTEST %-16s now_rel_ns=%lld next_rel_ns=%lld base_rel_ns=%lld " "backoff_ms=%lld attempts=%d\n", tag, rel(n), rel(w.dbgNextAllowed()), rel(w.dbgBaseline()), w.dbgBackoffMs(), w.attemptsThisStall()); }; std::fprintf(stderr, "WDTEST t0_ns=%lld timeout_ms=%lld max_ms=%lld\n", (long long)std::chrono::duration_cast(t0.time_since_epoch()).count(), w.dbgTimeoutMs(), w.dbgMaxBackoffMs()); dump("init", t0); // First attempt at the threshold. auto now = t0 + std::chrono::milliseconds(2000); dump("pre-attempt1", now); ST_ASSERT(w.poll(now)); dump("post-attempt1", now); ST_ASSERT_EQ(w.attemptsThisStall(), 1); // A genuinely gone publisher: no frame ever comes back. Immediately // asking again (the naive "retry every poll interval forever" a // watchdog without backoff would do) must NOT fire -- that is precisely // the "hammered every 2 seconds forever" this backoff exists to avoid. ST_ASSERT(!w.poll(now + std::chrono::milliseconds(250))); ST_ASSERT(!w.poll(now + std::chrono::milliseconds(1999))); // The backoff after attempt 1 is the base timeout (2000ms): the second // attempt is allowed at +2000ms from the first, not before. now += std::chrono::milliseconds(2000); dump("pre-attempt2", now); ST_ASSERT(w.poll(now)); dump("post-attempt2", now); ST_ASSERT_EQ(w.attemptsThisStall(), 2); // Backoff doubles: the third attempt needs a 4000ms gap, not 2000ms. ST_ASSERT(!w.poll(now + std::chrono::milliseconds(3999))); now += std::chrono::milliseconds(4000); dump("pre-attempt3", now); ST_ASSERT(w.poll(now)); dump("post-attempt3", now); ST_ASSERT_EQ(w.attemptsThisStall(), 3); // ... and again to 8000ms, and again to 16000ms. ST_ASSERT(!w.poll(now + std::chrono::milliseconds(7999))); now += std::chrono::milliseconds(8000); dump("pre-attempt4", now); ST_ASSERT(w.poll(now)); dump("post-attempt4", now); ST_ASSERT_EQ(w.attemptsThisStall(), 4); ST_ASSERT(!w.poll(now + std::chrono::milliseconds(15999))); now += std::chrono::milliseconds(16000); dump("pre-attempt5", now); ST_ASSERT(w.poll(now)); dump("post-attempt5", now); 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) { std::fprintf(stderr, "WDTEST ---- loop iteration %d ----\n", i); dump("loop-pre-29999", now + std::chrono::milliseconds(29999)); ST_ASSERT(!w.poll(now + std::chrono::milliseconds(29999))); now += std::chrono::milliseconds(30000); dump("loop-pre-fire", now); ST_ASSERT(w.poll(now)); dump("loop-post-fire", now); } // A frame finally arrives: the stall is over, and the NEXT one (a fresh // stall, not a continuation) starts back at the base cadence rather // than staying parked at the 30s ceiling forever. w.onFrameDelivered(now); ST_ASSERT_EQ(w.attemptsThisStall(), 0); ST_ASSERT(!w.poll(now + std::chrono::milliseconds(1999))); ST_ASSERT(w.poll(now + std::chrono::milliseconds(2000))); ST_ASSERT_EQ(w.attemptsThisStall(), 1); } // TEMPORARY DEBUG -- remove before merge. Same sequence as the test above, // but every number is RECORDED into plain locals and printed only after the // loop finishes, so the loop body stays as close to the original as // possible (an fprintf inside poll() made this pass on Windows, so the // observation itself perturbs the thing being observed). void testStallWatchdogTrace() { struct Rec { long long now_ns; long long next_ns; long long base_ns; long long backoff_ms; int attempts; int early; // result of poll(now + 29999) -- expected 0 int fire; // result of poll(now) -- expected 1 }; Rec recs[16]; int n = 0; const auto t0 = std::chrono::steady_clock::now(); StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000)); w.setExpectingFrames(true, t0); auto rel = [t0](std::chrono::steady_clock::time_point tp) { return (long long)std::chrono::duration_cast(tp - t0).count(); }; auto rec = [&](std::chrono::steady_clock::time_point now, int early, int fire) { if (n < 16) { recs[n].now_ns = rel(now); recs[n].next_ns = rel(w.dbgNextAllowed()); recs[n].base_ns = rel(w.dbgBaseline()); recs[n].backoff_ms = w.dbgBackoffMs(); recs[n].attempts = w.attemptsThisStall(); recs[n].early = early; recs[n].fire = fire; ++n; } }; auto now = t0 + std::chrono::milliseconds(2000); int f = w.poll(now) ? 1 : 0; rec(now, -1, f); const int gaps[] = {2000, 4000, 8000, 16000}; for (int i = 0; i < 4; ++i) { now += std::chrono::milliseconds(gaps[i]); f = w.poll(now) ? 1 : 0; rec(now, -1, f); } for (int i = 0; i < 6; ++i) { const int e = w.poll(now + std::chrono::milliseconds(29999)) ? 1 : 0; now += std::chrono::milliseconds(30000); f = w.poll(now) ? 1 : 0; rec(now, e, f); } std::fprintf(stderr, "WDTRACE t0_ns=%lld timeout_ms=%lld max_ms=%lld tick_num=%lld tick_den=%lld rep_bytes=%d\n", (long long)std::chrono::duration_cast(t0.time_since_epoch()).count(), w.dbgTimeoutMs(), w.dbgMaxBackoffMs(), (long long)std::chrono::steady_clock::period::num, (long long)std::chrono::steady_clock::period::den, (int)sizeof(std::chrono::steady_clock::rep)); for (int i = 0; i < n; ++i) { std::fprintf(stderr, "WDTRACE step=%2d now_rel_ns=%lld next_rel_ns=%lld base_rel_ns=%lld backoff_ms=%lld attempts=%d early=%d fire=%d\n", i, recs[i].now_ns, recs[i].next_ns, recs[i].base_ns, recs[i].backoff_ms, recs[i].attempts, recs[i].early, recs[i].fire); } } // TEMPORARY DEBUG -- remove before merge. // // A non-inlined wrapper, so the time point can be read on the CALLEE side: // whatever `received_ns` holds is what actually crossed the call boundary, // independent of what the caller's abstract-machine value was. #if defined(_MSC_VER) #define ST_NOINLINE __declspec(noinline) #else #define ST_NOINLINE __attribute__((noinline)) #endif ST_NOINLINE bool st_dbg_poll(StallWatchdog &w, std::chrono::steady_clock::time_point n, long long *received_ns) { *received_ns = (long long)n.time_since_epoch().count(); return w.poll(n); } // Shapes, each run twice: // mode 0 accumulate `now += 30000ms`, pass `now` (known-bad) // mode 1 accumulate, pass a named copy of `now` (known-good) // mode 2 accumulate, sum the ns the CALLER holds // mode 3 accumulate, through st_dbg_poll (callee-side capture) void testStallWatchdogQuiet() { struct Trial { int mode; int early[6]; int fire[6]; long long fire_arg_sum_ms; long long recv_rel_ms[6]; long long t0_ns; long long next_rel_ns; int attempts_after; }; const int kTrials = 8; Trial trials[kTrials]; for (int t = 0; t < kTrials; ++t) { const int mode = t % 4; const auto t0 = std::chrono::steady_clock::now(); StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000)); w.setExpectingFrames(true, t0); int early[6] = {0}; int fire[6] = {0}; long long fsum = 0; long long recv[6] = {0, 0, 0, 0, 0, 0}; auto now = t0 + std::chrono::milliseconds(2000); (void)w.poll(now); now += std::chrono::milliseconds(2000); (void)w.poll(now); now += std::chrono::milliseconds(4000); (void)w.poll(now); now += std::chrono::milliseconds(8000); (void)w.poll(now); now += std::chrono::milliseconds(16000); (void)w.poll(now); if (mode == 0) { for (int i = 0; i < 6; ++i) { early[i] = w.poll(now + std::chrono::milliseconds(29999)) ? 1 : 0; now += std::chrono::milliseconds(30000); fire[i] = w.poll(now) ? 1 : 0; } } else if (mode == 1) { for (int i = 0; i < 6; ++i) { early[i] = w.poll(now + std::chrono::milliseconds(29999)) ? 1 : 0; now += std::chrono::milliseconds(30000); const auto fire_arg = now; fire[i] = w.poll(fire_arg) ? 1 : 0; } } else if (mode == 2) { for (int i = 0; i < 6; ++i) { early[i] = w.poll(now + std::chrono::milliseconds(29999)) ? 1 : 0; now += std::chrono::milliseconds(30000); fsum += std::chrono::duration_cast(now - t0).count(); fire[i] = w.poll(now) ? 1 : 0; } } else { for (int i = 0; i < 6; ++i) { long long junk = 0; early[i] = st_dbg_poll(w, now + std::chrono::milliseconds(29999), &junk) ? 1 : 0; now += std::chrono::milliseconds(30000); fire[i] = st_dbg_poll(w, now, &recv[i]) ? 1 : 0; } } trials[t].mode = mode; for (int i = 0; i < 6; ++i) { trials[t].early[i] = early[i]; trials[t].fire[i] = fire[i]; trials[t].recv_rel_ms[i] = recv[i] ? (recv[i] - (long long)t0.time_since_epoch().count()) / 1000000LL : -1; } trials[t].fire_arg_sum_ms = fsum; trials[t].t0_ns = (long long)t0.time_since_epoch().count(); trials[t].next_rel_ns = (long long)std::chrono::duration_cast(w.dbgNextAllowed() - t0).count(); trials[t].attempts_after = w.attemptsThisStall(); } for (int t = 0; t < kTrials; ++t) { std::fprintf(stderr, "WDQUIET trial=%d mode=%d after=%d next_rel_ns=%lld fsum=%lld(want 822000) recv_ms=%lld,%lld,%lld,%lld,%lld,%lld early=%d%d%d%d%d%d fire=%d%d%d%d%d%d\n", t, trials[t].mode, trials[t].attempts_after, trials[t].next_rel_ns, trials[t].fire_arg_sum_ms, trials[t].recv_rel_ms[0], trials[t].recv_rel_ms[1], trials[t].recv_rel_ms[2], trials[t].recv_rel_ms[3], trials[t].recv_rel_ms[4], trials[t].recv_rel_ms[5], trials[t].early[0], trials[t].early[1], trials[t].early[2], trials[t].early[3], trials[t].early[4], trials[t].early[5], trials[t].fire[0], trials[t].fire[1], trials[t].fire[2], trials[t].fire[3], trials[t].fire[4], trials[t].fire[5]); } for (int t = 0; t < kTrials; ++t) { for (int i = 0; i < 6; ++i) { ST_ASSERT(trials[t].early[i] == 0); ST_ASSERT(trials[t].fire[i] == 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); auto now = t0 + std::chrono::milliseconds(1000); ST_ASSERT(w.poll(now)); // 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(now + std::chrono::milliseconds(gap - 1))); now += std::chrono::milliseconds(gap); ST_ASSERT(w.poll(now)); ST_ASSERT_EQ(w.attemptsThisStall(), ++attempt); } } // --------------------------------------------------------------------------- // Real SDK, failure paths only (no LiveKit server available headlessly) // --------------------------------------------------------------------------- void testConnectRejectsIncompleteConfig() { LiveKitSession session; std::atomic state_calls{0}; session.setStateHandler([&](SessionState, const std::string &) { state_calls.fetch_add(1); }); SessionConfig config; config.ws_url = ""; config.token = "t"; config.participant_identity = "cam1"; ST_ASSERT(!session.connect(config)); ST_ASSERT(session.state() == SessionState::Failed); ST_ASSERT(!session.stateDetail().empty()); config.ws_url = "ws://127.0.0.1:1"; config.token = ""; ST_ASSERT(!session.connect(config)); ST_ASSERT(session.state() == SessionState::Failed); config.token = "t"; config.participant_identity = ""; ST_ASSERT(!session.connect(config)); ST_ASSERT(session.state() == SessionState::Failed); // The state handler fired for each attempt (Connecting + Failed). ST_ASSERT(state_calls.load() >= 6); // Frame counters stay at zero and nothing crashes on teardown. ST_ASSERT_EQ(session.videoFrameCount(), std::uint64_t(0)); ST_ASSERT_EQ(session.audioFrameCount(), std::uint64_t(0)); session.disconnect(); session.disconnect(); // idempotent ST_ASSERT(session.state() == SessionState::Failed || session.state() == SessionState::Disconnected); } void testConnectToUnreachableServerFailsCleanly() { // Port 1 on loopback: nothing is listening, and the connection is // refused immediately rather than hanging. This exercises the real // livekit::Room::connect() failure path, with a real (garbage) token. LiveKitSession session; std::atomic video_frames{0}; session.setVideoHandler([&](const VideoFrameData &) { video_frames.fetch_add(1); }); SessionConfig config; config.ws_url = "ws://127.0.0.1:1"; config.token = "not.a.real.token"; config.participant_identity = "cam1"; config.connect_timeout_ms = 3000; const auto start = std::chrono::steady_clock::now(); const bool ok = session.connect(config); const auto elapsed = std::chrono::steady_clock::now() - start; ST_ASSERT(!ok); ST_ASSERT(session.state() == SessionState::Failed); ST_ASSERT(!session.stateDetail().empty()); ST_ASSERT_EQ(video_frames.load(), 0); // Must not sit on the caller's thread indefinitely -- this runs on an OBS // thread in the real adapter. ST_ASSERT(std::chrono::duration_cast(elapsed).count() < 60); session.disconnect(); } void testConnectToNonLiveKitServerFailsCleanly() { // A URL that resolves and connects but is not a LiveKit signalling // endpoint. The realistic operator mistake: pasting the app URL. LiveKitSession session; SessionConfig config; config.ws_url = "ws://127.0.0.1:1/rtc"; config.token = "eyJhbGciOiJIUzI1NiJ9.bm90YXRva2Vu.x"; config.participant_identity = "cam1"; config.connect_timeout_ms = 3000; ST_ASSERT(!session.connect(config)); ST_ASSERT(session.state() == SessionState::Failed); session.disconnect(); } void testDestroyWithoutDisconnect() { // The OBS adapter destroys sources without necessarily having called // disconnect() first (an OBS shutdown mid-connect, say). The destructor // must join every thread it started rather than terminating. { LiveKitSession session; SessionConfig config; config.ws_url = "ws://127.0.0.1:1"; config.token = "t"; config.participant_identity = "cam1"; config.connect_timeout_ms = 2000; (void)session.connect(config); } ST_ASSERT(true); // reaching here at all is the assertion } void testGlobalInitIsReferenceCounted() { // Several OBS sources may each hold the SDK open; the last one out turns // the lights off, and an unbalanced extra shutdown must not underflow. LiveKitSession::globalInitialize(); LiveKitSession::globalInitialize(); LiveKitSession::globalShutdown(); LiveKitSession::globalShutdown(); LiveKitSession::globalShutdown(); // extra, must be harmless LiveKitSession::globalInitialize(); LiveKitSession::globalShutdown(); ST_ASSERT(true); } } // namespace int main() { testFrameGeometry(); testTrackSelection(); testStateMachineHappyPath(); testPublisherSwapIsNotAnError(); testReconnect(); testFailureAndRecovery(); testStallWatchdogFiresAfterThreshold(); testStallWatchdogDoesNotFireWhenNotExpectingFrames(); testStallWatchdogBacksOffRatherThanLooping(); testStallWatchdogTrace(); testStallWatchdogQuiet(); testStallWatchdogNeverWaitsLongerThanTheCeiling(); LiveKitSession::globalInitialize(); testConnectRejectsIncompleteConfig(); testConnectToUnreachableServerFailsCleanly(); testConnectToNonLiveKitServerFailsCleanly(); testDestroyWithoutDisconnect(); LiveKitSession::globalShutdown(); testGlobalInitIsReferenceCounted(); return st_test_report("session"); }