diff --git a/include/livekit/local_track_publication.h b/include/livekit/local_track_publication.h index 412b4fb9..bae23d94 100644 --- a/include/livekit/local_track_publication.h +++ b/include/livekit/local_track_publication.h @@ -31,6 +31,9 @@ class LIVEKIT_API LocalTrackPublication : public TrackPublication { /// Note, this LocalTrackPublication is constructed internally only; /// safe to accept proto::OwnedTrackPublication. explicit LocalTrackPublication(const proto::OwnedTrackPublication& owned); + +private: + friend class Room; }; } // namespace livekit diff --git a/src/ffi_client.cpp b/src/ffi_client.cpp index 6c924560..867eba94 100644 --- a/src/ffi_client.cpp +++ b/src/ffi_client.cpp @@ -103,6 +103,8 @@ std::optional ExtractAsyncId(const proto::FfiEvent& event) { return event.get_stats().async_id(); case E::kGetSessionStats: return event.get_session_stats().async_id(); + case E::kSimulateScenario: + return event.simulate_scenario().async_id(); case E::kPublishSipDtmf: return event.publish_sip_dtmf().async_id(); case E::kChatMessage: @@ -655,6 +657,42 @@ std::future FfiClient::getSessionStatsAsync(uintptr_t room_handle) return fut; } +std::future FfiClient::simulateScenarioAsync(uintptr_t room_handle, int scenario) { + const AsyncId async_id = generateAsyncId(); + + auto fut = registerAsync( + async_id, + [async_id](const proto::FfiEvent& event) { + return event.has_simulate_scenario() && event.simulate_scenario().async_id() == async_id; + }, + [](const proto::FfiEvent& event, std::promise& pr) { + const auto& cb = event.simulate_scenario(); + if (cb.has_error() && !cb.error().empty()) { + pr.set_exception(std::make_exception_ptr(std::runtime_error(cb.error()))); + return; + } + pr.set_value(); + }); + + proto::FfiRequest req; + auto* msg = req.mutable_simulate_scenario(); + msg->set_room_handle(room_handle); + msg->set_scenario(static_cast(scenario)); + msg->set_request_async_id(async_id); + + try { + const proto::FfiResponse resp = sendRequest(req); + if (!resp.has_simulate_scenario()) { + logAndThrow("FfiResponse missing simulate_scenario"); + } + } catch (...) { + cancelPendingByAsyncId(async_id); + throw; + } + + return fut; +} + // Participant APIs Implementation std::future FfiClient::publishTrackAsync(std::uint64_t local_participant_handle, std::uint64_t track_handle, diff --git a/src/ffi_client.h b/src/ffi_client.h index c0835027..0104de76 100644 --- a/src/ffi_client.h +++ b/src/ffi_client.h @@ -106,6 +106,8 @@ class LIVEKIT_INTERNAL_API FfiClient { std::future getSessionStatsAsync(uintptr_t room_handle); + std::future simulateScenarioAsync(uintptr_t room_handle, int scenario); + // Participant APIs std::future publishTrackAsync(std::uint64_t local_participant_handle, std::uint64_t track_handle, diff --git a/src/room.cpp b/src/room.cpp index f37c5642..ff1dcd19 100644 --- a/src/room.cpp +++ b/src/room.cpp @@ -20,7 +20,10 @@ #include "ffi_client.h" #include "livekit/audio_stream.h" #include "livekit/e2ee.h" +#include "livekit/local_audio_track.h" #include "livekit/local_participant.h" +#include "livekit/local_track_publication.h" +#include "livekit/local_video_track.h" #include "livekit/remote_audio_track.h" #include "livekit/remote_data_track.h" #include "livekit/remote_participant.h" @@ -47,6 +50,33 @@ using proto::FfiResponse; namespace { +std::shared_ptr localTrackPublication(const std::shared_ptr& track) { + if (!track) { + return nullptr; + } + if (auto video = std::dynamic_pointer_cast(track)) { + return video->publication(); + } + if (auto audio = std::dynamic_pointer_cast(track)) { + return audio->publication(); + } + return nullptr; +} + +void updateLocalTrackPublicationInfo(LocalTrackPublication& publication, const proto::TrackPublicationInfo& info) { + publication.sid_ = info.sid(); + publication.name_ = info.name(); + publication.kind_ = fromProto(info.kind()); + publication.source_ = fromProto(info.source()); + publication.simulcasted_ = info.simulcasted(); + publication.width_ = info.width(); + publication.height_ = info.height(); + publication.mime_type_ = info.mime_type(); + publication.muted_ = info.muted(); + publication.encryption_type_ = static_cast(info.encryption_type()); + publication.audio_features_ = convertAudioFeatures(info.audio_features()); +} + std::shared_ptr createRemoteParticipant(const proto::OwnedParticipant& owned) { const auto& pinfo = owned.info(); std::unordered_map attrs; @@ -593,6 +623,38 @@ void Room::onEvent(const FfiEvent& event) { } break; } + case proto::RoomEvent::kLocalTrackRepublished: { + const std::scoped_lock guard(lock_); + if (!local_participant_) { + LK_LOG_ERROR("kLocalTrackRepublished: local_participant_ is nullptr"); + break; + } + const auto& ltr = re.local_track_republished(); + const std::string& previous_sid = ltr.previous_sid(); + + auto& published = local_participant_->published_tracks_by_sid_; + auto it = published.find(previous_sid); + if (it == published.end()) { + LK_LOG_WARN("local_track_republished for unknown previous sid: {}", previous_sid); + break; + } + auto track = it->second.lock(); + if (!track) { + published.erase(it); + LK_LOG_WARN("local_track_republished for expired previous sid: {}", previous_sid); + break; + } + auto publication = localTrackPublication(track); + if (!publication) { + LK_LOG_WARN("local_track_republished missing publication for sid: {}", previous_sid); + break; + } + updateLocalTrackPublicationInfo(*publication, ltr.info()); + published.erase(it); + published[publication->sid()] = track; + track->setPublication(publication); + break; + } case proto::RoomEvent::kLocalTrackSubscribed: { LocalTrackSubscribedEvent ev; { diff --git a/src/tests/common/room_test_access.h b/src/tests/common/room_test_access.h new file mode 100644 index 00000000..9509b00c --- /dev/null +++ b/src/tests/common/room_test_access.h @@ -0,0 +1,49 @@ +/* + * Copyright 2026 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. + */ + +#pragma once + +#include + +#include +#include +#include + +#include "ffi_client.h" +#include "room.pb.h" + +namespace livekit { + +struct RoomTestAccess { + static int listenerId(const Room& room) { + const std::scoped_lock guard(room.lock_); + return room.listener_id_; + } + + static void simulateScenario(Room& room, proto::SimulateScenarioKind scenario) { + std::shared_ptr handle; + { + const std::scoped_lock guard(room.lock_); + handle = room.room_handle_; + } + if (!handle) { + throw std::runtime_error("cannot simulate scenario for a disconnected room"); + } + FfiClient::instance().simulateScenarioAsync(handle->get(), static_cast(scenario)).get(); + } +}; + +} // namespace livekit diff --git a/src/tests/integration/test_local_track_publish_sid.cpp b/src/tests/integration/test_local_track_publish_sid.cpp index 91fbf33b..ee767450 100644 --- a/src/tests/integration/test_local_track_publish_sid.cpp +++ b/src/tests/integration/test_local_track_publish_sid.cpp @@ -16,15 +16,21 @@ #include +#include #include #include +#include #include "../common/audio_utils.h" +#include "../common/room_test_access.h" #include "../common/test_common.h" +#include "room.pb.h" namespace livekit::test { namespace { +using namespace std::chrono_literals; + void expectTrackSidAssigned(const Track& track, const LocalTrackPublication& publication) { const std::string& track_sid = track.sid(); const std::string& publication_sid = publication.sid(); @@ -33,6 +39,17 @@ void expectTrackSidAssigned(const Track& track, const LocalTrackPublication& pub EXPECT_EQ(track_sid, publication_sid); } +bool waitForSidChange(const Track& track, const std::string& previous_sid, std::chrono::milliseconds timeout) { + const auto deadline = std::chrono::steady_clock::now() + timeout; + while (std::chrono::steady_clock::now() < deadline) { + if (track.sid() != previous_sid && track.sid() != "TR_unknown") { + return true; + } + std::this_thread::sleep_for(50ms); + } + return track.sid() != previous_sid && track.sid() != "TR_unknown"; +} + } // namespace class LocalTrackPublishSidTest : public LiveKitTestBase {}; @@ -73,4 +90,34 @@ TEST_F(LocalTrackPublishSidTest, PublishAudioTrackAssignsSid) { lockLocalParticipant(room)->unpublishTrack(track->publication()->sid()); } +TEST_F(LocalTrackPublishSidTest, FullReconnectUpdatesPublishedSid) { + failIfNotConfigured(); + + Room room; + const RoomOptions room_options; + ASSERT_TRUE(room.connect(config_.url, config_.token_a, room_options)); + + auto source = std::make_shared(VideoCodec::H264, 16, 16); + std::shared_ptr track; + ASSERT_NO_THROW( + track = lockLocalParticipant(room)->publishVideoTrack("republish-sid-check", source, TrackSource::SOURCE_CAMERA)); + ASSERT_NE(track, nullptr); + ASSERT_NE(track->publication(), nullptr); + expectTrackSidAssigned(*track, *track->publication()); + + const std::string previous_sid = track->sid(); + ASSERT_NO_THROW(RoomTestAccess::simulateScenario(room, proto::SIMULATE_FULL_RECONNECT)); + ASSERT_TRUE(waitForSidChange(*track, previous_sid, 30s)) << "Timed out waiting for republished track SID"; + + ASSERT_NE(track->publication(), nullptr); + expectTrackSidAssigned(*track, *track->publication()); + EXPECT_NE(track->sid(), previous_sid); + + const auto pubs = lockLocalParticipant(room)->trackPublications(); + EXPECT_EQ(pubs.count(previous_sid), 0u); + EXPECT_EQ(pubs.count(track->sid()), 1u); + + lockLocalParticipant(room)->unpublishTrack(track->sid()); +} + } // namespace livekit::test diff --git a/src/tests/integration/test_room.cpp b/src/tests/integration/test_room.cpp index 942d8939..b6f3c80a 100644 --- a/src/tests/integration/test_room.cpp +++ b/src/tests/integration/test_room.cpp @@ -29,21 +29,11 @@ #include #include +#include "../common/room_test_access.h" #include "../common/test_common.h" using namespace std::chrono_literals; -namespace livekit { - -struct RoomTestAccess { - static int listenerId(const Room& room) { - const std::scoped_lock guard(room.lock_); - return room.listener_id_; - } -}; - -} // namespace livekit - namespace livekit::test { // Server-dependent tests - require LIVEKIT_URL and LIVEKIT_TOKEN_A env vars