/* * Copyright 2025 LiveKit * * 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 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an “AS IS” BASIS, * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * See the License for the specific language governing permissions and * limitations under the License. */ #include "livekit/room.h" #include "ffi.pb.h" #include "ffi_client.h" #include "livekit/audio_stream.h" #include "livekit/e2ee.h" #include "livekit/local_participant.h" #include "livekit/remote_audio_track.h" #include "livekit/remote_data_track.h" #include "livekit/remote_participant.h" #include "livekit/remote_track_publication.h" #include "livekit/remote_video_track.h" #include "livekit/room_delegate.h" #include "livekit/room_event_types.h" #include "livekit_ffi.h" #include "lk_log.h" #include "room.pb.h" #include "room_proto_converter.h" #include "trace/trace_event.h" #include "track.pb.h" #include "track_proto_converter.h" namespace livekit { using proto::ConnectCallback; using proto::ConnectRequest; using proto::FfiEvent; using proto::FfiRequest; using proto::FfiResponse; namespace { std::shared_ptr createRemoteParticipant(const proto::OwnedParticipant& owned) { const auto& pinfo = owned.info(); std::unordered_map attrs; attrs.reserve(pinfo.attributes_size()); for (const auto& kv : pinfo.attributes()) { attrs.emplace(kv.first, kv.second); } auto kind = livekit::fromProto(pinfo.kind()); auto reason = livekit::toDisconnectReason(pinfo.disconnect_reason()); livekit::FfiHandle handle(static_cast(owned.handle().id())); return std::make_shared(std::move(handle), pinfo.sid(), pinfo.name(), pinfo.identity(), pinfo.metadata(), std::move(attrs), kind, reason); } void readyForRoomEvent(std::uint64_t room_handle) { FfiRequest req; req.mutable_ready_for_room_event()->set_room_handle(room_handle); const auto resp = FfiClient::instance().sendRequest(req); if (!resp.has_ready_for_room_event()) { throw std::runtime_error("FfiResponse missing ready_for_room_event"); } } } // namespace Room::Room() : subscription_thread_dispatcher_(std::make_unique()) {} Room::~Room() { // disconnect() handles all destruction/graceful teardown functionality, simply call it here try { (void)disconnect(); // Don't need return value } catch (const std::exception& e) { LK_LOG_ERROR("Room::~Room: graceful disconnect failed: {}", e.what()); } catch (...) { LK_LOG_ERROR("Room::~Room: graceful disconnect failed: unknown exception"); } } void Room::setDelegate(RoomDelegate* delegate) { const std::scoped_lock g(lock_); delegate_ = delegate; } bool Room::connect(const std::string& url, const std::string& token, const RoomOptions& options) { TRACE_EVENT0("livekit", "Room::connect"); if (!FfiClient::instance().isInitialized()) { LK_LOG_ERROR("Room::connect failed: LiveKit is not initialized"); return false; } { const std::scoped_lock g(lock_); if (connection_state_ != ConnectionState::Disconnected) { throw std::runtime_error("already connected"); } connection_state_ = ConnectionState::Reconnecting; } FfiClient::ListenerId listenerId = 0; try { listenerId = FfiClient::instance().addListener([this](const proto::FfiEvent& e) { onEvent(e); }); { const std::scoped_lock g(lock_); listener_id_ = listenerId; } auto fut = FfiClient::instance().connectAsync(url, token, options); auto connectCb = fut.get(); // fut will throw if it fails to connect to the room const auto& owned_room = connectCb.result().room(); auto new_room_handle = std::make_shared(owned_room.handle().id()); const auto room_handle_id = static_cast(new_room_handle->get()); auto new_room_info = fromProto(owned_room.info()); // Setup local particpant std::shared_ptr new_local_participant; { const auto& owned_local = connectCb.result().local_participant(); const auto& pinfo = owned_local.info(); // Build attributes map std::unordered_map attrs; for (const auto& kv : pinfo.attributes()) { attrs.emplace(kv.first, kv.second); } auto kind = fromProto(pinfo.kind()); auto reason = toDisconnectReason(pinfo.disconnect_reason()); // Participant base stores a weak_ptr, so share the room handle FfiHandle participant_handle(static_cast(owned_local.handle().id())); new_local_participant = std::make_shared(std::move(participant_handle), pinfo.sid(), pinfo.name(), pinfo.identity(), pinfo.metadata(), std::move(attrs), kind, reason); } // Setup remote participants std::unordered_map> new_remote_participants; { const auto& participants = connectCb.result().participants(); const std::scoped_lock g(lock_); for (const auto& pt : participants) { const auto& owned = pt.participant(); auto rp = createRemoteParticipant(owned); // Add the initial remote participant tracks for (const auto& owned_publication_info : pt.publications()) { auto publication = std::make_shared(owned_publication_info); rp->mutableTrackPublications().emplace(publication->sid(), std::move(publication)); } new_remote_participants.emplace(rp->identity(), std::move(rp)); } } // Setup e2eeManager std::shared_ptr new_e2ee_manager; if (options.encryption) { LK_LOG_INFO("creating E2eeManager"); new_e2ee_manager = std::shared_ptr(new E2EEManager(new_room_handle->get(), options.encryption.value())); } // Publish all state atomically under lock { const std::scoped_lock g(lock_); room_handle_ = std::move(new_room_handle); room_info_ = std::move(new_room_info); local_participant_ = std::move(new_local_participant); remote_participants_ = std::move(new_remote_participants); e2ee_manager_ = std::move(new_e2ee_manager); connection_state_ = ConnectionState::Connected; } readyForRoomEvent(room_handle_id); return true; } catch (const std::exception& e) { int listener_to_remove = 0; std::shared_ptr local_participant_to_cleanup; { const std::scoped_lock g(lock_); connection_state_ = ConnectionState::Disconnected; if (listener_id_ == listenerId) { listener_to_remove = listener_id_; listener_id_ = 0; } local_participant_to_cleanup = std::move(local_participant_); remote_participants_.clear(); room_handle_.reset(); e2ee_manager_.reset(); text_stream_readers_.clear(); byte_stream_readers_.clear(); } if (local_participant_to_cleanup) { local_participant_to_cleanup->shutdown(); } if (listener_to_remove != 0) { FfiClient::instance().removeListener(listener_to_remove); } LK_LOG_ERROR("Room::connect failed: {}", e.what()); return false; } } bool Room::disconnect(DisconnectReason reason) { TRACE_EVENT0("livekit", "Room::disconnect"); return shutdown(true, reason, true); } bool Room::shutdown(bool disconnect_ffi, DisconnectReason reason, bool notify_delegate) { std::shared_ptr handle; RoomDelegate* delegate_snapshot = nullptr; std::shared_ptr local_participant_to_cleanup; std::unordered_map> remote_participants_to_clear; std::shared_ptr e2ee_manager_to_clear; std::unordered_map> text_stream_readers_to_clear; std::unordered_map> byte_stream_readers_to_clear; int listener_to_remove = 0; bool claimed_disconnect = false; { const std::scoped_lock g(lock_); const bool has_room_state = connection_state_ != ConnectionState::Disconnected || listener_id_ != 0 || room_handle_ || local_participant_ || !remote_participants_.empty(); // Return false for a no-op so callers can tell whether this call claimed the // room state and performed cleanup. Matches disconnect()'s documented contract. if (!has_room_state) { return false; } // The state transition determines which racing path owns the FFI request // and delegate notification. Remaining room state is still claimed here so // EOS or destruction can finish local cleanup after a server disconnect. claimed_disconnect = connection_state_ != ConnectionState::Disconnected; handle = std::move(room_handle_); delegate_snapshot = delegate_; local_participant_to_cleanup = std::move(local_participant_); remote_participants_to_clear = std::move(remote_participants_); e2ee_manager_to_clear = std::move(e2ee_manager_); text_stream_readers_to_clear = std::move(text_stream_readers_); byte_stream_readers_to_clear = std::move(byte_stream_readers_); listener_to_remove = listener_id_; listener_id_ = 0; connection_state_ = ConnectionState::Disconnected; } bool shutdown_ok = true; if (local_participant_to_cleanup) { try { local_participant_to_cleanup->shutdown(); } catch (const std::exception& e) { LK_LOG_ERROR("Room shutdown: local participant shutdown failed: {}", e.what()); shutdown_ok = false; } catch (...) { LK_LOG_ERROR("Room shutdown: local participant shutdown failed: unknown exception"); shutdown_ok = false; } } if (disconnect_ffi && claimed_disconnect && handle && handle->valid()) { try { FfiClient::instance().disconnectAsync(handle->get(), reason).get(); } catch (const std::exception& e) { LK_LOG_ERROR("Room shutdown: FFI disconnect failed (continuing local shutdown): {}", e.what()); shutdown_ok = false; } catch (...) { LK_LOG_ERROR("Room shutdown: FFI disconnect failed (continuing local shutdown): unknown exception"); shutdown_ok = false; } } if (subscription_thread_dispatcher_) { try { subscription_thread_dispatcher_->stopAll(); } catch (const std::exception& e) { LK_LOG_ERROR("Room shutdown: subscription shutdown failed: {}", e.what()); shutdown_ok = false; } catch (...) { LK_LOG_ERROR("Room shutdown: subscription shutdown failed: unknown exception"); shutdown_ok = false; } } if (listener_to_remove != 0) { try { FfiClient::instance().removeListener(listener_to_remove); } catch (const std::exception& e) { LK_LOG_ERROR("Room shutdown: listener removal failed: {}", e.what()); shutdown_ok = false; } catch (...) { LK_LOG_ERROR("Room shutdown: listener removal failed: unknown exception"); shutdown_ok = false; } } local_participant_to_cleanup.reset(); remote_participants_to_clear.clear(); e2ee_manager_to_clear.reset(); text_stream_readers_to_clear.clear(); byte_stream_readers_to_clear.clear(); handle.reset(); if (notify_delegate && claimed_disconnect && delegate_snapshot) { DisconnectedEvent ev; ev.reason = reason; try { delegate_snapshot->onDisconnected(*this, ev); } catch (const std::exception& e) { LK_LOG_ERROR("Room::disconnect: onDisconnected threw: {}", e.what()); } catch (...) { LK_LOG_ERROR("Room::disconnect: onDisconnected threw: unknown exception"); } } return claimed_disconnect && shutdown_ok; } RoomInfoData Room::roomInfo() const { const std::scoped_lock g(lock_); return room_info_; } std::weak_ptr Room::localParticipant() const { const std::scoped_lock g(lock_); return local_participant_; } std::weak_ptr Room::remoteParticipant(const std::string& identity) const { const std::scoped_lock g(lock_); auto it = remote_participants_.find(identity); return it == remote_participants_.end() ? std::weak_ptr{} : it->second; } std::vector> Room::remoteParticipants() const { const std::scoped_lock guard(lock_); std::vector> out; out.reserve(remote_participants_.size()); for (const auto& kv : remote_participants_) { out.push_back(kv.second); } return out; } ConnectionState Room::connectionState() const { const std::scoped_lock g(lock_); return connection_state_; } std::future Room::getStats() const { std::shared_ptr handle; { const std::scoped_lock g(lock_); handle = room_handle_; } if (!handle) { throw std::runtime_error("Room::getStats called on a disconnected room"); } return FfiClient::instance().getSessionStatsAsync(handle->get()); } std::weak_ptr Room::e2eeManager() const { const std::scoped_lock g(lock_); return e2ee_manager_; } void Room::registerTextStreamHandler(const std::string& topic, TextStreamHandler handler) { const std::scoped_lock g(lock_); auto [it, inserted] = text_stream_handlers_.emplace(topic, std::move(handler)); if (!inserted) { throw std::runtime_error("text stream handler for topic '" + topic + "' already set"); } } void Room::unregisterTextStreamHandler(const std::string& topic) { const std::scoped_lock g(lock_); text_stream_handlers_.erase(topic); } void Room::registerByteStreamHandler(const std::string& topic, ByteStreamHandler handler) { const std::scoped_lock g(lock_); auto [it, inserted] = byte_stream_handlers_.emplace(topic, std::move(handler)); if (!inserted) { throw std::runtime_error("byte stream handler for topic '" + topic + "' already set"); } } void Room::unregisterByteStreamHandler(const std::string& topic) { const std::scoped_lock g(lock_); byte_stream_handlers_.erase(topic); } // ------------------------------------------------------------------- // Frame callback registration // ------------------------------------------------------------------- void Room::setOnAudioFrameCallback(const std::string& participant_identity, const std::string& track_name, AudioFrameCallback callback, const AudioStream::Options& opts) { if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->setOnAudioFrameCallback(participant_identity, track_name, std::move(callback), opts); } } void Room::setOnVideoFrameCallback(const std::string& participant_identity, const std::string& track_name, VideoFrameCallback callback, const VideoStream::Options& opts) { if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->setOnVideoFrameCallback(participant_identity, track_name, std::move(callback), opts); } } void Room::setOnVideoFrameEventCallback(const std::string& participant_identity, const std::string& track_name, VideoFrameEventCallback callback, const VideoStream::Options& opts) { if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->setOnVideoFrameEventCallback(participant_identity, track_name, std::move(callback), opts); } } void Room::clearOnAudioFrameCallback(const std::string& participant_identity, const std::string& track_name) { if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->clearOnAudioFrameCallback(participant_identity, track_name); } } void Room::clearOnVideoFrameCallback(const std::string& participant_identity, const std::string& track_name) { if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->clearOnVideoFrameCallback(participant_identity, track_name); } } DataFrameCallbackId Room::addOnDataFrameCallback(const std::string& participant_identity, const std::string& track_name, DataFrameCallback callback) { if (subscription_thread_dispatcher_) { return subscription_thread_dispatcher_->addOnDataFrameCallback(participant_identity, track_name, std::move(callback)); } return std::numeric_limits::max(); } void Room::removeOnDataFrameCallback(DataFrameCallbackId id) { if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->removeOnDataFrameCallback(id); } } void Room::onEvent(const FfiEvent& event) { // Take a snapshot of the delegate under lock, but do NOT call it under the // lock. RoomDelegate* delegate_snapshot = nullptr; { const std::scoped_lock guard(lock_); delegate_snapshot = delegate_; } // First, handle RPC method invocations (not part of RoomEvent). if (event.message_case() == FfiEvent::kRpcMethodInvocation) { const auto& rpc = event.rpc_method_invocation(); LocalParticipant* lp = nullptr; { const std::scoped_lock guard(lock_); if (!local_participant_) { return; } auto local_handle = local_participant_->ffiHandleId(); if (local_handle == INVALID_HANDLE || rpc.local_participant_handle() != static_cast(local_handle)) { // RPC is not targeted at this room's local participant; ignore. return; } lp = local_participant_.get(); } // Call outside the lock to avoid deadlocks / re-entrancy issues. lp->handleRpcMethodInvocation(rpc.invocation_id(), rpc.method(), rpc.request_id(), rpc.caller_identity(), rpc.payload(), static_cast(rpc.response_timeout_ms()) / 1000.0); return; } switch (event.message_case()) { case FfiEvent::kRoomEvent: { const proto::RoomEvent& re = event.room_event(); // Check if this event is for our room handle { const std::scoped_lock guard(lock_); if (!room_handle_ || re.room_handle() != static_cast(room_handle_->get())) { return; } } switch (re.message_case()) { case proto::RoomEvent::kParticipantConnected: { std::shared_ptr new_participant; { const std::scoped_lock guard(lock_); const auto& owned = re.participant_connected().info(); // createRemoteParticipant takes proto::OwnedParticipant new_participant = createRemoteParticipant(owned); remote_participants_.emplace(new_participant->identity(), new_participant); } ParticipantConnectedEvent ev; ev.participant = new_participant.get(); if (delegate_snapshot) { delegate_snapshot->onParticipantConnected(*this, ev); } break; } case proto::RoomEvent::kParticipantDisconnected: { std::shared_ptr removed; DisconnectReason reason = DisconnectReason::Unknown; { const std::scoped_lock guard(lock_); const auto& pd = re.participant_disconnected(); const std::string& identity = pd.participant_identity(); reason = toDisconnectReason(pd.disconnect_reason()); auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { removed = it->second; remote_participants_.erase(it); } else { // We saw a disconnect event for a participant we don't track // internally. This can happen on races or if we never created a // RemoteParticipant LK_LOG_WARN("participant_disconnected for unknown identity: {}", identity); } } if (removed) { ParticipantDisconnectedEvent ev; ev.participant = removed.get(); ev.reason = reason; if (delegate_snapshot) { delegate_snapshot->onParticipantDisconnected(*this, ev); } } break; } case proto::RoomEvent::kLocalTrackPublished: { LocalTrackPublishedEvent ev; { const std::scoped_lock guard(lock_); if (!local_participant_) { LK_LOG_ERROR("kLocalTrackPublished: local_participant_ is nullptr"); break; } const auto& ltp = re.local_track_published(); const std::string& sid = ltp.track_sid(); const auto pubs = local_participant_->trackPublications(); auto it = pubs.find(sid); if (it == pubs.end()) { LK_LOG_WARN("local_track_published for unknown sid: {}", sid); break; } ev.publication = it->second; ev.track = ev.publication ? ev.publication->track() : nullptr; } if (delegate_snapshot) { delegate_snapshot->onLocalTrackPublished(*this, ev); } break; } case proto::RoomEvent::kLocalTrackUnpublished: { LocalTrackUnpublishedEvent ev; { const std::scoped_lock guard(lock_); if (!local_participant_) { LK_LOG_ERROR("kLocalTrackUnpublished: local_participant_ is nullptr"); break; } const auto& ltu = re.local_track_unpublished(); const std::string& pub_sid = ltu.publication_sid(); const auto pubs = local_participant_->trackPublications(); auto it = pubs.find(pub_sid); if (it == pubs.end()) { LK_LOG_WARN("local_track_unpublished for unknown publication sid: {}", pub_sid); break; } ev.publication = it->second; } if (delegate_snapshot) { delegate_snapshot->onLocalTrackUnpublished(*this, ev); } break; } case proto::RoomEvent::kLocalTrackSubscribed: { LocalTrackSubscribedEvent ev; { const std::scoped_lock guard(lock_); if (!local_participant_) { break; } const auto& lts = re.local_track_subscribed(); const std::string& sid = lts.track_sid(); const auto pubs = local_participant_->trackPublications(); auto it = pubs.find(sid); if (it == pubs.end()) { LK_LOG_WARN("local_track_subscribed for unknown sid: {}", sid); break; } auto publication = it->second; ev.track = publication ? publication->track() : nullptr; } if (delegate_snapshot) { delegate_snapshot->onLocalTrackSubscribed(*this, ev); } break; } case proto::RoomEvent::kTrackPublished: { TrackPublishedEvent ev; { const std::scoped_lock guard(lock_); const auto& tp = re.track_published(); const std::string& identity = tp.participant_identity(); auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { RemoteParticipant* rparticipant = it->second.get(); const auto& owned_publication = tp.publication(); auto rpublication = std::make_shared(owned_publication); // Store it on the participant, keyed by SID rparticipant->mutableTrackPublications().emplace(rpublication->sid(), std::move(rpublication)); ev.participant = rparticipant; ev.publication = rpublication; } else { // Optional: log if we get a track for an unknown participant LK_LOG_WARN("track_published for unknown participant: {}", identity); // Don't emit the break; } } if (delegate_snapshot) { delegate_snapshot->onTrackPublished(*this, ev); } break; } case proto::RoomEvent::kTrackUnpublished: { TrackUnpublishedEvent ev; { const std::scoped_lock guard(lock_); const auto& tu = re.track_unpublished(); const std::string& identity = tu.participant_identity(); const std::string& pub_sid = tu.publication_sid(); auto pit = remote_participants_.find(identity); if (pit == remote_participants_.end()) { LK_LOG_WARN("track_unpublished for unknown participant: {}", identity); break; } RemoteParticipant* rparticipant = pit->second.get(); auto& pubs = rparticipant->mutableTrackPublications(); auto it = pubs.find(pub_sid); if (it == pubs.end()) { LK_LOG_WARN( "track_unpublished for unknown publication sid {} " "(participant {})", pub_sid, identity); break; } ev.participant = rparticipant; ev.publication = it->second; pubs.erase(it); } if (delegate_snapshot) { delegate_snapshot->onTrackUnpublished(*this, ev); } break; } case proto::RoomEvent::kTrackSubscribed: { const auto& ts = re.track_subscribed(); const std::string& identity = ts.participant_identity(); const auto& owned_track = ts.track(); const auto& track_info = owned_track.info(); std::shared_ptr rpublication; RemoteParticipant* rparticipant = nullptr; std::shared_ptr remote_track; { const std::scoped_lock guard(lock_); // Find participant auto pit = remote_participants_.find(identity); if (pit == remote_participants_.end()) { LK_LOG_WARN("track_subscribed for unknown participant: {}", identity); break; } rparticipant = pit->second.get(); // Find existing publication by track SID (from track_published) auto& pubs = rparticipant->mutableTrackPublications(); auto pubIt = pubs.find(track_info.sid()); if (pubIt == pubs.end()) { LK_LOG_WARN( "track_subscribed for unknown publication sid {} " "(participant {})", track_info.sid(), identity); break; } rpublication = pubIt->second; // Create RemoteVideoTrack / RemoteAudioTrack if (track_info.kind() == proto::TrackKind::KIND_VIDEO) { remote_track = std::make_shared(owned_track); } else if (track_info.kind() == proto::TrackKind::KIND_AUDIO) { remote_track = std::make_shared(owned_track); } else { LK_LOG_WARN("track_subscribed with unsupported kind: {}", static_cast(track_info.kind())); break; } // Attach to publication, mark subscribed rpublication->setTrack(remote_track); rpublication->setSubscribed(true); } // Emit remote track_subscribed-style callback TrackSubscribedEvent ev; ev.track = remote_track; ev.publication = rpublication; ev.participant = rparticipant; if (delegate_snapshot) { delegate_snapshot->onTrackSubscribed(*this, ev); } if (subscription_thread_dispatcher_ && remote_track && rpublication) { subscription_thread_dispatcher_->handleTrackSubscribed(identity, rpublication->name(), remote_track); } break; } case proto::RoomEvent::kTrackUnsubscribed: { TrackUnsubscribedEvent ev; TrackSource unsub_source = TrackSource::SOURCE_UNKNOWN; std::string unsub_identity; { const std::scoped_lock guard(lock_); const auto& tu = re.track_unsubscribed(); unsub_identity = tu.participant_identity(); const std::string& track_sid = tu.track_sid(); auto pit = remote_participants_.find(unsub_identity); if (pit == remote_participants_.end()) { LK_LOG_WARN("track_unsubscribed for unknown participant: {}", unsub_identity); break; } RemoteParticipant* rparticipant = pit->second.get(); auto& pubs = rparticipant->mutableTrackPublications(); auto pubIt = pubs.find(track_sid); if (pubIt == pubs.end()) { LK_LOG_WARN( "track_unsubscribed for unknown publication sid {} " "(participant {})", track_sid, unsub_identity); break; } auto publication = pubIt->second; unsub_source = publication->source(); auto track = publication->track(); publication->setTrack(nullptr); publication->setSubscribed(false); ev.participant = rparticipant; ev.publication = publication; ev.track = track; } if (delegate_snapshot) { delegate_snapshot->onTrackUnsubscribed(*this, ev); } if (subscription_thread_dispatcher_ && unsub_source != TrackSource::SOURCE_UNKNOWN) { subscription_thread_dispatcher_->handleTrackUnsubscribed(unsub_identity, unsub_source, ev.publication ? ev.publication->name() : ""); } break; } case proto::RoomEvent::kTrackSubscriptionFailed: { TrackSubscriptionFailedEvent ev; { const std::scoped_lock guard(lock_); const auto& tsf = re.track_subscription_failed(); const std::string& identity = tsf.participant_identity(); auto pit = remote_participants_.find(identity); if (pit == remote_participants_.end()) { LK_LOG_WARN("track_subscription_failed for unknown participant: {}", identity); break; } ev.participant = pit->second.get(); ev.track_sid = tsf.track_sid(); ev.error = tsf.error(); } if (delegate_snapshot) { delegate_snapshot->onTrackSubscriptionFailed(*this, ev); } break; } case proto::RoomEvent::kDataTrackPublished: { const auto& rdtp = re.data_track_published(); auto remote_track = std::shared_ptr(new RemoteDataTrack(rdtp.track())); if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->handleDataTrackPublished(remote_track); } DataTrackPublishedEvent ev; ev.track = remote_track; if (delegate_snapshot) { delegate_snapshot->onDataTrackPublished(*this, ev); } break; } case proto::RoomEvent::kDataTrackUnpublished: { const auto& dtu = re.data_track_unpublished(); if (subscription_thread_dispatcher_) { subscription_thread_dispatcher_->handleDataTrackUnpublished(dtu.sid()); } DataTrackUnpublishedEvent ev; ev.sid = dtu.sid(); if (delegate_snapshot) { delegate_snapshot->onDataTrackUnpublished(*this, ev); } break; } case proto::RoomEvent::kTrackMuted: { TrackMutedEvent ev; bool success = false; { const std::scoped_lock guard(lock_); const auto& tm = re.track_muted(); const std::string& identity = tm.participant_identity(); const std::string& sid = tm.track_sid(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto pit = remote_participants_.find(identity); if (pit != remote_participants_.end()) { participant = pit->second.get(); } } if (!participant) { LK_LOG_WARN("track_muted for unknown participant: {}", identity); break; } auto pub = participant->findTrackPublication(sid); if (!pub) { LK_LOG_WARN("track_muted for unknown track sid: {}", sid); } else { pub->setMuted(true); if (auto t = pub->track()) { t->setMuted(true); } ev.participant = participant; ev.publication = pub; success = true; } } if (success && delegate_snapshot) { delegate_snapshot->onTrackMuted(*this, ev); } break; } case proto::RoomEvent::kTrackUnmuted: { TrackUnmutedEvent ev; bool success = false; { const std::scoped_lock guard(lock_); const auto& tu = re.track_unmuted(); const std::string& identity = tu.participant_identity(); const std::string& sid = tu.track_sid(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto pit = remote_participants_.find(identity); if (pit != remote_participants_.end()) { participant = pit->second.get(); } } if (!participant) { LK_LOG_WARN("track_unmuted for unknown participant: {}", identity); break; } auto pub = participant->findTrackPublication(sid); if (!pub) { LK_LOG_WARN("track_unmuted for unknown track sid: {}", sid); } else { pub->setMuted(false); if (auto t = pub->track()) { t->setMuted(false); } ev.participant = participant; ev.publication = pub; success = true; } ev.participant = participant; ev.publication = pub; } if (success && delegate_snapshot) { delegate_snapshot->onTrackUnmuted(*this, ev); } break; } case proto::RoomEvent::kActiveSpeakersChanged: { ActiveSpeakersChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& asc = re.active_speakers_changed(); for (const auto& identity : asc.participant_identities()) { // Appears to be clang-tidy false positive // NOLINTNEXTLINE(misc-const-correctness) Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto pit = remote_participants_.find(identity); if (pit != remote_participants_.end()) { participant = pit->second.get(); } } if (participant) { ev.speakers.push_back(participant); } } } if (delegate_snapshot) { delegate_snapshot->onActiveSpeakersChanged(*this, ev); } break; } case proto::RoomEvent::kRoomMetadataChanged: { RoomMetadataChangedEvent ev; { const std::scoped_lock guard(lock_); const auto old_metadata = room_info_.metadata; room_info_.metadata = re.room_metadata_changed().metadata(); ev.old_metadata = old_metadata; ev.new_metadata = room_info_.metadata; } if (delegate_snapshot) { delegate_snapshot->onRoomMetadataChanged(*this, ev); } break; } case proto::RoomEvent::kRoomSidChanged: { RoomSidChangedEvent ev; { const std::scoped_lock guard(lock_); room_info_.sid = re.room_sid_changed().sid(); ev.sid = room_info_.sid.value_or(std::string{}); } if (delegate_snapshot) { delegate_snapshot->onRoomSidChanged(*this, ev); } break; } case proto::RoomEvent::kParticipantMetadataChanged: { ParticipantMetadataChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& pm = re.participant_metadata_changed(); const std::string& identity = pm.participant_identity(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN("participant_metadata_changed for unknown participant: {}", identity); break; } const std::string old_metadata = participant->metadata(); participant->metadata_ = pm.metadata(); ev.participant = participant; ev.old_metadata = old_metadata; ev.new_metadata = participant->metadata(); } if (delegate_snapshot) { delegate_snapshot->onParticipantMetadataChanged(*this, ev); } break; } case proto::RoomEvent::kParticipantNameChanged: { ParticipantNameChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& pn = re.participant_name_changed(); const std::string& identity = pn.participant_identity(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN("participant_name_changed for unknown participant: {}", identity); break; } const std::string old_name = participant->name(); participant->name_ = pn.name(); ev.participant = participant; ev.old_name = old_name; ev.new_name = participant->name(); } if (delegate_snapshot) { delegate_snapshot->onParticipantNameChanged(*this, ev); } break; } case proto::RoomEvent::kParticipantAttributesChanged: { ParticipantAttributesChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& pa = re.participant_attributes_changed(); const std::string& identity = pa.participant_identity(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN("participant_attributes_changed for unknown participant: {}", identity); break; } // Build full attributes map std::unordered_map attrs; for (const auto& entry : pa.attributes()) { attrs.emplace(entry.key(), entry.value()); } participant->attributes_ = std::move(attrs); // Build changed_attributes map for (const auto& entry : pa.changed_attributes()) { ev.changed_attributes.emplace_back(entry.key(), entry.value()); } ev.participant = participant; } if (delegate_snapshot) { delegate_snapshot->onParticipantAttributesChanged(*this, ev); } break; } case proto::RoomEvent::kParticipantEncryptionStatusChanged: { ParticipantEncryptionStatusChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& pe = re.participant_encryption_status_changed(); const std::string& identity = pe.participant_identity(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN( "participant_encryption_status_changed for unknown " "participant: {}", identity); break; } ev.participant = participant; ev.is_encrypted = pe.is_encrypted(); } if (delegate_snapshot) { delegate_snapshot->onParticipantEncryptionStatusChanged(*this, ev); } break; } case proto::RoomEvent::kConnectionQualityChanged: { ConnectionQualityChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& cq = re.connection_quality_changed(); const std::string& identity = cq.participant_identity(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN("connection_quality_changed for unknown participant: {}", identity); break; } ev.participant = participant; ev.quality = static_cast(cq.quality()); } if (delegate_snapshot) { delegate_snapshot->onConnectionQualityChanged(*this, ev); } break; } // ------------------------------------------------------------------------ // Data packets: user vs SIP DTMF // ------------------------------------------------------------------------ case proto::RoomEvent::kDataPacketReceived: { const auto& dp = re.data_packet_received(); RemoteParticipant* rp = nullptr; { const std::scoped_lock guard(lock_); auto it = remote_participants_.find(dp.participant_identity()); if (it != remote_participants_.end()) { rp = it->second.get(); } } const auto which_val = dp.value_case(); if (which_val == proto::DataPacketReceived::kUser && delegate_snapshot) { const UserDataPacketEvent ev = userDataPacketFromProto(dp, rp); delegate_snapshot->onUserPacketReceived(*this, ev); } else if (which_val == proto::DataPacketReceived::kSipDtmf && delegate_snapshot) { const SipDtmfReceivedEvent ev = sipDtmfFromProto(dp, rp); delegate_snapshot->onSipDtmfReceived(*this, ev); } break; } // ------------------------------------------------------------------------ // E2EE state // ------------------------------------------------------------------------ case proto::RoomEvent::kE2EeStateChanged: { E2eeStateChangedEvent ev; { LK_LOG_DEBUG("e2ee_state_changed for participant"); const std::scoped_lock guard(lock_); const auto& es = re.e2ee_state_changed(); const std::string& identity = es.participant_identity(); Participant* participant = nullptr; if (local_participant_ && local_participant_->identity() == identity) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN("e2ee_state_changed for unknown participant: {}", identity); break; } ev.participant = participant; ev.state = static_cast(es.state()); } if (delegate_snapshot) { delegate_snapshot->onE2eeStateChanged(*this, ev); } break; } // ------------------------------------------------------------------------ // Connection state / lifecycle // ------------------------------------------------------------------------ case proto::RoomEvent::kConnectionStateChanged: { ConnectionStateChangedEvent ev; { const std::scoped_lock guard(lock_); const auto& cs = re.connection_state_changed(); // TODO, maybe we should update our |connection_state_| // correspoindingly, but the this kConnectionStateChanged event is never // triggered in my local test. LK_LOG_DEBUG("cs.state() is {} connection_state_ is {}", static_cast(cs.state()), static_cast(connection_state_)); ev.state = static_cast(cs.state()); } if (delegate_snapshot) { delegate_snapshot->onConnectionStateChanged(*this, ev); } break; } case proto::RoomEvent::kDisconnected: { bool should_notify = false; { const std::scoped_lock guard(lock_); // Local shutdown marks the state before awaiting the FFI response // and notifies the delegate itself. Suppress that duplicate while // passing server-initiated disconnects through unchanged. should_notify = connection_state_ != ConnectionState::Disconnected; connection_state_ = ConnectionState::Disconnected; } if (should_notify && delegate_snapshot) { DisconnectedEvent ev; ev.reason = toDisconnectReason(re.disconnected().reason()); delegate_snapshot->onDisconnected(*this, ev); } break; } case proto::RoomEvent::kReconnecting: { const ReconnectingEvent ev; if (delegate_snapshot) { delegate_snapshot->onReconnecting(*this, ev); } break; } case proto::RoomEvent::kReconnected: { const ReconnectedEvent ev; if (delegate_snapshot) { delegate_snapshot->onReconnected(*this, ev); } break; } case proto::RoomEvent::kTokenRefreshed: { const TokenRefreshedEvent ev = fromProto(re.token_refreshed()); if (delegate_snapshot) { delegate_snapshot->onTokenRefreshed(*this, ev); } break; } case proto::RoomEvent::kEos: { (void)shutdown(false, DisconnectReason::Unknown, false); const RoomEosEvent ev; if (delegate_snapshot) { delegate_snapshot->onRoomEos(*this, ev); } break; } case proto::RoomEvent::kChatMessage: { // Deprecated event, do nothing. break; } case proto::RoomEvent::kStreamHeaderReceived: { const auto& sh = re.stream_header_received(); const auto& header = sh.header(); const std::string& participant_identity = sh.participant_identity(); // Snapshot handler + create reader without holding lock during user // callback TextStreamHandler text_cb; ByteStreamHandler byte_cb; std::shared_ptr text_reader; std::shared_ptr byte_reader; { const std::scoped_lock guard(lock_); // Determine stream type from oneof in protobuf // Adjust these names if your generated C++ uses different ones const auto stream_type = header.content_header_case(); if (stream_type == proto::DataStream::Header::kTextHeader) { auto it = text_stream_handlers_.find(header.topic()); if (it == text_stream_handlers_.end()) { // Ignore if no callback attached break; } text_cb = it->second; const TextStreamInfo info = makeTextInfo(header); text_reader = std::make_shared(info); text_stream_readers_[header.stream_id()] = text_reader; } else if (stream_type == proto::DataStream::Header::kByteHeader) { auto it = byte_stream_handlers_.find(header.topic()); if (it == byte_stream_handlers_.end()) { break; } byte_cb = it->second; const ByteStreamInfo info = makeByteInfo(header); byte_reader = std::make_shared(info); byte_stream_readers_[header.stream_id()] = byte_reader; } else { // unknown header type: ignore break; } } // Invoke user callback outside lock (very important) if (text_reader) { text_cb(text_reader, participant_identity); } else if (byte_reader) { byte_cb(byte_reader, participant_identity); } break; } case proto::RoomEvent::kStreamChunkReceived: { const auto& sc = re.stream_chunk_received(); const auto& chunk = sc.chunk(); std::shared_ptr text_reader; std::shared_ptr byte_reader; { const std::scoped_lock guard(lock_); auto itT = text_stream_readers_.find(chunk.stream_id()); if (itT != text_stream_readers_.end()) { text_reader = itT->second; } else { auto itB = byte_stream_readers_.find(chunk.stream_id()); if (itB != byte_stream_readers_.end()) { byte_reader = itB->second; } } } if (text_reader) { // chunk.content() is bytes; treat as UTF-8 string. text_reader->onChunkUpdate(chunk.content()); } else if (byte_reader) { // Convert string bytes -> vector const std::string& s = chunk.content(); const std::vector bytes(s.begin(), s.end()); byte_reader->onChunkUpdate(bytes); } break; } case proto::RoomEvent::kStreamTrailerReceived: { const auto& st = re.stream_trailer_received(); const auto& trailer = st.trailer(); std::shared_ptr text_reader; std::shared_ptr byte_reader; std::map trailer_attrs; for (const auto& kv : trailer.attributes()) { trailer_attrs.emplace(kv.first, kv.second); } { const std::scoped_lock guard(lock_); auto itT = text_stream_readers_.find(trailer.stream_id()); if (itT != text_stream_readers_.end()) { text_reader = itT->second; text_stream_readers_.erase(itT); } else { auto itB = byte_stream_readers_.find(trailer.stream_id()); if (itB != byte_stream_readers_.end()) { byte_reader = itB->second; byte_stream_readers_.erase(itB); } } } if (text_reader) { text_reader->onStreamClose(trailer_attrs); } else if (byte_reader) { byte_reader->onStreamClose(trailer_attrs); } break; } case proto::RoomEvent::kDataChannelLowThresholdChanged: { auto ev = fromProto(re.data_channel_low_threshold_changed()); if (delegate_snapshot) { delegate_snapshot->onDataChannelBufferedAmountLowThresholdChanged(*this, ev); } break; } case proto::RoomEvent::kByteStreamOpened: { auto ev = fromProto(re.byte_stream_opened()); if (delegate_snapshot) { delegate_snapshot->onByteStreamOpened(*this, ev); } break; } case proto::RoomEvent::kTextStreamOpened: { auto ev = fromProto(re.text_stream_opened()); if (delegate_snapshot) { delegate_snapshot->onTextStreamOpened(*this, ev); } break; } case proto::RoomEvent::kRoomUpdated: { auto ev = roomUpdatedFromProto(re.room_updated()); if (delegate_snapshot) { delegate_snapshot->onRoomUpdated(*this, ev); } break; } case proto::RoomEvent::kMoved: { auto ev = roomMovedFromProto(re.moved()); if (delegate_snapshot) { delegate_snapshot->onRoomMoved(*this, ev); } break; } case proto::RoomEvent::kParticipantsUpdated: { ParticipantsUpdatedEvent ev; { const std::scoped_lock guard(lock_); const auto& pu = re.participants_updated(); for (const auto& info : pu.participants()) { const std::string& identity = info.identity(); Participant* participant = nullptr; if (local_participant_ && identity == local_participant_->identity()) { participant = local_participant_.get(); } else { auto it = remote_participants_.find(identity); if (it != remote_participants_.end()) { participant = it->second.get(); } } if (!participant) { LK_LOG_WARN("kParticipantsUpdated: participant does not exist: {}", identity); continue; } participant->name_ = info.name(); participant->metadata_ = info.metadata(); std::unordered_map attrs; attrs.reserve(info.attributes_size()); for (const auto& kv : info.attributes()) { attrs.emplace(kv.first, kv.second); } participant->attributes_ = std::move(attrs); participant->kind_ = fromProto(info.kind()); participant->reason_ = toDisconnectReason(info.disconnect_reason()); ev.participants.push_back(participant); } } if (delegate_snapshot) { delegate_snapshot->onParticipantsUpdated(*this, ev); } break; } case proto::RoomEvent::MESSAGE_NOT_SET: default: break; } break; } default: break; } } } // namespace livekit