Agentes ambientales
Ejecutar un agente ambiental
Se requiere un controlador de activación. La ruta compatible es with_invoker, que acepta cualquier elemento
que implemente adk_core::AgentInvoker; Runner lo hace:
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 controla tres aspectos:
| Método | Propósito | Predeterminado |
|---|---|---|
new(user_id) | Identidad con la que se registran las ejecuciones. Un activador no tiene un usuario interactivo, así que usa "system" o el nombre de una cuenta de servicio. | obligatorio |
with_session_policy | PerTrigger proporciona a cada evento su propia sesión; Shared(id) reutiliza una. | PerTrigger |
with_prompt | Convierte el evento en texto de prompt. | indica el origen y serializa la carga útil |
PerTrigger es el valor predeterminado porque una ejecución programada cada minuto en una única sesión compartida hace crecer sin límite el historial de esa sesión —y el coste en tokens de cada ejecución posterior—.
Runner serializa los turnos invocados externamente que apuntan a la misma sesión compartida hasta que finaliza cada flujo de eventos, mientras que las sesiones con distintos ID aún pueden ejecutarse simultáneamente.
AgentInvoker::invoke crea la sesión cuando no existe. Runner::run no lo hace: resuelve una sesión existente y, de lo contrario, produce session.not_found a través del flujo, algo que una ejecución activada externamente no tiene oportunidad de registrar previamente.
Cuando el invocador expone su raíz ejecutable, como hace Runner, with_invoker usa ese agente para el registro ambiental y los diagnósticos. Esto evita que la telemetría indique un agente mientras otro es el que realmente gestiona el activador.
Proporcionar un controlador directamente
with_trigger_handler sigue disponible para los llamadores que controlan algo distinto de un Runner. Recibe el evento y el agente, y debe devolver el flujo de eventos; por tanto, crear la sesión es responsabilidad del controlador.
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 falla sin un controlador:
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.
Importante: iniciar sin un controlador antes funcionaba correctamente y después registraba cada activador, por lo que
AmbientAgent::new(..).start()parecía estar ejecutando un agente que nunca se ejecutaba.
Salida y concurrencia
| Comportamiento | Control |
|---|---|
| Los eventos y errores que produce el agente se entregan a un canal | take_output(capacity) |
| Desencadenadores gestionados de inmediato | with_max_concurrent_triggers (4 de forma predeterminada, cero se trata como uno) |
Los eventos producidos se registraban anteriormente en el nivel de depuración y se descartaban, por lo que quien llamaba no podía observar qué había hecho una ejecución ni si había fallado. Los activadores también se gestionaban estrictamente de uno en uno: el bucle agotaba todo el flujo de eventos de un controlador antes de volver a consultar el origen, por lo que un activador lento bloqueaba todos los posteriores.
Nota: el límite controla la concurrencia, no el paralelismo. Los controladores comparten la tarea ambiental, por lo que un controlador que bloquea el subproceso sigue deteniendo el bucle. Usa
tokio::task::spawn_blockingpara esos casos. Los desplazamientos duraderos de los activadores, la gestión de mensajes no entregables y los reintentos siguen siendo responsabilidad de quien realiza la llamada. Los agentes ambientales reaccionan a eventos en lugar de a un turno del usuario. UnEventSourceproduceTriggerEvents y el agente se ejecuta para cada uno.
Requiere la característica ambient:
[dependencies]
adk-agent = { version = "2.1.0", features = ["ambient"] }
Orígenes de eventos
| Fuente | Se activa con | Principal |
|---|---|---|
CronTrigger | Una programación cron | None — una programación no tiene invocador |
FileWatchTrigger | Un cambio coincidente en el sistema de archivos | None — un cambio de archivo no tiene invocador |
WebhookTrigger | Un POST autorizado HTTP | El resultado del verificador |
TriggerEvent::principal permite que un controlador distinga un activador autorizado de uno
anónimo, en lugar de tratar cada evento como igualmente confiable.
Intervalos omitidos
CronTrigger::subscribe calcula el siguiente intervalo desde el momento en que se invoca. Por lo tanto, un activador que
se reinicia después de un periodo de inactividad o se ejecuta en un host que se suspende, reanuda su ejecución en el siguiente
intervalo futuro y descarta todos los intervalos que vencieron entretanto.
MissedTickPolicy decide qué ocurre con ese periodo:
| Política | Comportamiento | Usar para |
|---|---|---|
Skip | Descarta los intervalos transcurridos y espera al siguiente programado. La opción predeterminada. | Programaciones en las que una ejecución tardía no aporta valor |
CoalesceOne | Emite un evento que cubre todo el intervalo transcurrido. | Barridos en los que solo importa el estado actual |
All | Emite un evento por cada intervalo transcurrido, comenzando por el más antiguo. | Programaciones en las que cada ocurrencia tiene su propio trabajo |
Una política por sí sola solo cubre las brechas dentro de una suscripción. Detectar una brecha que abarque un reinicio del proceso requiere un TickWatermark para registrar dónde quedó el horario:
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 almacena un cursor RFC 3339. Escribe mediante un archivo temporal hermano único, lo sincroniza y reemplaza atómicamente el destino en Unix y Windows. Implementa TickWatermark para otros almacenes subyacentes.
Limitar una reproducción
All en un horario frecuente puede dejar miles de ejecuciones pendientes después de una interrupción prolongada. with_max_catch_up limita cuántas se reproducen en una pasada (64 de forma predeterminada); una vez alcanzado el límite, el resto de la brecha se descarta, el cursor persistente avanza más allá de ella y el activador se reanuda en la siguiente ejecución futura, registrando cuántas ejecuciones se descartaron. Reiniciar antes de la siguiente ejecución ordinaria no recupera el resto descartado.
Contrato de entrega
La marca de agua avanza cuando el activador emite una ejecución, no cuando el consumidor termina de actuar sobre ella. Un fallo entre la emisión y la finalización descarta esa ejecución en lugar de repetirla: como máximo una vez, no al menos una vez. Esto evita que un consumidor que deja de consultar reproduzca la misma brecha en cada reinicio. Los consumidores cuyo trabajo deba sobrevivir a un fallo durante la ejecución deben registrar su propio estado de finalización. Si no se puede persistir una marca de agua configurada, el flujo de cron se detiene antes de emitir el evento afectado, en lugar de debilitar silenciosamente esta garantía.
Activadores de webhook
Un webhook accesible es un punto de entrada remoto a la lógica de la aplicación, por lo que WebhookTrigger usa loopback de forma predeterminada y se niega a servir en una dirección más amplia sin un verificador.
Desarrollo local
use adk_agent::ambient::WebhookTrigger;
// Binds 127.0.0.1 — reachable only from this host.
let trigger = WebhookTrigger::new(8080, "/webhook");
Los listeners expuestos requieren un verificador
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())));
Suscribirse a una dirección que no sea de bucle invertido sin un verificador falla con
agent.ambient.webhook_unauthenticated. La comprobación se realiza al suscribirse, cuando el
error todavía es fácil de corregir.
Importante: un verificador recibe el cuerpo sin procesar mediante
WebhookRequest::body, porque los esquemas de firma se calculan sobre los bytes exactos recibidos. Verifique antes de confiar en cualquier forma analizada de la solicitud.
Gestión de solicitudes
| Condición | Respuesta | Evento |
|---|---|---|
Cuerpo por encima del límite (1 MiB predeterminado, with_max_body_bytes) | 413 | ninguno |
| El verificador rechaza | 401 | ninguno |
| El cuerpo no es válido JSON | 400 | ninguno |
El cuerpo no es válido JSON, conjunto accept_non_json() | 200 | la carga útil es una cadena JSON |
| El suscriptor se ha desconectado | 503 | ninguno |
| De lo contrario | 200 | la carga útil es el JSON analizado |
Un 401 no contiene detalles sobre qué parte de una credencial falló; en su lugar, se registra el motivo, por lo que el endpoint no puede utilizarse para sondear credenciales.
accept_non_json está desactivado de forma predeterminada porque un cuerpo malformado envuelto como cadena produce un evento de activación indistinguible de uno deliberado.
Ciclo de vida
El listener HTTP pertenece al flujo devuelto por subscribe. Al descartar el flujo, el servidor se apaga correctamente y libera el puerto, por lo que se puede volver a vincular el mismo puerto:
let stream = trigger.subscribe().await?;
// ... consume events ...
drop(stream); // the listener stops and the port is free
Nota: anteriormente, el listener sobrevivía a su consumidor. Al descartar el flujo, el servidor permanecía vinculado y seguía aceptando solicitudes que no podía entregar, por lo que un reinicio en el mismo puerto fallaba.