関数型 API

通常の非同期 Rust 関数としてエージェントのワークフローを記述し、自動チェックポイント、型付き状態リデューサー、割り込み/再開サポートを利用できます。

概要

API(functionaladk-graph 機能)は、明示的なグラフ/ノード/エッジ構築に代わる、より高レベルな代替手段を提供します。StateGraph を手動で構築する代わりに、#[entrypoint]#[task] のマクロで関数に注釈を付け、標準の Rust の制御フローを使用します。

主な利点:

  • ワークフローを通常の非同期 Rust として記述できる — グラフ DSL は不要
  • 各タスク後に自動でチェックポイントを作成し、クラッシュ復旧を支援
  • 標準の Rust 制御フロー(if/else、for、loop、match)
  • 永続化保証付きの型付き状態コンテナ
  • 既存の StreamEvent および Checkpointer インフラストラクチャと互換

はじめに

[dependencies]
adk-graph = { version = "2.0.0", features = ["functional"] }
adk-rust-macros = "2.0.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
}

バックグラウンド実行

backgroundadk-server 機能は、非同期ワークフロー実行のための REST エンドポイントを追加します。

エンドポイント

メソッドパス説明
POST/runsバックグラウンド実行を送信する
GET/runs/{run_id}実行ステータスをポーリングする
DELETE/runs/{run_id}実行をキャンセルする

使用方法

run は workflowId に名前を付けます。その名前を作業に変換するものが必要なので、WorkflowExecutor を登録してください。これがないと、送信された run は完了を報告する代わりに 失敗 します:

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 を送信すると、実行不能な run をキューに入れるのではなく 404 が返されます。WorkflowExecutor を直接実装して、adk-graph、関数型の API、または独自のディスパッチャーへ接続してください。run を失敗させるには Err を返します。これにより、リトライ予算に意味が生まれます。

ステータスのライフサイクル

queued → running → completed
                 → failed (retries if configured)
                 → cancelled (via DELETE)

リトライは、元の入力を使ってワークフローを 最初から 再実行します。チェックポイントを認識しないため、副作用のあるワークフローは冪等にするか、自身の進行を保護する必要があります。

重要: run レコードはインメモリストアに保存されます。プロセスを再起動すると、進行中の run を含め、すべてのレコードが失われます。これらのエンドポイントは永続的なジョブキューではなく、非同期実行のための表面として扱ってください。

Cron スケジューリング

background 機能には、cron ジョブ管理も含まれます。

エンドポイント

メソッドパス説明
POST/croncron job を作成する
GET/cronすべての job を一覧表示する
GET/cron/{job_id}1件のジョブを取得
PATCH/cron/{job_id}一時停止/再開
DELETE/cron/{job_id}ジョブを削除

同時実行ポリシー

  • skip: 直前の実行がまだアクティブな場合、その発生をスキップします(デフォルト)
  • allow: 同時実行を許可します
  • queue: 発生をキューに入れ、アクティブな実行が終了したときに実行します

各スケジュール発生は 1 回だけ取得されるため、実行がアクティブな間に到着した発生は、スケジューラのポーリングごとではなく、ちょうど 1 つのキュー済み実行を生成します。キュー済み実行は順番に処理され、直接トリガーされた実行と同じ方法で監視されるため、実行が失敗したりキャンセルされたりしても、アクティブ数はゼロに戻ります。

Cron ジョブはバックグラウンド実行と同じ executor を使用します。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

機能フラグ

機能クレート追加内容
functionaladk-graphTaskContext、typed reducers、schema validation、proc macros
backgroundadk-serverバックグラウンド実行エンドポイント、cronスケジューリング、スケジューラループ

Previous: ← Graph Agents | Next: Realtime Agents →

割り込みと型付き再開

TaskContext::interrupt<T> はワークフローを一時停止し、後の実行で、その箇所に対して渡された値を返します。

各割り込み箇所には、実行内での位置から continuation key が割り当てられます。interrupt-1interrupt-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?;

その key は、continuation_key の下で割り込みチェックポイントにも書き込まれます。

状況結果
キーに値がないFunctionalError::Suspended { continuation_key, message }
T にデシリアライズされる値Ok(value)
T にデシリアライズされない値サイトに名前を付ける FunctionalError::InterruptTypeMismatch

重要: interrupt は以前、InterruptTypeMismatch毎回 "workflow interrupted" とともに返しており、メソッド外で resume 値を消費するものはありませんでした。呼び出し側は「入力が必要」なのか「指定した値の型が間違っていた」のかを判別できず、どのキーの下に値を渡せばよいかも分からず、呼び出し箇所で型付きの値を受け取ることもありませんでした。

注: From<FunctionalError> for GraphErrorGraphError::Other(String) に平坦化されるため、 構造的に一致を行う呼び出し側は、変換前に FunctionalError で一致させる必要があります。