Skip to main content

adk_managed/
default_runtime.rs

1//! Default implementation of the [`ManagedAgentRuntime`] trait.
2//!
3//! [`DefaultManagedAgentRuntime`] composes existing ADK crates (`Runner`,
4//! `SessionService`, optional sandbox and memory) behind the unified lifecycle
5//! trait. It manages active sessions as supervised background tasks with
6//! durable checkpointing, event streaming, and custom tool parking.
7//!
8//! # Architecture
9//!
10//! The runtime is a library, not a service. The platform hosts it:
11//!
12//! - **Testable in isolation**: Zero HTTP/auth/billing dependencies
13//! - **Embeddable**: Self-hosted deployments use the runtime trait directly
14//! - **Swappable platform**: Different platforms can host the same runtime
15//! - **Provider-neutral**: Identical event sequences regardless of model provider
16//!
17//! # Example
18//!
19//! ```rust,ignore
20//! use std::sync::Arc;
21//! use adk_managed::default_runtime::DefaultManagedAgentRuntime;
22//! use adk_managed::resolver::DefaultModelResolver;
23//! use adk_session::InMemorySessionService;
24//!
25//! let resolver = Arc::new(DefaultModelResolver::new());
26//! let sessions = Arc::new(InMemorySessionService::new());
27//!
28//! let runtime = DefaultManagedAgentRuntime::new(resolver, sessions);
29//! ```
30
31use 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// ─── ActiveSession ───────────────────────────────────────────────────────────
60
61/// The addressing a managed session's conversation was persisted under.
62///
63/// Deletion has to remove what creation wrote. Keeping the triple rather than rebuilding it
64/// means a change to how sessions are addressed cannot leave orphaned conversation data
65/// behind in the configured backend.
66#[derive(Debug, Clone, PartialEq, Eq)]
67pub(crate) struct PersistedIdentity {
68    /// The app name the session was created under.
69    pub(crate) app_name: String,
70    /// The user the session was created for.
71    pub(crate) user_id: String,
72    /// The session's own identifier.
73    pub(crate) session_id: String,
74}
75
76/// Internal state for an active (or recently active) session.
77///
78/// Each session spawns a background task running the [`SessionLoop`](crate::session_loop::SessionLoop).
79/// This struct holds the communication handles and control primitives needed
80/// to interact with that background task from the runtime methods.
81#[allow(dead_code)] // Fields are retained for the session lifecycle
82pub(crate) struct ActiveSession {
83    /// The built agent driving this session.
84    pub(crate) agent: Arc<dyn Agent>,
85    /// Sender for user events into the session loop.
86    pub(crate) event_tx: mpsc::Sender<crate::types::UserEvent>,
87    /// Broadcast sender for session events (fan-out to stream subscribers).
88    pub(crate) broadcast_tx: broadcast::Sender<crate::types::SessionEvent>,
89    /// Cancellation token for interrupt handling.
90    pub(crate) cancel_token: CancellationToken,
91    /// Pause flag — when true, the session loop parks until resumed.
92    pub(crate) pause_flag: Arc<Mutex<bool>>,
93    /// Notify used to wake the session loop after resume.
94    pub(crate) pause_notify: Arc<Notify>,
95    /// Current session status (shared with the session loop).
96    pub(crate) status: Arc<RwLock<SessionStatus>>,
97    /// The identity this session's conversation is persisted under.
98    ///
99    /// Recorded at creation so deletion removes exactly what creation wrote, rather than
100    /// re-deriving an identity and risking a mismatch that silently leaves data behind.
101    pub(crate) persisted_as: PersistedIdentity,
102    /// Checkpoint manager for durable state.
103    pub(crate) checkpoint: Arc<RwLock<CheckpointManager>>,
104}
105
106// ─── DefaultManagedAgentRuntime ──────────────────────────────────────────────
107
108/// Default implementation of the managed agent runtime.
109///
110/// Composed from a [`ModelResolver`] + a pluggable [`SessionService`] +
111/// optional sandbox factory and memory service. Has no platform dependencies —
112/// the platform injects its own implementations of these traits.
113///
114/// # Fields
115///
116/// - `model_resolver` — resolves [`ModelRef`](crate::types::ModelRef) into `Arc<dyn Llm>`
117/// - `session_service` — persistent session storage backend
118/// - `sandbox_factory` — optional sandbox for built-in tool execution
119/// - `memory` — optional cross-session memory service
120/// - `sessions` — active session registry
121///
122/// # Example
123///
124/// ```rust,ignore
125/// use std::sync::Arc;
126/// use adk_managed::default_runtime::DefaultManagedAgentRuntime;
127/// use adk_managed::resolver::DefaultModelResolver;
128/// use adk_session::InMemorySessionService;
129///
130/// // Minimal runtime with defaults
131/// let runtime = DefaultManagedAgentRuntime::new(
132///     Arc::new(DefaultModelResolver::new()),
133///     Arc::new(InMemorySessionService::new()),
134/// );
135///
136/// // With sandbox and memory (feature-gated)
137/// let runtime = DefaultManagedAgentRuntime::new(
138///     Arc::new(DefaultModelResolver::new()),
139///     Arc::new(InMemorySessionService::new()),
140/// )
141/// .with_sandbox(my_sandbox)
142/// .with_memory(my_memory_service);
143/// ```
144pub struct DefaultManagedAgentRuntime {
145    /// Resolves ModelRef → `Arc<dyn Llm>`.
146    model_resolver: Arc<dyn ModelResolver>,
147    /// Persistent session storage.
148    session_service: Arc<dyn SessionService>,
149    /// Optional sandbox backend for isolated built-in tool execution.
150    ///
151    /// When set, built-in tools (bash, code_execution, etc.) execute inside
152    /// this sandbox. When `None`, built-in tools execute in-process.
153    #[cfg(feature = "sandbox")]
154    sandbox: Option<Arc<dyn SandboxBackend>>,
155    /// Optional memory service for cross-session persistent memory.
156    ///
157    /// Passed to the Runner's `memory_service` field so agents can search
158    /// and store semantic memories across sessions.
159    #[cfg(feature = "memory")]
160    memory: Option<Arc<dyn Memory>>,
161    /// Registered agents keyed by agent handle ID.
162    agents: Arc<RwLock<HashMap<String, RegisteredAgent>>>,
163    /// Active session registry keyed by session ID.
164    sessions: Arc<RwLock<HashMap<String, ActiveSession>>>,
165}
166
167/// Internal state for a registered agent.
168#[allow(dead_code)] // `def` is retained for future session creation
169struct RegisteredAgent {
170    /// The built agent instance.
171    agent: Arc<dyn Agent>,
172    /// The original definition (retained for session creation).
173    def: ManagedAgentDef,
174}
175
176impl DefaultManagedAgentRuntime {
177    /// Create a new `DefaultManagedAgentRuntime` with injected services.
178    ///
179    /// # Arguments
180    ///
181    /// * `model_resolver` - Resolves `ModelRef` declarations into callable LLM instances.
182    /// * `session_service` - Persistent storage backend for sessions and checkpoints.
183    ///
184    /// Use `.with_sandbox()` and `.with_memory()` builder methods to inject
185    /// optional sandbox and memory services (feature-gated).
186    ///
187    /// # Example
188    ///
189    /// ```rust,ignore
190    /// use std::sync::Arc;
191    /// use adk_managed::default_runtime::DefaultManagedAgentRuntime;
192    /// use adk_managed::resolver::DefaultModelResolver;
193    /// use adk_session::InMemorySessionService;
194    ///
195    /// let runtime = DefaultManagedAgentRuntime::new(
196    ///     Arc::new(DefaultModelResolver::new()),
197    ///     Arc::new(InMemorySessionService::new()),
198    /// );
199    /// ```
200    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    /// Set the sandbox backend for isolated built-in tool execution.
217    #[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    /// Set the memory service for cross-session persistent memory.
224    #[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    /// Get a reference to the model resolver.
231    pub fn model_resolver(&self) -> &Arc<dyn ModelResolver> {
232        &self.model_resolver
233    }
234
235    /// Get a reference to the session service.
236    pub fn session_service(&self) -> &Arc<dyn SessionService> {
237        &self.session_service
238    }
239
240    /// Get a reference to the optional sandbox backend.
241    #[cfg(feature = "sandbox")]
242    pub fn sandbox(&self) -> Option<&Arc<dyn SandboxBackend>> {
243        self.sandbox.as_ref()
244    }
245
246    /// Get a reference to the optional memory service.
247    #[cfg(feature = "memory")]
248    pub fn memory(&self) -> Option<&Arc<dyn Memory>> {
249        self.memory.as_ref()
250    }
251
252    /// Get a reference to the active sessions map.
253    #[cfg(test)]
254    pub(crate) fn sessions(&self) -> &Arc<RwLock<HashMap<String, ActiveSession>>> {
255        &self.sessions
256    }
257}
258
259// ─── Channel and timeout defaults ────────────────────────────────────────────
260
261/// Default capacity for the user event mpsc channel.
262const DEFAULT_EVENT_CHANNEL_CAPACITY: usize = 64;
263
264/// Default capacity for the session event broadcast channel.
265const DEFAULT_BROADCAST_CHANNEL_CAPACITY: usize = 256;
266
267/// Default timeout for custom tool parking (5 minutes).
268const DEFAULT_PARKING_TIMEOUT: Duration = Duration::from_secs(300);
269
270// ─── ManagedAgentRuntime implementation ──────────────────────────────────────
271
272#[async_trait]
273impl ManagedAgentRuntime for DefaultManagedAgentRuntime {
274    /// Create a managed agent from a declarative definition.
275    ///
276    /// Resolves the `ModelRef` into an `Arc<dyn Llm>`, builds a runnable agent,
277    /// stores it in the internal registry, and returns an opaque handle.
278    async fn create(&self, def: ManagedAgentDef) -> Result<AgentHandle, RuntimeError> {
279        // 1. Resolve model
280        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        // 2. Build agent from definition
288        #[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        // 3. Generate handle ID
300        let handle_id = uuid::Uuid::new_v4().to_string();
301
302        info!(agent_handle = %handle_id, agent_name = %def.name, "agent created");
303
304        // 4. Store in registry
305        let registered = RegisteredAgent { agent, def };
306        self.agents.write().await.insert(handle_id.clone(), registered);
307
308        Ok(AgentHandle(handle_id))
309    }
310
311    /// Start a new session for the given agent.
312    ///
313    /// Creates internal communication channels, spawns the session loop as a
314    /// background task, and stores the active session handle. Initial status
315    /// is `Queued`.
316    async fn start_session(
317        &self,
318        agent: &AgentHandle,
319        owner: &ManagedOwner,
320        env: Option<EnvironmentConfig>,
321    ) -> Result<SessionHandle, RuntimeError> {
322        // Environment configuration was accepted and discarded — the parameter was named
323        // `_env`. Applying `env_vars` or `working_dir` to an in-process session loop would
324        // mutate process-global state shared with every other session, so the runtime refuses
325        // rather than pretending. A sandboxed execution boundary is what would make this
326        // honourable.
327        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        // 1. Look up agent from registry
341        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(&registered.agent);
346        drop(agents);
347
348        // 2. Generate session ID
349        let session_id = uuid::Uuid::new_v4().to_string();
350
351        // 3. Create mpsc channel for user events
352        let (event_tx, event_rx) = mpsc::channel(DEFAULT_EVENT_CHANNEL_CAPACITY);
353
354        // 4. Create broadcast channel for session events
355        let (broadcast_tx, _) = broadcast::channel(DEFAULT_BROADCAST_CHANNEL_CAPACITY);
356
357        // 5. Create control primitives
358        let cancel_token = CancellationToken::new();
359        let pause_flag = Arc::new(Mutex::new(false));
360        let pause_notify = Arc::new(Notify::new());
361
362        // 6. Create ToolParkingLot and CheckpointManager
363        let parking = Arc::new(ToolParkingLot::new(DEFAULT_PARKING_TIMEOUT));
364        let checkpoint = Arc::new(RwLock::new(CheckpointManager::new(session_id.clone())));
365
366        // 7. Seed the session in the SessionService.
367        //    The Runner's run() calls session_service.get() which requires the
368        //    session to exist. We create it here with the same triple
369        //    (owner app, owner user, session_id) that
370        //    build_runner/run_str use in the session loop.
371        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        // 8. Spawn SessionLoop as background task, sharing the status the caller observes
388        //    so normal queued → running → idle transitions are visible through
389        //    `ManagedAgentRuntime::status`.
390        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        // 10. Create and store ActiveSession
426        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    /// Send a user event to the session.
446    ///
447    /// Dispatches the event to the session loop's input channel.
448    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    /// Subscribe to the session's event stream.
468    ///
469    /// If `from_seq` is provided, replays historical events first, then attaches
470    /// to the live broadcast.
471    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        // Subscribe to broadcast channel
482        let broadcast_rx = active.broadcast_tx.subscribe();
483
484        // Read checkpoint for replay
485        let checkpoint = active.checkpoint.read().await;
486        let stream = create_event_stream(&checkpoint, broadcast_rx, from_seq);
487
488        Ok(stream)
489    }
490
491    /// Interrupt the session at the next safe boundary.
492    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    /// Pause the session, checkpointing current state.
505    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    /// Resume a paused session.
519    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    /// Query the current status of a session.
534    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    /// Archive a session (terminal state).
544    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    /// Delete a session and its associated data.
558    async fn delete_session(&self, session: &SessionHandle) -> Result<(), RuntimeError> {
559        // First archive (set terminal state and cancel loop)
560        {
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        // Remove from sessions map
569        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        // Deleting the handle is the control plane only. `start_session` seeded a persistent
575        // session and the Runner appended every turn to it, so stopping here reported
576        // deletion while the conversation stayed in the configured backend — surviving the
577        // process that "deleted" it. Delete under the identity creation used.
578        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                // The handle is already gone, so the caller must be told the data is not.
588                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    /// A minimal in-memory session service for testing.
617    /// Uses adk-session's InMemorySessionService.
618    fn mock_session_service() -> Arc<dyn SessionService> {
619        Arc::new(adk_session::InMemorySessionService::new())
620    }
621
622    /// Mock LLM for testing the full runtime lifecycle.
623    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    /// A mock resolver that returns a MockLlm for any model ref.
658    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                        // Read and write isolation are reported separately since #485.
717                        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        // Verify we get references back (type-level verification)
773        let _r: &Arc<dyn ModelResolver> = runtime.model_resolver();
774        let _s: &Arc<dyn SessionService> = runtime.session_service();
775    }
776
777    // ─── Task 7.2: create() method tests ─────────────────────────────────────
778
779    #[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    // ─── Task 7.3: start_session() method tests ──────────────────────────────
844
845    #[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    // ─── Task 7.4: send_event() method tests ─────────────────────────────────
905
906    #[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    // ─── Task 7.5: stream_events() method tests ──────────────────────────────
948
949    #[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        // Subscribe to stream
972        let mut stream = runtime.stream_events(&session, None).await.unwrap();
973
974        // Send a message (the session loop will process it and emit events)
975        let event =
976            UserEvent::Message { content: vec![ContentBlock::Text { text: "Test".to_string() }] };
977        runtime.send_event(&session, event).await.unwrap();
978
979        // We should receive at least a StatusRunning event
980        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    // ─── Task 7.6: interrupt/pause/resume/status/archive/delete tests ────────
1001
1002    #[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        // Session should no longer be accessible
1136        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}