Skip to repository content1098 lines · 36.6 KB · rust
tenant.openagents/omega
No repository description is available.
OpenAgents Git authority 2026-07-28T03:59:01.911Z Public web read
NIP-34 coordinate
30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omegaMaintainersHidden in public view
References2 branches · 1 tag
Read-only clone
git clone https://openagents.com/git/tenant.openagents/omega.gitBrowse files
test.rs
1use crate::{AudioStream, Participant, RemoteTrack, RoomEvent, TrackPublication};
2
3use crate::mock_client::{participant::*, publication::*, track::*};
4use anyhow::{Context as _, Result};
5use async_trait::async_trait;
6use collections::{BTreeMap, HashMap, HashSet, btree_map::Entry as BTreeEntry, hash_map::Entry};
7use gpui::{App, AsyncApp, BackgroundExecutor};
8use livekit_api::{proto, token};
9use parking_lot::Mutex;
10use postage::{mpsc, sink::Sink};
11use std::sync::{
12 Arc, Weak,
13 atomic::{AtomicBool, AtomicU64, Ordering::SeqCst},
14};
15
16#[derive(Clone, Debug, Eq, Hash, PartialEq, PartialOrd, Ord)]
17pub struct ParticipantIdentity(pub String);
18
19#[derive(Clone, Debug, Eq, Hash, PartialEq, PartialOrd, Ord)]
20pub struct TrackSid(pub(crate) String);
21
22impl std::fmt::Display for TrackSid {
23 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
24 self.0.fmt(f)
25 }
26}
27
28impl TryFrom<String> for TrackSid {
29 type Error = anyhow::Error;
30
31 fn try_from(value: String) -> Result<Self, Self::Error> {
32 Ok(TrackSid(value))
33 }
34}
35
36#[derive(Copy, Clone, Debug, Eq, Hash, PartialEq, PartialOrd, Ord)]
37#[non_exhaustive]
38pub enum ConnectionState {
39 Connected,
40 Disconnected,
41}
42
43#[derive(Clone, Debug, Default)]
44pub struct SessionStats {
45 pub publisher_stats: Vec<RtcStats>,
46 pub subscriber_stats: Vec<RtcStats>,
47}
48
49#[derive(Clone, Debug)]
50pub enum RtcStats {}
51
52static SERVERS: Mutex<BTreeMap<String, Arc<TestServer>>> = Mutex::new(BTreeMap::new());
53
54pub struct TestServer {
55 pub url: String,
56 pub api_key: String,
57 pub secret_key: String,
58 rooms: Mutex<HashMap<String, TestServerRoom>>,
59 executor: BackgroundExecutor,
60 timestamp_source: Arc<dyn token::UnixTimestampSource>,
61}
62
63pub struct ManualUnixTimestampSource(AtomicU64);
64
65impl ManualUnixTimestampSource {
66 pub fn new(timestamp: u64) -> Self {
67 Self(AtomicU64::new(timestamp))
68 }
69
70 pub fn advance(&self) {
71 self.0.fetch_add(1, SeqCst);
72 }
73}
74
75impl token::UnixTimestampSource for ManualUnixTimestampSource {
76 fn unix_timestamp(&self) -> Result<u64> {
77 Ok(self.0.load(SeqCst))
78 }
79}
80
81impl TestServer {
82 pub fn create(
83 url: String,
84 api_key: String,
85 secret_key: String,
86 executor: BackgroundExecutor,
87 ) -> Result<Arc<TestServer>> {
88 Self::create_with_timestamp_source(
89 url,
90 api_key,
91 secret_key,
92 executor,
93 Arc::new(token::SystemUnixTimestampSource),
94 )
95 }
96
97 pub fn create_with_timestamp_source(
98 url: String,
99 api_key: String,
100 secret_key: String,
101 executor: BackgroundExecutor,
102 timestamp_source: Arc<dyn token::UnixTimestampSource>,
103 ) -> Result<Arc<TestServer>> {
104 let mut servers = SERVERS.lock();
105 if let BTreeEntry::Vacant(e) = servers.entry(url.clone()) {
106 let server = Arc::new(TestServer {
107 url,
108 api_key,
109 secret_key,
110 rooms: Default::default(),
111 executor,
112 timestamp_source,
113 });
114 e.insert(server.clone());
115 Ok(server)
116 } else {
117 anyhow::bail!("a server with url {url:?} already exists");
118 }
119 }
120
121 fn get(url: &str) -> Result<Arc<TestServer>> {
122 Ok(SERVERS
123 .lock()
124 .get(url)
125 .context("no server found for url")?
126 .clone())
127 }
128
129 pub fn teardown(&self) -> Result<()> {
130 SERVERS
131 .lock()
132 .remove(&self.url)
133 .with_context(|| format!("server with url {:?} does not exist", self.url))?;
134 Ok(())
135 }
136
137 pub fn create_api_client(&self) -> TestApiClient {
138 TestApiClient {
139 url: self.url.clone(),
140 }
141 }
142
143 #[cfg(any(test, feature = "test-support"))]
144 fn validate_token<'a>(&self, token: &'a str) -> Result<token::ClaimGrants<'a>> {
145 token::validate_with_timestamp_source(
146 token,
147 &self.secret_key,
148 self.timestamp_source.as_ref(),
149 )
150 }
151
152 #[cfg(not(any(test, feature = "test-support")))]
153 fn validate_token<'a>(&self, token: &'a str) -> Result<token::ClaimGrants<'a>> {
154 token::validate(token, &self.secret_key)
155 }
156
157 pub async fn create_room(&self, room: String) -> Result<()> {
158 self.simulate_random_delay().await;
159
160 let mut server_rooms = self.rooms.lock();
161 if let Entry::Vacant(e) = server_rooms.entry(room.clone()) {
162 e.insert(Default::default());
163 Ok(())
164 } else {
165 anyhow::bail!("{room:?} already exists");
166 }
167 }
168
169 async fn delete_room(&self, room: String) -> Result<()> {
170 self.simulate_random_delay().await;
171
172 let mut server_rooms = self.rooms.lock();
173 server_rooms
174 .remove(&room)
175 .with_context(|| format!("room {room:?} does not exist"))?;
176 Ok(())
177 }
178
179 async fn join_room(&self, token: String, client_room: Room) -> Result<ParticipantIdentity> {
180 self.simulate_random_delay().await;
181
182 let claims = self.validate_token(&token)?;
183 let identity = ParticipantIdentity(
184 claims
185 .sub
186 .context("missing participant identity")?
187 .to_string(),
188 );
189 let room_name = claims.video.room.context("missing room name")?.to_string();
190 let mut server_rooms = self.rooms.lock();
191 let room = (*server_rooms).entry(room_name.clone()).or_default();
192 if let Some(revoked_before) = room.token_revocations.get(&identity) {
193 anyhow::ensure!(claims.nbf >= *revoked_before, "invalid token: revoked");
194 }
195
196 if let Entry::Vacant(e) = room.client_rooms.entry(identity.clone()) {
197 for server_track in &room.video_tracks {
198 let track = RemoteTrack::Video(RemoteVideoTrack {
199 server_track: server_track.clone(),
200 _room: client_room.downgrade(),
201 });
202 client_room
203 .0
204 .lock()
205 .updates_tx
206 .blocking_send(RoomEvent::TrackSubscribed {
207 track: track.clone(),
208 publication: RemoteTrackPublication {
209 sid: server_track.sid.clone(),
210 room: client_room.downgrade(),
211 track,
212 },
213 participant: RemoteParticipant {
214 room: client_room.downgrade(),
215 identity: server_track.publisher_id.clone(),
216 },
217 })
218 .unwrap();
219 }
220 for server_track in &room.audio_tracks {
221 let track = RemoteTrack::Audio(RemoteAudioTrack {
222 server_track: server_track.clone(),
223 room: client_room.downgrade(),
224 });
225 client_room
226 .0
227 .lock()
228 .updates_tx
229 .blocking_send(RoomEvent::TrackSubscribed {
230 track: track.clone(),
231 publication: RemoteTrackPublication {
232 sid: server_track.sid.clone(),
233 room: client_room.downgrade(),
234 track,
235 },
236 participant: RemoteParticipant {
237 room: client_room.downgrade(),
238 identity: server_track.publisher_id.clone(),
239 },
240 })
241 .unwrap();
242 }
243 e.insert(client_room);
244 Ok(identity)
245 } else {
246 anyhow::bail!("{identity:?} attempted to join room {room_name:?} twice");
247 }
248 }
249
250 async fn leave_room(&self, token: String) -> Result<()> {
251 self.simulate_random_delay().await;
252
253 let claims = self.validate_token(&token)?;
254 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
255 let room_name = claims.video.room.unwrap();
256 let mut server_rooms = self.rooms.lock();
257 let room = server_rooms
258 .get_mut(&*room_name)
259 .with_context(|| format!("room {room_name:?} does not exist"))?;
260 room.client_rooms.remove(&identity).with_context(|| {
261 format!("{identity:?} attempted to leave room {room_name:?} before joining it")
262 })?;
263 Ok(())
264 }
265
266 fn remote_participants(
267 &self,
268 token: String,
269 ) -> Result<HashMap<ParticipantIdentity, RemoteParticipant>> {
270 let claims = self.validate_token(&token)?;
271 let local_identity = ParticipantIdentity(claims.sub.unwrap().to_string());
272 let room_name = claims.video.room.unwrap().to_string();
273
274 if let Some(server_room) = self.rooms.lock().get(&room_name) {
275 let room = server_room
276 .client_rooms
277 .get(&local_identity)
278 .unwrap()
279 .downgrade();
280 Ok(server_room
281 .client_rooms
282 .iter()
283 .filter(|(identity, _)| *identity != &local_identity)
284 .map(|(identity, _)| {
285 (
286 identity.clone(),
287 RemoteParticipant {
288 room: room.clone(),
289 identity: identity.clone(),
290 },
291 )
292 })
293 .collect())
294 } else {
295 Ok(Default::default())
296 }
297 }
298
299 async fn remove_participant(
300 &self,
301 room_name: String,
302 identity: ParticipantIdentity,
303 ) -> Result<()> {
304 self.simulate_random_delay().await;
305 let revoked_before = self.timestamp_source.unix_timestamp()?;
306
307 let mut server_rooms = self.rooms.lock();
308 let room = server_rooms
309 .get_mut(&room_name)
310 .with_context(|| format!("room {room_name} does not exist"))?;
311 let removed_room = room
312 .client_rooms
313 .remove(&identity)
314 .with_context(|| format!("participant {identity:?} did not join room {room_name:?}"))?;
315 room.token_revocations.insert(identity, revoked_before);
316 let mut removed_room = removed_room.0.lock();
317 removed_room.connection_state = ConnectionState::Disconnected;
318 removed_room
319 .updates_tx
320 .blocking_send(RoomEvent::Disconnected {
321 reason: "PARTICIPANT_REMOVED",
322 })
323 .ok();
324 Ok(())
325 }
326
327 async fn update_participant(
328 &self,
329 room_name: String,
330 identity: String,
331 permission: proto::ParticipantPermission,
332 ) -> Result<()> {
333 self.simulate_random_delay().await;
334 let revoked_before = self.timestamp_source.unix_timestamp()?;
335
336 let mut server_rooms = self.rooms.lock();
337 let room = server_rooms
338 .get_mut(&room_name)
339 .with_context(|| format!("room {room_name} does not exist"))?;
340 let identity = ParticipantIdentity(identity);
341 room.participant_permissions
342 .insert(identity.clone(), permission);
343 // Permission changes in LiveKit Cloud invalidate existing participant
344 // tokens, so the mock needs to reject tokens minted before the update.
345 room.token_revocations.insert(identity, revoked_before);
346 Ok(())
347 }
348
349 pub async fn disconnect_client(&self, client_identity: String) {
350 let client_identity = ParticipantIdentity(client_identity);
351
352 self.simulate_random_delay().await;
353
354 let mut server_rooms = self.rooms.lock();
355 for room in server_rooms.values_mut() {
356 if let Some(room) = room.client_rooms.remove(&client_identity) {
357 let mut room = room.0.lock();
358 room.connection_state = ConnectionState::Disconnected;
359 room.updates_tx
360 .blocking_send(RoomEvent::Disconnected {
361 reason: "SIGNAL_CLOSED",
362 })
363 .ok();
364 }
365 }
366 }
367
368 pub(crate) async fn publish_video_track(
369 &self,
370 token: String,
371 _local_track: LocalVideoTrack,
372 ) -> Result<TrackSid> {
373 self.simulate_random_delay().await;
374
375 let claims = self.validate_token(&token)?;
376 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
377 let room_name = claims.video.room.unwrap();
378
379 let mut server_rooms = self.rooms.lock();
380 let room = server_rooms
381 .get_mut(&*room_name)
382 .with_context(|| format!("room {room_name} does not exist"))?;
383
384 let can_publish = room
385 .participant_permissions
386 .get(&identity)
387 .map(|permission| permission.can_publish)
388 .or(claims.video.can_publish)
389 .unwrap_or(true);
390
391 anyhow::ensure!(can_publish, "user is not allowed to publish");
392
393 let sid: TrackSid = format!("TR_{}", nanoid::nanoid!(17)).try_into().unwrap();
394 let server_track = Arc::new(TestServerVideoTrack {
395 sid: sid.clone(),
396 publisher_id: identity.clone(),
397 });
398
399 room.video_tracks.push(server_track.clone());
400
401 for (room_identity, client_room) in &room.client_rooms {
402 if *room_identity != identity {
403 let track = RemoteTrack::Video(RemoteVideoTrack {
404 server_track: server_track.clone(),
405 _room: client_room.downgrade(),
406 });
407 let publication = RemoteTrackPublication {
408 sid: sid.clone(),
409 room: client_room.downgrade(),
410 track: track.clone(),
411 };
412 let participant = RemoteParticipant {
413 identity: identity.clone(),
414 room: client_room.downgrade(),
415 };
416 client_room
417 .0
418 .lock()
419 .updates_tx
420 .blocking_send(RoomEvent::TrackSubscribed {
421 track,
422 publication,
423 participant,
424 })
425 .unwrap();
426 }
427 }
428
429 Ok(sid)
430 }
431
432 pub(crate) async fn publish_audio_track(
433 &self,
434 token: String,
435 _local_track: &LocalAudioTrack,
436 ) -> Result<TrackSid> {
437 self.simulate_random_delay().await;
438
439 let claims = self.validate_token(&token)?;
440 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
441 let room_name = claims.video.room.unwrap();
442
443 let mut server_rooms = self.rooms.lock();
444 let room = server_rooms
445 .get_mut(&*room_name)
446 .with_context(|| format!("room {room_name} does not exist"))?;
447
448 let can_publish = room
449 .participant_permissions
450 .get(&identity)
451 .map(|permission| permission.can_publish)
452 .or(claims.video.can_publish)
453 .unwrap_or(true);
454
455 anyhow::ensure!(can_publish, "user is not allowed to publish");
456
457 let sid: TrackSid = format!("TR_{}", nanoid::nanoid!(17)).try_into().unwrap();
458 let server_track = Arc::new(TestServerAudioTrack {
459 sid: sid.clone(),
460 publisher_id: identity.clone(),
461 muted: AtomicBool::new(false),
462 });
463
464 room.audio_tracks.push(server_track.clone());
465
466 for (room_identity, client_room) in &room.client_rooms {
467 if *room_identity != identity {
468 let track = RemoteTrack::Audio(RemoteAudioTrack {
469 server_track: server_track.clone(),
470 room: client_room.downgrade(),
471 });
472 let publication = RemoteTrackPublication {
473 sid: sid.clone(),
474 room: client_room.downgrade(),
475 track: track.clone(),
476 };
477 let participant = RemoteParticipant {
478 identity: identity.clone(),
479 room: client_room.downgrade(),
480 };
481 client_room
482 .0
483 .lock()
484 .updates_tx
485 .blocking_send(RoomEvent::TrackSubscribed {
486 track,
487 publication,
488 participant,
489 })
490 .ok();
491 }
492 }
493
494 Ok(sid)
495 }
496
497 pub(crate) async fn unpublish_track(&self, token: String, track_sid: &TrackSid) -> Result<()> {
498 let claims = self.validate_token(&token)?;
499 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
500 let room_name = claims.video.room.unwrap();
501
502 let mut server_rooms = self.rooms.lock();
503 let room = server_rooms
504 .get_mut(&*room_name)
505 .with_context(|| format!("room {room_name} does not exist"))?;
506
507 if let Some(video_to_unpublish) = room.video_tracks.iter().position(|t| t.sid == *track_sid)
508 {
509 let video_to_unpublish = room.video_tracks.remove(video_to_unpublish);
510 for client_room in room
511 .client_rooms
512 .iter()
513 .filter(|(id, _)| **id != identity)
514 .map(|(_, room)| room)
515 {
516 let track = RemoteTrack::Video(RemoteVideoTrack {
517 server_track: video_to_unpublish.clone(),
518 _room: client_room.downgrade(),
519 });
520 let publication = RemoteTrackPublication {
521 sid: track_sid.clone(),
522 room: client_room.downgrade(),
523 track: track.clone(),
524 };
525 let participant = RemoteParticipant {
526 identity: identity.clone(),
527 room: client_room.downgrade(),
528 };
529 let event = RoomEvent::TrackUnsubscribed {
530 track,
531 publication,
532 participant,
533 };
534
535 client_room.0.lock().updates_tx.blocking_send(event).ok();
536 }
537 }
538
539 if let Some(audio_to_unpublish) = room.audio_tracks.iter().position(|t| t.sid == *track_sid)
540 {
541 let audio_to_unpublish = room.audio_tracks.remove(audio_to_unpublish);
542 for client_room in room
543 .client_rooms
544 .iter()
545 .filter(|(id, _)| **id != identity)
546 .map(|(_, room)| room)
547 {
548 let track = RemoteTrack::Audio(RemoteAudioTrack {
549 server_track: audio_to_unpublish.clone(),
550 room: client_room.downgrade(),
551 });
552 let publication = RemoteTrackPublication {
553 sid: track_sid.clone(),
554 room: client_room.downgrade(),
555 track: track.clone(),
556 };
557 let participant = RemoteParticipant {
558 identity: identity.clone(),
559 room: client_room.downgrade(),
560 };
561 let event = RoomEvent::TrackUnsubscribed {
562 track,
563 publication,
564 participant,
565 };
566
567 client_room.0.lock().updates_tx.blocking_send(event).ok();
568 }
569 }
570
571 Ok(())
572 }
573
574 pub(crate) fn set_track_muted(
575 &self,
576 token: &str,
577 track_sid: &TrackSid,
578 muted: bool,
579 ) -> Result<()> {
580 let claims = self.validate_token(token)?;
581 let room_name = claims.video.room.unwrap();
582 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
583 let mut server_rooms = self.rooms.lock();
584 let room = server_rooms
585 .get_mut(&*room_name)
586 .with_context(|| format!("room {room_name} does not exist"))?;
587 if let Some(track) = room
588 .audio_tracks
589 .iter_mut()
590 .find(|track| track.sid == *track_sid)
591 {
592 track.muted.store(muted, SeqCst);
593 for (id, client_room) in room.client_rooms.iter() {
594 if *id != identity {
595 let participant = Participant::Remote(RemoteParticipant {
596 identity: identity.clone(),
597 room: client_room.downgrade(),
598 });
599 let track = RemoteTrack::Audio(RemoteAudioTrack {
600 server_track: track.clone(),
601 room: client_room.downgrade(),
602 });
603 let publication = TrackPublication::Remote(RemoteTrackPublication {
604 sid: track_sid.clone(),
605 room: client_room.downgrade(),
606 track,
607 });
608
609 let event = if muted {
610 RoomEvent::TrackMuted {
611 participant,
612 publication,
613 }
614 } else {
615 RoomEvent::TrackUnmuted {
616 participant,
617 publication,
618 }
619 };
620
621 client_room
622 .0
623 .lock()
624 .updates_tx
625 .blocking_send(event)
626 .unwrap();
627 }
628 }
629 }
630 Ok(())
631 }
632
633 pub(crate) fn is_track_muted(&self, token: &str, track_sid: &TrackSid) -> Option<bool> {
634 let claims = self.validate_token(token).ok()?;
635 let room_name = claims.video.room.unwrap();
636
637 let mut server_rooms = self.rooms.lock();
638 let room = server_rooms.get_mut(&*room_name)?;
639 room.audio_tracks.iter().find_map(|track| {
640 if track.sid == *track_sid {
641 Some(track.muted.load(SeqCst))
642 } else {
643 None
644 }
645 })
646 }
647
648 pub(crate) fn video_tracks(&self, token: String) -> Result<Vec<RemoteVideoTrack>> {
649 let claims = self.validate_token(&token)?;
650 let room_name = claims.video.room.unwrap();
651 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
652
653 let mut server_rooms = self.rooms.lock();
654 let room = server_rooms
655 .get_mut(&*room_name)
656 .with_context(|| format!("room {room_name} does not exist"))?;
657 let client_room = room
658 .client_rooms
659 .get(&identity)
660 .context("not a participant in room")?;
661 Ok(room
662 .video_tracks
663 .iter()
664 .map(|track| RemoteVideoTrack {
665 server_track: track.clone(),
666 _room: client_room.downgrade(),
667 })
668 .collect())
669 }
670
671 pub(crate) fn audio_tracks(&self, token: String) -> Result<Vec<RemoteAudioTrack>> {
672 let claims = self.validate_token(&token)?;
673 let room_name = claims.video.room.unwrap();
674 let identity = ParticipantIdentity(claims.sub.unwrap().to_string());
675
676 let mut server_rooms = self.rooms.lock();
677 let room = server_rooms
678 .get_mut(&*room_name)
679 .with_context(|| format!("room {room_name} does not exist"))?;
680 let client_room = room
681 .client_rooms
682 .get(&identity)
683 .context("not a participant in room")?;
684 Ok(room
685 .audio_tracks
686 .iter()
687 .map(|track| RemoteAudioTrack {
688 server_track: track.clone(),
689 room: client_room.downgrade(),
690 })
691 .collect())
692 }
693
694 async fn simulate_random_delay(&self) {
695 #[cfg(any(test, feature = "test-support"))]
696 self.executor.simulate_random_delay().await;
697 }
698}
699
700#[derive(Default, Debug)]
701struct TestServerRoom {
702 client_rooms: HashMap<ParticipantIdentity, Room>,
703 video_tracks: Vec<Arc<TestServerVideoTrack>>,
704 audio_tracks: Vec<Arc<TestServerAudioTrack>>,
705 participant_permissions: HashMap<ParticipantIdentity, proto::ParticipantPermission>,
706 token_revocations: HashMap<ParticipantIdentity, u64>,
707}
708
709#[derive(Debug)]
710pub(crate) struct TestServerVideoTrack {
711 pub(crate) sid: TrackSid,
712 pub(crate) publisher_id: ParticipantIdentity,
713 // frames_rx: async_broadcast::Receiver<Frame>,
714}
715
716#[derive(Debug)]
717pub(crate) struct TestServerAudioTrack {
718 pub(crate) sid: TrackSid,
719 pub(crate) publisher_id: ParticipantIdentity,
720 pub(crate) muted: AtomicBool,
721}
722
723pub struct TestApiClient {
724 url: String,
725}
726
727#[async_trait]
728impl livekit_api::Client for TestApiClient {
729 fn url(&self) -> &str {
730 &self.url
731 }
732
733 async fn create_room(&self, name: String) -> Result<()> {
734 let server = TestServer::get(&self.url)?;
735 server.create_room(name).await?;
736 Ok(())
737 }
738
739 async fn delete_room(&self, name: String) -> Result<()> {
740 let server = TestServer::get(&self.url)?;
741 server.delete_room(name).await?;
742 Ok(())
743 }
744
745 async fn remove_participant(&self, room: String, identity: String) -> Result<()> {
746 let server = TestServer::get(&self.url)?;
747 server
748 .remove_participant(room, ParticipantIdentity(identity))
749 .await?;
750 Ok(())
751 }
752
753 async fn update_participant(
754 &self,
755 room: String,
756 identity: String,
757 permission: livekit_api::proto::ParticipantPermission,
758 ) -> Result<()> {
759 let server = TestServer::get(&self.url)?;
760 server
761 .update_participant(room, identity, permission)
762 .await?;
763 Ok(())
764 }
765
766 fn room_token(&self, room: &str, identity: &str) -> Result<String> {
767 let server = TestServer::get(&self.url)?;
768 token::create_with_timestamp_source(
769 &server.api_key,
770 &server.secret_key,
771 Some(identity),
772 token::VideoGrant::to_join(room),
773 server.timestamp_source.as_ref(),
774 )
775 }
776
777 fn guest_token(&self, room: &str, identity: &str) -> Result<String> {
778 let server = TestServer::get(&self.url)?;
779 token::create_with_timestamp_source(
780 &server.api_key,
781 &server.secret_key,
782 Some(identity),
783 token::VideoGrant::for_guest(room),
784 server.timestamp_source.as_ref(),
785 )
786 }
787}
788
789pub(crate) struct RoomState {
790 pub(crate) url: String,
791 pub(crate) token: String,
792 pub(crate) local_identity: ParticipantIdentity,
793 pub(crate) connection_state: ConnectionState,
794 pub(crate) paused_audio_tracks: HashSet<TrackSid>,
795 pub(crate) updates_tx: mpsc::Sender<RoomEvent>,
796}
797
798#[derive(Clone, Debug)]
799pub struct Room(pub(crate) Arc<Mutex<RoomState>>);
800
801#[derive(Clone, Debug)]
802pub(crate) struct WeakRoom(Weak<Mutex<RoomState>>);
803
804impl std::fmt::Debug for RoomState {
805 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
806 f.debug_struct("Room")
807 .field("url", &self.url)
808 .field("token", &self.token)
809 .field("local_identity", &self.local_identity)
810 .field("connection_state", &self.connection_state)
811 .field("paused_audio_tracks", &self.paused_audio_tracks)
812 .finish()
813 }
814}
815
816impl Room {
817 pub(crate) fn downgrade(&self) -> WeakRoom {
818 WeakRoom(Arc::downgrade(&self.0))
819 }
820
821 pub fn connection_state(&self) -> ConnectionState {
822 self.0.lock().connection_state
823 }
824
825 pub fn local_participant(&self) -> LocalParticipant {
826 let identity = self.0.lock().local_identity.clone();
827 LocalParticipant {
828 identity,
829 room: self.clone(),
830 }
831 }
832
833 pub async fn connect(
834 url: String,
835 token: String,
836 _cx: &mut AsyncApp,
837 ) -> Result<(Self, mpsc::Receiver<RoomEvent>)> {
838 let server = TestServer::get(&url)?;
839 let (updates_tx, updates_rx) = mpsc::channel(1024);
840 let this = Self(Arc::new(Mutex::new(RoomState {
841 local_identity: ParticipantIdentity(String::new()),
842 url: url.to_string(),
843 token: token.to_string(),
844 connection_state: ConnectionState::Disconnected,
845 paused_audio_tracks: Default::default(),
846 updates_tx,
847 })));
848
849 let identity = server
850 .join_room(token.to_string(), this.clone())
851 .await
852 .context("room join")?;
853 {
854 let mut state = this.0.lock();
855 state.local_identity = identity;
856 state.connection_state = ConnectionState::Connected;
857 }
858
859 Ok((this, updates_rx))
860 }
861
862 pub fn remote_participants(&self) -> HashMap<ParticipantIdentity, RemoteParticipant> {
863 self.test_server()
864 .remote_participants(self.0.lock().token.clone())
865 .unwrap()
866 }
867
868 pub(crate) fn test_server(&self) -> Arc<TestServer> {
869 TestServer::get(&self.0.lock().url).unwrap()
870 }
871
872 pub(crate) fn token(&self) -> String {
873 self.0.lock().token.clone()
874 }
875
876 pub fn name(&self) -> String {
877 "test_room".to_string()
878 }
879
880 pub async fn sid(&self) -> String {
881 "RM_test_session".to_string()
882 }
883
884 pub fn play_remote_audio_track(
885 &self,
886 _track: &RemoteAudioTrack,
887 _cx: &App,
888 ) -> anyhow::Result<AudioStream> {
889 Ok(AudioStream {})
890 }
891
892 pub async fn unpublish_local_track(&self, sid: TrackSid, cx: &mut AsyncApp) -> Result<()> {
893 self.local_participant().unpublish_track(sid, cx).await
894 }
895
896 pub async fn publish_local_microphone_track(
897 &self,
898 _track_name: String,
899 _is_staff: bool,
900 cx: &mut AsyncApp,
901 ) -> Result<(LocalTrackPublication, AudioStream, Arc<AtomicU64>)> {
902 self.local_participant().publish_microphone_track(cx).await
903 }
904
905 pub async fn get_stats(&self) -> Result<SessionStats> {
906 Ok(SessionStats::default())
907 }
908
909 pub fn stats_task(&self, _cx: &impl gpui::AppContext) -> gpui::Task<Result<SessionStats>> {
910 gpui::Task::ready(Ok(SessionStats::default()))
911 }
912}
913
914impl Drop for RoomState {
915 fn drop(&mut self) {
916 if self.connection_state == ConnectionState::Connected
917 && let Ok(server) = TestServer::get(&self.url)
918 {
919 let executor = server.executor.clone();
920 let token = self.token.clone();
921 executor
922 .spawn(async move { server.leave_room(token).await.ok() })
923 .detach();
924 }
925 }
926}
927
928impl WeakRoom {
929 pub(crate) fn upgrade(&self) -> Option<Room> {
930 self.0.upgrade().map(Room)
931 }
932}
933
934#[cfg(test)]
935mod tests {
936 use super::*;
937 use gpui::TestAppContext;
938 use livekit_api::Client as _;
939 use std::{ops::Deref, sync::atomic::AtomicUsize};
940
941 struct TestServerGuard {
942 server: Arc<TestServer>,
943 timestamp_source: Arc<ManualUnixTimestampSource>,
944 }
945
946 impl TestServerGuard {
947 fn advance_timestamp(&self) {
948 self.timestamp_source.advance();
949 }
950 }
951
952 impl Deref for TestServerGuard {
953 type Target = TestServer;
954
955 fn deref(&self) -> &Self::Target {
956 self.server.as_ref()
957 }
958 }
959
960 impl Drop for TestServerGuard {
961 fn drop(&mut self) {
962 self.server.teardown().ok();
963 }
964 }
965
966 fn create_test_server(name: &str, executor: BackgroundExecutor) -> TestServerGuard {
967 static NEXT_SERVER_ID: AtomicUsize = AtomicUsize::new(0);
968 let server_id = NEXT_SERVER_ID.fetch_add(1, SeqCst);
969 let timestamp_source = Arc::new(ManualUnixTimestampSource::new(1_234_567));
970 let server = TestServer::create_with_timestamp_source(
971 format!("http://livekit-{name}-{server_id}.test"),
972 format!("api-key-{server_id}"),
973 format!("secret-key-{server_id}"),
974 executor,
975 timestamp_source.clone(),
976 )
977 .expect("create LiveKit test server");
978 TestServerGuard {
979 server,
980 timestamp_source,
981 }
982 }
983
984 async fn assert_token_was_revoked(server: &TestServer, token: String, cx: &mut TestAppContext) {
985 match Room::connect(server.url.clone(), token, &mut cx.to_async()).await {
986 Ok(_) => panic!("revoked token unexpectedly connected"),
987 Err(error) => {
988 let error = format!("{error:#}");
989 assert!(
990 error.contains("invalid token: revoked"),
991 "expected revoked token error, got {error}"
992 );
993 }
994 }
995 }
996
997 #[gpui::test]
998 async fn token_created_after_participant_removal_can_join(
999 executor: BackgroundExecutor,
1000 cx: &mut TestAppContext,
1001 ) {
1002 let server = create_test_server("room-token", executor);
1003 server
1004 .create_room("room".into())
1005 .await
1006 .expect("create LiveKit test room");
1007 let api_client = server.create_api_client();
1008
1009 let initial_token = api_client
1010 .room_token("room", "participant")
1011 .expect("create initial room token");
1012 let (initial_room, _) = Room::connect(
1013 server.url.clone(),
1014 initial_token.clone(),
1015 &mut cx.to_async(),
1016 )
1017 .await
1018 .expect("connect with initial room token");
1019
1020 server.advance_timestamp();
1021 api_client
1022 .remove_participant("room".into(), "participant".into())
1023 .await
1024 .expect("remove participant");
1025
1026 assert_eq!(
1027 initial_room.connection_state(),
1028 ConnectionState::Disconnected
1029 );
1030 assert_token_was_revoked(&server, initial_token, cx).await;
1031
1032 let fresh_token = api_client
1033 .room_token("room", "participant")
1034 .expect("create fresh room token");
1035 let (fresh_room, _) = Room::connect(server.url.clone(), fresh_token, &mut cx.to_async())
1036 .await
1037 .expect("connect with fresh room token");
1038
1039 assert_eq!(fresh_room.connection_state(), ConnectionState::Connected);
1040 }
1041
1042 #[gpui::test]
1043 async fn guest_token_created_after_permission_update_can_join(
1044 executor: BackgroundExecutor,
1045 cx: &mut TestAppContext,
1046 ) {
1047 let server = create_test_server("guest-token", executor);
1048 server
1049 .create_room("room".into())
1050 .await
1051 .expect("create LiveKit test room");
1052 let api_client = server.create_api_client();
1053
1054 let initial_token = api_client
1055 .guest_token("room", "participant")
1056 .expect("create initial guest token");
1057 let (initial_room, _) = Room::connect(
1058 server.url.clone(),
1059 initial_token.clone(),
1060 &mut cx.to_async(),
1061 )
1062 .await
1063 .expect("connect with initial guest token");
1064
1065 server.advance_timestamp();
1066 api_client
1067 .update_participant(
1068 "room".into(),
1069 "participant".into(),
1070 proto::ParticipantPermission {
1071 can_subscribe: true,
1072 can_publish: true,
1073 can_publish_data: true,
1074 hidden: false,
1075 recorder: false,
1076 },
1077 )
1078 .await
1079 .expect("update participant permissions");
1080 assert_token_was_revoked(&server, initial_token, cx).await;
1081
1082 server.disconnect_client("participant".into()).await;
1083 assert_eq!(
1084 initial_room.connection_state(),
1085 ConnectionState::Disconnected
1086 );
1087
1088 let fresh_token = api_client
1089 .guest_token("room", "participant")
1090 .expect("create fresh guest token");
1091 let (fresh_room, _) = Room::connect(server.url.clone(), fresh_token, &mut cx.to_async())
1092 .await
1093 .expect("connect with fresh guest token");
1094
1095 assert_eq!(fresh_room.connection_state(), ConnectionState::Connected);
1096 }
1097}
1098