फ़ंक्शनल API
स्वचालित चेकपॉइंटिंग, टाइप्ड स्थिति रिड्यूसर और इंटरप्ट/रिज़्यूम समर्थन के साथ agent वर्कफ़्लो को सामान्य async Rust फ़ंक्शन के रूप में लिखें।
अवलोकन
फ़ंक्शनल API (functional में adk-graph सुविधा) स्पष्ट ग्राफ़/नोड/एज निर्माण का उच्च-स्तरीय विकल्प प्रदान करता है। मैन्युअल रूप से StateGraph बनाने के बजाय, आप फ़ंक्शन पर #[entrypoint] और #[task] मैक्रो लागू करते हैं और मानक Rust नियंत्रण प्रवाह का उपयोग करते हैं।
मुख्य लाभ:
- वर्कफ़्लो को सामान्य async Rust के रूप में लिखें — किसी ग्राफ़ DSL की आवश्यकता नहीं
- क्रैश से पुनर्प्राप्ति के लिए प्रत्येक टास्क के बाद स्वचालित चेकपॉइंटिंग
- मानक Rust नियंत्रण प्रवाह (if/else, for, loop, match)
- स्थायित्व की गारंटी वाले टाइप्ड स्थिति कंटेनर
- मौजूदा StreamEvent और Checkpointer इन्फ्रास्ट्रक्चर के साथ संगत
शुरुआत करना
[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 सुविधा async वर्कफ़्लो निष्पादन के लिए 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)
पुनःप्रयास मूल इनपुट के साथ वर्कफ़्लो को शुरुआत से फिर निष्पादित करता है। यह चेकपॉइंट-जागरूक नहीं है, इसलिए साइड इफ़ेक्ट वाले वर्कफ़्लो को इडेम्पोटेंट होना चाहिए या अपनी प्रगति की स्वयं सुरक्षा करनी चाहिए।
महत्वपूर्ण: रन रिकॉर्ड इन-मेमोरी स्टोर में रहते हैं। प्रक्रिया के पुनः आरंभ होने पर सभी रिकॉर्ड, जिनमें चल रहे रन भी शामिल हैं, खो जाते हैं। इन एंडपॉइंट्स को स्थायी जॉब कतार के बजाय असिंक्रोनस निष्पादन सतह के रूप में मानें।
क्रॉन शेड्यूलिंग
background सुविधा में क्रॉन जॉब प्रबंधन भी शामिल है।
एंडपॉइंट્સ
| विधि | पथ | विवरण |
|---|---|---|
| POST | /cron | क्रॉन जॉब बनाएँ |
| GET | /cron | सभी जॉब सूचीबद्ध करें |
| GET | /cron/{job_id} | एक जॉब प्राप्त करें |
| PATCH | /cron/{job_id} | रोकें/फिर से शुरू करें |
| DELETE | /cron/{job_id} | जॉब हटाएँ |
समवर्ती नीतियाँ
- skip: यदि पिछला रन अभी भी सक्रिय है, तो इस घटना को छोड़ दें (डिफ़ॉल्ट)
- allow: समवर्ती निष्पादन की अनुमति दें
- queue: इस घटना को कतार में रखें और सक्रिय रन पूरा होने पर इसे चलाएँ
प्रत्येक शेड्यूल घटना का दावा केवल एक बार किया जाता है, इसलिए सक्रिय रन के दौरान आने वाली घटना प्रत्येक शेड्यूलर पोल के लिए एक रन के बजाय ठीक एक कतारबद्ध रन उत्पन्न करती है। कतारबद्ध रन क्रम से पूरे होते हैं, और प्रत्येक की निगरानी उसी तरह की जाती है जैसे सीधे ट्रिगर किए गए रन की, इसलिए रन विफल होने या रद्द किए जाने पर भी सक्रिय संख्या शून्य पर लौट आती है।
क्रॉन जॉब बैकग्राउंड रन के समान 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
फ़ीचर फ़्लैग-ಗಳು
| सुविधा | क्रेट | जोड़ता है |
|---|---|---|
functional | adk-graph | TaskContext, टाइप किए गए रिड्यूसर, स्कीमा सत्यापन, proc macros |
background | adk-server | बैकग्राउंड रन एंडपॉइंट, क्रॉन शेड्यूलिंग, शेड्यूलर लूप |
पिछला: ← Graph Agents | अगला: Realtime Agents →
इंटरप्ट और टाइपयुक्त पुनरारंभ
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पहलेInterruptTypeMismatchको हर कॉल पर "workflow interrupted" के साथ लौटाता था, और विधि के बाहर कुछ भी resume value का उपयोग नहीं करता था। कॉलर यह नहीं बता सकता था कि "इनपुट आवश्यक है" या "आपका मान गलत प्रकार का था", उसके पास किसी मान को उपलब्ध कराने के लिए कोई कुंजी नहीं थी, और कॉल साइट पर उसे कभी typed value प्राप्त नहीं होती थी।
नोट:
From<FunctionalError> for GraphError,GraphError::Other(String)में flatten होता है, इसलिए संरचनात्मक रूप से मिलान करने वाले कॉलर को conversion से पहलेFunctionalErrorपर मिलान करना चाहिए।