Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 30 additions & 2 deletions include/livekit/local_data_track.h
Original file line number Diff line number Diff line change
Expand Up @@ -67,17 +67,45 @@ class LIVEKIT_API LocalDataTrack {

/// Try to push a frame to all subscribers of this track.
///
/// @return success on delivery acceptance, or a typed error describing why
/// An empty payload succeeds without sending a frame to subscribers.
///
/// @param frame The frame to push.
/// @return Success on delivery acceptance, or a typed error describing why
/// the frame could not be queued.
Result<void, LocalDataTrackTryPushError> tryPush(const DataTrackFrame& frame);

/// Try to push a frame to all subscribers of this track.
///
/// @return success on delivery acceptance, or a typed error describing why
/// An empty payload succeeds without sending a frame to subscribers.
///
/// @param payload The payload to push.
/// @param user_timestamp Optional application-defined timestamp. The unit is
/// caller-defined.
/// @return Success on delivery acceptance, or a typed error describing why
/// the frame could not be queued.
Result<void, LocalDataTrackTryPushError> tryPush(std::vector<std::uint8_t>&& payload,
std::optional<std::uint64_t> user_timestamp = std::nullopt);

/// Try to push a frame from a borrowed byte buffer.
///
/// Copies @p size bytes from @p data into an FFI request before returning;
/// the SDK does not retain the buffer.
/// An empty payload succeeds without sending a frame to subscribers.
///
/// @note This avoids an intermediate copy when the caller already owns a byte
/// buffer, but it is not a zero-copy send.
///
/// @param data Pointer to @p size payload bytes. Must be non-null when
/// @p size is non-zero.
/// @param size Number of bytes at @p data. A value of zero is a successful
/// no-op.
/// @param user_timestamp Optional application-defined timestamp. The unit is
/// caller-defined.
/// @return Success on delivery acceptance, or a typed error describing why
/// the frame could not be queued.
Result<void, LocalDataTrackTryPushError> tryPush(const std::uint8_t* data, std::size_t size,
std::optional<std::uint64_t> user_timestamp = std::nullopt);

/// Whether the track is still published in the room.
bool isPublished() const;

Expand Down
32 changes: 21 additions & 11 deletions src/local_data_track.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,19 +29,37 @@ LocalDataTrack::LocalDataTrack(const proto::OwnedLocalDataTrack& owned)
: handle_(static_cast<uintptr_t>(owned.handle().id())), info_(fromProto(owned.info())) {}

Result<void, LocalDataTrackTryPushError> LocalDataTrack::tryPush(const DataTrackFrame& frame) {
return tryPush(frame.payload.data(), frame.payload.size(), frame.user_timestamp);
}

Result<void, LocalDataTrackTryPushError> LocalDataTrack::tryPush(std::vector<std::uint8_t>&& payload,
std::optional<std::uint64_t> user_timestamp) {
const DataTrackFrame frame(std::move(payload), user_timestamp);
return tryPush(frame);
}

Result<void, LocalDataTrackTryPushError> LocalDataTrack::tryPush(const std::uint8_t* data, std::size_t size,
std::optional<std::uint64_t> user_timestamp) {
if (!handle_.valid()) {
return Result<void, LocalDataTrackTryPushError>::failure(LocalDataTrackTryPushError{
LocalDataTrackTryPushErrorCode::INVALID_HANDLE, "LocalDataTrack::tryPush: invalid FFI handle"});
}
if (size == 0) {
return Result<void, LocalDataTrackTryPushError>::success();
}
if (data == nullptr) {
return Result<void, LocalDataTrackTryPushError>::failure(LocalDataTrackTryPushError{
LocalDataTrackTryPushErrorCode::INTERNAL, "LocalDataTrack::tryPush: payload pointer is null"});
}

try {
proto::FfiRequest req;
auto* msg = req.mutable_local_data_track_try_push();
msg->set_track_handle(static_cast<uint64_t>(handle_.get()));
auto* pf = msg->mutable_frame();
pf->set_payload(frame.payload.data(), frame.payload.size());
if (frame.user_timestamp.has_value()) {
pf->set_user_timestamp(frame.user_timestamp.value());
pf->set_payload(data, size);
if (user_timestamp.has_value()) {
pf->set_user_timestamp(user_timestamp.value());
}

const proto::FfiResponse resp = FfiClient::instance().sendRequest(req);
Expand All @@ -56,14 +74,6 @@ Result<void, LocalDataTrackTryPushError> LocalDataTrack::tryPush(const DataTrack
}
}

Result<void, LocalDataTrackTryPushError> LocalDataTrack::tryPush(std::vector<std::uint8_t>&& payload,
std::optional<std::uint64_t> user_timestamp) {
DataTrackFrame frame;
frame.payload = std::move(payload);
frame.user_timestamp = user_timestamp;
return tryPush(frame);
}

bool LocalDataTrack::isPublished() const {
if (!handle_.valid()) {
return false;
Expand Down
109 changes: 109 additions & 0 deletions src/tests/integration/test_data_track.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -955,6 +955,115 @@ TEST_F(DataTrackE2ETest, PreservesUserTimestampEndToEnd) {
EXPECT_EQ(frame.user_timestamp.value(), sent_timestamp);
}

TEST_F(DataTrackE2ETest, RejectsNullBorrowedPayloadWithNonZeroSize) {
const auto track_name = makeTrackName("null_borrowed_payload");
auto rooms = testRooms(1);
auto local_track = requirePublishedTrack(rooms[0]->localParticipant(), track_name);

const auto push_result = local_track->tryPush(static_cast<const std::uint8_t*>(nullptr), 4);
ASSERT_FALSE(push_result);
EXPECT_EQ(push_result.error().code, LocalDataTrackTryPushErrorCode::INTERNAL);
EXPECT_FALSE(push_result.error().message.empty());

local_track->unpublishDataTrack();
}

TEST_F(DataTrackE2ETest, AcceptsEmptyPayload) {
const auto track_name = makeTrackName("empty_borrowed_payload");
auto rooms = testRooms(1);
auto local_track = requirePublishedTrack(rooms[0]->localParticipant(), track_name);

const std::uint8_t payload = 0;
EXPECT_TRUE(local_track->tryPush(&payload, 0));
EXPECT_TRUE(local_track->tryPush(nullptr, 0));
EXPECT_TRUE(local_track->tryPush(std::vector<std::uint8_t>{}));

local_track->unpublishDataTrack();
}

TEST_F(DataTrackE2ETest, CopiesBorrowedPayloadBeforeTryPushReturns) {
const auto track_name = makeTrackName("borrowed_payload");
const auto sent_timestamp = getTimestampUs();
constexpr std::uint8_t kOriginal = 0xA5;
constexpr std::uint8_t kMutated = 0x5A;
constexpr std::size_t kPayloadSize = 64;

DataTrackPublishedDelegate subscriber_delegate;
std::vector<TestRoomConnectionOptions> room_configs(2);
room_configs[1].delegate = &subscriber_delegate;

auto rooms = testRooms(room_configs);
auto& publisher_room = rooms[0];

auto publish_result = lockLocalParticipant(*publisher_room)->publishDataTrack(track_name);
if (!publish_result) {
FAIL() << describeDataTrackError(publish_result.error());
}
auto local_track = publish_result.value();
ASSERT_TRUE(local_track->isPublished());

auto remote_track = subscriber_delegate.waitForTrack(kTrackWaitTimeout);
ASSERT_NE(remote_track, nullptr) << "Timed out waiting for remote data track";

auto subscribe_result = remote_track->subscribe();
if (!subscribe_result) {
FAIL() << describeDataTrackError(subscribe_result.error());
}
auto subscription = subscribe_result.value();

std::promise<DataTrackFrame> frame_promise;
auto frame_future = frame_promise.get_future();
std::thread reader([&]() {
try {
DataTrackFrame frame;
if (!subscription->read(frame)) {
throw std::runtime_error("Subscription ended before borrowed-payload frame arrived");
}
frame_promise.set_value(std::move(frame));
} catch (...) {
frame_promise.set_exception(std::current_exception());
}
});

bool pushed = false;
for (int attempt = 0; attempt < kTimestampFrameAttempts; ++attempt) {
std::array<std::uint8_t, kPayloadSize> payload{};
payload.fill(kOriginal);
auto push_result = local_track->tryPush(payload.data(), payload.size(), sent_timestamp);
payload.fill(kMutated);
pushed = static_cast<bool>(push_result) || pushed;
if (frame_future.wait_for(25ms) == std::future_status::ready) {
break;
}
}
const auto frame_status = frame_future.wait_for(5s);

if (frame_status != std::future_status::ready) {
subscription->close();
}

subscription->close();
reader.join();
local_track->unpublishDataTrack();

ASSERT_TRUE(pushed) << "Failed to push borrowed data frame";
ASSERT_EQ(frame_status, std::future_status::ready) << "Timed out waiting for borrowed-payload frame";

DataTrackFrame frame;
try {
frame = frame_future.get();
} catch (const std::exception& e) {
FAIL() << e.what();
}

ASSERT_EQ(frame.payload.size(), kPayloadSize);
EXPECT_TRUE(std::all_of(frame.payload.begin(), frame.payload.end(), [expected = kOriginal](std::uint8_t byte) {
return byte == expected;
})) << "Received payload reflects caller mutation after tryPush returned";
ASSERT_TRUE(frame.user_timestamp.has_value());
EXPECT_EQ(frame.user_timestamp.value(), sent_timestamp);
}

TEST_F(DataTrackE2ETest, PublishesAndReceivesEncryptedFramesEndToEnd) {
runEncryptedDataTrackRoundTrip(kDefaultKeyDerivationFunction, "e2ee_transport");
}
Expand Down
Loading