环境代理

运行环境代理

需要一个触发器处理程序。受支持的路径是 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_policyPerTrigger 为每个事件创建独立会话;Shared(id) 重用同一个会话。PerTrigger
with_prompt将事件转换为提示文本。说明来源并序列化有效负载

PerTrigger 是默认设置,因为每分钟触发一次并写入同一个共享会话,会使该会话的历史记录以及之后每次运行的令牌成本无限增长。 Runner 会将面向同一个共享会话的外部调用轮次串行化,直到每个事件流完成;而不同的会话 ID 仍可并发运行。

如果会话不存在,AgentInvoker::invoke 会创建该会话。Runner::run 则不会:它会解析一个已有的会话,否则会通过流产生 session.not_found,而外部触发的运行没有机会预先注册会话。

当调用方公开其可执行根目录时(正如 Runner 所做的那样),with_invoker 会使用该 agent 进行环境日志记录和诊断。这样可以避免遥测数据中显示的是一个 agent,而实际处理触发器的却是另一个 agent。

直接提供处理程序

对于驱动非 Runner 对象的调用方,with_trigger_handler 仍然可用。它接收事件和 agent,并且必须返回事件流;因此,创建会话的责任由处理程序承担。

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() 看起来像是在运行一个实际上从未运行的 agent。

输出与并发

行为控制
agent 产生的事件和错误会发送到一个通道take_output(capacity)
立即处理触发器with_max_concurrent_triggers(默认为 4,零按 1 处理)

生成的事件此前以调试级别记录后被丢弃,因此调用方无法观察一次运行执行了什么,也无法得知运行是否失败。触发器也严格按顺序逐个处理——循环会先耗尽某个处理程序的整个事件流,然后才再次轮询源——因此,一个缓慢的触发器会阻塞其后的所有触发器。

注意: 此限制约束的是并发数,而不是并行性。处理程序共享环境任务,因此阻塞线程的处理程序仍会阻塞循环。对于这类情况,请使用 tokio::task::spawn_blocking。持久化触发器偏移量、死信处理和重试仍由调用方负责。 环境代理会对事件作出反应,而不是对用户回合作出反应。一个 EventSource 会生成 TriggerEvent,代理会针对每个事件运行一次。

需要启用 ambient 功能:

[dependencies]
adk-agent = { version = "2.1.0", features = ["ambient"] }

事件源

来源触发于主体
CronTriggerCron 计划None — 计划没有调用方
FileWatchTrigger匹配的文件系统更改None — 文件更改没有调用方
WebhookTrigger已授权的 HTTP POST验证器的结果

TriggerEvent::principal 让处理程序能够区分经过授权的触发器和匿名触发器,而不是将每个事件都视为同等可信。

错过的 tick

CronTrigger::subscribe 根据其被调用的时刻计算下一个 tick。因此,停机后重新启动的触发器,或运行在发生挂起的主机上的触发器,会在下一个未来的 tick 恢复运行,而期间到期的每个 tick 都会被丢弃。

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 默认绑定到回环地址,并且在没有验证器的情况下拒绝监听更广泛的地址。

本地开发

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_bytes413
验证器拒绝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

注意: 此前监听器的生命周期长于其使用者。丢弃流会使服务器仍保持绑定状态,继续接受它无法传递的请求,并导致在同一端口上重启失败。

环境代理 - ADK-Rust 文档 | ADK-Rust