Round 1's in-poll fprintf made Windows pass, so poll() is back to its pristine form and the numbers are recorded into locals and printed after the fact. The quiet copy runs the same sequence with no observation at all, to tell codegen apart from data. TEMPORARY. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
660 lines
26 KiB
C++
660 lines
26 KiB
C++
/*
|
|
streamer-tools OBS Camera Plugin - session wrapper tests
|
|
Copyright (C) 2026 CyberCoveLLC <jknapp85@gmail.com>
|
|
|
|
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 <atomic>
|
|
#include <cstdio>
|
|
#include <cstdio>
|
|
#include <chrono>
|
|
#include <string>
|
|
#include <thread>
|
|
|
|
#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<std::chrono::nanoseconds>(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<std::chrono::nanoseconds>(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<std::chrono::nanoseconds>(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<std::chrono::nanoseconds>(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 byte-for-byte quiet copy of the
|
|
// loop from testStallWatchdogBacksOffRatherThanLooping: no accessor calls,
|
|
// no printing inside the loop. Results are collected into an int array and
|
|
// reported afterwards, so if the loud trace above passes while this one
|
|
// fails, the observation is the cure and the fault is codegen, not data.
|
|
void testStallWatchdogQuiet()
|
|
{
|
|
int early[6] = {0};
|
|
int fire[6] = {0};
|
|
|
|
const auto t0 = std::chrono::steady_clock::now();
|
|
StallWatchdog w(std::chrono::milliseconds(2000), std::chrono::milliseconds(30000));
|
|
w.setExpectingFrames(true, t0);
|
|
|
|
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);
|
|
const int attempts_before_loop = w.attemptsThisStall();
|
|
|
|
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;
|
|
}
|
|
|
|
std::fprintf(stderr, "WDQUIET attempts_before_loop=%d\n", attempts_before_loop);
|
|
for (int i = 0; i < 6; ++i)
|
|
std::fprintf(stderr, "WDQUIET i=%d early=%d(want 0) fire=%d(want 1)\n", i, early[i], fire[i]);
|
|
for (int i = 0; i < 6; ++i) {
|
|
ST_ASSERT(early[i] == 0);
|
|
ST_ASSERT(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<int> 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<int> 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<std::chrono::seconds>(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");
|
|
}
|