الوكلاء المحيطيون
تشغيل وكيل محيطي
يلزم وجود معالج للمحفز. المسار المدعوم هو with_invoker، والذي يقبل أي شيء
ينفّذ adk_core::AgentInvoker — ويقوم Runner بذلك:
use adk_agent::ambient::{AmbientAgent, RunnerTriggerConfig};
use std::sync::Arc;
let mut ambient = AmbientAgent::new(agent, source)
.with_invoker(runner, RunnerTriggerConfig::new("system"))
.with_max_concurrent_triggers(4);
let mut outputs = ambient.take_output(64);
ambient.start().await?;
يتحكم RunnerTriggerConfig في ثلاثة أشياء:
| الطريقة | الغرض | الافتراضي |
|---|---|---|
new(user_id) | الهوية التي تُسجَّل العمليات باسمها. لا يملك المشغّل مستخدمًا تفاعليًا، لذا استخدم "system" أو اسم حساب خدمة. | مطلوب |
with_session_policy | ينشئ PerTrigger جلسة خاصة به لكل حدث؛ بينما يعيد Shared(id) استخدام جلسة واحدة. | PerTrigger |
with_prompt | يحوّل الحدث إلى نص مطالبة. | يحدّد المصدر ويُسلسل الحمولة |
PerTrigger هو الإعداد الافتراضي، لأن تشغيل جدول زمني كل دقيقة في جلسة مشتركة واحدة يؤدي إلى نمو سجل تلك الجلسة — وتكلفة الرموز لكل تشغيل لاحق — بلا حدود.
تعمل Runner على تسلسل الأدوار المستدعاة خارجيًا التي تستهدف الجلسة المشتركة نفسها إلى أن ينتهي كل تدفق أحداث، بينما تظل معرّفات الجلسات المختلفة قادرة على العمل بالتزامن.
تنشئ AgentInvoker::invoke الجلسة عندما لا تكون موجودة. أما Runner::run فلا تفعل ذلك: إذ تحلّ جلسة موجودة وتُخرج session.not_found عبر التدفق بخلاف ذلك، وهو ما لا يتوفر للتشغيل الذي يُشغَّل خارجيًا أي فرصة لتسجيله مسبقًا.
عندما يكشف المستدعي عن جذر التنفيذ الخاص به، كما تفعل Runner، تستخدم with_invoker ذلك الوكيل للتسجيل والتشخيص التلقائيين. ويمنع هذا بيانات القياس عن تسمية وكيل واحد بينما يتولى وكيل آخر معالجة المشغّل فعليًا.
توفير معالج مباشرةً
تظل with_trigger_handler متاحة للمتصلين الذين يقودون شيئًا غير Runner. وهي تستقبل الحدث والوكيل، ويجب أن تُعيد تدفق الأحداث؛ وعندئذٍ تصبح مسؤولية إنشاء الجلسة على عاتق المعالج.
use adk_agent::ambient::{AmbientAgent, TriggerHandler};
use std::sync::Arc;
let handler: TriggerHandler = Arc::new(move |event, agent| {
let backend = backend.clone();
Box::pin(async move { backend.dispatch(event, agent).await })
});
let mut ambient = AmbientAgent::new(agent, source).with_trigger_handler(handler);
تفشل start من دون معالج:
AmbientAgent has no trigger handler, so starting it would log trigger events without ever
invoking the agent. Call `with_trigger_handler` with a closure that drives the agent through a
Runner.
مهم: كان البدء من دون معالج ينجح سابقًا، ثم يسجّل كل مشغّل، لذلك بدت
AmbientAgent::new(..).start()وكأنها تشغّل وكيلًا لم يُشغَّل قط.
المخرجات والتزامن
| السلوك | التحكم |
|---|---|
| تُسلَّم الأحداث والأخطاء التي ينتجها الوكيل إلى قناة | take_output(capacity) |
| المحفزات التي تتم معالجتها دفعة واحدة | with_max_concurrent_triggers (الافتراضي 4، ويُعامل الصفر على أنه واحد) |
كانت الأحداث المُنتجة تُسجَّل سابقًا بمستوى تصحيح الأخطاء ثم تُهمَل، لذلك لم يكن بإمكان المستدعي معرفة ما فعلته عملية التشغيل أو ما إذا كانت قد فشلت. كما كانت المُشغِّلات تُعالَج واحدًا تلو الآخر بدقة — إذ كانت الحلقة تستنزف مجرى الأحداث بالكامل لأحد المعالجات قبل استطلاع المصدر مجددًا — ولذلك كان مُشغِّل بطيء واحد يحجب جميع المُشغِّلات اللاحقة.
ملاحظة: يحدّد الحدّ التزامن، وليس التوازي. تشترك المعالجات في المهمة المحيطة، لذا فإن المعالج الذي يحجب الخيط يظل يوقف الحلقة. استخدم
tokio::task::spawn_blockingلمثل هذه الحالات. تظل إزاحات المُشغِّلات الدائمة، ومعالجة الرسائل الميتة، وإعادة المحاولة من مسؤولية المستدعي.
تستجيب الوكلاء المحيطة للأحداث بدلًا من استجابتها لدور مستخدم. ينتج EventSource
TriggerEvents، ويعمل الوكيل من أجل كل واحد منها.
يتطلب ميزة ambient:
[dependencies]
adk-agent = { version = "2.1.0", features = ["ambient"] }
مصادر الأحداث
| المصدر | يعمل عند | الجهة الرئيسية |
|---|---|---|
CronTrigger | جدول زمني لـ cron | None — لا يوجد مستدعٍ للجدول الزمني |
FileWatchTrigger | تغيير مطابق في نظام الملفات | None — لا يوجد مستدعٍ لتغيير الملف |
WebhookTrigger | طلب POST مُصرَّح به من HTTP | نتيجة أداة التحقق |
TriggerEvent::principal يتيح للمعالج التمييز بين مُشغّل مُصرَّح به وآخر مجهول، بدلًا من التعامل مع كل حدث على أنه موثوق بالقدر نفسه.
النبضات الفائتة
يحسب CronTrigger::subscribe النبضة التالية انطلاقًا من اللحظة التي يتم استدعاؤه فيها. ولذلك، فإن المُشغّل الذي يُعاد تشغيله بعد فترة توقف، أو الذي يعمل على مضيف يدخل في وضع التعليق، يستأنف عند النبضة المستقبلية التالية، وتُهمَل كل نبضة كان موعدها قد حلّ في الأثناء.
يحدّد MissedTickPolicy ما يحدث خلال تلك الفترة:
| السياسة | السلوك | الاستخدام |
|---|---|---|
Skip | تجاهل النبضات المنقضية وانتظر النبضة المجدولة التالية. هذا هو الإعداد الافتراضي. | الجداول الزمنية التي لا قيمة فيها للتنفيذ المتأخر |
CoalesceOne | إصدار حدث واحد يغطي كامل الفترة المنقضية. | عمليات المسح التي تكون فيها الحالة الحالية وحدها مهمة |
All | إصدار حدث واحد لكل نبضة زمنية منقضية، بدءًا بالأقدم. | الجداول الزمنية التي يكون لكل تكرار فيها عمله الخاص |
لا تغطي السياسة وحدها سوى الفجوات داخل اشتراك واحد. يتطلب اكتشاف فجوة تمتد عبر إعادة تشغيل العملية وجود TickWatermark لتسجيل الموضع الذي توقف عنده الجدول:
use std::sync::Arc;
use adk_agent::ambient::{CronTrigger, FileTickWatermark, MissedTickPolicy};
let trigger = CronTrigger::new("0 */5 * * * *")?
.with_missed_tick_policy(MissedTickPolicy::CoalesceOne)
.with_watermark(Arc::new(FileTickWatermark::new("/var/lib/my-agent/sweep.tick")));
يخزّن FileTickWatermark مؤشرًا واحدًا بتنسيق RFC 3339. ويكتب عبر ملف مؤقت شقيق وفريد،
ويزامنه، ويستبدل الوجهة ذريًا على Unix وWindows. نفّذ TickWatermark لمخازن الدعم الأخرى.
تحديد إعادة التشغيل
يمكن أن يترك All وفق جدول متكرر آلاف النبضات المعلّقة بعد انقطاع طويل.
يحدّد with_max_catch_up عدد مرات إعادة التشغيل التي ينفذها تمرير واحد (القيمة الافتراضية 64)؛ وبمجرد بلوغ الحد،
يُهمل الجزء المتبقي من الفجوة، ويتقدم المؤشر الدائم إلى ما بعده، ويستأنف المشغّل عند
النبضة المستقبلية التالية، مع تسجيل عدد النبضات التي أُسقطت. لا تؤدي إعادة التشغيل قبل النبضة العادية
التالية إلى استعادة الجزء المتبقي الذي أُهمل.
عقد التسليم
يتقدم المؤشر المائي عندما يُصدر المشغّل نبضة، وليس عندما ينتهي المستهلك من التعامل معها. يؤدي التعطل بين الإصدار والإكمال إلى إسقاط ذلك التشغيل بدلًا من تكراره — مرة واحدة على الأكثر، وليس مرة واحدة على الأقل. وهذا ما يمنع المستهلك الذي يتوقف عن الاستطلاع من إعادة تشغيل الفجوة نفسها عند كل إعادة تشغيل. يجب على المستهلكين الذين ينبغي أن يستمر عملهم بعد تعطل في منتصف التشغيل تسجيل حالة إكمالهم الخاصة. وإذا تعذر استمرار مؤشر مائي مُكوَّن، يتوقف تدفق cron قبل إصدار الحدث المتأثر بدلًا من إضعاف هذا الضمان بصمت.
مشغلات Webhook
إن webhook القابل للوصول هو نقطة دخول بعيدة إلى منطق التطبيق، ولذلك يضبط WebhookTrigger
افتراضيًا على loopback ويرفض تقديم الخدمة على عنوان أوسع دون أداة تحقق.
التطوير المحلي
use adk_agent::ambient::WebhookTrigger;
// Binds 127.0.0.1 — reachable only from this host.
let trigger = WebhookTrigger::new(8080, "/webhook");
المستمعون المكشوفون يتطلبون أداة تحقق
use adk_agent::ambient::{WebhookRequest, WebhookTrigger, WebhookVerifier};
use std::net::SocketAddr;
use std::sync::Arc;
#[derive(Debug)]
struct SharedToken(String);
impl WebhookVerifier for SharedToken {
fn verify(&self, request: &WebhookRequest<'_>) -> Result<String, String> {
match request.header("x-webhook-token") {
Some(value) if value == self.0 => Ok("ci-system".to_string()),
Some(_) => Err("token mismatch".to_string()),
None => Err("missing x-webhook-token".to_string()),
}
}
}
let address: SocketAddr = "0.0.0.0:8080".parse().unwrap();
let trigger = WebhookTrigger::new(8080, "/webhook")
.with_bind_address(address)
.with_verifier(Arc::new(SharedToken(std::env::var("WEBHOOK_TOKEN").unwrap_or_default())));
يفشل الاشتراك بعنوان غير استرجاعي من دون أداة تحقق مع
agent.ambient.webhook_unauthenticated. يحدث الفحص وقت الاشتراك، حين يكون تصحيح
الخطأ لا يزال قليل التكلفة.
مهم: تتلقى أداة التحقق النص الأساسي الخام عبر
WebhookRequest::body، لأن مخططات التوقيع تُحتسب على وحدات البايت الدقيقة المستلمة. تحقّق قبل الوثوق بأي شكل محلّل من الطلب.
معالجة الطلب
| الشرط | الاستجابة | الحدث |
|---|---|---|
تجاوز الجسم الحد الأقصى (1 MiB افتراضيًا، with_max_body_bytes) | 413 | لا شيء |
| يرفض المدقّق | 401 | لا شيء |
| النص الأساسي غير صالح JSON | 400 | لا شيء |
النص الأساسي غير صالح JSON، تم تعيين accept_non_json() | 200 | الحمولة عبارة عن سلسلة نصية JSON |
| غادر المشترك | 503 | لا شيء |
| بخلاف ذلك | 200 | الحمولة هي JSON المُحلَّلة |
لا يحمل 401 أي تفاصيل حول الجزء من بيانات الاعتماد الذي فشل؛ بل يُسجَّل السبب
بدلًا من ذلك، ولذلك لا يمكن استخدام نقطة النهاية لفحص بيانات الاعتماد.
يكون accept_non_json معطّلًا افتراضيًا لأن الجسم غير الصالح المضمَّن كسلسلة نصية ينتج
حدث تشغيل لا يمكن تمييزه عن حدث متعمَّد.
مدة الصلاحية
ينتمي مستمع HTTP إلى الدفق الذي يعيده subscribe. يؤدي إسقاط الدفق إلى إيقاف تشغيل
الخادم بسلاسة وتحرير المنفذ، ولذلك يمكن ربط المنفذ نفسه مجددًا:
let stream = trigger.subscribe().await?;
// ... consume events ...
drop(stream); // the listener stops and the port is free
ملاحظة: كان المستمع سابقًا يتجاوز عمر المستهلك. أدى إسقاط الدفق إلى إبقاء الخادم مرتبطًا، مع استمراره في قبول الطلبات التي لم يتمكن من تسليمها، وفشل إعادة التشغيل على المنفذ نفسه.