diff --git a/.changeset/quiet-ducks-yell.md b/.changeset/quiet-ducks-yell.md new file mode 100644 index 00000000..424930c7 --- /dev/null +++ b/.changeset/quiet-ducks-yell.md @@ -0,0 +1,5 @@ +--- +'@livekit/rtc-node': patch +--- + +Convert data streams to use livekit-ffi exposed data streams interface diff --git a/packages/livekit-rtc/package.json b/packages/livekit-rtc/package.json index b4ce7c47..a53495fa 100644 --- a/packages/livekit-rtc/package.json +++ b/packages/livekit-rtc/package.json @@ -32,7 +32,7 @@ "dependencies": { "@datastructures-js/deque": "1.0.8", "@livekit/mutex": "^1.0.0", - "@livekit/rtc-ffi-bindings": "0.12.60", + "@livekit/rtc-ffi-bindings": "0.12.71", "@livekit/typed-emitter": "^3.0.0", "pino": "^9.0.0", "pino-pretty": "^13.0.0" diff --git a/packages/livekit-rtc/src/data_streams/stream_reader.ts b/packages/livekit-rtc/src/data_streams/stream_reader.ts index 7067d69d..e6ffd63f 100644 --- a/packages/livekit-rtc/src/data_streams/stream_reader.ts +++ b/packages/livekit-rtc/src/data_streams/stream_reader.ts @@ -66,7 +66,10 @@ export class ByteStreamReader extends BaseStreamReader { // consumer never calls return() (e.g. breaking out of for-await). reader.releaseLock(); log.error('error processing stream update: %s', error); - return { done: true, value: undefined as unknown }; + // Propagate abnormal termination (e.g. remote abort, payload over + // the receiver's size limit) instead of presenting the truncated + // payload as a clean EOF. + throw error; } }, @@ -153,7 +156,10 @@ export class TextStreamReader extends BaseStreamReader { reader.releaseLock(); receivedChunks.clear(); log.error('error processing stream update: %s', error); - return { done: true, value: undefined }; + // Propagate abnormal termination (e.g. remote abort, payload over + // the receiver's size limit) instead of presenting the truncated + // payload as a clean EOF. + throw error; } }, diff --git a/packages/livekit-rtc/src/data_streams/types.ts b/packages/livekit-rtc/src/data_streams/types.ts index 5f7c9ff4..24f5c3be 100644 --- a/packages/livekit-rtc/src/data_streams/types.ts +++ b/packages/livekit-rtc/src/data_streams/types.ts @@ -42,6 +42,12 @@ export interface TextStreamOptions extends DataStreamOptions { export interface ByteStreamOptions extends DataStreamOptions { name?: string; onProgress?: (progress: number) => void; + /** + * Whether the payload may be compressed on the wire. Defaults to true; + * compression is only applied when it actually reduces the payload size + * and every recipient supports it. + */ + compress?: boolean; } export type ByteStreamHandler = ( diff --git a/packages/livekit-rtc/src/participant.ts b/packages/livekit-rtc/src/participant.ts index 87372e8f..4b55e4ef 100644 --- a/packages/livekit-rtc/src/participant.ts +++ b/packages/livekit-rtc/src/participant.ts @@ -1,7 +1,7 @@ // SPDX-FileCopyrightText: 2024 LiveKit, Inc. // // SPDX-License-Identifier: Apache-2.0 -import { Mutex } from '@livekit/mutex'; +import { type Mutex } from '@livekit/mutex'; import { DisconnectReason, type OwnedParticipant, @@ -10,13 +10,16 @@ import { type ParticipantKindDetail, } from '@livekit/rtc-ffi-bindings'; import { + type ByteStreamOpenCallback, + ByteStreamOpenRequest, + type ByteStreamOpenResponse, + type ByteStreamWriterCloseCallback, + ByteStreamWriterCloseRequest, + type ByteStreamWriterCloseResponse, + type ByteStreamWriterWriteCallback, + ByteStreamWriterWriteRequest, + type ByteStreamWriterWriteResponse, ChatMessage as ChatMessageModel, - DataStream_ByteHeader, - DataStream_Chunk, - DataStream_Header, - DataStream_OperationType, - DataStream_TextHeader, - DataStream_Trailer, EditChatMessageRequest, TranscriptionSegment as ProtoTranscriptionSegment, type PublishDataCallback, @@ -34,15 +37,6 @@ import { type SendChatMessageCallback, SendChatMessageRequest, type SendChatMessageResponse, - type SendStreamChunkCallback, - SendStreamChunkRequest, - type SendStreamChunkResponse, - type SendStreamHeaderCallback, - SendStreamHeaderRequest, - type SendStreamHeaderResponse, - type SendStreamTrailerCallback, - SendStreamTrailerRequest, - type SendStreamTrailerResponse, type SetLocalAttributesCallback, SetLocalAttributesRequest, type SetLocalAttributesResponse, @@ -52,6 +46,23 @@ import { type SetLocalNameCallback, SetLocalNameRequest, type SetLocalNameResponse, + StreamByteOptions, + type StreamSendFileCallback, + StreamSendFileRequest, + type StreamSendFileResponse, + type StreamSendTextCallback, + StreamSendTextRequest, + type StreamSendTextResponse, + StreamTextOptions, + type TextStreamOpenCallback, + TextStreamOpenRequest, + type TextStreamOpenResponse, + type TextStreamWriterCloseCallback, + TextStreamWriterCloseRequest, + type TextStreamWriterCloseResponse, + type TextStreamWriterWriteCallback, + TextStreamWriterWriteRequest, + type TextStreamWriterWriteResponse, type TrackPublishOptions, type UnpublishTrackCallback, UnpublishTrackRequest, @@ -71,12 +82,10 @@ import { UnregisterRpcMethodRequest, } from '@livekit/rtc-ffi-bindings'; import type { PathLike } from 'node:fs'; -import { open, stat } from 'node:fs/promises'; +import { fileURLToPath } from 'node:url'; import { - type ByteStreamInfo, type ByteStreamOptions, ByteStreamWriter, - type TextStreamInfo, TextStreamWriter, } from './data_streams/index.js'; import { FfiClient, FfiHandle } from './ffi_client.js'; @@ -87,9 +96,8 @@ import type { RemoteTrackPublication, TrackPublication } from './track_publicati import { LocalTrackPublication } from './track_publication.js'; import type { Transcription } from './transcription.js'; import type { ChatMessage } from './types.js'; -import { numberToBigInt, splitUtf8 } from './utils.js'; - -const STREAM_CHUNK_SIZE = 15_000; +import { byteStreamInfoFromProto, textStreamInfoFromProto } from './utils.js'; +import { numberToBigInt } from './utils.js'; export abstract class Participant { /** @internal */ @@ -278,92 +286,54 @@ export class LocalParticipant extends Participant { streamId?: string; senderIdentity?: string; }): Promise { - const senderIdentity = options?.senderIdentity ?? this.identity; - const streamId = options?.streamId ?? crypto.randomUUID(); - const destinationIdentities = options?.destinationIdentities; - - const info: TextStreamInfo = { - streamId: streamId, - mimeType: 'text/plain', - topic: options?.topic ?? '', - timestamp: Date.now(), - }; - - const headerReq = new SendStreamHeaderRequest({ - senderIdentity, - destinationIdentities, + const req = new TextStreamOpenRequest({ localParticipantHandle: this.ffi_handle.handle, - header: new DataStream_Header({ - streamId, - mimeType: info.mimeType, - topic: info.topic, - timestamp: numberToBigInt(info.timestamp), + options: new StreamTextOptions({ + topic: options?.topic ?? '', attributes: options?.attributes, - contentHeader: { - case: 'textHeader', - value: new DataStream_TextHeader({ - operationType: DataStream_OperationType.CREATE, - version: 0, - replyToStreamId: '', - generated: false, - }), - }, + destinationIdentities: options?.destinationIdentities, + id: options?.streamId, + senderIdentity: options?.senderIdentity ?? this.identity, }), }); - await this.sendStreamHeader(headerReq); + const res = FfiClient.instance.request({ + message: { case: 'textStreamOpen', value: req }, + }); + + const cb = await FfiClient.instance.waitFor( + (ev) => ev.message.case == 'textStreamOpen' && ev.message.value.asyncId == res.asyncId, + { signal: this.disconnectSignal }, + ); + + if (cb.result.case !== 'writer') { + throw new Error(cb.result.case === 'error' ? cb.result.value.description : 'unknown error'); + } - let nextChunkId = 0; - const localHandle = this.ffi_handle.handle; - const sendTrailer = this.sendStreamTrailer; - const sendChunk = this.sendStreamChunk; + const writerHandle = cb.result.value.handle!.id!; + const info = textStreamInfoFromProto(cb.result.value.info!); + const writeText = this.textStreamWriterWrite; + const closeWriter = this.textStreamWriterClose; const writableStream = new WritableStream({ // Implement the sink async write(text) { - for (const textByteChunk of splitUtf8(text, STREAM_CHUNK_SIZE)) { - const chunkRequest = new SendStreamChunkRequest({ - senderIdentity, - localParticipantHandle: localHandle, - destinationIdentities, - chunk: new DataStream_Chunk({ - content: textByteChunk, - streamId, - chunkIndex: numberToBigInt(nextChunkId), - }), - }); - - await sendChunk(chunkRequest); - nextChunkId += 1; - } + await writeText(new TextStreamWriterWriteRequest({ writerHandle, text })); }, async close() { - const trailerReq = new SendStreamTrailerRequest({ - senderIdentity, - localParticipantHandle: localHandle, - destinationIdentities, - trailer: new DataStream_Trailer({ - streamId, - reason: '', - }), - }); - await sendTrailer(trailerReq); + await closeWriter(new TextStreamWriterCloseRequest({ writerHandle })); }, - // Send a trailer with the error reason so the remote side's stream + // Close the stream with the error reason so the remote side's stream // controller is closed instead of waiting for data that won't arrive. async abort(err) { log.error(err, 'Sink Error'); try { - const trailerReq = new SendStreamTrailerRequest({ - senderIdentity, - localParticipantHandle: localHandle, - destinationIdentities, - trailer: new DataStream_Trailer({ - streamId, + await closeWriter( + new TextStreamWriterCloseRequest({ + writerHandle, reason: err instanceof Error ? err.message : String(err ?? ''), }), - }); - await sendTrailer(trailerReq); + ); } catch { // Best-effort: the connection may already be gone. } @@ -382,12 +352,41 @@ export class LocalParticipant extends Participant { attributes?: Record; destinationIdentities?: Array; streamId?: string; + /** + * Whether the payload may be compressed on the wire. Defaults to true; + * compression is only applied when it actually reduces the payload size + * and every recipient supports it. + */ + compress?: boolean; }, ) { - const writer = await this.streamText(options); - await writer.write(text); - await writer.close(); - return writer.info; + const req = new StreamSendTextRequest({ + localParticipantHandle: this.ffi_handle.handle, + options: new StreamTextOptions({ + topic: options?.topic ?? '', + attributes: options?.attributes, + destinationIdentities: options?.destinationIdentities, + id: options?.streamId, + senderIdentity: this.identity, + compress: options?.compress, + }), + text, + }); + + const res = FfiClient.instance.request({ + message: { case: 'sendText', value: req }, + }); + + const cb = await FfiClient.instance.waitFor( + (ev) => ev.message.case == 'sendText' && ev.message.value.asyncId == res.asyncId, + { signal: this.disconnectSignal }, + ); + + if (cb.result.case !== 'info') { + throw new Error(cb.result.case === 'error' ? cb.result.value.description : 'unknown error'); + } + + return textStreamInfoFromProto(cb.result.value); } async streamBytes(options?: { @@ -399,102 +398,56 @@ export class LocalParticipant extends Participant { mimeType?: string; totalSize?: number; }) { - const senderIdentity = this.identity; - const streamId = options?.streamId ?? crypto.randomUUID(); - const destinationIdentities = options?.destinationIdentities; - - const info: ByteStreamInfo = { - streamId: streamId, - mimeType: options?.mimeType ?? 'application/octet-stream', - topic: options?.topic ?? '', - timestamp: Date.now(), - attributes: options?.attributes, - totalSize: options?.totalSize, - name: options?.name ?? 'unknown', - }; - - const headerReq = new SendStreamHeaderRequest({ - senderIdentity, - destinationIdentities, + const req = new ByteStreamOpenRequest({ localParticipantHandle: this.ffi_handle.handle, - header: new DataStream_Header({ - streamId, - mimeType: info.mimeType, - topic: info.topic, - timestamp: numberToBigInt(info.timestamp), - attributes: info.attributes, - totalLength: numberToBigInt(info.totalSize), - contentHeader: { - case: 'byteHeader', - value: new DataStream_ByteHeader({ - name: info.name, - }), - }, + options: new StreamByteOptions({ + topic: options?.topic ?? '', + attributes: options?.attributes, + destinationIdentities: options?.destinationIdentities, + id: options?.streamId, + name: options?.name ?? 'unknown', + mimeType: options?.mimeType ?? 'application/octet-stream', + totalLength: numberToBigInt(options?.totalSize), + senderIdentity: this.identity, }), }); - await this.sendStreamHeader(headerReq); + const res = FfiClient.instance.request({ + message: { case: 'byteStreamOpen', value: req }, + }); + + const cb = await FfiClient.instance.waitFor( + (ev) => ev.message.case == 'byteStreamOpen' && ev.message.value.asyncId == res.asyncId, + { signal: this.disconnectSignal }, + ); + + if (cb.result.case !== 'writer') { + throw new Error(cb.result.case === 'error' ? cb.result.value.description : 'unknown error'); + } - let chunkId = 0; - const localHandle = this.ffi_handle.handle; - const sendTrailer = this.sendStreamTrailer; - const sendChunk = this.sendStreamChunk; - const writeMutex = new Mutex(); + const writerHandle = cb.result.value.handle!.id!; + const info = byteStreamInfoFromProto(cb.result.value.info!); + const writeBytes = this.byteStreamWriterWrite; + const closeWriter = this.byteStreamWriterClose; const writableStream = new WritableStream({ async write(chunk) { - const unlock = await writeMutex.lock(); - - let byteOffset = 0; - try { - while (byteOffset < chunk.byteLength) { - const subChunk = chunk.slice(byteOffset, byteOffset + STREAM_CHUNK_SIZE); - const chunkRequest = new SendStreamChunkRequest({ - senderIdentity, - localParticipantHandle: localHandle, - destinationIdentities, - chunk: new DataStream_Chunk({ - content: subChunk, - streamId, - chunkIndex: numberToBigInt(chunkId), - }), - }); - - await sendChunk(chunkRequest); - chunkId += 1; - byteOffset += subChunk.byteLength; - } - } finally { - unlock(); - } + await writeBytes(new ByteStreamWriterWriteRequest({ writerHandle, bytes: chunk })); }, async close() { - const trailerReq = new SendStreamTrailerRequest({ - senderIdentity, - localParticipantHandle: localHandle, - destinationIdentities, - trailer: new DataStream_Trailer({ - streamId, - reason: '', - }), - }); - await sendTrailer(trailerReq); + await closeWriter(new ByteStreamWriterCloseRequest({ writerHandle })); }, - // Send a trailer with the error reason so the remote side's stream + // Close the stream with the error reason so the remote side's stream // controller is closed instead of waiting for data that won't arrive. async abort(err) { log.error(err, 'Sink error'); try { - const trailerReq = new SendStreamTrailerRequest({ - senderIdentity, - localParticipantHandle: localHandle, - destinationIdentities, - trailer: new DataStream_Trailer({ - streamId, + await closeWriter( + new ByteStreamWriterCloseRequest({ + writerHandle, reason: err instanceof Error ? err.message : String(err ?? ''), }), - }); - await sendTrailer(trailerReq); + ); } catch { // Best-effort: the connection may already be gone. } @@ -508,77 +461,97 @@ export class LocalParticipant extends Participant { /** Sends a file provided as PathLike to specified recipients */ async sendFile(path: PathLike, options?: ByteStreamOptions) { - const fileStats = await stat(path); - const file = await open(path); - try { - const stream: ReadableStream = file.readableWebStream(); - const streamId = crypto.randomUUID(); - const destinationIdentities = options?.destinationIdentities; - - const writer = await this.streamBytes({ - streamId: streamId, - name: options?.name, - totalSize: fileStats.size, - destinationIdentities, - topic: options?.topic, - mimeType: options?.mimeType, + const filePath = path instanceof URL ? fileURLToPath(path) : path.toString(); + + const req = new StreamSendFileRequest({ + localParticipantHandle: this.ffi_handle.handle, + options: new StreamByteOptions({ + topic: options?.topic ?? '', attributes: options?.attributes, - }); + destinationIdentities: options?.destinationIdentities, + name: options?.name ?? 'unknown', + mimeType: options?.mimeType ?? 'application/octet-stream', + senderIdentity: this.identity, + compress: options?.compress, + }), + filePath, + }); - for await (const chunk of stream) { - await writer.write(chunk); - } - await writer.close(); - } finally { - await file.close(); + const res = FfiClient.instance.request({ + message: { case: 'sendFile', value: req }, + }); + + const cb = await FfiClient.instance.waitFor( + (ev) => ev.message.case == 'sendFile' && ev.message.value.asyncId == res.asyncId, + { signal: this.disconnectSignal }, + ); + + if (cb.result.case === 'error') { + throw new Error(cb.result.value.description); } } - private async sendStreamHeader(req: SendStreamHeaderRequest) { - const type = 'sendStreamHeader'; - const res = FfiClient.instance.request({ - message: { case: type, value: req }, + private textStreamWriterWrite = async (req: TextStreamWriterWriteRequest) => { + const res = FfiClient.instance.request({ + message: { case: 'textStreamWrite', value: req }, }); - const cb = await FfiClient.instance.waitFor( - (ev) => ev.message.case == type && ev.message.value.asyncId == res.asyncId, + const cb = await FfiClient.instance.waitFor( + (ev) => + ev.message.case == 'textStreamWriterWrite' && ev.message.value.asyncId == res.asyncId, { signal: this.disconnectSignal }, ); if (cb.error) { - throw new Error(cb.error); + throw new Error(cb.error.description); } - } + }; - private sendStreamChunk = async (req: SendStreamChunkRequest) => { - const type = 'sendStreamChunk'; - const res = FfiClient.instance.request({ - message: { case: type, value: req }, + private textStreamWriterClose = async (req: TextStreamWriterCloseRequest) => { + const res = FfiClient.instance.request({ + message: { case: 'textStreamClose', value: req }, }); - const cb = await FfiClient.instance.waitFor( - (ev) => ev.message.case == type && ev.message.value.asyncId == res.asyncId, + const cb = await FfiClient.instance.waitFor( + (ev) => + ev.message.case == 'textStreamWriterClose' && ev.message.value.asyncId == res.asyncId, { signal: this.disconnectSignal }, ); if (cb.error) { - throw new Error(cb.error); + throw new Error(cb.error.description); } }; - private sendStreamTrailer = async (req: SendStreamTrailerRequest) => { - const type = 'sendStreamTrailer'; - const res = FfiClient.instance.request({ - message: { case: type, value: req }, + private byteStreamWriterWrite = async (req: ByteStreamWriterWriteRequest) => { + const res = FfiClient.instance.request({ + message: { case: 'byteStreamWrite', value: req }, }); - const cb = await FfiClient.instance.waitFor( - (ev) => ev.message.case == type && ev.message.value.asyncId == res.asyncId, + const cb = await FfiClient.instance.waitFor( + (ev) => + ev.message.case == 'byteStreamWriterWrite' && ev.message.value.asyncId == res.asyncId, { signal: this.disconnectSignal }, ); if (cb.error) { - throw new Error(cb.error); + throw new Error(cb.error.description); + } + }; + + private byteStreamWriterClose = async (req: ByteStreamWriterCloseRequest) => { + const res = FfiClient.instance.request({ + message: { case: 'byteStreamClose', value: req }, + }); + + const cb = await FfiClient.instance.waitFor( + (ev) => + ev.message.case == 'byteStreamWriterClose' && ev.message.value.asyncId == res.asyncId, + { signal: this.disconnectSignal }, + ); + + if (cb.error) { + throw new Error(cb.error.description); } }; diff --git a/packages/livekit-rtc/src/room.ts b/packages/livekit-rtc/src/room.ts index 6ee23de0..62e71d18 100644 --- a/packages/livekit-rtc/src/room.ts +++ b/packages/livekit-rtc/src/room.ts @@ -9,12 +9,10 @@ import type { GetSessionStatsResponse, } from '@livekit/rtc-ffi-bindings'; import { DisconnectReason, type OwnedParticipant } from '@livekit/rtc-ffi-bindings'; -import type { - DataStream_Trailer, - DisconnectCallback, - TrackPublicationInfo, -} from '@livekit/rtc-ffi-bindings'; +import { type DisconnectCallback, type TrackPublicationInfo } from '@livekit/rtc-ffi-bindings'; import { + ByteStreamReaderReadIncrementalRequest, + type ByteStreamReaderReadIncrementalResponse, type ConnectCallback, ConnectRequest, type ConnectResponse, @@ -22,29 +20,27 @@ import { ConnectionState, ContinualGatheringPolicy, type DataPacketKind, - type DataStream_Chunk, - type DataStream_Header, + DataStream_Chunk, type DisconnectResponse, + RoomDataStreamOptions as FfiRoomDataStreamOptions, RoomOptions as FfiRoomOptions, type IceServer, IceTransportType, + type OwnedByteStreamReader, + type OwnedTextStreamReader, type ReadyForRoomEventResponse, type RoomInfo, type SimulateScenarioCallback, type SimulateScenarioKind, type SimulateScenarioResponse, + TextStreamReaderReadIncrementalRequest, + type TextStreamReaderReadIncrementalResponse, } from '@livekit/rtc-ffi-bindings'; import { TrackKind } from '@livekit/rtc-ffi-bindings'; import type { TypedEventEmitter as TypedEmitter } from '@livekit/typed-emitter'; import EventEmitter from 'events'; import { ByteStreamReader, TextStreamReader } from './data_streams/stream_reader.js'; -import type { - ByteStreamHandler, - ByteStreamInfo, - StreamController, - TextStreamHandler, - TextStreamInfo, -} from './data_streams/types.js'; +import { type ByteStreamHandler, type TextStreamHandler } from './data_streams/types.js'; import type { E2EEOptions } from './e2ee.js'; import { E2EEManager, defaultE2EEOptions } from './e2ee.js'; import { FfiClient, FfiClientEvent, FfiHandle } from './ffi_client.js'; @@ -56,7 +52,12 @@ import { RemoteAudioTrack, RemoteVideoTrack } from './track.js'; import type { LocalTrackPublication, TrackPublication } from './track_publication.js'; import { RemoteTrackPublication } from './track_publication.js'; import type { ChatMessage } from './types.js'; -import { bigIntToNumber } from './utils.js'; +import { + bigIntToNumber, + byteStreamInfoFromProto, + numberToBigInt, + textStreamInfoFromProto, +} from './utils.js'; export interface RtcConfiguration { iceTransportType: IceTransportType; @@ -70,6 +71,15 @@ export const defaultRtcConfiguration: RtcConfiguration = { iceServers: [], }; +export interface RoomDataStreamOptions { + /** + * Maximum decompressed payload size in bytes accepted for a single incoming + * data stream. Incoming streams exceeding this limit terminate with an error + * on the receiving side. Unset falls back to the SDK default. + */ + maxPayloadByteLength?: number; +} + export interface RoomOptions { autoSubscribe: boolean; dynacast: boolean; @@ -79,6 +89,8 @@ export interface RoomOptions { e2ee?: E2EEOptions; encryption?: E2EEOptions; rtcConfig?: RtcConfiguration; + /** Options controlling data stream behavior for this room. */ + dataStream?: RoomDataStreamOptions; } export const defaultRoomOptions = new FfiRoomOptions({ @@ -101,8 +113,15 @@ export class Room extends (EventEmitter as new () => TypedEmitter */ private ffiEventLock = new Mutex(); - private byteStreamControllers = new Map>(); - private textStreamControllers = new Map>(); + // Active incoming stream readers keyed by their FFI reader handle id. Each entry + // owns the FfiClient listener feeding reader events into the ReadableStream. + private streamReaders = new Map< + bigint, + { + controller: ReadableStreamDefaultController; + listener: (ev: FfiEvent) => void; + } + >(); private byteStreamHandlers = new Map(); private textStreamHandlers = new Map(); @@ -272,7 +291,10 @@ export class Room extends (EventEmitter as new () => TypedEmitter * @throws ConnectError - if connection fails */ async connect(url: string, token: string, opts?: RoomOptions) { - const options = { ...defaultRoomOptions, ...opts }; + // dataStream is a plain user-facing object; convert it to its proto + // message separately instead of spreading it into the FFI options. + const { dataStream, ...restOpts } = opts ?? {}; + const options = { ...defaultRoomOptions, ...restOpts }; const e2eeEnabled = options.encryption || options.e2ee; const e2eeOptions = options.encryption ? { ...defaultE2EEOptions, ...options.encryption } @@ -281,7 +303,14 @@ export class Room extends (EventEmitter as new () => TypedEmitter const req = new ConnectRequest({ url: url, token: token, - options, + options: { + ...options, + dataStream: dataStream + ? new FfiRoomDataStreamOptions({ + maxPayloadByteLength: numberToBigInt(dataStream.maxPayloadByteLength), + }) + : undefined, + }, }); FfiClient.instance.on(FfiClientEvent.FfiEvent, this.onFfiEvent); @@ -421,28 +450,20 @@ export class Room extends (EventEmitter as new () => TypedEmitter if (this.hasCleanedUp) return; this.hasCleanedUp = true; - // Error all in-progress stream controllers to prevent FD leaks. - // Streams that were receiving data but never got a trailer (e.g. the sender - // disconnected mid-transfer) would otherwise keep their ReadableStream open - // indefinitely, leaking the underlying controller and any buffered chunks. + // Error all in-progress stream readers to prevent FD leaks. + // Streams that were receiving data but never reached end-of-stream (e.g. the + // sender disconnected mid-transfer) would otherwise keep their ReadableStream + // open indefinitely, leaking the underlying controller and any buffered chunks. // Using error() instead of close() signals an abnormal termination to consumers. - for (const [, streamController] of this.byteStreamControllers) { - try { - streamController.controller.error(new Error('Disconnected while receiving')); - } catch { - // controller may already be closed or errored - } - } - this.byteStreamControllers.clear(); - - for (const [, streamController] of this.textStreamControllers) { + for (const [, { controller, listener }] of this.streamReaders) { + FfiClient.instance.off(FfiClientEvent.FfiEvent, listener); try { - streamController.controller.error(new Error('Disconnected while receiving')); + controller.error(new Error('Disconnected while receiving')); } catch { // controller may already be closed or errored } } - this.textStreamControllers.clear(); + this.streamReaders.clear(); // Detach every track from this room so attached FrameProcessors receive // onStreamInfoCleared/onCredentialsCleared and the tokenRefreshed listener @@ -836,12 +857,10 @@ export class Room extends (EventEmitter as new () => TypedEmitter this.emit(RoomEvent.Reconnected); } else if (ev.case == 'roomSidChanged') { this.emit(RoomEvent.RoomSidChanged, ev.value.sid!); - } else if (ev.case === 'streamHeaderReceived' && ev.value.header) { - this.handleStreamHeader(ev.value.header, ev.value.participantIdentity!); - } else if (ev.case === 'streamChunkReceived' && ev.value.chunk) { - this.handleStreamChunk(ev.value.chunk); - } else if (ev.case === 'streamTrailerReceived' && ev.value.trailer) { - this.handleStreamTrailer(ev.value.trailer); + } else if (ev.case === 'byteStreamOpened' && ev.value.reader) { + this.handleByteStreamOpened(ev.value.reader, ev.value.participantIdentity!); + } else if (ev.case === 'textStreamOpened' && ev.value.reader) { + this.handleTextStreamOpened(ev.value.reader, ev.value.participantIdentity!); } else if (ev.case === 'roomUpdated') { this.info = ev.value; this.emit(RoomEvent.RoomUpdated); @@ -926,104 +945,135 @@ export class Room extends (EventEmitter as new () => TypedEmitter return participant; } - private handleStreamHeader(streamHeader: DataStream_Header, participantIdentity: string) { - if (streamHeader.contentHeader.case === 'byteHeader') { - const streamHandlerCallback = this.byteStreamHandlers.get(streamHeader.topic ?? ''); + private handleByteStreamOpened(ownedReader: OwnedByteStreamReader, participantIdentity: string) { + const info = byteStreamInfoFromProto(ownedReader.info!); + const readerHandle = ownedReader.handle!.id!; - if (!streamHandlerCallback) { - log.debug( - 'ignoring incoming byte stream due to no handler for topic: %s', - streamHeader.topic ?? 'undefined', - ); - return; - } - let streamController: ReadableStreamDefaultController; - const stream = new ReadableStream({ - start: (controller) => { - streamController = controller; - this.byteStreamControllers.set(streamHeader.streamId!, { - header: streamHeader, - controller: streamController, - startTime: Date.now(), - }); - }, - }); - const info: ByteStreamInfo = { - streamId: streamHeader.streamId!, - name: streamHeader.contentHeader.value.name ?? 'unknown', - mimeType: streamHeader.mimeType!, - totalSize: streamHeader.totalLength ? Number(streamHeader.totalLength) : undefined, - topic: streamHeader.topic!, - timestamp: bigIntToNumber(streamHeader.timestamp!), - attributes: streamHeader.attributes, - }; - streamHandlerCallback( - new ByteStreamReader(info, stream, bigIntToNumber(streamHeader.totalLength)), - { identity: participantIdentity }, - ); - } else if (streamHeader.contentHeader.case === 'textHeader') { - const streamHandlerCallback = this.textStreamHandlers.get(streamHeader.topic ?? ''); - - if (!streamHandlerCallback) { - log.debug( - 'ignoring incoming text stream due to no handler for topic: %s', - streamHeader.topic ?? 'undefined', - ); - return; - } - let streamController: ReadableStreamDefaultController; - const stream = new ReadableStream({ - start: (controller) => { - streamController = controller; - this.textStreamControllers.set(streamHeader.streamId!, { - header: streamHeader, - controller: streamController, - startTime: Date.now(), - }); - }, - }); - const info: TextStreamInfo = { - streamId: streamHeader.streamId!, - mimeType: streamHeader.mimeType!, - totalSize: streamHeader.totalLength ? Number(streamHeader.totalLength) : undefined, - topic: streamHeader.topic!, - timestamp: Number(streamHeader.timestamp), - attributes: streamHeader.attributes, - }; - streamHandlerCallback( - new TextStreamReader(info, stream, bigIntToNumber(streamHeader.totalLength)), - { identity: participantIdentity }, - ); + const streamHandlerCallback = this.byteStreamHandlers.get(info.topic); + if (!streamHandlerCallback) { + log.debug('ignoring incoming byte stream due to no handler for topic: %s', info.topic); + // The native reader was never consumed — dispose it to free the handle + // and any chunks it has buffered. + new FfiHandle(readerHandle).dispose(); + return; } - } - private handleStreamChunk(chunk: DataStream_Chunk) { - const fileBuffer = this.byteStreamControllers.get(chunk.streamId!); - if (fileBuffer) { - if (chunk.content!.length > 0) { - fileBuffer.controller.enqueue(chunk); - } - } - const textBuffer = this.textStreamControllers.get(chunk.streamId!); - if (textBuffer) { - if (chunk.content!.length > 0) { - textBuffer.controller.enqueue(chunk); - } - } + const stream = new ReadableStream({ + start: (controller) => { + let nextChunkIndex = 0; + const listener = (ev: FfiEvent) => { + if ( + ev.message.case !== 'byteStreamReaderEvent' || + ev.message.value.readerHandle !== readerHandle + ) { + return; + } + const detail = ev.message.value.detail; + if (detail.case === 'chunkReceived') { + const content = detail.value.content!; + if (content.length > 0) { + controller.enqueue( + new DataStream_Chunk({ content, chunkIndex: numberToBigInt(nextChunkIndex++) }), + ); + } + } else if (detail.case === 'eos') { + FfiClient.instance.off(FfiClientEvent.FfiEvent, listener); + this.streamReaders.delete(readerHandle); + const error = detail.value.error; + if (error) { + // Abnormal termination (e.g. remote abort, payload over the + // receiver's size limit): surface the error to the consumer + // instead of presenting a truncated payload as a clean EOF. + controller.error(new Error(error.description ?? 'stream terminated')); + } else { + controller.close(); + } + } + }; + FfiClient.instance.on(FfiClientEvent.FfiEvent, listener); + this.streamReaders.set(readerHandle, { controller, listener }); + + // Start the incremental read; the native reader buffers chunks received + // before this request, so nothing is lost. This consumes the FFI handle. + FfiClient.instance.request({ + message: { + case: 'byteReadIncremental', + value: new ByteStreamReaderReadIncrementalRequest({ readerHandle }), + }, + }); + }, + }); + + streamHandlerCallback( + new ByteStreamReader(info, stream, bigIntToNumber(ownedReader.info!.totalLength)), + { identity: participantIdentity }, + ); } - private handleStreamTrailer(trailer: DataStream_Trailer) { - const streamId = trailer.streamId!; - const fileBuffer = this.byteStreamControllers.get(streamId); - if (fileBuffer) { - fileBuffer.controller.close(); - this.byteStreamControllers.delete(streamId); - } - const textBuffer = this.textStreamControllers.get(streamId); - if (textBuffer) { - textBuffer.controller.close(); - this.textStreamControllers.delete(streamId); + private handleTextStreamOpened(ownedReader: OwnedTextStreamReader, participantIdentity: string) { + const info = textStreamInfoFromProto(ownedReader.info!); + const readerHandle = ownedReader.handle!.id!; + + const streamHandlerCallback = this.textStreamHandlers.get(info.topic); + if (!streamHandlerCallback) { + log.debug('ignoring incoming text stream due to no handler for topic: %s', info.topic); + // The native reader was never consumed — dispose it to free the handle + // and any chunks it has buffered. + new FfiHandle(readerHandle).dispose(); + return; } + + const textEncoder = new TextEncoder(); + const stream = new ReadableStream({ + start: (controller) => { + let nextChunkIndex = 0; + const listener = (ev: FfiEvent) => { + if ( + ev.message.case !== 'textStreamReaderEvent' || + ev.message.value.readerHandle !== readerHandle + ) { + return; + } + const detail = ev.message.value.detail; + if (detail.case === 'chunkReceived') { + const content = textEncoder.encode(detail.value.content!); + if (content.length > 0) { + controller.enqueue( + new DataStream_Chunk({ content, chunkIndex: numberToBigInt(nextChunkIndex++) }), + ); + } + } else if (detail.case === 'eos') { + FfiClient.instance.off(FfiClientEvent.FfiEvent, listener); + this.streamReaders.delete(readerHandle); + const error = detail.value.error; + if (error) { + // Abnormal termination (e.g. remote abort, payload over the + // receiver's size limit): surface the error to the consumer + // instead of presenting a truncated payload as a clean EOF. + controller.error(new Error(error.description ?? 'stream terminated')); + } else { + controller.close(); + } + } + }; + FfiClient.instance.on(FfiClientEvent.FfiEvent, listener); + this.streamReaders.set(readerHandle, { controller, listener }); + + // Start the incremental read; the native reader buffers chunks received + // before this request, so nothing is lost. This consumes the FFI handle. + FfiClient.instance.request({ + message: { + case: 'textReadIncremental', + value: new TextStreamReaderReadIncrementalRequest({ readerHandle }), + }, + }); + }, + }); + + streamHandlerCallback( + new TextStreamReader(info, stream, bigIntToNumber(ownedReader.info!.totalLength)), + { identity: participantIdentity }, + ); } } diff --git a/packages/livekit-rtc/src/utils.ts b/packages/livekit-rtc/src/utils.ts index 7bc5a650..090fabc0 100644 --- a/packages/livekit-rtc/src/utils.ts +++ b/packages/livekit-rtc/src/utils.ts @@ -1,6 +1,11 @@ // SPDX-FileCopyrightText: 2024 LiveKit, Inc. // // SPDX-License-Identifier: Apache-2.0 +import type { + ByteStreamInfo as ProtoByteStreamInfo, + TextStreamInfo as ProtoTextStreamInfo, +} from '@livekit/rtc-ffi-bindings'; +import type { ByteStreamInfo, TextStreamInfo } from './data_streams/types.js'; /** convert bigints to numbers preserving undefined values */ export function bigIntToNumber( @@ -40,3 +45,28 @@ export function splitUtf8(s: string, n: number): Uint8Array[] { } return result; } + +/** @internal */ +export function textStreamInfoFromProto(info: ProtoTextStreamInfo): TextStreamInfo { + return { + streamId: info.streamId!, + mimeType: info.mimeType!, + topic: info.topic!, + timestamp: bigIntToNumber(info.timestamp!), + totalSize: info.totalLength !== undefined ? bigIntToNumber(info.totalLength) : undefined, + attributes: info.attributes, + }; +} + +/** @internal */ +export function byteStreamInfoFromProto(info: ProtoByteStreamInfo): ByteStreamInfo { + return { + streamId: info.streamId!, + name: info.name ?? 'unknown', + mimeType: info.mimeType!, + topic: info.topic!, + timestamp: bigIntToNumber(info.timestamp!), + totalSize: info.totalLength !== undefined ? bigIntToNumber(info.totalLength) : undefined, + attributes: info.attributes, + }; +}