/* streamer-tools OBS Camera Plugin - LiveKit session wrapper Copyright (C) 2026 CyberCoveLLC This program is free software; you can redistribute it and/or modify it under the terms of the GNU General Public License as published by the Free Software Foundation; either version 2 of the License, or (at your option) any later version. This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License for more details. You should have received a copy of the GNU General Public License along with this program. If not, see */ #include "stplugin/session.h" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include 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 track; }; livekit::Room room; SessionConfig config; mutable std::mutex state_mutex; SessionStateMachine machine; VideoFrameHandler on_video; AudioFrameHandler on_audio; SessionStateHandler on_state; std::atomic video_frames{0}; std::atomic audio_frames{0}; std::atomic 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 queue; std::thread worker; bool worker_running = false; // Owned exclusively by the worker thread. std::shared_ptr video_stream; std::thread video_thread; std::shared_ptr audio_stream; std::thread audio_thread; bool connected = false; ~Impl() override = default; // --- state helpers ----------------------------------------------------- template void mutateState(Fn &&fn) { SessionState state; std::string detail; SessionStateHandler handler; { std::lock_guard 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 track = nullptr) { { std::lock_guard guard(queue_mutex); if (!worker_running) return; queue.push_back(Command{type, std::move(track)}); } queue_cv.notify_one(); } // --- 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)) post(CommandType::AttachVideo, event.track); 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 guard(queue_mutex); queue.clear(); worker_running = true; } worker = std::thread([this] { workerLoop(); }); } void stopWorker() { { std::lock_guard 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 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 &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 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(video_stream); video_stream.reset(); if (had) mutateState([](SessionStateMachine &m) { m.onVideoDetached(); }); } void attachAudio(const std::shared_ptr &track) { if (!track) return; detachAudio(); livekit::AudioStream::Options options; options.capacity = config.audio_queue_capacity; std::shared_ptr 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(audio_stream); audio_stream.reset(); if (had) mutateState([](SessionStateMachine &m) { m.onAudioDetached(); }); } void videoReaderLoop(std::shared_ptr stream) { VideoFrameHandler handler; { std::lock_guard 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 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(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(width) * 4u; out.planes[0].size = static_cast(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 stream) { AudioFrameHandler handler; { std::lock_guard 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 &publication = entry.second; if (!publication) continue; const std::shared_ptr 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)) post(CommandType::AttachVideo, track); 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 guard(impl_->state_mutex); impl_->on_video = std::move(handler); } void LiveKitSession::setAudioHandler(AudioFrameHandler handler) { std::lock_guard guard(impl_->state_mutex); impl_->on_audio = std::move(handler); } void LiveKitSession::setStateHandler(SessionStateHandler handler) { std::lock_guard 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. // // 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; 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 guard(impl_->state_mutex); return impl_->machine.state(); } std::string LiveKitSession::stateDetail() const { std::lock_guard guard(impl_->state_mutex); return impl_->machine.detail(); } bool LiveKitSession::hasVideo() const { std::lock_guard guard(impl_->state_mutex); return impl_->machine.hasVideo(); } bool LiveKitSession::hasAudio() const { std::lock_guard guard(impl_->state_mutex); return impl_->machine.hasAudio(); } bool LiveKitSession::waitingForCamera() const { std::lock_guard 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 guard(globalMutex()); if (globalRefCount()++ == 0) livekit::initialize(livekit::LogLevel::Warn); } void LiveKitSession::globalShutdown() { std::lock_guard guard(globalMutex()); if (globalRefCount() > 0 && --globalRefCount() == 0) livekit::shutdown(); } } // namespace stplugin