परिवेश एजेंट
परिवेश एजेंट चलाना
एक ट्रिगर हैंडलर आवश्यक है। समर्थित पथ 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 | इवेंट को प्रॉम्प्ट टेक्स्ट में बदलता है। | स्रोत बताता है और payload को serialize करता है |
PerTrigger डिफ़ॉल्ट है क्योंकि हर मिनट एक साझा सत्र में सक्रिय होने वाला शेड्यूल उस सत्र का इतिहास — और बाद के प्रत्येक रन की टोकन लागत — को बिना सीमा के बढ़ाता रहता है।
Runner उसी साझा सत्र को लक्षित करने वाले बाहरी रूप से 호출ित टर्न को तब तक क्रमबद्ध करता है, जब तक प्रत्येक इवेंट स्ट्रीम समाप्त नहीं हो जाती, जबकि अलग-अलग सत्र ID अभी भी समवर्ती रूप से चल सकती हैं।
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, शून्य को एक माना जाता है) |
उत्पादित इवेंट पहले debug स्तर पर लॉग किए जाते थे और फिर छोड़ दिए जाते थे, इसलिए कोई कॉलर यह नहीं देख सकता था कि किसी रन ने क्या किया या वह विफल हुआ या नहीं। ट्रिगर भी सख्ती से एक-एक करके संभाले जाते थे — लूप स्रोत को फिर से पोल करने से पहले किसी हैंडलर की पूरी इवेंट स्ट्रीम को समाप्त कर देता था — इसलिए एक धीमा ट्रिगर उसके बाद आने वाले सभी ट्रिगर को रोक देता था।
नोट: यह सीमा समवर्तीता को नियंत्रित करती है, समानांतरता को नहीं। हैंडलर परिवेशीय टास्क साझा करते हैं, इसलिए जो हैंडलर थ्रेड को ब्लॉक करता है, वह लूप को अब भी रोक देता है। ऐसे हैंडलरों के लिए
tokio::task::spawn_blockingका उपयोग करें। टिकाऊ ट्रिगर ऑफ़सेट, डेड-लेटर प्रबंधन और पुनःप्रयास की ज़िम्मेदारी कॉलर की बनी रहती है। परिवेशीय एजेंट उपयोगकर्ता के टर्न के बजाय इवेंट पर प्रतिक्रिया देते हैं। एकEventSourceTriggerEvents उत्पन्न करता है और एजेंट प्रत्येक के लिए चलता है।
ambient फ़ीचर आवश्यक है:
[dependencies]
adk-agent = { version = "2.1.0", features = ["ambient"] }
इवेंट स्रोत
| स्रोत | इस पर सक्रिय होता है | प्रधान |
|---|---|---|
CronTrigger | क्रॉन शेड्यूल | None — शेड्यूल का कोई कॉलर नहीं होता |
FileWatchTrigger | मेल खाता फ़ाइल सिस्टम परिवर्तन | None — फ़ाइल परिवर्तन का कोई कॉलर नहीं होता |
WebhookTrigger | एक अधिकृत HTTP POST | सत्यापक का परिणाम |
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); सीमा पहुँचने पर अंतराल का शेष भाग छोड़ दिया जाता है, स्थायी कर्सर उसके आगे बढ़ जाता है, और ट्रिगर अगले भविष्य के टिक से फिर शुरू होता है तथा छोड़े गए टिकों की संख्या लॉग करता है। अगले सामान्य टिक से पहले पुनः आरंभ करने पर छोड़ा गया शेष भाग पुनर्प्राप्त नहीं होता।
डिलीवरी अनुबंध
वॉटरमार्क तब आगे बढ़ता है जब ट्रिगर कोई टिक उत्सर्जित करता है, न कि तब जब उपभोक्ता उस पर कार्रवाई पूरी करता है। उत्सर्जन और पूर्णता के बीच क्रैश होने पर वह रन दोहराने के बजाय छूट जाता है — अधिकतम एक बार, न्यूनतम एक बार नहीं। यही बात किसी ऐसे उपभोक्ता को, जो पोलिंग करना बंद कर देता है, हर पुनः आरंभ पर उसी अंतराल को दोबारा चलाने से रोकती है। जिन उपभोक्ताओं का कार्य बीच में क्रैश होने पर भी सुरक्षित रहना आवश्यक है, उन्हें अपनी पूर्णता स्थिति स्वयं दर्ज करनी चाहिए। यदि कॉन्फ़िगर किया गया वॉटरमार्क स्थायी रूप से संग्रहीत नहीं किया जा सकता, तो क्रोन स्ट्रीम प्रभावित इवेंट को उत्सर्जित करने से पहले रुक जाती है, बजाय इसके कि यह गारंटी चुपचाप कमजोर कर दी जाए।
वेबहुक ट्रिगर
किसी पहुँच योग्य वेबहुक से एप्लिकेशन लॉजिक में दूरस्थ प्रवेश-बिंदु मिलता है, इसलिए WebhookTrigger डिफ़ॉल्ट रूप से लूपबैक पर रहता है और सत्यापक के बिना व्यापक पते पर सेवा देने से इनकार करता है।
स्थानीय विकास
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())));
किसी सत्यापक के बिना non-loopback पते पर सदस्यता लेने पर 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
नोट: पहले लिस्नर अपने उपभोक्ता से अधिक समय तक जीवित रहता था। स्ट्रीम को ड्रॉप करने पर सर्वर बाइंड रहता था और ऐसे अनुरोध स्वीकार करता रहता था जिन्हें वह वितरित नहीं कर सकता था, तथा उसी पोर्ट पर पुनः आरंभ करना विफल हो जाता था।