API الوظيفية
اكتب سير عمل الوكلاء كدوال Rust عادية غير متزامنة مع إنشاء نقاط تحقق تلقائيًا، ومختزلات حالة مُنَمّطة، ودعم المقاطعة والاستئناف.
نظرة عامة
توفر API (ميزة functional في adk-graph) بديلاً عالي المستوى لإنشاء الرسم البياني/العقد/الحواف بشكل صريح. بدلاً من إنشاء StateGraph يدويًا، أضف التعليقات التوضيحية إلى الدوال باستخدام وحدات الماكرو #[entrypoint] و#[task]، واستخدم تدفق التحكم القياسي في Rust.
الفوائد الرئيسية:
- كتابة سير العمل كدوال Rust عادية غير متزامنة — دون الحاجة إلى لغة خاصة بالرسوم البيانية
- إنشاء نقاط تحقق تلقائيًا بعد كل مهمة لاستعادة العمل بعد الأعطال
- تدفق تحكم قياسي في 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
حاوية لرسائل المحادثة مع إزالة التكرارات تلقائيًا استنادًا إلى معرّفات الرسائل.
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
}
عمليات التشغيل في الخلفية
تضيف ميزة background في adk-server نقاط نهاية 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 | إنشاء مهمة مجدولة |
| GET | /cron | سرد جميع المهام |
| GET | /cron/{job_id} | الحصول على مهمة واحدة |
| PATCH | /cron/{job_id} | إيقاف مؤقت/استئناف |
| DELETE | /cron/{job_id} | حذف مهمة |
سياسات التزامن
- skip: تخطَّ التكرار إذا كان تشغيل سابق لا يزال نشطًا (افتراضيًا)
- allow: اسمح بالتنفيذات المتزامنة
- queue: ضع التكرار في قائمة الانتظار وشغّله عند انتهاء التشغيل النشط
يتم المطالبة بكل تكرار للجدول مرة واحدة، لذا فإن التكرار الذي يصل أثناء نشاط تشغيل
ينتج تشغيلًا واحدًا بالضبط في قائمة الانتظار، وليس تشغيلًا لكل عملية استطلاع للجدولة.
تُفرَّغ التشغيلات الموجودة في قائمة الانتظار بالترتيب، وتتم مراقبة كل منها بالطريقة نفسها
التي تتم بها مراقبة تشغيل مُشغَّل مباشرةً، ولذلك يعود العدد النشط إلى الصفر حتى إذا فشل
تشغيل أو أُلغي.
تستخدم وظائف 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) |
قيمة لا يمكن فك تسلسلها إلى T | FunctionalError::InterruptTypeMismatch تسمّي الموقع |
مهم: كان
interruptيُرجع سابقًاInterruptTypeMismatchمع عبارة "تمت مقاطعة سير العمل" في كل استدعاء، ولم يكن أي شيء خارج الطريقة يستهلك قيمة الاستئناف. لم يكن بإمكان المستدعي التمييز بين "يتطلب إدخالًا" و"نوع قيمتك غير صحيح"، ولم يكن لديه مفتاح يقدّم القيمة تحته، ولم يتلقَّ قيمة مكتوبة النوع في موضع الاستدعاء.
ملاحظة: يُبسّط
From<FunctionalError> for GraphErrorإلىGraphError::Other(String)، لذا ينبغي على المستدعي الذي يطابق بنيويًا أن يطابقFunctionalErrorقبل التحويل.