Runner
Die Ausführungsumgebung von adk-runner, die die Agentenausführung orchestriert.
Übersicht
Der Runner verwaltet den vollständigen Lebenszyklus der Agentenausführung:
- Sitzungsverwaltung (Sitzungen erstellen/abrufen)
- Speicherinjektion (relevante Speicher suchen und einfügen)
- Artefaktverwaltung (bereichsbezogener Zugriff auf Artefakte)
- Ereignis-Streaming (Ereignisse verarbeiten und weiterleiten)
- Agentenübertragungen (Multi-Agenten-Handoffs behandeln)
Installation
[dependencies]
adk-runner = "2.0.0"
RunnerConfig
Konfiguriere den Runner mit den erforderlichen Diensten:
use adk_runner::{Runner, RunnerConfig};
use adk_session::InMemorySessionService;
use adk_artifact::InMemoryArtifactService;
use std::sync::Arc;
let config = RunnerConfig {
app_name: "my_app".to_string(),
agent: Arc::new(my_agent),
session_service: Arc::new(InMemorySessionService::new()),
artifact_service: Some(Arc::new(InMemoryArtifactService::new())),
memory_service: None,
plugin_manager: None,
run_config: None,
compaction_config: None,
context_cache_config: None,
cache_capable: None,
request_context: None,
cancellation_token: None,
};
let runner = Runner::new(config)?;
RunnerConfigBuilder (Empfohlen)
Verwende den Typestaten-Builder, um einen Runner zu konstruieren. Der Builder erzwingt erforderliche Felder zur Compile-Zeit und setzt alle optionalen Felder auf Standardwerte, sodass das Hinzufügen neuer Felder in zukünftigen Releases deinen Code nicht bricht:
use adk_runner::Runner;
let runner = Runner::builder()
.app_name("my_app")
.agent(Arc::new(my_agent))
.session_service(Arc::new(InMemorySessionService::new()))
// Optional fields — only set what you need
.artifact_service(Arc::new(InMemoryArtifactService::new()))
.build()?;
Der Builder erfordert drei Felder: app_name, agent und session_service. Alles andere ist optional und hat sinnvolle Standardwerte. Die build()-Methode ist erst verfügbar, wenn alle drei erforderlichen Felder gesetzt sind — ein fehlendes Feld ist ein Compile-Zeit-Fehler, kein Laufzeitfehler.
Konfigurationsfelder
| Feld | Typ | Erforderlich | Beschreibung |
|---|---|---|---|
app_name | String | Ja | Anwendungskennung |
agent | Arc<dyn Agent> | Ja | Auszuführender Stamm-Agent |
session_service | Arc<dyn SessionService> | Ja | Sitzungs-Speicher-Backend |
artifact_service | Option<Arc<dyn ArtifactService>> | Nein | Artefakt-Speicher |
memory_service | Option<Arc<dyn Memory>> | Nein | Langzeitgedächtnis |
plugin_manager | Option<Arc<PluginManager>> | Nein | Plugin-Lebenszyklus-Hooks |
compaction_config | Option<EventsCompactionConfig> | Nein | Einstellungen zur Kontextkompaktion |
run_config | Option<RunConfig> | Nein | Ausführungsoptionen |
context_cache_config | Option<ContextCacheConfig> | Nein | Lebenszyklus des Kontext-Cache auf Runner-Ebene (experimentell — siehe unten) |
cache_capable | Option<Arc<dyn CacheCapable>> | Nein | Cache-fähige Modellreferenz (experimentell — siehe unten) |
request_context | Option<RequestContext> | Nein | Kontext der Auth-Middleware |
cancellation_token | Option<CancellationToken> | Nein | Kooperative Abbruchsteuerung |
Prompt-Caching
Caching ist eine Angelegenheit auf Provider-Ebene und erfordert keine Runner-Konfiguration. Jede Provider-Integration behandelt es dort, wo die Anfrage zusammengesetzt wird:
| Anbieter | Mechanismus | Standard |
|---|---|---|
| Anthropic / Bedrock | cache_control-Haltepunkte | aktiv (AnthropicConfig::prompt_caching, abwählbar mit with_prompt_caching(false)) |
| OpenAI | serverseitiges Prompt-Caching, PromptCacheRetention zur Aufbewahrung | automatisch |
| Gemini | implizites Caching auf 2.5/3.x — ein gemeinsames Präfix erhält einen Rabatt ohne Codeänderung | automatisch |
Cache-Treffern sind ohne zusätzliche Verkabelung beobachtbar: Die Gemini-Integration
zeichnet cachedContentTokenCount bei jeder Antwort auf.
context_cache_configundcache_capablesind experimentell und sollten unset gesetzt bleiben. Sie steuern Gemini's explizitecachedContentsAPI von der Runner. Dieses API erfordert, dass der Cachesystem_instruction,toolsundtool_configersetzt — das Senden eines Caches zusammen mit einem davon wird mitINVALID_ARGUMENTabgelehnt. Der Runner wählt einen Cache aus, bevor der Agent seine Tools auflöst, daher kann er diese Anfrage nicht zusammenstellen, und das Aktivieren dieser Felder führt derzeit nicht zu Cache-Treffern. Garantiertes (statt Best-Effort-) Caching für Gemini gehört in die Modellintegration, zusammen mit der Art und Weise, wie die anderen Anbieter es handhaben.
Agents ausführen
Einen Agenten mit Benutzereingabe ausführen:
use adk_core::{Content, SessionId, UserId};
use futures::StreamExt;
let user_content = Content::new("user").with_text("Hello!");
let mut stream = runner.run(
UserId::new("user-123")?,
SessionId::new("session-456")?,
user_content,
).await?;
while let Some(event) = stream.next().await {
match event {
Ok(e) => {
if let Some(content) = e.content() {
for part in &content.parts {
if let Some(text) = part.text() {
print!("{}", text);
}
}
}
}
Err(e) => eprintln!("Error: {}", e),
}
}
String-Komfortmethode
Für einfache Aufrufstellen akzeptiert run_str() einfache &str-Argumente und übernimmt die Newtype-Konvertierung intern:
let mut stream = runner.run_str(
"user-123",
"session-456",
Content::new("user").with_text("Hello!"),
).await?;
Wenn die Zeichenkette die Validierung nicht besteht (leer, enthält Null-Bytes oder überschreitet die Längenbeschränkung), gibt run_str() einen Fehler zurück, bevor die Agentenschleife gestartet wird. Die bestehende run()-Methode mit typisierten UserId/SessionId bleibt unverändert.
Unterbrechung und Run-Isolation
Ein Run wird registriert, sobald run() seinen Stream zurückgibt, und deregistriert, wenn
dieser Stream entfernt wird — auch wenn er entfernt wird, ohne jemals abgefragt worden zu sein.
| Methode | Bereich |
|---|---|
interrupt(session_id) | Bricht jede laufende Ausführung für diese Session-ID ab, über Apps und Benutzer hinweg |
interrupt_identity(app_name, user_id, session_id) | Bricht Ausführungen für eine exakt übereinstimmende Identität ab |
active_runs() | Die Identität jedes laufenden Runs; eine wiederholte Identität bedeutet gleichzeitige Runs |
active_session_ids() | Deduplizierte Sitzungs-IDs laufender Runs |
// Cancel one tenant's run without touching another that shares the session ID
let cancelled = runner.interrupt_identity("my-app", "user-1", "session-1");
Eine Session-ID ist nur innerhalb einer App und eines Benutzers eindeutig, daher ist interrupt(session_id) die weit gefasste Form und interrupt_identity die präzise. Bevorzuge interrupt_identity, wenn eine einzelne Runner mehr als eine App oder einen Benutzer bedient.
Runs werden durch eine eindeutige Run-ID und nicht durch eine Session-ID verfolgt, daher werden zwei Runs für dieselbe Identität getrennt verfolgt, und jeder deregistriert nur sich selbst.
Persistence Is Identity-Bound
Jedes Ereignis, das der Runner persistiert — Benutzerturns, Modellantworten, Transferereignisse, Plugin-Ereignisse und Kompaktierungsereignisse — wird über SessionService::append_event_for_identity mit dem vollständigen
(app_name, user_id, session_id)-Triple geschrieben. Ein SessionService, dessen natürlicher Schlüssel
zusammengesetzt ist, kann daher jedes Ereignis an seinen Mandanten binden und den
rohen Session-ID-append_event-Pfad vollständig ablehnen oder ignorieren.
Execution Flow
┌─────────────────────────────────────────────────────────────┐
│ Runner.run() │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 1. Session Retrieval │
│ │
│ SessionService.get(app_name, user_id, session_id) │
│ → Creates new session if not exists │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 2. Agent Selection │
│ │
│ Check session state for active agent │
│ → Use root agent or transferred agent │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 3. Context Creation │
│ │
│ InvocationContext with: │
│ - Session (mutable) │
│ - Artifacts (scoped to session) │
│ - Memory (if configured) │
│ - Run config │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 4. Agent Execution │
│ │
│ agent.run(ctx) → EventStream │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 5. Event Processing │
│ │
│ For each event: │
│ - Update session state │
│ - Handle transfers │
│ - Forward to caller │
└─────────────────────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────┐
│ 6. Session Save │
│ │
│ SessionService.append_event(session, events) │
└─────────────────────────────────────────────────────────────┘
InvocationContext
Der Kontext, der Agenten während der Ausführung bereitgestellt wird:
pub trait InvocationContext: CallbackContext {
/// The agent being executed
fn agent(&self) -> Arc<dyn Agent>;
/// Memory service (if configured)
fn memory(&self) -> Option<Arc<dyn Memory>>;
/// Current session
fn session(&self) -> &dyn Session;
/// Execution configuration
fn run_config(&self) -> &RunConfig;
/// Signal end of invocation
fn end_invocation(&self);
/// Check if invocation has ended
fn ended(&self) -> bool;
}
RunConfig
Ausführungsoptionen:
pub struct RunConfig {
/// Streaming mode for responses
pub streaming_mode: StreamingMode,
// ... other fields (tool_confirmation_decisions, cached_content, etc.)
}
ToolExecutionStrategy
Steuert, wie mehrere Tool-Aufrufe aus einer einzelnen LLM-Antwort verteilt werden:
| Strategie | Verhalten |
|---|---|
Sequential (Standard) | Führe Tools nacheinander in der von LLM zurückgegebenen Reihenfolge aus |
Parallel | Führe alle Tools gleichzeitig aus; der Aufrufer ist für die Sicherheit verantwortlich |
Auto | Führe die sichere schreibgeschützte Teilmenge gleichzeitig aus, dann alle verbleibenden Aufrufe nacheinander |
Set per-agent via LlmAgentBuilder:
use adk_core::ToolExecutionStrategy;
let agent = LlmAgentBuilder::new("fast_agent")
.model(model)
.tool_execution_strategy(ToolExecutionStrategy::Auto)
.tool(Arc::new(
search_tool
.with_read_only(true)
.with_concurrency_safe(true),
))
.tool(Arc::new(save_tool)) // runs after the concurrent safe subset
.build()?;
Im Auto-Modus fragt die Dispatch-Schleife sowohl is_read_only() als auch is_concurrency_safe() ab. Aufrufe, deren ausgewählte Tools für beide Methoden true zurückgeben, werden zuerst parallel ausgeführt; alle verbleibenden Aufrufe werden danach sequentiell ausgeführt. Parallel umgeht diese Metadatenprüfungen als explizite Überschreibung durch den Aufrufer. Ergebnisse werden immer wieder in der ursprünglichen von LLM zurückgegebenen Reihenfolge zusammengesetzt, unabhängig von der Strategie. Fehlgeschlagene Tools erzeugen eine JSON-Fehlermeldung, ohne den Batch abzubrechen.
pub enum StreamingMode {
/// No streaming, return complete response
None,
/// Server-Sent Events (default)
SSE,
/// Bidirectional streaming (realtime)
Bidi,
}
Agenten-Übergaben
Der Runner behandelt Übergaben zwischen mehreren Agenten automatisch:
// In an agent's tool or callback
if should_transfer {
// Set transfer in event actions
ctx.set_actions(EventActions {
transfer_to_agent: Some("specialist_agent".to_string()),
..Default::default()
});
}
Der Runner wird:
- Die Übergabeforderung im Event erkennen
- Den Ziel-Agenten in sub_agents finden
- Den Sitzungsstatus mit dem neuen aktiven Agenten aktualisieren
- Die Ausführung mit dem neuen Agenten fortsetzen
Kontextkomprimierung
Für lang laufende Sitzungen aktivieren Sie die automatische Kontextkomprimierung, um das LLM-Kontextfenster begrenzt zu halten:
use adk_runner::{Runner, RunnerConfig, EventsCompactionConfig};
use adk_agent::LlmEventSummarizer;
use std::sync::Arc;
let summarizer = LlmEventSummarizer::new(model.clone());
let config = RunnerConfig {
// ... other fields ...
compaction_config: Some(EventsCompactionConfig {
compaction_interval: 3, // Compact every 3 invocations
overlap_size: 1, // Keep 1 event overlap for continuity
summarizer: Arc::new(summarizer),
}),
// ...
};
Wenn die Komprimierung ausgelöst wird, werden ältere Events durch ein Zusammenfassungs-Event ersetzt. conversation_history() verwendet automatisch die Zusammenfassung anstelle der ursprünglichen Events.
Siehe Kontextkomprimierung für die vollständige Dokumentation.
Integration mit Launcher
Der Launcher verwendet intern Runner:
// Launcher creates Runner with default services
Launcher::new(agent)
.app_name("my_app")
.run()
.await?;
// Equivalent to using the builder:
let runner = Runner::builder()
.app_name("my_app")
.agent(agent)
.session_service(Arc::new(InMemorySessionService::new()))
.build()?;
Benutzerdefinierte Runner-Verwendung
Für fortgeschrittene Szenarien verwenden Sie Runner direkt:
use adk_runner::Runner;
// Production configuration using the builder
let runner = Runner::builder()
.app_name("production_app")
.agent(my_agent)
.session_service(Arc::new(SqliteSessionService::new(db_pool)))
.artifact_service(Arc::new(S3ArtifactService::new(s3_client)))
.memory_service(Arc::new(QdrantMemoryService::new(qdrant_client)))
.build()?;
// Use in HTTP handler with run_str() for convenience
async fn chat_handler(runner: &Runner, request: ChatRequest) -> Response {
let stream = runner.run_str(
&request.user_id,
&request.session_id,
request.content,
).await?;
// Stream events to client
Response::sse(stream)
}
Vorherige: ← Kern-Typen | Nächste: Launcher →