Verwaltete Agent-Laufzeit

STABILITÄT: Experimentell — Diese Funktion ist additiv und durch managed-runtime funktionsgesteuert. Sie wirkt sich nicht auf bestehende Runner/LlmAgent APIs 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:

MethodeBeschreibung
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 Assistenten
  • agent.tool_use — Aufruf eines integrierten Tools
  • agent.custom_tool_use — Vom Client ausgeführtes benutzerdefiniertes Tool (Schleife pausiert)
  • agent.mcp_tool_use — Aufruf des MCP-Tools
  • status.running — Durchlauf gestartet
  • status.idle — Durchlauf abgeschlossen (mit stop_reason)
  • error — Ausführungsfehler

UserEvent

Ereignisse vom Client an den Agenten:

  • user.message — Inhalt an den Agenten senden
  • user.interrupt — Aktuellen Durchlauf stoppen
  • user.tool_confirmation — Tool-Ausführung erlauben/ablehnen
  • user.custom_tool_result — Ergebnisse des benutzerdefinierten Tools zurückgeben
  • user.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:

HaltbarkeitBedeutung
ProcessLocalWiedergabe und Fortsetzung funktionieren, solange der Prozess läuft. Bei einem Absturz geht der Zustand verloren, und ein anderer Prozess kann die Sitzung nicht fortsetzen.
CrashDurableDer 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:

ÜbergangUrsache
QueuedRunningEin Durchlauf beginnt
RunningIdleDer Durchlauf wird abgeschlossen und die Nutzung wird aufgezeichnet
beliebig → Pausedpause
beliebig → Archivedarchive oder delete_session

Hinweis: Zuvor war dies ein gemeinsam genutztes Handle, status meldete Queued wä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:

  1. Setzt die Sitzung auf beendet und bricht ihre Schleife ab.
  2. Entfernt das Laufzeit-Handle.
  3. Löscht die persistierte Konversation über das injizierte SessionService unter derselben Identität start_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: checkpoint wurde 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 findet restore in 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_user persistiert, 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 _env und wurde verworfen, sodass ein Aufrufer, der Umgebungskonfiguration übergab, eine Sitzung erhielt, die sie stillschweigend ignorierte.