Skip to content

Commit fb3e47d

Browse files
Add a request_async_id to the async requests
1 parent 2975ff4 commit fb3e47d

11 files changed

Lines changed: 84 additions & 40 deletions

File tree

livekit-ffi/protocol/audio_frame.proto

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -56,11 +56,12 @@ message NewAudioSourceRequest {
5656
}
5757
message NewAudioSourceResponse { required OwnedAudioSource source = 1; }
5858

59-
// Push a frame to an AudioSource
59+
// Push a frame to an AudioSource
6060
// The data provided must be available as long as the client receive the callback.
61-
message CaptureAudioFrameRequest {
61+
message CaptureAudioFrameRequest {
6262
required uint64 source_handle = 1;
6363
required AudioFrameBufferInfo buffer = 2;
64+
optional uint64 request_async_id = 3;
6465
}
6566
message CaptureAudioFrameResponse {
6667
required uint64 async_id = 1;

livekit-ffi/protocol/data_stream.proto

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ message TextStreamReaderReadIncrementalResponse {}
3737
// Reads an incoming text stream in its entirety.
3838
message TextStreamReaderReadAllRequest {
3939
required uint64 reader_handle = 1;
40+
optional uint64 request_async_id = 2;
4041
}
4142
message TextStreamReaderReadAllResponse {
4243
required uint64 async_id = 1;
@@ -82,6 +83,7 @@ message ByteStreamReaderReadIncrementalResponse {}
8283
// Reads an incoming byte stream in its entirety.
8384
message ByteStreamReaderReadAllRequest {
8485
required uint64 reader_handle = 1;
86+
optional uint64 request_async_id = 2;
8587
}
8688
message ByteStreamReaderReadAllResponse {
8789
required uint64 async_id = 1;
@@ -97,6 +99,7 @@ message ByteStreamReaderReadAllCallback {
9799
// Writes data from an incoming stream to a file as it arrives.
98100
message ByteStreamReaderWriteToFileRequest {
99101
required uint64 reader_handle = 1;
102+
optional uint64 request_async_id = 2;
100103

101104
// Directory to write the file in (must be writable by the current process).
102105
// If not provided, the file will be written to the system's temp directory.
@@ -145,6 +148,8 @@ message StreamSendFileRequest {
145148

146149
// Path of the file to send (must be readable by the current process).
147150
required string file_path = 3;
151+
152+
optional uint64 request_async_id = 4;
148153
}
149154
message StreamSendFileResponse {
150155
required uint64 async_id = 1;
@@ -167,6 +172,8 @@ message StreamSendBytesRequest {
167172

168173
// Bytes to send.
169174
required bytes bytes = 3;
175+
176+
optional uint64 request_async_id = 4;
170177
}
171178
message StreamSendBytesResponse {
172179
required uint64 async_id = 1;
@@ -189,6 +196,8 @@ message StreamSendTextRequest {
189196

190197
// Text to send.
191198
required string text = 3;
199+
200+
optional uint64 request_async_id = 4;
192201
}
193202
message StreamSendTextResponse {
194203
required uint64 async_id = 1;
@@ -215,6 +224,8 @@ message ByteStreamOpenRequest {
215224

216225
// Options to use for opening the stream.
217226
required StreamByteOptions options = 2;
227+
228+
optional uint64 request_async_id = 3;
218229
}
219230
message ByteStreamOpenResponse {
220231
required uint64 async_id = 1;
@@ -231,6 +242,7 @@ message ByteStreamOpenCallback {
231242
message ByteStreamWriterWriteRequest {
232243
required uint64 writer_handle = 1;
233244
required bytes bytes = 2;
245+
optional uint64 request_async_id = 3;
234246
}
235247
message ByteStreamWriterWriteResponse {
236248
required uint64 async_id = 1;
@@ -244,6 +256,7 @@ message ByteStreamWriterWriteCallback {
244256
message ByteStreamWriterCloseRequest {
245257
required uint64 writer_handle = 1;
246258
optional string reason = 2;
259+
optional uint64 request_async_id = 3;
247260
}
248261
message ByteStreamWriterCloseResponse {
249262
required uint64 async_id = 1;
@@ -267,6 +280,8 @@ message TextStreamOpenRequest {
267280

268281
// Options to use for opening the stream.
269282
required StreamTextOptions options = 2;
283+
284+
optional uint64 request_async_id = 3;
270285
}
271286
message TextStreamOpenResponse {
272287
required uint64 async_id = 1;
@@ -283,6 +298,7 @@ message TextStreamOpenCallback {
283298
message TextStreamWriterWriteRequest {
284299
required uint64 writer_handle = 1;
285300
required string text = 2;
301+
optional uint64 request_async_id = 3;
286302
}
287303
message TextStreamWriterWriteResponse {
288304
required uint64 async_id = 1;
@@ -296,6 +312,7 @@ message TextStreamWriterWriteCallback {
296312
message TextStreamWriterCloseRequest {
297313
required uint64 writer_handle = 1;
298314
optional string reason = 2;
315+
optional uint64 request_async_id = 3;
299316
}
300317
message TextStreamWriterCloseResponse {
301318
required uint64 async_id = 1;

livekit-ffi/protocol/room.proto

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ message ConnectRequest {
3030
required string url = 1;
3131
required string token = 2;
3232
required RoomOptions options = 3;
33+
optional uint64 request_async_id = 4;
3334
}
3435
message ConnectResponse {
3536
required uint64 async_id = 1;
@@ -58,7 +59,10 @@ message ConnectCallback {
5859
}
5960

6061
// Disconnect from the a room
61-
message DisconnectRequest { required uint64 room_handle = 1; }
62+
message DisconnectRequest {
63+
required uint64 room_handle = 1;
64+
optional uint64 request_async_id = 2;
65+
}
6266
message DisconnectResponse { required uint64 async_id = 1; }
6367
message DisconnectCallback { required uint64 async_id = 1; }
6468

@@ -67,6 +71,7 @@ message PublishTrackRequest {
6771
required uint64 local_participant_handle = 1;
6872
required uint64 track_handle = 2;
6973
required TrackPublishOptions options = 3;
74+
optional uint64 request_async_id = 4;
7075
}
7176
message PublishTrackResponse {
7277
required uint64 async_id = 1;
@@ -85,6 +90,7 @@ message UnpublishTrackRequest {
8590
required uint64 local_participant_handle = 1;
8691
required string track_sid = 2;
8792
required bool stop_on_unpublish = 3;
93+
optional uint64 request_async_id = 4;
8894
}
8995
message UnpublishTrackResponse {
9096
required uint64 async_id = 1;
@@ -103,6 +109,7 @@ message PublishDataRequest {
103109
repeated string destination_sids = 5 [deprecated=true];
104110
optional string topic = 6;
105111
repeated string destination_identities = 7;
112+
optional uint64 request_async_id = 8;
106113
}
107114
message PublishDataResponse {
108115
required uint64 async_id = 1;
@@ -118,6 +125,7 @@ message PublishTranscriptionRequest {
118125
required string participant_identity = 2;
119126
required string track_id = 3;
120127
repeated TranscriptionSegment segments = 4;
128+
optional uint64 request_async_id = 5;
121129
}
122130
message PublishTranscriptionResponse {
123131
required uint64 async_id = 1;
@@ -133,6 +141,7 @@ message PublishSipDtmfRequest {
133141
required uint32 code = 2;
134142
required string digit = 3;
135143
repeated string destination_identities = 4;
144+
optional uint64 request_async_id = 5;
136145
}
137146
message PublishSipDtmfResponse {
138147
required uint64 async_id = 1;
@@ -146,6 +155,7 @@ message PublishSipDtmfCallback {
146155
message SetLocalMetadataRequest {
147156
required uint64 local_participant_handle = 1;
148157
required string metadata = 2;
158+
optional uint64 request_async_id = 3;
149159
}
150160
message SetLocalMetadataResponse {
151161
required uint64 async_id = 1;
@@ -160,13 +170,15 @@ message SendChatMessageRequest {
160170
required string message = 2;
161171
repeated string destination_identities = 3;
162172
optional string sender_identity = 4;
173+
optional uint64 request_async_id = 5;
163174
}
164175
message EditChatMessageRequest {
165176
required uint64 local_participant_handle = 1;
166177
required string edit_text = 2;
167178
required ChatMessage original_message = 3;
168179
repeated string destination_identities = 4;
169180
optional string sender_identity = 5;
181+
optional uint64 request_async_id = 6;
170182
}
171183
message SendChatMessageResponse {
172184
required uint64 async_id = 1;
@@ -183,6 +195,7 @@ message SendChatMessageCallback {
183195
message SetLocalAttributesRequest {
184196
required uint64 local_participant_handle = 1;
185197
repeated AttributesEntry attributes = 2;
198+
optional uint64 request_async_id = 3;
186199
}
187200

188201
message AttributesEntry {
@@ -202,6 +215,7 @@ message SetLocalAttributesCallback {
202215
message SetLocalNameRequest {
203216
required uint64 local_participant_handle = 1;
204217
required string name = 2;
218+
optional uint64 request_async_id = 3;
205219
}
206220
message SetLocalNameResponse {
207221
required uint64 async_id = 1;
@@ -220,6 +234,7 @@ message SetSubscribedResponse {}
220234

221235
message GetSessionStatsRequest {
222236
required uint64 room_handle = 1;
237+
optional uint64 request_async_id = 2;
223238
}
224239
message GetSessionStatsResponse {
225240
required uint64 async_id = 1;
@@ -644,20 +659,23 @@ message SendStreamHeaderRequest {
644659
required DataStream.Header header = 2;
645660
repeated string destination_identities = 3;
646661
required string sender_identity = 4;
662+
optional uint64 request_async_id = 5;
647663
}
648664

649665
message SendStreamChunkRequest {
650666
required uint64 local_participant_handle = 1;
651667
required DataStream.Chunk chunk = 2;
652668
repeated string destination_identities = 3;
653669
required string sender_identity = 4;
670+
optional uint64 request_async_id = 5;
654671
}
655672

656673
message SendStreamTrailerRequest {
657674
required uint64 local_participant_handle = 1;
658675
required DataStream.Trailer trailer = 2;
659676
repeated string destination_identities = 3;
660677
required string sender_identity = 4;
678+
optional uint64 request_async_id = 5;
661679
}
662680

663681
message SendStreamHeaderResponse {

livekit-ffi/protocol/rpc.proto

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ message PerformRpcRequest {
3030
required string method = 3;
3131
required string payload = 4;
3232
optional uint32 response_timeout_ms = 5;
33+
optional uint64 request_async_id = 6;
3334
}
3435

3536
message RegisterRpcMethodRequest {

livekit-ffi/protocol/track.proto

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ message CreateAudioTrackResponse {
4141

4242
message GetStatsRequest {
4343
required uint64 track_handle = 1;
44+
optional uint64 request_async_id = 2;
4445
}
4546
message GetStatsResponse {
4647
required uint64 async_id = 1;

livekit-ffi/src/server/audio_source.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,7 +75,7 @@ impl FfiAudioSource {
7575
let buffer = capture.buffer;
7676

7777
let source = self.source.clone();
78-
let async_id = server.next_id();
78+
let async_id = server.resolve_async_id(capture.request_async_id);
7979

8080
let data = unsafe {
8181
let len = buffer.num_channels * buffer.samples_per_channel;

livekit-ffi/src/server/data_stream.rs

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -94,9 +94,9 @@ impl FfiByteStreamReader {
9494
pub fn read_all(
9595
self,
9696
server: &'static FfiServer,
97-
_request: proto::ByteStreamReaderReadAllRequest,
97+
request: proto::ByteStreamReaderReadAllRequest,
9898
) -> FfiResult<proto::ByteStreamReaderReadAllResponse> {
99-
let async_id = server.next_id();
99+
let async_id = server.resolve_async_id(request.request_async_id);
100100
let handle = server.async_runtime.spawn(async move {
101101
let result = self.inner.read_all().await.into();
102102
let callback =
@@ -112,7 +112,7 @@ impl FfiByteStreamReader {
112112
server: &'static FfiServer,
113113
request: proto::ByteStreamReaderWriteToFileRequest,
114114
) -> FfiResult<proto::ByteStreamReaderWriteToFileResponse> {
115-
let async_id = server.next_id();
115+
let async_id = server.resolve_async_id(request.request_async_id);
116116

117117
let handle = server.async_runtime.spawn(async move {
118118
let result = self
@@ -174,9 +174,9 @@ impl FfiTextStreamReader {
174174
pub fn read_all(
175175
self,
176176
server: &'static FfiServer,
177-
_request: proto::TextStreamReaderReadAllRequest,
177+
request: proto::TextStreamReaderReadAllRequest,
178178
) -> FfiResult<proto::TextStreamReaderReadAllResponse> {
179-
let async_id = server.next_id();
179+
let async_id = server.resolve_async_id(request.request_async_id);
180180
let handle = server.async_runtime.spawn(async move {
181181
let result = self.inner.read_all().await.into();
182182
let callback =
@@ -208,7 +208,7 @@ impl FfiByteStreamWriter {
208208
server: &'static FfiServer,
209209
request: proto::ByteStreamWriterWriteRequest,
210210
) -> FfiResult<proto::ByteStreamWriterWriteResponse> {
211-
let async_id = server.next_id();
211+
let async_id = server.resolve_async_id(request.request_async_id);
212212
let inner = self.inner.clone();
213213
let handle = server.async_runtime.spawn(async move {
214214
let result = inner.write(&request.bytes).await;
@@ -227,7 +227,7 @@ impl FfiByteStreamWriter {
227227
server: &'static FfiServer,
228228
request: proto::ByteStreamWriterCloseRequest,
229229
) -> FfiResult<proto::ByteStreamWriterCloseResponse> {
230-
let async_id = server.next_id();
230+
let async_id = server.resolve_async_id(request.request_async_id);
231231
let handle = server.async_runtime.spawn(async move {
232232
let result = match request.reason {
233233
Some(reason) => self.inner.close_with_reason(&reason).await,
@@ -264,7 +264,7 @@ impl FfiTextStreamWriter {
264264
server: &'static FfiServer,
265265
request: proto::TextStreamWriterWriteRequest,
266266
) -> FfiResult<proto::TextStreamWriterWriteResponse> {
267-
let async_id = server.next_id();
267+
let async_id = server.resolve_async_id(request.request_async_id);
268268
let inner = self.inner.clone();
269269
let handle = server.async_runtime.spawn(async move {
270270
let result = inner.write(&request.text).await;
@@ -283,7 +283,7 @@ impl FfiTextStreamWriter {
283283
server: &'static FfiServer,
284284
request: proto::TextStreamWriterCloseRequest,
285285
) -> FfiResult<proto::TextStreamWriterCloseResponse> {
286-
let async_id = server.next_id();
286+
let async_id = server.resolve_async_id(request.request_async_id);
287287
let handle = server.async_runtime.spawn(async move {
288288
let result = match request.reason {
289289
Some(reason) => self.inner.close_with_reason(&reason).await,

livekit-ffi/src/server/mod.rs

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -179,6 +179,12 @@ impl FfiServer {
179179
self.next_id.fetch_add(1, Ordering::Relaxed)
180180
}
181181

182+
/// Resolves the async_id to use for a request.
183+
/// Uses the client-provided ID if available, otherwise generates a new one.
184+
pub fn resolve_async_id(&self, request_async_id: Option<u64>) -> FfiHandleId {
185+
request_async_id.unwrap_or_else(|| self.next_id())
186+
}
187+
182188
pub fn store_handle<T>(&self, id: FfiHandleId, handle: T)
183189
where
184190
T: FfiHandle,

0 commit comments

Comments
 (0)