2026-09-06 21:39:49 -07:00
|
|
|
/*
|
|
|
|
|
streamer-tools OBS Camera Plugin - LiveKit session wrapper
|
|
|
|
|
Copyright (C) 2026 CyberCoveLLC <jknapp85@gmail.com>
|
|
|
|
|
|
2026-09-07 04:44:16 -07:00
|
|
|
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
|
2026-09-06 21:39:49 -07:00
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
#include "stplugin/session.h"
|
|
|
|
|
|
|
|
|
|
#include <atomic>
|
|
|
|
|
#include <chrono>
|
|
|
|
|
#include <condition_variable>
|
|
|
|
|
#include <deque>
|
|
|
|
|
#include <exception>
|
|
|
|
|
#include <mutex>
|
|
|
|
|
#include <thread>
|
|
|
|
|
#include <utility>
|
|
|
|
|
#include <vector>
|
|
|
|
|
|
|
|
|
|
#include <livekit/audio_frame.h>
|
|
|
|
|
#include <livekit/audio_stream.h>
|
|
|
|
|
#include <livekit/livekit.h>
|
|
|
|
|
#include <livekit/remote_participant.h>
|
|
|
|
|
#include <livekit/remote_track_publication.h>
|
|
|
|
|
#include <livekit/room.h>
|
|
|
|
|
#include <livekit/room_delegate.h>
|
|
|
|
|
#include <livekit/room_event_types.h>
|
|
|
|
|
#include <livekit/track.h>
|
|
|
|
|
#include <livekit/video_frame.h>
|
|
|
|
|
#include <livekit/video_stream.h>
|
|
|
|
|
|
|
|
|
|
namespace stplugin {
|
|
|
|
|
|
|
|
|
|
namespace {
|
|
|
|
|
|
|
|
|
|
// --- LiveKit <-> plugin type conversion ------------------------------------
|
|
|
|
|
|
|
|
|
|
MediaKind toMediaKind(livekit::TrackKind kind)
|
|
|
|
|
{
|
|
|
|
|
switch (kind) {
|
|
|
|
|
case livekit::TrackKind::KIND_AUDIO: return MediaKind::Audio;
|
|
|
|
|
case livekit::TrackKind::KIND_VIDEO: return MediaKind::Video;
|
|
|
|
|
case livekit::TrackKind::KIND_UNKNOWN: break;
|
|
|
|
|
}
|
|
|
|
|
return MediaKind::Unknown;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
MediaSource toMediaSource(livekit::TrackSource source)
|
|
|
|
|
{
|
|
|
|
|
switch (source) {
|
|
|
|
|
case livekit::TrackSource::SOURCE_CAMERA: return MediaSource::Camera;
|
|
|
|
|
case livekit::TrackSource::SOURCE_MICROPHONE: return MediaSource::Microphone;
|
|
|
|
|
case livekit::TrackSource::SOURCE_SCREENSHARE: return MediaSource::Screenshare;
|
|
|
|
|
case livekit::TrackSource::SOURCE_SCREENSHARE_AUDIO: return MediaSource::ScreenshareAudio;
|
|
|
|
|
case livekit::TrackSource::SOURCE_UNKNOWN: break;
|
|
|
|
|
}
|
|
|
|
|
return MediaSource::Unknown;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
livekit::VideoBufferType toLiveKitBufferType(PixelFormat format)
|
|
|
|
|
{
|
|
|
|
|
switch (format) {
|
|
|
|
|
case PixelFormat::I420: return livekit::VideoBufferType::I420;
|
|
|
|
|
case PixelFormat::NV12: return livekit::VideoBufferType::NV12;
|
|
|
|
|
case PixelFormat::BGRA: return livekit::VideoBufferType::BGRA;
|
|
|
|
|
}
|
|
|
|
|
return livekit::VideoBufferType::I420;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Returns false when the SDK handed us a format the OBS adapter cannot
|
|
|
|
|
/// consume, in which case the caller converts.
|
|
|
|
|
bool fromLiveKitBufferType(livekit::VideoBufferType type, PixelFormat &out)
|
|
|
|
|
{
|
|
|
|
|
switch (type) {
|
|
|
|
|
case livekit::VideoBufferType::I420: out = PixelFormat::I420; return true;
|
|
|
|
|
case livekit::VideoBufferType::NV12: out = PixelFormat::NV12; return true;
|
|
|
|
|
case livekit::VideoBufferType::BGRA: out = PixelFormat::BGRA; return true;
|
|
|
|
|
default: return false;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Which disconnect reasons are worth telling the operator "this will not fix
|
|
|
|
|
/// itself" about. Everything else is reported as an ordinary disconnect,
|
|
|
|
|
/// because the SDK's own reconnect logic covers it.
|
|
|
|
|
bool isFatalDisconnect(livekit::DisconnectReason reason)
|
|
|
|
|
{
|
|
|
|
|
switch (reason) {
|
|
|
|
|
case livekit::DisconnectReason::DuplicateIdentity:
|
|
|
|
|
case livekit::DisconnectReason::ParticipantRemoved:
|
|
|
|
|
case livekit::DisconnectReason::RoomDeleted:
|
|
|
|
|
case livekit::DisconnectReason::JoinFailure:
|
|
|
|
|
case livekit::DisconnectReason::UserRejected:
|
|
|
|
|
return true;
|
|
|
|
|
default:
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const char *describeDisconnectReason(livekit::DisconnectReason reason)
|
|
|
|
|
{
|
|
|
|
|
switch (reason) {
|
|
|
|
|
case livekit::DisconnectReason::Unknown: return "connection lost";
|
|
|
|
|
case livekit::DisconnectReason::ClientInitiated: return "disconnected";
|
|
|
|
|
case livekit::DisconnectReason::DuplicateIdentity: return "another client joined with the same identity";
|
|
|
|
|
case livekit::DisconnectReason::ServerShutdown: return "the LiveKit server is shutting down";
|
|
|
|
|
case livekit::DisconnectReason::ParticipantRemoved: return "removed from the room";
|
|
|
|
|
case livekit::DisconnectReason::RoomDeleted: return "the room was deleted";
|
|
|
|
|
case livekit::DisconnectReason::StateMismatch: return "session could not be resumed";
|
|
|
|
|
case livekit::DisconnectReason::JoinFailure: return "could not join the room (token rejected or expired?)";
|
|
|
|
|
case livekit::DisconnectReason::Migration: return "migrating to another server";
|
|
|
|
|
case livekit::DisconnectReason::SignalClose: return "the signalling connection closed";
|
|
|
|
|
case livekit::DisconnectReason::RoomClosed: return "the room closed";
|
|
|
|
|
case livekit::DisconnectReason::UserUnavailable: return "user unavailable";
|
|
|
|
|
case livekit::DisconnectReason::UserRejected: return "connection rejected";
|
|
|
|
|
case livekit::DisconnectReason::SipTrunkFailure: return "SIP trunk failure";
|
|
|
|
|
case livekit::DisconnectReason::ConnectionTimeout: return "connection timed out";
|
|
|
|
|
case livekit::DisconnectReason::MediaFailure: return "media connection failed";
|
|
|
|
|
case livekit::DisconnectReason::AgentError: return "agent error";
|
|
|
|
|
}
|
|
|
|
|
return "disconnected";
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// --- Process-wide SDK lifetime ---------------------------------------------
|
|
|
|
|
|
|
|
|
|
std::mutex &globalMutex()
|
|
|
|
|
{
|
|
|
|
|
static std::mutex m;
|
|
|
|
|
return m;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
int &globalRefCount()
|
|
|
|
|
{
|
|
|
|
|
static int n = 0;
|
|
|
|
|
return n;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
} // namespace
|
|
|
|
|
|
|
|
|
|
// ---------------------------------------------------------------------------
|
|
|
|
|
// Impl
|
|
|
|
|
// ---------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
struct LiveKitSession::Impl : public livekit::RoomDelegate {
|
|
|
|
|
enum class CommandType { AttachVideo, DetachVideo, AttachAudio, DetachAudio, Stop };
|
|
|
|
|
|
|
|
|
|
struct Command {
|
|
|
|
|
CommandType type;
|
|
|
|
|
std::shared_ptr<livekit::Track> track;
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
livekit::Room room;
|
|
|
|
|
SessionConfig config;
|
|
|
|
|
|
|
|
|
|
mutable std::mutex state_mutex;
|
|
|
|
|
SessionStateMachine machine;
|
|
|
|
|
|
|
|
|
|
VideoFrameHandler on_video;
|
|
|
|
|
AudioFrameHandler on_audio;
|
|
|
|
|
SessionStateHandler on_state;
|
|
|
|
|
|
|
|
|
|
std::atomic<std::uint64_t> video_frames{0};
|
|
|
|
|
std::atomic<std::uint64_t> audio_frames{0};
|
|
|
|
|
std::atomic<std::uint64_t> dropped_frames{0};
|
|
|
|
|
|
|
|
|
|
// Command queue. Every interaction with livekit::VideoStream /
|
|
|
|
|
// 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).
|
|
|
|
|
std::mutex queue_mutex;
|
|
|
|
|
std::condition_variable queue_cv;
|
|
|
|
|
std::deque<Command> queue;
|
|
|
|
|
std::thread worker;
|
|
|
|
|
bool worker_running = false;
|
|
|
|
|
|
|
|
|
|
// Owned exclusively by the worker thread.
|
|
|
|
|
std::shared_ptr<livekit::VideoStream> video_stream;
|
|
|
|
|
std::thread video_thread;
|
|
|
|
|
std::shared_ptr<livekit::AudioStream> audio_stream;
|
|
|
|
|
std::thread audio_thread;
|
|
|
|
|
|
|
|
|
|
bool connected = false;
|
|
|
|
|
|
|
|
|
|
~Impl() override = default;
|
|
|
|
|
|
|
|
|
|
// --- state helpers -----------------------------------------------------
|
|
|
|
|
|
|
|
|
|
template<typename Fn> void mutateState(Fn &&fn)
|
|
|
|
|
{
|
|
|
|
|
SessionState state;
|
|
|
|
|
std::string detail;
|
|
|
|
|
SessionStateHandler handler;
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(state_mutex);
|
|
|
|
|
fn(machine);
|
|
|
|
|
state = machine.state();
|
|
|
|
|
detail = machine.detail();
|
|
|
|
|
handler = on_state;
|
|
|
|
|
}
|
|
|
|
|
// Notified outside the lock: the handler is OBS adapter code and must
|
|
|
|
|
// never be able to deadlock against a concurrent state query.
|
|
|
|
|
if (handler)
|
|
|
|
|
handler(state, detail);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void post(CommandType type, std::shared_ptr<livekit::Track> track = nullptr)
|
|
|
|
|
{
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(queue_mutex);
|
|
|
|
|
if (!worker_running)
|
|
|
|
|
return;
|
|
|
|
|
queue.push_back(Command{type, std::move(track)});
|
|
|
|
|
}
|
|
|
|
|
queue_cv.notify_one();
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-07 10:43:54 -07:00
|
|
|
// 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:
|
|
|
|
|
//
|
|
|
|
|
// - 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).
|
|
|
|
|
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);
|
|
|
|
|
} catch (const std::exception &) {
|
|
|
|
|
// Best-effort: worst case this track keeps whatever quality
|
|
|
|
|
// it already had, which is the pre-existing behaviour.
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
post(CommandType::AttachVideo, track);
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-06 21:39:49 -07:00
|
|
|
// --- RoomDelegate ------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
void onTrackSubscribed(livekit::Room &, const livekit::TrackSubscribedEvent &event) override
|
|
|
|
|
{
|
|
|
|
|
if (!event.participant || !event.track)
|
|
|
|
|
return;
|
|
|
|
|
const std::string identity = event.participant->identity();
|
|
|
|
|
const MediaKind kind = toMediaKind(event.track->kind());
|
|
|
|
|
const MediaSource source =
|
|
|
|
|
event.publication ? toMediaSource(event.publication->source()) : MediaSource::Unknown;
|
|
|
|
|
|
|
|
|
|
if (isWantedVideoTrack(config.participant_identity, identity, kind, source))
|
2026-09-07 10:43:54 -07:00
|
|
|
handleWantedVideoTrack(event.track, event.publication);
|
2026-09-06 21:39:49 -07:00
|
|
|
else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source))
|
|
|
|
|
post(CommandType::AttachAudio, event.track);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void onTrackUnsubscribed(livekit::Room &, const livekit::TrackUnsubscribedEvent &event) override
|
|
|
|
|
{
|
|
|
|
|
if (!event.participant || !event.track)
|
|
|
|
|
return;
|
|
|
|
|
if (event.participant->identity() != config.participant_identity)
|
|
|
|
|
return;
|
|
|
|
|
const MediaKind kind = toMediaKind(event.track->kind());
|
|
|
|
|
if (kind == MediaKind::Video)
|
|
|
|
|
post(CommandType::DetachVideo);
|
|
|
|
|
else if (kind == MediaKind::Audio)
|
|
|
|
|
post(CommandType::DetachAudio);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void onParticipantDisconnected(livekit::Room &, const livekit::ParticipantDisconnectedEvent &event) override
|
|
|
|
|
{
|
|
|
|
|
if (!event.participant || event.participant->identity() != config.participant_identity)
|
|
|
|
|
return;
|
|
|
|
|
// The slot went away entirely. This is the placeholder state, not an
|
|
|
|
|
// error: the operator's room is fine, the camera just left.
|
|
|
|
|
post(CommandType::DetachVideo);
|
|
|
|
|
post(CommandType::DetachAudio);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void onReconnecting(livekit::Room &, const livekit::ReconnectingEvent &) override
|
|
|
|
|
{
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onReconnecting(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void onReconnected(livekit::Room &, const livekit::ReconnectedEvent &) override
|
|
|
|
|
{
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onReconnected(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void onDisconnected(livekit::Room &, const livekit::DisconnectedEvent &event) override
|
|
|
|
|
{
|
|
|
|
|
const std::string reason = describeDisconnectReason(event.reason);
|
|
|
|
|
const bool fatal = isFatalDisconnect(event.reason);
|
|
|
|
|
post(CommandType::DetachVideo);
|
|
|
|
|
post(CommandType::DetachAudio);
|
|
|
|
|
mutateState([&](SessionStateMachine &m) { m.onRoomEnded(reason, fatal); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void onRoomEos(livekit::Room &, const livekit::RoomEosEvent &) override
|
|
|
|
|
{
|
|
|
|
|
post(CommandType::DetachVideo);
|
|
|
|
|
post(CommandType::DetachAudio);
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onRoomEnded("the room session ended", false); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// --- worker ------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
void startWorker()
|
|
|
|
|
{
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(queue_mutex);
|
|
|
|
|
queue.clear();
|
|
|
|
|
worker_running = true;
|
|
|
|
|
}
|
|
|
|
|
worker = std::thread([this] { workerLoop(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void stopWorker()
|
|
|
|
|
{
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(queue_mutex);
|
|
|
|
|
if (!worker_running)
|
|
|
|
|
return;
|
|
|
|
|
queue.push_back(Command{CommandType::Stop, nullptr});
|
|
|
|
|
worker_running = false;
|
|
|
|
|
}
|
|
|
|
|
queue_cv.notify_one();
|
|
|
|
|
if (worker.joinable())
|
|
|
|
|
worker.join();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void workerLoop()
|
|
|
|
|
{
|
|
|
|
|
for (;;) {
|
|
|
|
|
Command command{CommandType::Stop, nullptr};
|
|
|
|
|
{
|
|
|
|
|
std::unique_lock<std::mutex> lock(queue_mutex);
|
|
|
|
|
queue_cv.wait(lock, [this] { return !queue.empty(); });
|
|
|
|
|
command = std::move(queue.front());
|
|
|
|
|
queue.pop_front();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
switch (command.type) {
|
|
|
|
|
case CommandType::AttachVideo:
|
|
|
|
|
attachVideo(command.track);
|
|
|
|
|
break;
|
|
|
|
|
case CommandType::DetachVideo:
|
|
|
|
|
detachVideo();
|
|
|
|
|
break;
|
|
|
|
|
case CommandType::AttachAudio:
|
|
|
|
|
attachAudio(command.track);
|
|
|
|
|
break;
|
|
|
|
|
case CommandType::DetachAudio:
|
|
|
|
|
detachAudio();
|
|
|
|
|
break;
|
|
|
|
|
case CommandType::Stop:
|
|
|
|
|
detachVideo();
|
|
|
|
|
detachAudio();
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void attachVideo(const std::shared_ptr<livekit::Track> &track)
|
|
|
|
|
{
|
|
|
|
|
if (!track)
|
|
|
|
|
return;
|
|
|
|
|
// Replacing an existing stream is the publisher-swap path: tear the
|
|
|
|
|
// old reader all the way down first so no frame from the previous
|
|
|
|
|
// publisher can arrive after the new one starts.
|
|
|
|
|
detachVideo();
|
|
|
|
|
|
|
|
|
|
livekit::VideoStream::Options options;
|
|
|
|
|
options.capacity = config.video_queue_capacity;
|
|
|
|
|
options.format = toLiveKitBufferType(config.video_format);
|
|
|
|
|
|
|
|
|
|
std::shared_ptr<livekit::VideoStream> stream;
|
|
|
|
|
try {
|
|
|
|
|
stream = livekit::VideoStream::fromTrack(track, options);
|
|
|
|
|
} catch (const std::exception &e) {
|
|
|
|
|
mutateState([&](SessionStateMachine &m) {
|
|
|
|
|
m.onRoomEnded(std::string("could not open the video stream: ") + e.what(), true);
|
|
|
|
|
});
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
if (!stream)
|
|
|
|
|
return;
|
|
|
|
|
|
|
|
|
|
video_stream = stream;
|
|
|
|
|
video_thread = std::thread([this, stream] { videoReaderLoop(stream); });
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onVideoAttached(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void detachVideo()
|
|
|
|
|
{
|
|
|
|
|
if (video_stream)
|
|
|
|
|
video_stream->close(); // wakes the blocking read()
|
|
|
|
|
if (video_thread.joinable())
|
|
|
|
|
video_thread.join();
|
|
|
|
|
const bool had = static_cast<bool>(video_stream);
|
|
|
|
|
video_stream.reset();
|
|
|
|
|
if (had)
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onVideoDetached(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void attachAudio(const std::shared_ptr<livekit::Track> &track)
|
|
|
|
|
{
|
|
|
|
|
if (!track)
|
|
|
|
|
return;
|
|
|
|
|
detachAudio();
|
|
|
|
|
|
|
|
|
|
livekit::AudioStream::Options options;
|
|
|
|
|
options.capacity = config.audio_queue_capacity;
|
|
|
|
|
|
|
|
|
|
std::shared_ptr<livekit::AudioStream> stream;
|
|
|
|
|
try {
|
|
|
|
|
stream = livekit::AudioStream::fromTrack(track, options);
|
|
|
|
|
} catch (const std::exception &) {
|
|
|
|
|
// Audio is not worth failing the whole source over: a camera with
|
|
|
|
|
// no usable audio track is still a usable camera.
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
if (!stream)
|
|
|
|
|
return;
|
|
|
|
|
|
|
|
|
|
audio_stream = stream;
|
|
|
|
|
audio_thread = std::thread([this, stream] { audioReaderLoop(stream); });
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onAudioAttached(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void detachAudio()
|
|
|
|
|
{
|
|
|
|
|
if (audio_stream)
|
|
|
|
|
audio_stream->close();
|
|
|
|
|
if (audio_thread.joinable())
|
|
|
|
|
audio_thread.join();
|
|
|
|
|
const bool had = static_cast<bool>(audio_stream);
|
|
|
|
|
audio_stream.reset();
|
|
|
|
|
if (had)
|
|
|
|
|
mutateState([](SessionStateMachine &m) { m.onAudioDetached(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void videoReaderLoop(std::shared_ptr<livekit::VideoStream> stream)
|
|
|
|
|
{
|
|
|
|
|
VideoFrameHandler handler;
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(state_mutex);
|
|
|
|
|
handler = on_video;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
livekit::VideoFrameEvent event;
|
|
|
|
|
while (stream->read(event)) {
|
|
|
|
|
if (!handler)
|
|
|
|
|
continue;
|
|
|
|
|
deliverVideoFrame(event, handler);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void deliverVideoFrame(livekit::VideoFrameEvent &event, const VideoFrameHandler &handler)
|
|
|
|
|
{
|
|
|
|
|
PixelFormat format;
|
|
|
|
|
const livekit::VideoFrame *frame = &event.frame;
|
|
|
|
|
livekit::VideoFrame converted;
|
|
|
|
|
|
|
|
|
|
if (!fromLiveKitBufferType(frame->type(), format)) {
|
|
|
|
|
// The SDK gave us something the adapter cannot hand to OBS.
|
|
|
|
|
// convert() is a full CPU repack, so this is a fallback, not the
|
|
|
|
|
// normal path -- the normal path is the format we asked for.
|
|
|
|
|
try {
|
|
|
|
|
converted = frame->convert(toLiveKitBufferType(config.video_format));
|
|
|
|
|
} catch (const std::exception &) {
|
|
|
|
|
dropped_frames.fetch_add(1);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
frame = &converted;
|
|
|
|
|
format = config.video_format;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const int width = frame->width();
|
|
|
|
|
const int height = frame->height();
|
|
|
|
|
const std::size_t expected = expectedFrameBytes(format, width, height);
|
|
|
|
|
if (expected == 0 || frame->dataSize() < expected) {
|
|
|
|
|
// Geometry that does not match the buffer would make OBS read off
|
|
|
|
|
// the end of it. Drop rather than trust.
|
|
|
|
|
dropped_frames.fetch_add(1);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
VideoFrameData out;
|
|
|
|
|
out.width = width;
|
|
|
|
|
out.height = height;
|
|
|
|
|
out.format = format;
|
|
|
|
|
out.data = frame->data();
|
|
|
|
|
out.size = frame->dataSize();
|
|
|
|
|
out.timestamp_us = event.timestamp_us;
|
|
|
|
|
|
|
|
|
|
const std::vector<livekit::VideoPlaneInfo> planes = frame->planeInfos();
|
|
|
|
|
const int wanted_planes = planeCount(format);
|
|
|
|
|
int count = 0;
|
|
|
|
|
for (const livekit::VideoPlaneInfo &plane : planes) {
|
|
|
|
|
if (count >= 4)
|
|
|
|
|
break;
|
|
|
|
|
out.planes[count].data = reinterpret_cast<const std::uint8_t *>(plane.data_ptr);
|
|
|
|
|
out.planes[count].stride = plane.stride;
|
|
|
|
|
out.planes[count].size = plane.size;
|
|
|
|
|
++count;
|
|
|
|
|
}
|
|
|
|
|
if (count == 0 && wanted_planes == 1) {
|
|
|
|
|
// planeInfos() documents that packed formats may return an empty
|
|
|
|
|
// list rather than one plane. Synthesise it from the frame buffer
|
|
|
|
|
// instead of dropping a perfectly good BGRA frame.
|
|
|
|
|
out.planes[0].data = frame->data();
|
|
|
|
|
out.planes[0].stride = static_cast<std::uint32_t>(width) * 4u;
|
|
|
|
|
out.planes[0].size = static_cast<std::uint32_t>(frame->dataSize());
|
|
|
|
|
count = 1;
|
|
|
|
|
}
|
|
|
|
|
if (count != wanted_planes) {
|
|
|
|
|
dropped_frames.fetch_add(1);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
out.plane_count = count;
|
|
|
|
|
|
|
|
|
|
video_frames.fetch_add(1);
|
|
|
|
|
handler(out);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void audioReaderLoop(std::shared_ptr<livekit::AudioStream> stream)
|
|
|
|
|
{
|
|
|
|
|
AudioFrameHandler handler;
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(state_mutex);
|
|
|
|
|
handler = on_audio;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
livekit::AudioFrameEvent event;
|
|
|
|
|
while (stream->read(event)) {
|
|
|
|
|
if (!handler)
|
|
|
|
|
continue;
|
|
|
|
|
const livekit::AudioFrame &frame = event.frame;
|
|
|
|
|
if (frame.numChannels() <= 0 || frame.samplesPerChannel() <= 0 || frame.sampleRate() <= 0)
|
|
|
|
|
continue;
|
|
|
|
|
AudioFrameData out;
|
|
|
|
|
out.samples = frame.data().data();
|
|
|
|
|
out.sample_count = frame.totalSamples();
|
|
|
|
|
out.sample_rate = frame.sampleRate();
|
|
|
|
|
out.channels = frame.numChannels();
|
|
|
|
|
out.samples_per_channel = frame.samplesPerChannel();
|
|
|
|
|
audio_frames.fetch_add(1);
|
|
|
|
|
handler(out);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// 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.
|
|
|
|
|
void attachExistingTracks()
|
|
|
|
|
{
|
|
|
|
|
auto participant = room.remoteParticipant(config.participant_identity).lock();
|
|
|
|
|
if (!participant)
|
|
|
|
|
return;
|
|
|
|
|
const std::string identity = participant->identity();
|
|
|
|
|
for (const auto &entry : participant->trackPublications()) {
|
|
|
|
|
const std::shared_ptr<livekit::RemoteTrackPublication> &publication = entry.second;
|
|
|
|
|
if (!publication)
|
|
|
|
|
continue;
|
|
|
|
|
const std::shared_ptr<livekit::Track> track = publication->track();
|
|
|
|
|
if (!track)
|
|
|
|
|
continue; // published but not subscribed yet
|
|
|
|
|
const MediaKind kind = toMediaKind(track->kind());
|
|
|
|
|
const MediaSource source = toMediaSource(publication->source());
|
|
|
|
|
if (isWantedVideoTrack(config.participant_identity, identity, kind, source))
|
2026-09-07 10:43:54 -07:00
|
|
|
handleWantedVideoTrack(track, publication);
|
2026-09-06 21:39:49 -07:00
|
|
|
else if (config.subscribe_audio && isWantedAudioTrack(config.participant_identity, identity, kind, source))
|
|
|
|
|
post(CommandType::AttachAudio, track);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
// ---------------------------------------------------------------------------
|
|
|
|
|
// LiveKitSession
|
|
|
|
|
// ---------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
LiveKitSession::LiveKitSession() : impl_(new Impl()) {}
|
|
|
|
|
|
|
|
|
|
LiveKitSession::~LiveKitSession()
|
|
|
|
|
{
|
|
|
|
|
disconnect();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void LiveKitSession::setVideoHandler(VideoFrameHandler handler)
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
impl_->on_video = std::move(handler);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void LiveKitSession::setAudioHandler(AudioFrameHandler handler)
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
impl_->on_audio = std::move(handler);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void LiveKitSession::setStateHandler(SessionStateHandler handler)
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
impl_->on_state = std::move(handler);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
bool LiveKitSession::connect(const SessionConfig &config)
|
|
|
|
|
{
|
|
|
|
|
if (impl_->connected)
|
|
|
|
|
disconnect();
|
|
|
|
|
|
|
|
|
|
impl_->config = config;
|
|
|
|
|
impl_->video_frames.store(0);
|
|
|
|
|
impl_->audio_frames.store(0);
|
|
|
|
|
impl_->dropped_frames.store(0);
|
|
|
|
|
|
|
|
|
|
impl_->mutateState([](SessionStateMachine &m) { m.onConnectRequested(); });
|
|
|
|
|
|
|
|
|
|
if (config.ws_url.empty() || config.token.empty() || config.participant_identity.empty()) {
|
|
|
|
|
impl_->mutateState(
|
|
|
|
|
[](SessionStateMachine &m) { m.onConnectFailed("missing LiveKit URL, token or camera selection"); });
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl_->startWorker();
|
|
|
|
|
|
|
|
|
|
livekit::RoomOptions options;
|
|
|
|
|
// auto_subscribe is what makes track_subscribed events (and therefore any
|
|
|
|
|
// media at all) happen; the SDK is emphatic about this.
|
2026-09-06 22:59:57 -07:00
|
|
|
//
|
|
|
|
|
// 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.
|
2026-09-06 21:39:49 -07:00
|
|
|
options.auto_subscribe = true;
|
|
|
|
|
options.dynacast = false;
|
|
|
|
|
// This client never publishes, so a single peer connection is all it
|
|
|
|
|
// needs.
|
|
|
|
|
options.single_peer_connection = true;
|
|
|
|
|
options.connect_timeout = std::chrono::milliseconds(config.connect_timeout_ms);
|
|
|
|
|
|
|
|
|
|
impl_->room.setDelegate(impl_.get());
|
|
|
|
|
|
|
|
|
|
bool ok = false;
|
|
|
|
|
try {
|
|
|
|
|
ok = impl_->room.connect(config.ws_url, config.token, options);
|
|
|
|
|
} catch (const std::exception &e) {
|
|
|
|
|
ok = false;
|
|
|
|
|
impl_->mutateState([&](SessionStateMachine &m) { m.onConnectFailed(e.what()); });
|
|
|
|
|
impl_->stopWorker();
|
|
|
|
|
impl_->room.setDelegate(nullptr);
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (!ok) {
|
|
|
|
|
impl_->mutateState([](SessionStateMachine &m) {
|
|
|
|
|
m.onConnectFailed("could not connect to LiveKit (check the server URL, or the token may have expired)");
|
|
|
|
|
});
|
|
|
|
|
impl_->stopWorker();
|
|
|
|
|
impl_->room.setDelegate(nullptr);
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl_->connected = true;
|
|
|
|
|
impl_->mutateState([](SessionStateMachine &m) { m.onConnectSucceeded(); });
|
|
|
|
|
impl_->attachExistingTracks();
|
|
|
|
|
return true;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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.
|
|
|
|
|
impl_->stopWorker();
|
|
|
|
|
|
|
|
|
|
if (impl_->connected) {
|
|
|
|
|
impl_->connected = false;
|
|
|
|
|
try {
|
|
|
|
|
impl_->room.disconnect(livekit::DisconnectReason::ClientInitiated);
|
|
|
|
|
} catch (const std::exception &) {
|
|
|
|
|
// Best effort: a failed graceful disconnect must not stop the
|
|
|
|
|
// OBS source from being destroyed.
|
|
|
|
|
}
|
|
|
|
|
impl_->mutateState([](SessionStateMachine &m) { m.onLocalDisconnect(); });
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl_->room.setDelegate(nullptr);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
SessionState LiveKitSession::state() const
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
return impl_->machine.state();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
std::string LiveKitSession::stateDetail() const
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
return impl_->machine.detail();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
bool LiveKitSession::hasVideo() const
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
return impl_->machine.hasVideo();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
bool LiveKitSession::hasAudio() const
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
return impl_->machine.hasAudio();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
bool LiveKitSession::waitingForCamera() const
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(impl_->state_mutex);
|
|
|
|
|
return impl_->machine.waitingForCamera();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
std::uint64_t LiveKitSession::videoFrameCount() const
|
|
|
|
|
{
|
|
|
|
|
return impl_->video_frames.load();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
std::uint64_t LiveKitSession::audioFrameCount() const
|
|
|
|
|
{
|
|
|
|
|
return impl_->audio_frames.load();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void LiveKitSession::globalInitialize()
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(globalMutex());
|
|
|
|
|
if (globalRefCount()++ == 0)
|
|
|
|
|
livekit::initialize(livekit::LogLevel::Warn);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void LiveKitSession::globalShutdown()
|
|
|
|
|
{
|
|
|
|
|
std::lock_guard<std::mutex> guard(globalMutex());
|
|
|
|
|
if (globalRefCount() > 0 && --globalRefCount() == 0)
|
|
|
|
|
livekit::shutdown();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
} // namespace stplugin
|