Files
obs-streamer-tools-plugin/core/src/session.cpp
T

736 lines
25 KiB
C++
Raw Normal View History

/*
streamer-tools OBS Camera Plugin - LiveKit session wrapper
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
*/
#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();
}
// --- 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<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))
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<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.
//
// 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<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