Skip to repository content

tenant.openagents/omega

No repository description is available.

OpenAgents Git authority 2026-07-28T04:59:56.292Z Public web read
NIP-34 coordinate30617:7649603503856e5148d571eac2766b288a8ff1e9e35d380337a1d2b0015b4f92:omega
MaintainersHidden in public view
References2 branches · 1 tag
Read-only clonegit clone https://openagents.com/git/tenant.openagents/omega.git
Browse files

test.rs

1098 lines · 36.6 KB · rust
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
Served at tenant.openagents/omega Member data and write actions are omitted.