الوكلاء المحيطيون

تشغيل وكيل محيطي

يلزم وجود معالج للمحفز. المسار المدعوم هو 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جدول زمني لـ cronNone — لا يوجد مستدعٍ للجدول الزمني
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لا شيء
النص الأساسي غير صالح JSON400لا شيء
النص الأساسي غير صالح 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

ملاحظة: كان المستمع سابقًا يتجاوز عمر المستهلك. أدى إسقاط الدفق إلى إبقاء الخادم مرتبطًا، مع استمراره في قبول الطلبات التي لم يتمكن من تسليمها، وفشل إعادة التشغيل على المنفذ نفسه.