기능적 API

자동 체크포인트, 타입이 지정된 상태 리듀서, 인터럽트/재개 지원을 통해 일반적인 async Rust 함수로 에이전트 워크플로를 작성할 수 있습니다.

개요

기능적 API(adk-graphfunctional 기능)는 명시적인 그래프/노드/엣지 구성에 대한 더 높은 수준의 대안을 제공합니다. 수동으로 StateGraph을 구축하는 대신, 함수에 #[entrypoint]#[task] 매크로를 적용하고 표준 Rust 제어 흐름을 사용합니다.

주요 이점:

  • 그래프 DSL 없이 일반적인 async Rust로 워크플로 작성
  • 충돌 복구를 위해 각 작업 후 자동 체크포인트 생성
  • 표준 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-serverbackground 기능은 비동기 워크플로 실행을 위한 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/croncron 작업 생성
GET/cron모든 작업 나열
GET/cron/{job_id}작업 하나 가져오기
PATCH/cron/{job_id}일시중지/재개
DELETE/cron/{job_id}작업 삭제

동시성 정책

  • skip: 이전 실행이 아직 활성 상태이면 해당 발생을 건너뜁니다(기본값).
  • allow: 동시 실행을 허용합니다.
  • queue: 해당 발생을 대기열에 추가하고 활성 실행이 완료되면 실행합니다.

각 일정 발생은 한 번만 처리되므로, 실행이 활성 상태일 때 도착한 발생은 스케줄러 폴링마다 하나씩이 아니라 정확히 하나의 대기 실행을 생성합니다. 대기 중인 실행은 순서대로 처리되며, 직접 트리거된 실행과 동일한 방식으로 각각 모니터링됩니다. 따라서 실행이 실패하거나 취소되더라도 활성 실행 수는 0으로 돌아갑니다.

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, 타입이 지정된 리듀서, 스키마 검증, 프로시저 매크로
backgroundadk-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)
T로 역직렬화되지 않는 값사이트를 지정하는 FunctionalError::InterruptTypeMismatch

중요: interrupt는 이전에 모든 호출에서 "workflow interrupted"와 함께 InterruptTypeMismatch를 반환했으며, 메서드 외부에서는 재개 값을 소비하지 않았습니다. 호출자는 "입력이 필요함"과 "값의 형식이 잘못됨"을 구분할 수 없었고, 값을 제공할 키도 없었으며, 호출 지점에서 타입이 지정된 값을 받지도 못했습니다.

참고: From<FunctionalError> for GraphErrorGraphError::Other(String)로 평탄화되므로, 구조적으로 일치시키는 호출자는 변환 전에 FunctionalError에 대해 일치시켜야 합니다.