use adk_core::{Content, FinishReason, Llm, LlmRequest, LlmResponse, LlmResponseStream};
use adk_managed::resolver::{ModelResolver, ResolverResult};
use adk_managed::types::{ContentBlock, ManagedAgentDef, ModelRef, SessionStatus, UserEvent};
use adk_managed::{
DefaultManagedAgentRuntime, EnvironmentConfig, ManagedAgentRuntime, ManagedOwner,
};
use adk_session::InMemorySessionService;
use adk_session::service::{GetRequest, SessionService};
use async_trait::async_trait;
use std::sync::Arc;
struct SilentModel;
#[async_trait]
impl Llm for SilentModel {
fn name(&self) -> &str {
"silent"
}
async fn generate_content(
&self,
_request: LlmRequest,
_stream: bool,
) -> adk_core::Result<LlmResponseStream> {
let response = LlmResponse {
content: Some(Content::new("model").with_text("ok")),
partial: false,
turn_complete: true,
finish_reason: Some(FinishReason::Stop),
..Default::default()
};
Ok(Box::pin(async_stream::stream! { yield Ok(response); }))
}
}
struct SilentResolver;
#[async_trait]
impl ModelResolver for SilentResolver {
async fn resolve(&self, _model: &ModelRef) -> ResolverResult<Arc<dyn Llm>> {
Ok(Arc::new(SilentModel) as Arc<dyn Llm>)
}
}
const MANAGED_APP: &str = "managed";
const MANAGED_USER: &str = "managed_user";
fn runtime_with_service() -> (DefaultManagedAgentRuntime, Arc<InMemorySessionService>) {
let service = Arc::new(InMemorySessionService::new());
let runtime = DefaultManagedAgentRuntime::new(
Arc::new(SilentResolver) as Arc<dyn ModelResolver>,
Arc::clone(&service) as Arc<dyn SessionService>,
);
(runtime, service)
}
fn test_owner() -> ManagedOwner {
ManagedOwner::new(MANAGED_APP, MANAGED_USER).expect("valid owner")
}
fn agent_def(name: &str) -> ManagedAgentDef {
ManagedAgentDef::new(name, ModelRef::Shorthand("silent".to_string()))
}
async fn session_exists(service: &InMemorySessionService, session_id: &str) -> bool {
service
.get(GetRequest {
app_name: MANAGED_APP.to_string(),
user_id: MANAGED_USER.to_string(),
session_id: session_id.to_string(),
num_recent_events: None,
after: None,
})
.await
.map(|_| true)
.unwrap_or(false)
}
#[tokio::test]
async fn deleting_a_managed_session_removes_its_persisted_conversation() {
let (runtime, service) = runtime_with_service();
let agent = runtime.create(agent_def("deleter")).await.expect("agent");
let session = runtime.start_session(&agent, &test_owner(), None).await.expect("session");
assert!(
session_exists(&service, session.0.as_str()).await,
"start_session must seed a persistent session for the Runner to append to"
);
runtime.delete_session(&session).await.expect("delete");
assert!(
!session_exists(&service, session.0.as_str()).await,
"the persisted conversation survived a reported deletion"
);
}
#[tokio::test]
async fn deleting_an_unknown_session_is_still_reported_as_not_found() {
let (runtime, _service) = runtime_with_service();
let agent = runtime.create(agent_def("deleter")).await.expect("agent");
let session = runtime.start_session(&agent, &test_owner(), None).await.expect("session");
runtime.delete_session(&session).await.expect("first delete");
assert!(
runtime.delete_session(&session).await.is_err(),
"deleting an already-deleted session must report not found"
);
}
#[tokio::test]
async fn a_new_session_reports_queued_through_the_public_handle() {
let (runtime, _service) = runtime_with_service();
let agent = runtime.create(agent_def("reporter")).await.expect("agent");
let session = runtime.start_session(&agent, &test_owner(), None).await.expect("session");
assert_eq!(runtime.status(&session).await.expect("status"), SessionStatus::Queued);
}
#[tokio::test]
async fn a_working_session_stops_reporting_queued() {
let (runtime, _service) = runtime_with_service();
let agent = runtime.create(agent_def("reporter")).await.expect("agent");
let session = runtime.start_session(&agent, &test_owner(), None).await.expect("session");
runtime
.send_event(
&session,
UserEvent::Message { content: vec![ContentBlock::Text { text: "hi".into() }] },
)
.await
.expect("send");
let moved_on = tokio::time::timeout(std::time::Duration::from_secs(5), async {
loop {
let status = runtime.status(&session).await.expect("status");
if status != SessionStatus::Queued {
return status;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await;
let status = moved_on.expect("a session that has done work must stop reporting Queued");
assert!(
matches!(status, SessionStatus::Running | SessionStatus::Idle),
"expected a normal working transition, saw {status:?}"
);
}
#[tokio::test]
async fn archive_is_visible_through_the_public_handle() {
let (runtime, _service) = runtime_with_service();
let agent = runtime.create(agent_def("reporter")).await.expect("agent");
let session = runtime.start_session(&agent, &test_owner(), None).await.expect("session");
runtime.archive(&session).await.expect("archive");
assert_eq!(runtime.status(&session).await.expect("status"), SessionStatus::Archived);
}
#[tokio::test]
async fn two_owners_persist_into_separate_namespaces() {
let (runtime, service) = runtime_with_service();
let agent = runtime.create(agent_def("shared")).await.expect("agent");
let alice = ManagedOwner::new("console", "alice").unwrap();
let bob = ManagedOwner::new("console", "bob").unwrap();
let alice_session = runtime.start_session(&agent, &alice, None).await.expect("alice");
let bob_session = runtime.start_session(&agent, &bob, None).await.expect("bob");
assert!(exists_for(&service, &alice, alice_session.0.as_str()).await);
assert!(exists_for(&service, &bob, bob_session.0.as_str()).await);
assert!(
!exists_for(&service, &bob, alice_session.0.as_str()).await,
"one owner's session must not be addressable as another's"
);
}
#[tokio::test]
async fn deleting_one_owners_session_leaves_the_others_intact() {
let (runtime, service) = runtime_with_service();
let agent = runtime.create(agent_def("shared")).await.expect("agent");
let alice = ManagedOwner::new("console", "alice").unwrap();
let bob = ManagedOwner::new("console", "bob").unwrap();
let alice_session = runtime.start_session(&agent, &alice, None).await.expect("alice");
let bob_session = runtime.start_session(&agent, &bob, None).await.expect("bob");
runtime.delete_session(&alice_session).await.expect("delete");
assert!(!exists_for(&service, &alice, alice_session.0.as_str()).await);
assert!(
exists_for(&service, &bob, bob_session.0.as_str()).await,
"deleting one owner's session must not remove another's"
);
}
#[tokio::test]
async fn an_owner_needs_both_components() {
assert!(
ManagedOwner::new("", "user").is_err(),
"a blank app name recreates a shared namespace"
);
assert!(ManagedOwner::new("app", "").is_err(), "a blank user id cannot be scoped to a caller");
assert!(ManagedOwner::new(" ", "user").is_err(), "whitespace is not an identity");
assert!(ManagedOwner::new("app", "user").is_ok());
}
#[tokio::test]
async fn environment_configuration_is_refused_rather_than_ignored() {
let (runtime, _service) = runtime_with_service();
let agent = runtime.create(agent_def("env")).await.expect("agent");
let owner = ManagedOwner::new("console", "alice").unwrap();
let mut env = EnvironmentConfig::default();
env.env_vars.insert("API_KEY".to_string(), "value".to_string());
let error = runtime
.start_session(&agent, &owner, Some(env))
.await
.expect_err("configuration the runtime cannot honour must not be silently discarded");
assert!(error.to_string().contains("in-process"), "{error}");
assert!(
runtime.start_session(&agent, &owner, Some(EnvironmentConfig::default())).await.is_ok()
);
}
async fn exists_for(
service: &InMemorySessionService,
owner: &ManagedOwner,
session_id: &str,
) -> bool {
service
.get(GetRequest {
app_name: owner.app_name().to_string(),
user_id: owner.user_id().to_string(),
session_id: session_id.to_string(),
num_recent_events: None,
after: None,
})
.await
.map(|_| true)
.unwrap_or(false)
}