関数型 API
自動チェックポイント、型付き状態リデューサー、割り込み/再開サポートを備えたエージェントワークフローを、通常の非同期 Rust 関数として記述します。
概要
関数型 API(adk-graph の functional 機能)は、明示的なグラフ/ノード/エッジ構築に代わる、より高レベルな方法を提供します。StateGraph を手動で構築する代わりに、関数に #[entrypoint] および #[task] マクロを付け、標準的な Rust の制御フローを使用します。
主な利点:
- 通常の非同期 Rust としてワークフローを記述できる — グラフ DSL は不要
- クラッシュ復旧のため、各タスク後に自動的にチェックポイントを作成
- 標準的な Rust の制御フロー(if/else、for、loop、match)
- 永続化が保証された型付き状態コンテナ
- 既存の StreamEvent およびチェックポインターインフラストラクチャと互換
はじめに
[dependencies]
adk-graph = { version = "2.1.0", features = ["functional"] }
adk-rust-macros = "2.1.0"
コア型
TaskContext
すべてのワークフロー関数に渡されるランタイムコンテキスト。状態、チェックポイント、割り込み、ストリーミングへのアクセスを提供します。
use adk_graph::functional::TaskContext;
// Read state
let count: Option<i64> = ctx.get("counter");
// Write state (uses configured reducer)
ctx.set("counter", serde_json::json!(count.unwrap_or(0) + 1));
// Emit progress events
ctx.emit(StreamEvent::custom("my_task", "progress", json!({"pct": 50})));
// Interrupt for human-in-the-loop
let approval: bool = ctx.interrupt("Please approve").await?;
ReducedValue<T>
チェックポイント間で永続化される追記専用の状態コンテナ。値は蓄積され、上書きされることはありません。
use adk_graph::functional::ReducedValue;
let mut results: ReducedValue<String> = ReducedValue::new();
results.push("step 1 output".to_string());
results.push("step 2 output".to_string());
assert_eq!(results.len(), 2);
assert_eq!(&results[0], "step 1 output");
UntrackedValue<T>
チェックポイントの永続化対象から除外される一時的なランタイム値。再開時にデフォルト値へリセットされます。
use adk_graph::functional::UntrackedValue;
let mut temp: UntrackedValue<Vec<u8>> = UntrackedValue::new();
temp.set(vec![1, 2, 3]);
// After checkpoint restore: temp.get() == &[]
MessagesValue
メッセージ ID に基づく自動重複排除機能を備えたチャットメッセージコンテナ。
use adk_graph::functional::{MessagesValue, ChatMessage, MessageRole};
let mut messages = MessagesValue::new();
messages.push(ChatMessage {
id: "msg-1".to_string(),
role: MessageRole::User,
content: "Hello".to_string(),
metadata: None,
});
// Pushing same ID replaces the message (upsert)
StateSchemaValidator
ワークフローの境界で状態の型を検証し、不一致を早期に検出します。
use adk_graph::functional::{StateSchemaValidator, ExpectedType};
use adk_graph::state::StateSchema;
let validator = StateSchemaValidator::new(schema)
.expect_type("counter", ExpectedType::Number)
.expect_type("status", ExpectedType::String)
.require_field("status");
validator.validate_state(&state)?; // Fails with descriptive error
ExecutionLog
再開時にタスクをスキップする動作のために、タスクの完了状況を追跡します。完了済みのタスクはワークフローの再開時にスキップされます。
use adk_graph::functional::ExecutionLog;
let mut log = ExecutionLog::new();
log.record_start("fetch");
log.record_completion("fetch", json!({"data": [1,2,3]}));
// On resume: skip completed tasks
if log.is_completed("fetch") {
let cached = log.get_result("fetch"); // Returns cached result
}
バックグラウンド実行
adk-server の background 機能は、非同期ワークフロー実行のための REST エンドポイントを追加します。
エンドポイント
| メソッド | パス | 説明 |
|---|---|---|
| POST | /runs | バックグラウンド実行を送信 |
| GET | /runs/{run_id} | 実行ステータスをポーリング |
| 削除 | /runs/{run_id} | 実行をキャンセル |
使用方法
実行は workflowId に名前を付けます。その名前を処理に変換するものが必要なため、
WorkflowExecutor を登録してください。登録しない場合、送信された実行は完了を報告するのではなく
失敗します。
use adk_server::background::{
BackgroundState, WorkflowRegistry, background_runs_router_with_state,
};
use serde_json::json;
use std::sync::Arc;
let registry = WorkflowRegistry::new().register("summarize", |input, cancel| async move {
if cancel.is_cancelled() {
return Err("cancelled before starting".to_string());
}
Ok(json!({ "summary": "…", "of": input.get("document") }))
});
let state = BackgroundState::new().with_executor(Arc::new(registry));
let app = axum::Router::new().merge(background_runs_router_with_state(state));
未登録の workflowId を送信すると、実行できない実行がキューに追加されるのではなく、
404 が返されます。WorkflowExecutor を直接実装して adk-graph、機能的な API、または独自のディスパッチャーに接続し、実行を失敗させるには Err を返してください。これにより、再試行予算が意味を持つようになります。
ステータスのライフサイクル
queued → running → completed
→ failed (retries if configured)
→ cancelled (via DELETE)
再試行では、元の入力を使用してワークフローを最初から再実行します。チェックポイントには対応していないため、副作用のあるワークフローはべき等にするか、進行状況を独自に保護する必要があります。
重要: 実行記録はメモリ内ストアに保存されます。プロセスを再起動すると、実行中の実行を含むすべての記録が失われます。これらのエンドポイントは永続的なジョブキューではなく、非同期実行のためのインターフェースとして扱ってください。
Cron スケジューリング
background 機能には、cron ジョブの管理も含まれています。
エンドポイント
| メソッド | パス | 説明 |
|---|---|---|
| POST | /cron | cron ジョブを作成 |
| GET | /cron | すべてのジョブを一覧表示 |
| GET | /cron/{job_id} | ジョブを1件取得 |
| PATCH | /cron/{job_id} | 一時停止/再開 |
| DELETE | /cron/{job_id} | ジョブを削除 |
並行実行ポリシー
- skip: 前回の実行がまだアクティブな場合は、その発生をスキップします(デフォルト)
- allow: 並行実行を許可します
- queue: その発生をキューに追加し、アクティブな実行が終了したときに実行します
各スケジュールの発生は一度だけ取得されるため、実行中に到着した発生は、スケジューラーのポーリングごとに1つずつではなく、正確に1つのキュー実行を生成します。キューに入った実行は順番に処理され、それぞれが直接トリガーされた実行と同じ方法で監視されるため、実行が失敗またはキャンセルされた場合でも、アクティブな数は最終的にゼロに戻ります。
Cron ジョブはバックグラウンド実行と同じエグゼキューターを使用します。workflowId が登録されていないジョブは、その実行に失敗します。
使用方法
use adk_server::background::{BackgroundState, CronState, cron_jobs_router_with_state, start_cron_scheduler};
let bg_state = BackgroundState::new();
let cron_state = CronState::new(bg_state);
let app = axum::Router::new().merge(cron_jobs_router_with_state(cron_state.clone()));
// Start the background scheduler
start_cron_scheduler(cron_state);
例
# Functional API (TaskContext, ReducedValue, MessagesValue, etc.)
cargo run --manifest-path examples/functional_workflow/Cargo.toml
# Background Runs (REST API with Axum)
cargo run --manifest-path examples/background_runs/Cargo.toml
# Cron Scheduling (full lifecycle demo)
cargo run --manifest-path examples/cron_scheduling/Cargo.toml
機能フラグ
| 機能 | クレート | 追加されるもの |
|---|---|---|
functional | adk-graph | TaskContext、型付きリデューサー、スキーマ検証、手続きマクロ |
background | adk-server | バックグラウンド実行エンドポイント、cronスケジューリング、スケジューラーループ |
前へ: ← グラフエージェント | 次へ: リアルタイムエージェント →
中断と型付き再開
TaskContext::interrupt<T> はワークフローを中断し、後の実行時に、そのサイトに提供された値を返します。
各中断サイトは、実行内の位置から 継続キー を取得します — interrupt-1、
interrupt-2 など — これにより、再実行されたワークフローは同じ順序で同じ中断に到達し、
渡された値を見つけられます。
// First run: no value supplied, so the workflow suspends.
let error = ctx.interrupt::<Approval>("approve the refund").await.unwrap_err();
// workflow suspended at interrupt 'interrupt-1': approve the refund
// Second run: supply the value under that key.
let ctx = ctx.with_resume_values(HashMap::from([
("interrupt-1".to_string(), serde_json::json!({ "approved": true, "approver": "alice" })),
]));
let approval: Approval = ctx.interrupt("approve the refund").await?;
このキーは、continuation_key の下にある中断チェックポイントにも書き込まれます。
| 状況 | 結果 |
|---|---|
| キーに値がない | FunctionalError::Suspended { continuation_key, message } |
T にデシリアライズされる値 | Ok(value) |
FunctionalError::InterruptTypeMismatch にデシリアライズされない値 | サイトを示す T |
重要: 以前の
interruptはすべての呼び出しで「ワークフローが中断されました」を示すInterruptTypeMismatchを返し、メソッドの外部で再開値を消費するものはありませんでした。呼び出し元は「入力が必要」なのか「値の型が正しくない」のかを判別できず、値を提供するためのキーもなく、呼び出し箇所で型付きの値を受け取ることもできませんでした。
注:
From<FunctionalError> for GraphErrorはGraphError::Other(String)にフラット化されるため、構造的にマッチする呼び出し元は変換前にFunctionalErrorに対してマッチさせる必要があります。