1use std::collections::HashMap;
32use std::sync::Arc;
33use std::time::Duration;
34
35use async_trait::async_trait;
36use futures::stream::BoxStream;
37use tokio::sync::{Mutex, Notify, RwLock, broadcast, mpsc};
38use tokio_util::sync::CancellationToken;
39use tracing::{debug, info};
40
41use adk_core::Agent;
42#[cfg(feature = "memory")]
43use adk_core::Memory;
44#[cfg(feature = "sandbox")]
45use adk_sandbox::SandboxBackend;
46use adk_session::service::{CreateRequest, SessionService};
47
48use crate::agent_builder::{BuildError, build_agent};
49use crate::checkpoint::CheckpointManager;
50use crate::parking::ToolParkingLot;
51use crate::replay::create_event_stream;
52use crate::resolver::ModelResolver;
53use crate::runtime::{
54 AgentHandle, EnvironmentConfig, ManagedAgentRuntime, ManagedOwner, SessionHandle,
55};
56use crate::session_loop::SessionLoop;
57use crate::types::{ManagedAgentDef, RuntimeError, SessionEvent, SessionStatus, UserEvent};
58
59#[derive(Debug, Clone, PartialEq, Eq)]
67pub(crate) struct PersistedIdentity {
68 pub(crate) app_name: String,
70 pub(crate) user_id: String,
72 pub(crate) session_id: String,
74}
75
76#[allow(dead_code)] pub(crate) struct ActiveSession {
83 pub(crate) agent: Arc<dyn Agent>,
85 pub(crate) event_tx: mpsc::Sender<crate::types::UserEvent>,
87 pub(crate) broadcast_tx: broadcast::Sender<crate::types::SessionEvent>,
89 pub(crate) cancel_token: CancellationToken,
91 pub(crate) pause_flag: Arc<Mutex<bool>>,
93 pub(crate) pause_notify: Arc<Notify>,
95 pub(crate) status: Arc<RwLock<SessionStatus>>,
97 pub(crate) persisted_as: PersistedIdentity,
102 pub(crate) checkpoint: Arc<RwLock<CheckpointManager>>,
104}
105
106pub struct DefaultManagedAgentRuntime {
145 model_resolver: Arc<dyn ModelResolver>,
147 session_service: Arc<dyn SessionService>,
149 #[cfg(feature = "sandbox")]
154 sandbox: Option<Arc<dyn SandboxBackend>>,
155 #[cfg(feature = "memory")]
160 memory: Option<Arc<dyn Memory>>,
161 agents: Arc<RwLock<HashMap<String, RegisteredAgent>>>,
163 sessions: Arc<RwLock<HashMap<String, ActiveSession>>>,
165}
166
167#[allow(dead_code)] struct RegisteredAgent {
170 agent: Arc<dyn Agent>,
172 def: ManagedAgentDef,
174}
175
176impl DefaultManagedAgentRuntime {
177 pub fn new(
201 model_resolver: Arc<dyn ModelResolver>,
202 session_service: Arc<dyn SessionService>,
203 ) -> Self {
204 Self {
205 model_resolver,
206 session_service,
207 #[cfg(feature = "sandbox")]
208 sandbox: None,
209 #[cfg(feature = "memory")]
210 memory: None,
211 agents: Arc::new(RwLock::new(HashMap::new())),
212 sessions: Arc::new(RwLock::new(HashMap::new())),
213 }
214 }
215
216 #[cfg(feature = "sandbox")]
218 pub fn with_sandbox(mut self, sandbox: Arc<dyn SandboxBackend>) -> Self {
219 self.sandbox = Some(sandbox);
220 self
221 }
222
223 #[cfg(feature = "memory")]
225 pub fn with_memory(mut self, memory: Arc<dyn Memory>) -> Self {
226 self.memory = Some(memory);
227 self
228 }
229
230 pub fn model_resolver(&self) -> &Arc<dyn ModelResolver> {
232 &self.model_resolver
233 }
234
235 pub fn session_service(&self) -> &Arc<dyn SessionService> {
237 &self.session_service
238 }
239
240 #[cfg(feature = "sandbox")]
242 pub fn sandbox(&self) -> Option<&Arc<dyn SandboxBackend>> {
243 self.sandbox.as_ref()
244 }
245
246 #[cfg(feature = "memory")]
248 pub fn memory(&self) -> Option<&Arc<dyn Memory>> {
249 self.memory.as_ref()
250 }
251
252 #[cfg(test)]
254 pub(crate) fn sessions(&self) -> &Arc<RwLock<HashMap<String, ActiveSession>>> {
255 &self.sessions
256 }
257}
258
259const DEFAULT_EVENT_CHANNEL_CAPACITY: usize = 64;
263
264const DEFAULT_BROADCAST_CHANNEL_CAPACITY: usize = 256;
266
267const DEFAULT_PARKING_TIMEOUT: Duration = Duration::from_secs(300);
269
270#[async_trait]
273impl ManagedAgentRuntime for DefaultManagedAgentRuntime {
274 async fn create(&self, def: ManagedAgentDef) -> Result<AgentHandle, RuntimeError> {
279 let model = self.model_resolver.resolve(&def.model).await.map_err(|e| {
281 RuntimeError::ProviderError {
282 provider: format!("{:?}", def.model),
283 message: e.to_string(),
284 }
285 })?;
286
287 #[cfg(feature = "sandbox")]
289 let agent = build_agent(&def, model, self.sandbox.clone()).map_err(|e| match e {
290 BuildError::InvalidDef(msg) => RuntimeError::invalid_request(msg),
291 BuildError::BuildFailed(msg) => RuntimeError::internal(msg),
292 })?;
293 #[cfg(not(feature = "sandbox"))]
294 let agent = build_agent(&def, model).map_err(|e| match e {
295 BuildError::InvalidDef(msg) => RuntimeError::invalid_request(msg),
296 BuildError::BuildFailed(msg) => RuntimeError::internal(msg),
297 })?;
298
299 let handle_id = uuid::Uuid::new_v4().to_string();
301
302 info!(agent_handle = %handle_id, agent_name = %def.name, "agent created");
303
304 let registered = RegisteredAgent { agent, def };
306 self.agents.write().await.insert(handle_id.clone(), registered);
307
308 Ok(AgentHandle(handle_id))
309 }
310
311 async fn start_session(
317 &self,
318 agent: &AgentHandle,
319 owner: &ManagedOwner,
320 env: Option<EnvironmentConfig>,
321 ) -> Result<SessionHandle, RuntimeError> {
322 if let Some(env) = &env
328 && (!env.env_vars.is_empty() || env.working_dir.is_some())
329 {
330 return Err(RuntimeError::InvalidRequest {
331 message: "EnvironmentConfig cannot be honoured by this runtime: sessions run \
332 in-process, so per-session environment variables and working \
333 directories would have to mutate process-global state shared with \
334 other sessions. Pass `None`, or configure a sandboxed runtime."
335 .to_string(),
336 param: Some("env".to_string()),
337 });
338 }
339
340 let agents = self.agents.read().await;
342 let registered = agents
343 .get(&agent.0)
344 .ok_or_else(|| RuntimeError::NotFound { session_id: agent.0.clone() })?;
345 let agent_arc = Arc::clone(®istered.agent);
346 drop(agents);
347
348 let session_id = uuid::Uuid::new_v4().to_string();
350
351 let (event_tx, event_rx) = mpsc::channel(DEFAULT_EVENT_CHANNEL_CAPACITY);
353
354 let (broadcast_tx, _) = broadcast::channel(DEFAULT_BROADCAST_CHANNEL_CAPACITY);
356
357 let cancel_token = CancellationToken::new();
359 let pause_flag = Arc::new(Mutex::new(false));
360 let pause_notify = Arc::new(Notify::new());
361
362 let parking = Arc::new(ToolParkingLot::new(DEFAULT_PARKING_TIMEOUT));
364 let checkpoint = Arc::new(RwLock::new(CheckpointManager::new(session_id.clone())));
365
366 let persisted_as = PersistedIdentity {
372 app_name: owner.app_name().to_string(),
373 user_id: owner.user_id().to_string(),
374 session_id: session_id.clone(),
375 };
376
377 self.session_service
378 .create(CreateRequest {
379 app_name: persisted_as.app_name.clone(),
380 user_id: persisted_as.user_id.clone(),
381 session_id: Some(session_id.clone()),
382 state: std::collections::HashMap::new(),
383 })
384 .await
385 .map_err(|e| RuntimeError::internal(format!("failed to seed session: {e}")))?;
386
387 let status = Arc::new(RwLock::new(SessionStatus::Queued));
391
392 #[cfg(feature = "memory")]
393 let session_loop = SessionLoop::with_pause_controls(
394 session_id.clone(),
395 event_rx,
396 broadcast_tx.clone(),
397 Arc::clone(&parking),
398 cancel_token.clone(),
399 Arc::clone(&pause_flag),
400 Arc::clone(&pause_notify),
401 Arc::clone(&checkpoint),
402 Arc::clone(&agent_arc),
403 Arc::clone(&self.session_service),
404 self.memory.clone(),
405 )
406 .with_shared_status(Arc::clone(&status))
407 .with_owner(owner.app_name(), owner.user_id());
408 #[cfg(not(feature = "memory"))]
409 let session_loop = SessionLoop::with_pause_controls(
410 session_id.clone(),
411 event_rx,
412 broadcast_tx.clone(),
413 Arc::clone(&parking),
414 cancel_token.clone(),
415 Arc::clone(&pause_flag),
416 Arc::clone(&pause_notify),
417 Arc::clone(&checkpoint),
418 Arc::clone(&agent_arc),
419 Arc::clone(&self.session_service),
420 )
421 .with_shared_status(Arc::clone(&status))
422 .with_owner(owner.app_name(), owner.user_id());
423 tokio::spawn(session_loop.run());
424
425 let active_session = ActiveSession {
427 agent: agent_arc,
428 persisted_as,
429 event_tx,
430 broadcast_tx,
431 cancel_token,
432 pause_flag,
433 pause_notify,
434 status,
435 checkpoint,
436 };
437
438 self.sessions.write().await.insert(session_id.clone(), active_session);
439
440 info!(session_id = %session_id, "session started");
441
442 Ok(SessionHandle(session_id))
443 }
444
445 async fn send_event(
449 &self,
450 session: &SessionHandle,
451 event: UserEvent,
452 ) -> Result<(), RuntimeError> {
453 let sessions = self.sessions.read().await;
454 let active = sessions
455 .get(&session.0)
456 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
457
458 active
459 .event_tx
460 .send(event)
461 .await
462 .map_err(|_| RuntimeError::conflict("session loop channel closed"))?;
463
464 Ok(())
465 }
466
467 async fn stream_events(
472 &self,
473 session: &SessionHandle,
474 from_seq: Option<u64>,
475 ) -> Result<BoxStream<'static, SessionEvent>, RuntimeError> {
476 let sessions = self.sessions.read().await;
477 let active = sessions
478 .get(&session.0)
479 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
480
481 let broadcast_rx = active.broadcast_tx.subscribe();
483
484 let checkpoint = active.checkpoint.read().await;
486 let stream = create_event_stream(&checkpoint, broadcast_rx, from_seq);
487
488 Ok(stream)
489 }
490
491 async fn interrupt(&self, session: &SessionHandle) -> Result<(), RuntimeError> {
493 let sessions = self.sessions.read().await;
494 let active = sessions
495 .get(&session.0)
496 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
497
498 debug!(session_id = %session.0, "interrupting session");
499 active.cancel_token.cancel();
500
501 Ok(())
502 }
503
504 async fn pause(&self, session: &SessionHandle) -> Result<(), RuntimeError> {
506 let sessions = self.sessions.read().await;
507 let active = sessions
508 .get(&session.0)
509 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
510
511 debug!(session_id = %session.0, "pausing session");
512 *active.pause_flag.lock().await = true;
513 *active.status.write().await = SessionStatus::Paused;
514
515 Ok(())
516 }
517
518 async fn resume(&self, session: &SessionHandle) -> Result<(), RuntimeError> {
520 let sessions = self.sessions.read().await;
521 let active = sessions
522 .get(&session.0)
523 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
524
525 debug!(session_id = %session.0, "resuming session");
526 *active.pause_flag.lock().await = false;
527 *active.status.write().await = SessionStatus::Running;
528 active.pause_notify.notify_one();
529
530 Ok(())
531 }
532
533 async fn status(&self, session: &SessionHandle) -> Result<SessionStatus, RuntimeError> {
535 let sessions = self.sessions.read().await;
536 let active = sessions
537 .get(&session.0)
538 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
539
540 Ok(*active.status.read().await)
541 }
542
543 async fn archive(&self, session: &SessionHandle) -> Result<(), RuntimeError> {
545 let sessions = self.sessions.read().await;
546 let active = sessions
547 .get(&session.0)
548 .ok_or_else(|| RuntimeError::NotFound { session_id: session.0.clone() })?;
549
550 debug!(session_id = %session.0, "archiving session");
551 *active.status.write().await = SessionStatus::Archived;
552 active.cancel_token.cancel();
553
554 Ok(())
555 }
556
557 async fn delete_session(&self, session: &SessionHandle) -> Result<(), RuntimeError> {
559 {
561 let sessions = self.sessions.read().await;
562 if let Some(active) = sessions.get(&session.0) {
563 *active.status.write().await = SessionStatus::Archived;
564 active.cancel_token.cancel();
565 }
566 }
567
568 let removed = self.sessions.write().await.remove(&session.0);
570 let Some(removed) = removed else {
571 return Err(RuntimeError::NotFound { session_id: session.0.clone() });
572 };
573
574 let identity = removed.persisted_as;
579 self.session_service
580 .delete(adk_session::DeleteRequest {
581 app_name: identity.app_name.clone(),
582 user_id: identity.user_id.clone(),
583 session_id: identity.session_id.clone(),
584 })
585 .await
586 .map_err(|e| {
587 RuntimeError::internal(format!(
589 "session {} was removed from the runtime but its persisted conversation \
590 could not be deleted: {e}. The data remains under app {} / user {} and \
591 needs manual cleanup.",
592 identity.session_id, identity.app_name, identity.user_id
593 ))
594 })?;
595
596 debug!(
597 session_id = %session.0,
598 app_name = %identity.app_name,
599 user_id = %identity.user_id,
600 "session deleted, including persisted conversation"
601 );
602 Ok(())
603 }
604}
605
606#[cfg(test)]
607mod tests {
608 use super::*;
609 use crate::resolver::DefaultModelResolver;
610 use crate::types::{ContentBlock, ModelRef};
611 use adk_core::{Content, FinishReason, Llm, LlmRequest, LlmResponse, LlmResponseStream};
612 use async_stream::stream;
613 use futures::StreamExt;
614 use std::time::Duration;
615
616 fn mock_session_service() -> Arc<dyn SessionService> {
619 Arc::new(adk_session::InMemorySessionService::new())
620 }
621
622 struct MockLlm {
624 name: String,
625 }
626
627 impl MockLlm {
628 fn new(name: &str) -> Self {
629 Self { name: name.to_string() }
630 }
631 }
632
633 #[async_trait]
634 impl Llm for MockLlm {
635 fn name(&self) -> &str {
636 &self.name
637 }
638
639 async fn generate_content(
640 &self,
641 _request: LlmRequest,
642 _stream: bool,
643 ) -> adk_core::Result<LlmResponseStream> {
644 let s = stream! {
645 yield Ok(LlmResponse {
646 content: Some(Content::new("model").with_text("Hello from mock")),
647 partial: false,
648 turn_complete: true,
649 finish_reason: Some(FinishReason::Stop),
650 ..Default::default()
651 });
652 };
653 Ok(Box::pin(s))
654 }
655 }
656
657 struct MockResolver;
659
660 #[async_trait]
661 impl ModelResolver for MockResolver {
662 async fn resolve(
663 &self,
664 _model_ref: &ModelRef,
665 ) -> crate::resolver::ResolverResult<Arc<dyn Llm>> {
666 Ok(Arc::new(MockLlm::new("mock-model")))
667 }
668 }
669
670 fn test_owner() -> ManagedOwner {
671 ManagedOwner::new("app", "user").expect("valid owner")
672 }
673
674 fn create_test_runtime() -> DefaultManagedAgentRuntime {
675 let resolver: Arc<dyn ModelResolver> = Arc::new(MockResolver);
676 let sessions = mock_session_service();
677 DefaultManagedAgentRuntime::new(resolver, sessions)
678 }
679
680 #[test]
681 fn test_new_with_minimal_config() {
682 let resolver = Arc::new(DefaultModelResolver::new());
683 let sessions = mock_session_service();
684
685 let _runtime = DefaultManagedAgentRuntime::new(resolver, sessions);
686
687 #[cfg(feature = "sandbox")]
688 assert!(_runtime.sandbox().is_none());
689 #[cfg(feature = "memory")]
690 assert!(_runtime.memory().is_none());
691 }
692
693 #[cfg(all(feature = "sandbox", feature = "memory"))]
694 #[test]
695 fn test_new_with_sandbox_and_memory() {
696 use adk_sandbox::{
697 BackendCapabilities, EnforcedLimits, ExecRequest, ExecResult, Language, SandboxBackend,
698 SandboxError,
699 };
700
701 struct FakeSandbox;
702
703 #[async_trait]
704 impl SandboxBackend for FakeSandbox {
705 fn name(&self) -> &str {
706 "fake"
707 }
708 fn capabilities(&self) -> BackendCapabilities {
709 BackendCapabilities {
710 supported_languages: vec![Language::Python],
711 isolation_class: "fake".to_string(),
712 enforced_limits: EnforcedLimits {
713 timeout: true,
714 memory: false,
715 network_isolation: false,
716 filesystem_write_isolation: false,
718 filesystem_read_isolation: false,
719 environment_isolation: false,
720 },
721 }
722 }
723 async fn execute(&self, _request: ExecRequest) -> Result<ExecResult, SandboxError> {
724 Ok(ExecResult {
725 stdout: "ok".to_string(),
726 stderr: String::new(),
727 exit_code: 0,
728 duration: std::time::Duration::from_millis(1),
729 })
730 }
731 }
732
733 struct FakeMemory;
734
735 #[async_trait]
736 impl adk_core::Memory for FakeMemory {
737 async fn search(&self, _query: &str) -> adk_core::Result<Vec<adk_core::MemoryEntry>> {
738 Ok(vec![])
739 }
740 }
741
742 let resolver = Arc::new(DefaultModelResolver::new());
743 let sessions = mock_session_service();
744
745 let runtime = DefaultManagedAgentRuntime::new(resolver, sessions)
746 .with_sandbox(Arc::new(FakeSandbox))
747 .with_memory(Arc::new(FakeMemory));
748
749 assert!(runtime.sandbox().is_some());
750 assert!(runtime.memory().is_some());
751 }
752
753 #[test]
754 fn test_sessions_map_starts_empty() {
755 let resolver = Arc::new(DefaultModelResolver::new());
756 let sessions = mock_session_service();
757
758 let runtime = DefaultManagedAgentRuntime::new(resolver, sessions);
759
760 let sessions = runtime.sessions().try_read().unwrap();
761 assert!(sessions.is_empty());
762 }
763
764 #[test]
765 fn test_accessors_return_injected_services() {
766 let resolver: Arc<dyn ModelResolver> = Arc::new(DefaultModelResolver::new());
767 let session_service = mock_session_service();
768
769 let runtime =
770 DefaultManagedAgentRuntime::new(Arc::clone(&resolver), Arc::clone(&session_service));
771
772 let _r: &Arc<dyn ModelResolver> = runtime.model_resolver();
774 let _s: &Arc<dyn SessionService> = runtime.session_service();
775 }
776
777 #[tokio::test]
780 async fn test_create_agent_returns_handle() {
781 let runtime = create_test_runtime();
782
783 let def = ManagedAgentDef {
784 name: "test-agent".to_string(),
785 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
786 system: Some("You are helpful.".to_string()),
787 description: None,
788 tools: vec![],
789 mcp_servers: vec![],
790 skills: vec![],
791 permission_policy: None,
792 metadata: None,
793 };
794
795 let handle = runtime.create(def).await.unwrap();
796 assert!(!handle.0.is_empty());
797 }
798
799 #[tokio::test]
800 async fn test_create_agent_stores_in_registry() {
801 let runtime = create_test_runtime();
802
803 let def = ManagedAgentDef {
804 name: "stored-agent".to_string(),
805 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
806 system: None,
807 description: None,
808 tools: vec![],
809 mcp_servers: vec![],
810 skills: vec![],
811 permission_policy: None,
812 metadata: None,
813 };
814
815 let handle = runtime.create(def).await.unwrap();
816 let agents = runtime.agents.read().await;
817 assert!(agents.contains_key(&handle.0));
818 }
819
820 #[tokio::test]
821 async fn test_create_multiple_agents() {
822 let runtime = create_test_runtime();
823
824 let make_def = |name: &str| ManagedAgentDef {
825 name: name.to_string(),
826 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
827 system: None,
828 description: None,
829 tools: vec![],
830 mcp_servers: vec![],
831 skills: vec![],
832 permission_policy: None,
833 metadata: None,
834 };
835
836 let h1 = runtime.create(make_def("agent-1")).await.unwrap();
837 let h2 = runtime.create(make_def("agent-2")).await.unwrap();
838
839 assert_ne!(h1.0, h2.0);
840 assert_eq!(runtime.agents.read().await.len(), 2);
841 }
842
843 #[tokio::test]
846 async fn test_start_session_returns_handle() {
847 let runtime = create_test_runtime();
848
849 let def = ManagedAgentDef {
850 name: "session-agent".to_string(),
851 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
852 system: None,
853 description: None,
854 tools: vec![],
855 mcp_servers: vec![],
856 skills: vec![],
857 permission_policy: None,
858 metadata: None,
859 };
860
861 let agent = runtime.create(def).await.unwrap();
862 let session = runtime
863 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
864 .await
865 .unwrap();
866 assert!(!session.0.is_empty());
867 }
868
869 #[tokio::test]
870 async fn test_start_session_initial_status_queued() {
871 let runtime = create_test_runtime();
872
873 let def = ManagedAgentDef {
874 name: "status-agent".to_string(),
875 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
876 system: None,
877 description: None,
878 tools: vec![],
879 mcp_servers: vec![],
880 skills: vec![],
881 permission_policy: None,
882 metadata: None,
883 };
884
885 let agent = runtime.create(def).await.unwrap();
886 let session = runtime
887 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
888 .await
889 .unwrap();
890
891 let status = runtime.status(&session).await.unwrap();
892 assert_eq!(status, SessionStatus::Queued);
893 }
894
895 #[tokio::test]
896 async fn test_start_session_unknown_agent_returns_error() {
897 let runtime = create_test_runtime();
898
899 let fake_agent = AgentHandle("nonexistent".to_string());
900 let result = runtime.start_session(&fake_agent, &test_owner(), None).await;
901 assert!(result.is_err());
902 }
903
904 #[tokio::test]
907 async fn test_send_event_message() {
908 let runtime = create_test_runtime();
909
910 let def = ManagedAgentDef {
911 name: "event-agent".to_string(),
912 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
913 system: None,
914 description: None,
915 tools: vec![],
916 mcp_servers: vec![],
917 skills: vec![],
918 permission_policy: None,
919 metadata: None,
920 };
921
922 let agent = runtime.create(def).await.unwrap();
923 let session = runtime
924 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
925 .await
926 .unwrap();
927
928 let event =
929 UserEvent::Message { content: vec![ContentBlock::Text { text: "Hello".to_string() }] };
930
931 let result = runtime.send_event(&session, event).await;
932 assert!(result.is_ok());
933 }
934
935 #[tokio::test]
936 async fn test_send_event_unknown_session_returns_error() {
937 let runtime = create_test_runtime();
938
939 let fake_session = SessionHandle("nonexistent".to_string());
940 let event =
941 UserEvent::Message { content: vec![ContentBlock::Text { text: "Hello".to_string() }] };
942
943 let result = runtime.send_event(&fake_session, event).await;
944 assert!(result.is_err());
945 }
946
947 #[tokio::test]
950 async fn test_stream_events_receives_broadcast() {
951 let runtime = create_test_runtime();
952
953 let def = ManagedAgentDef {
954 name: "stream-agent".to_string(),
955 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
956 system: None,
957 description: None,
958 tools: vec![],
959 mcp_servers: vec![],
960 skills: vec![],
961 permission_policy: None,
962 metadata: None,
963 };
964
965 let agent = runtime.create(def).await.unwrap();
966 let session = runtime
967 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
968 .await
969 .unwrap();
970
971 let mut stream = runtime.stream_events(&session, None).await.unwrap();
973
974 let event =
976 UserEvent::Message { content: vec![ContentBlock::Text { text: "Test".to_string() }] };
977 runtime.send_event(&session, event).await.unwrap();
978
979 let first_event = tokio::time::timeout(Duration::from_secs(2), stream.next())
981 .await
982 .expect("timed out waiting for event")
983 .expect("stream ended unexpectedly");
984
985 match first_event {
986 SessionEvent::StatusRunning { .. } => {}
987 other => panic!("expected StatusRunning, got: {other:?}"),
988 }
989 }
990
991 #[tokio::test]
992 async fn test_stream_events_unknown_session_returns_error() {
993 let runtime = create_test_runtime();
994
995 let fake_session = SessionHandle("nonexistent".to_string());
996 let result = runtime.stream_events(&fake_session, None).await;
997 assert!(result.is_err());
998 }
999
1000 #[tokio::test]
1003 async fn test_interrupt_cancels_session() {
1004 let runtime = create_test_runtime();
1005
1006 let def = ManagedAgentDef {
1007 name: "interrupt-agent".to_string(),
1008 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
1009 system: None,
1010 description: None,
1011 tools: vec![],
1012 mcp_servers: vec![],
1013 skills: vec![],
1014 permission_policy: None,
1015 metadata: None,
1016 };
1017
1018 let agent = runtime.create(def).await.unwrap();
1019 let session = runtime
1020 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
1021 .await
1022 .unwrap();
1023
1024 let result = runtime.interrupt(&session).await;
1025 assert!(result.is_ok());
1026 }
1027
1028 #[tokio::test]
1029 async fn test_pause_sets_paused_status() {
1030 let runtime = create_test_runtime();
1031
1032 let def = ManagedAgentDef {
1033 name: "pause-agent".to_string(),
1034 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
1035 system: None,
1036 description: None,
1037 tools: vec![],
1038 mcp_servers: vec![],
1039 skills: vec![],
1040 permission_policy: None,
1041 metadata: None,
1042 };
1043
1044 let agent = runtime.create(def).await.unwrap();
1045 let session = runtime
1046 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
1047 .await
1048 .unwrap();
1049
1050 runtime.pause(&session).await.unwrap();
1051 let status = runtime.status(&session).await.unwrap();
1052 assert_eq!(status, SessionStatus::Paused);
1053 }
1054
1055 #[tokio::test]
1056 async fn test_resume_clears_pause() {
1057 let runtime = create_test_runtime();
1058
1059 let def = ManagedAgentDef {
1060 name: "resume-agent".to_string(),
1061 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
1062 system: None,
1063 description: None,
1064 tools: vec![],
1065 mcp_servers: vec![],
1066 skills: vec![],
1067 permission_policy: None,
1068 metadata: None,
1069 };
1070
1071 let agent = runtime.create(def).await.unwrap();
1072 let session = runtime
1073 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
1074 .await
1075 .unwrap();
1076
1077 runtime.pause(&session).await.unwrap();
1078 assert_eq!(runtime.status(&session).await.unwrap(), SessionStatus::Paused);
1079
1080 runtime.resume(&session).await.unwrap();
1081 assert_eq!(runtime.status(&session).await.unwrap(), SessionStatus::Running);
1082 }
1083
1084 #[tokio::test]
1085 async fn test_archive_sets_archived_status() {
1086 let runtime = create_test_runtime();
1087
1088 let def = ManagedAgentDef {
1089 name: "archive-agent".to_string(),
1090 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
1091 system: None,
1092 description: None,
1093 tools: vec![],
1094 mcp_servers: vec![],
1095 skills: vec![],
1096 permission_policy: None,
1097 metadata: None,
1098 };
1099
1100 let agent = runtime.create(def).await.unwrap();
1101 let session = runtime
1102 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
1103 .await
1104 .unwrap();
1105
1106 runtime.archive(&session).await.unwrap();
1107 let status = runtime.status(&session).await.unwrap();
1108 assert_eq!(status, SessionStatus::Archived);
1109 }
1110
1111 #[tokio::test]
1112 async fn test_delete_session_removes_from_registry() {
1113 let runtime = create_test_runtime();
1114
1115 let def = ManagedAgentDef {
1116 name: "delete-agent".to_string(),
1117 model: ModelRef::Shorthand("gemini-2.5-flash".to_string()),
1118 system: None,
1119 description: None,
1120 tools: vec![],
1121 mcp_servers: vec![],
1122 skills: vec![],
1123 permission_policy: None,
1124 metadata: None,
1125 };
1126
1127 let agent = runtime.create(def).await.unwrap();
1128 let session = runtime
1129 .start_session(&agent, &ManagedOwner::new("app", "user").unwrap(), None)
1130 .await
1131 .unwrap();
1132
1133 runtime.delete_session(&session).await.unwrap();
1134
1135 let result = runtime.status(&session).await;
1137 assert!(result.is_err());
1138 }
1139
1140 #[tokio::test]
1141 async fn test_delete_nonexistent_session_returns_error() {
1142 let runtime = create_test_runtime();
1143
1144 let fake_session = SessionHandle("nonexistent".to_string());
1145 let result = runtime.delete_session(&fake_session).await;
1146 assert!(result.is_err());
1147 }
1148
1149 #[tokio::test]
1150 async fn test_interrupt_nonexistent_session_returns_error() {
1151 let runtime = create_test_runtime();
1152
1153 let fake_session = SessionHandle("nonexistent".to_string());
1154 let result = runtime.interrupt(&fake_session).await;
1155 assert!(result.is_err());
1156 }
1157}