Verwaltete Agent-Laufzeit
STABILITÄT: Experimentell — Diese Funktion ist additiv und durch
managed-runtimefunktionsgesteuert. Sie wirkt sich nicht auf bestehendeRunner/LlmAgentAPIs aus, wenn die Funktion deaktiviert ist. Die Oberfläche von API kann sich in zukünftigen Versionen ändern.
Übersicht
Die verwaltete Agent-Laufzeit (adk-managed) ist eine anbieterneutrale, dauerhafte und
fortsetzbare Ausführungs-Engine für Agenten. Sie übernimmt eine deklarative
ManagedAgentDef, erstellt daraus einen ausführbaren Agenten und betreibt ihn als
prüfpunktfortsetzbare, ereignisstreamende Hintergrundsitzung.
Die Laufzeit ist eine Bibliothek, kein Dienst. Die Plattform hostet sie. Das bedeutet:
- Isoliert testbar: Keine Abhängigkeiten von HTTP/Authentifizierung/Abrechnung
- Einbettbar: Selbst gehostete Bereitstellungen verwenden direkt dasselbe Laufzeit-Trait
- Austauschbare Plattform: Verschiedene Plattformen können dieselbe Laufzeit hosten
- Anbieterneutral: Identische Ereignisfolgen unabhängig vom Modellanbieter
Schnellstart
Füge die Funktion zu deiner Cargo.toml hinzu:
[dependencies]
adk-rust = { version = "2.1.0", features = ["managed-runtime"] }
Oder verwende direkt das Crate adk-managed:
[dependencies]
adk-managed = "2.1.0"
adk-session = "2.1.0"
Minimalbeispiel (ScriptedLlm — kein API-Schlüssel)
use std::sync::Arc;
use adk_managed::{
DefaultManagedAgentRuntime, ManagedAgentRuntime, ModelResolver,
ScriptedLlm, ScriptedTurn,
resolver::ResolverResult,
types::{ContentBlock, ManagedAgentDef, ModelRef, UserEvent},
};
use adk_session::InMemorySessionService;
use async_trait::async_trait;
use futures::StreamExt;
// A resolver that returns our scripted LLM
struct MockResolver { llm: Arc<dyn adk_core::Llm> }
#[async_trait]
impl ModelResolver for MockResolver {
async fn resolve(&self, _: &ModelRef) -> ResolverResult<Arc<dyn adk_core::Llm>> {
Ok(self.llm.clone())
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// 1. Create a scripted LLM (deterministic, offline, $0)
let llm = Arc::new(ScriptedLlm::new("test-model", vec![
ScriptedTurn { text: Some("Hello!".into()), tool_calls: vec![] },
]));
// 2. Build the runtime
let runtime = DefaultManagedAgentRuntime::new(
Arc::new(MockResolver { llm }),
Arc::new(InMemorySessionService::new()),
);
// 3. Create an agent
let def = ManagedAgentDef::new("my-agent", ModelRef::Shorthand("test-model".into()))
.with_system("You are helpful.");
let agent = runtime.create(def).await?;
// 4. Start a session
let session = runtime.start_session(&agent, None).await?;
// 5. Subscribe to events and send a message
let mut stream = runtime.stream_events(&session, None).await?;
runtime.send_event(&session, UserEvent::Message {
content: vec![ContentBlock::Text { text: "Hi!".into() }],
}).await?;
// 6. Collect events
while let Some(event) = stream.next().await {
println!("{event:?}");
}
Ok(())
}
Architektur
┌─────────────────────────────────────────────────────────────┐
│ Platform Layer (ep-* crates) │
│ HTTP Routes │ Auth │ Billing │ Multi-tenancy │
└──────────────────────────┬──────────────────────────────────┘
│ Rust trait calls (in-process)
▼
┌─────────────────────────────────────────────────────────────┐
│ Runtime Layer (adk-managed) │
│ │
│ ManagedAgentRuntime trait + DefaultManagedAgentRuntime │
│ ─────────────────────────────────────────────────── │
│ • Builds runnable agents from ManagedAgentDef │
│ • Runs supervised session loop (durable, resumable) │
│ • Emits provider-neutral SessionEvent stream │
│ • Manages custom tool parking, checkpoints, interrupts │
│ • Resolves ModelRef → Arc<dyn Llm> │
│ │
│ Composes existing crates: │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐ │
│ │adk-runner│ │adk-session│ │adk-model │ │adk-tool │ │
│ └──────────┘ └──────────┘ └──────────┘ └────────────┘ │
└─────────────────────────────────────────────────────────────┘
Zentrale Typen
ManagedAgentRuntime-Trait
Das zentrale asynchrone Trait, das den vollständigen Lebenszyklus des Agenten definiert:
| Methode | Beschreibung |
|---|---|
create(def) | Registriert eine Agentendefinition und gibt AgentHandle zurück |
start_session(agent, env?) | Startet eine neue Sitzung mit dem anfänglichen Status Queued |
send_event(session, event) | Eine UserEvent an die Sitzung senden |
stream_events(session, from_seq?) | Den Stream SessionEvent abonnieren |
interrupt(session) | An der nächsten Grenze anhalten, status.idle ausgeben |
pause(session) | Einen Prüfpunkt erstellen und die Verarbeitung pausieren |
resume(session) | Aus der Pause fortsetzen oder neu starten |
status(session) | Aktuellen SessionStatus abfragen |
archive(session) | Endzustand, Daten beibehalten |
delete_session(session) | Sitzungsdaten entfernen |
ManagedAgentDef
Deklarative Agent-Definition mit Builder API:
let def = ManagedAgentDef::new("my-agent", ModelRef::Shorthand("gemini-3.7-flash".into()))
.with_system("You are a helpful assistant.")
.with_description("Research agent with web search")
.with_tools(vec![ToolConfig::BuiltIn(ManagedBuiltinTool::WebSearch)]);
SessionEvent
Anbieterneutraler Ereignisstrom mit monoton steigenden Sequenznummern:
agent.message— Textinhalt des Assistentenagent.tool_use— Aufruf eines integrierten Toolsagent.custom_tool_use— Vom Client ausgeführtes benutzerdefiniertes Tool (Schleife pausiert)agent.mcp_tool_use— Aufruf des MCP-Toolsstatus.running— Durchlauf gestartetstatus.idle— Durchlauf abgeschlossen (mitstop_reason)error— Ausführungsfehler
UserEvent
Ereignisse vom Client an den Agenten:
user.message— Inhalt an den Agenten sendenuser.interrupt— Aktuellen Durchlauf stoppenuser.tool_confirmation— Tool-Ausführung erlauben/ablehnenuser.custom_tool_result— Ergebnisse des benutzerdefinierten Tools zurückgebenuser.tool_result— Ergebnis des integrierten Tools (nur selbst gehostet)user.define_outcome— Erfolgskriterien festlegen
ModelRef
Anbieterneutrale Modellreferenz zur Unterstützung aller Anbieter:
// Shorthand (provider inferred from name)
ModelRef::Shorthand("gemini-3.7-flash".into())
ModelRef::Shorthand("gpt-5.6-terra".into())
ModelRef::Shorthand("claude-sonnet-5".into())
// Structured (explicit provider)
ModelRef::Structured {
provider: Provider::OpenaiCompatible,
model: ModelConfig::Compatible {
model: "my-model".into(),
base_url: "https://my-endpoint.com/v1".into(),
api_key: "sk-...".into(),
},
speed: None,
}
Hauptfunktionen
Dauerhafte Sitzungen
Jedes Ereignis wird atomar als Prüfpunkt gespeichert. Bei einem Prozessabsturz stellt resume()
den Zustand vom letzten konsistenten Prüfpunkt ohne Ereignisverlust wieder her:
// Before crash: events 0..5 committed
// After restart:
runtime.resume(&session).await?;
// Continues from seq=5, no gap, no duplicate
Pausieren benutzerdefinierter Tools
Wenn der Agent agent.custom_tool_use ausgibt, pausiert die Schleife, bis der Client
Ergebnisse zurückgibt oder ein konfigurierbares Zeitlimit abläuft:
// Agent emits: agent.custom_tool_use { custom_tool_use_id: "ct_1", name: "deploy" }
// Client executes the tool, then:
runtime.send_event(&session, UserEvent::CustomToolResult {
custom_tool_use_id: "ct_1".into(),
content: vec![ContentBlock::Text { text: "Deployed successfully".into() }],
}).await?;
Ereigniswiedergabe
Unterstützung der Wiederverbindung von SSE Last-Event-ID durch sequenzbasierte Wiedergabe:
// Reconnect from seq 42 — replays events 43, 44, ... then live tail
let stream = runtime.stream_events(&session, Some(42)).await?;
Anbieterparität
Identisches ManagedAgentDef erzeugt byteidentische Sequenzen von Ereignistypen über
Gemini, OpenAI, Anthropic, Ollama und OpenAI-kompatible Anbieter hinweg (Fixture F-8).
Testen mit ScriptedLlm
ScriptedLlm ist ein deterministisches LLM-Double, das die vollständige Laufzeitpipeline durchläuft. Nur der Anbieter-API-Aufruf wird ersetzt:
use adk_managed::testing::{ScriptedLlm, ScriptedTurn, ScriptedToolCall};
use serde_json::json;
let llm = ScriptedLlm::new("test", vec![
ScriptedTurn {
text: Some("I'll search for that.".into()),
tool_calls: vec![ScriptedToolCall {
name: "web_search".into(),
input: json!({"query": "rust agents"}),
id: Some("tc_1".into()),
}],
},
ScriptedTurn {
text: Some("Here are the results...".into()),
tool_calls: vec![],
},
]);
API-Referenz
Die vollständige API-Dokumentation ist auf docs.rs verfügbar:
Beispiel für einen Rauchtest
Für Plattformteams steht ein eigenständiger Beispiel-Crate zur Verfügung:
cargo run --manifest-path examples/managed_runtime_hello/Cargo.toml
Damit wird das Fixture F-1 mit ScriptedLlm durchgängig ausgeführt (kein API-Schlüssel erforderlich).
Dauerhaftigkeit des verwalteten Zustands
Der verwaltete Sitzungszustand — das Ereignisprotokoll, die Sequenzposition, angehaltene Tool-Aufrufe und der Lebenszyklusstatus — befindet sich in einem ManagedStateStore. Der Speicher gibt seine eigene Garantie an:
| Haltbarkeit | Bedeutung |
|---|---|
ProcessLocal | Wiedergabe und Fortsetzung funktionieren, solange der Prozess läuft. Bei einem Absturz geht der Zustand verloren, und ein anderer Prozess kann die Sitzung nicht fortsetzen. |
CrashDurable | Der Zustand wird in einen zugrunde liegenden Speicher geschrieben, bevor der Schreibvorgang bestätigt wird, sodass ein anderer Prozess die Sitzung rekonstruieren kann. |
Nur InMemoryManagedStateStore wird ausgeliefert, und es ist ProcessLocal. Prüfen Sie die Garantie, anstatt sie aus dem Vorhandensein von Checkpointing abzuleiten:
use adk_managed::{Durability, InMemoryManagedStateStore, ManagedStateStore};
let store = InMemoryManagedStateStore::new();
assert_eq!(store.durability(), Durability::ProcessLocal);
assert!(!store.durability().survives_process_loss());
Checkpointing im Vergleich zum Flushen
CheckpointManager::checkpoint zeichnet ein Ereignis und den neuen Laufzustand zusammen auf, sodass bei der Wiedergabe nie das eine ohne das andere zu sehen ist. Dabei werden die eigenen Felder des Managers geschrieben. flush schreibt den Snapshot in den konfigurierten Speicher, und restore erstellt daraus einen Manager neu:
use adk_managed::{CheckpointManager, InMemoryManagedStateStore, ManagedStateStore};
use std::sync::Arc;
# async fn example() -> Result<(), adk_managed::types::RuntimeError> {
let store: Arc<dyn ManagedStateStore> = Arc::new(InMemoryManagedStateStore::new());
let manager = CheckpointManager::new("session-1".to_string()).with_store(Arc::clone(&store));
manager.flush().await?;
let restored = CheckpointManager::restore("session-1".to_string(), store).await?;
assert_eq!(restored.session_id(), "session-1");
# Ok(())
# }
Statusberichterstattung
ManagedAgentRuntime::status liest dasselbe Handle, in das die Sitzungsschleife schreibt, sodass normale Übergänge sichtbar sind und nicht nur Übergänge auf der Steuerungsebene:
| Übergang | Ursache |
|---|---|
Queued → Running | Ein Durchlauf beginnt |
Running → Idle | Der Durchlauf wird abgeschlossen und die Nutzung wird aufgezeichnet |
beliebig → Paused | pause |
beliebig → Archived | archive oder delete_session |
Hinweis: Zuvor war dies ein gemeinsam genutztes Handle,
statusmeldeteQueuedwährend der gesamten Lebensdauer einer Sitzung, auch während der Ausführung von Durchläufen. Übergänge der Steuerungsebene (Pause, Fortsetzen, Archivieren) waren sichtbar, da sie direkt in das Handle schrieben.
Löschsemantik
delete_session entfernt beide Ebenen:
- Setzt die Sitzung auf beendet und bricht ihre Schleife ab.
- Entfernt das Laufzeit-Handle.
- Löscht die persistierte Konversation über das injizierte
SessionServiceunter derselben Identitätstart_session, mit der sie erstellt wurde.
Falls Schritt 3 fehlschlägt, gibt delete_session einen Fehler zurück, der die App, den Benutzer und die Sitzung nennt, in denen die Daten weiterhin vorhanden sind — das Handle ist zu diesem Zeitpunkt bereits entfernt, daher muss dem Aufrufer mitgeteilt werden, was manuell bereinigt werden muss, anstatt ihm zu erlauben, von einem Erfolg auszugehen.
runtime.delete_session(&session).await?;
// The handle is gone and the conversation is no longer in the session backend.
Wichtig: Das Löschen entfernt die Konversation für die Sitzung unter deren Besitzer. Siehe Sitzungsbesitz.
Sitzungsbesitz
start_session erfordert ein ManagedOwner. Die Sitzung wird unter dieser Identität persistiert, und jeder Runner-Aufruf, den die Sitzungsschleife ausführt, verwendet sie:
use adk_managed::{ManagedAgentRuntime, ManagedOwner};
# async fn start(runtime: &dyn ManagedAgentRuntime, agent: &adk_managed::AgentHandle)
# -> Result<(), adk_managed::RuntimeError> {
let owner = ManagedOwner::new("support-console", "user-42")?;
let session = runtime.start_session(agent, &owner, None).await?;
# Ok(())
# }
Wichtig:
checkpointwurde als „atomar persistieren“ dokumentiert, mit der Garantie, dass „die Wiedergabe nach jedem Absturz eine konsistente Ansicht sieht“, und das Laden wurde als Rückgabe von „allem beschrieben, was erforderlich ist, um eine Sitzung nach einem Neustart wiederherzustellen“. Beides traf nicht zu: Beide Vorgänge arbeiteten mit In-Memory-Feldern, ohne Transaktion gegenüber einem persistenten Speicher. Mit dem mitgelieferten Speicher findetrestorein einem neuen Prozess nichts.
Beide Komponenten sind erforderlich und dürfen nicht leer sein. Sitzungen verschiedener Besitzer werden separat adressiert, sodass Suche und Löschung auf einen Besitzer beschränkt sind und nicht auf die Daten eines anderen Besitzers zugreifen können.
Hinweis: Jede verwaltete Sitzung wurde zuvor unter den Konstanten
managed/managed_userpersistiert, sodass alle denselben logischen Namensraum gemeinsam nutzten: Nichts konnte auf einen Aufrufer begrenzt werden, und keine Sitzung konnte einem Aufrufer zugeordnet werden.
Umgebungskonfiguration
EnvironmentConfig enthält env_vars und working_dir. Diese Laufzeit weist eine
Konfiguration zurück, die eines von beiden anfordert:
invalid request: EnvironmentConfig cannot be honoured by this runtime: sessions run
in-process, so per-session environment variables and working directories would have to mutate
process-global state shared with other sessions. Pass `None`, or configure a sandboxed runtime.
Sitzungen laufen innerhalb des Prozesses, sodass das Anwenden sitzungsspezifischer Umgebungsvariablen oder eines Arbeitsverzeichnisses den Zustand verändern würde, der von jeder anderen Sitzung gemeinsam genutzt wird. Die Ablehnung ist das ehrliche Ergebnis; eine sandboxbasierte Ausführungsgrenze würde die Anfrage erfüllbar machen.
Hinweis: Das Argument hieß zuvor
_envund wurde verworfen, sodass ein Aufrufer, der Umgebungskonfiguration übergab, eine Sitzung erhielt, die sie stillschweigend ignorierte.