Agentes Ambientais
Executando um agente ambiental
É necessário um manipulador de gatilho. O caminho compatível é with_invoker, que aceita qualquer coisa
que implemente adk_core::AgentInvoker — Runner faz o seguinte:
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 três coisas:
| Método | Finalidade | Padrão |
|---|---|---|
new(user_id) | Identidade sob a qual as execuções são registradas. Um gatilho não tem um usuário interativo, portanto use "system" ou o nome de uma conta de serviço. | obrigatório |
with_session_policy | PerTrigger fornece a cada evento sua própria sessão; Shared(id) reutiliza uma. | PerTrigger |
with_prompt | Converte o evento em texto de prompt. | indica a origem e serializa o payload |
PerTrigger é o padrão porque um agendamento disparado a cada minuto em uma única sessão compartilhada faz o histórico dessa sessão — e o custo de tokens de cada execução posterior — crescer sem limite.
Runner serializa as execuções invocadas externamente que têm como destino a mesma sessão compartilhada até que cada fluxo de eventos seja concluído, enquanto IDs de sessão diferentes ainda podem ser executados simultaneamente.
AgentInvoker::invoke cria a sessão quando ela não existe. Runner::run não faz isso: ele resolve uma sessão existente e, caso contrário, produz session.not_found pelo fluxo, algo que uma execução acionada externamente não tem a oportunidade de pré-registrar.
Quando o invocador expõe sua raiz executável, como Runner faz, with_invoker usa esse agente para registro e diagnósticos no contexto. Isso impede que a telemetria identifique um agente enquanto outro realmente trata o acionamento.
Fornecendo um manipulador diretamente
with_trigger_handler continua disponível para chamadores que controlam algo diferente de um Runner. Ele recebe o evento e o agente e deve retornar o fluxo de eventos; nesse caso, a criação da sessão fica sob responsabilidade do manipulador.
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 falha sem um manipulador:
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 sem um manipulador anteriormente funcionava e depois registrava cada acionamento, fazendo com que
AmbientAgent::new(..).start()parecesse estar executando um agente que nunca era executado.
Saída e simultaneidade
| Comportamento | Controle |
|---|---|
| Eventos e erros produzidos pelo agente são entregues a um canal | take_output(capacity) |
| Gatilhos tratados de uma vez | with_max_concurrent_triggers (padrão 4, zero considerado como um) |
Os eventos produzidos anteriormente eram registrados no nível de depuração e descartados, portanto o chamador não podia observar o que uma execução havia feito ou se ela havia falhado. Os gatilhos também eram tratados estritamente um de cada vez — o loop consumia todo o fluxo de eventos de um manipulador antes de consultar a fonte novamente — portanto, um gatilho lento bloqueava todos os posteriores.
Observação: o limite controla a simultaneidade, não o paralelismo. Os manipuladores compartilham a tarefa ambiente, portanto um manipulador que bloqueia a thread ainda interrompe o loop. Use
tokio::task::spawn_blockingpara esses casos. Os deslocamentos duráveis dos gatilhos, o tratamento de mensagens não entregues e as novas tentativas continuam sendo responsabilidade do chamador. Os agentes de ambiente reagem a eventos em vez de a uma interação do usuário. UmEventSourceproduzTriggerEvents, e o agente é executado para cada um deles.
Requer o recurso ambient:
[dependencies]
adk-agent = { version = "2.1.0", features = ["ambient"] }
Fontes de eventos
| Fonte | Dispara em | Principal |
|---|---|---|
CronTrigger | Uma agenda cron | None — uma agenda não tem chamador |
FileWatchTrigger | Uma alteração correspondente no sistema de arquivos | None — uma alteração de arquivo não tem chamador |
WebhookTrigger | Um POST autorizado por HTTP | O resultado do verificador |
TriggerEvent::principal permite que um manipulador diferencie um acionador autorizado de um anônimo, em vez de tratar todos os eventos como igualmente confiáveis.
Tiques perdidos
CronTrigger::subscribe calcula o próximo tique a partir do momento em que é chamado. Portanto, um acionador que reinicia após um período de inatividade ou é executado em um host que entra em suspensão retoma no próximo tique futuro, e todos os tiques que ocorreram nesse intervalo são descartados.
MissedTickPolicy decide o que acontece com esse intervalo:
| Política | Comportamento | Usar para |
|---|---|---|
Skip | Descartar os intervalos transcorridos e aguardar o próximo agendado. O padrão. | Agendamentos em que uma execução atrasada não tem valor |
CoalesceOne | Emitir um evento abrangendo todo o intervalo transcorrido. | Varreduras em que apenas o estado atual importa |
All | Emite um evento por intervalo transcorrido, começando pelo mais antigo. | Agendamentos em que cada ocorrência tem seu próprio trabalho |
Uma política, por si só, cobre apenas lacunas dentro de uma assinatura. Detectar uma lacuna que abrange uma reinicialização do processo requer um TickWatermark para registrar onde o agendamento foi interrompido:
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 armazena um cursor RFC 3339. Ele grava por meio de um arquivo temporário irmão exclusivo, sincroniza esse arquivo e substitui atomicamente o destino no Unix e no Windows. Implemente TickWatermark para outros armazenamentos subjacentes.
Limitando uma reprodução
All em um agendamento frequente pode deixar milhares de marcações pendentes após uma longa interrupção. with_max_catch_up limita quantas uma única passagem reproduz (o padrão é 64); quando o limite é atingido, o restante da lacuna é descartado, o cursor persistente avança além dela e o gatilho retoma na próxima marcação futura, registrando quantas marcações foram descartadas. Reiniciar antes da próxima marcação normal não recupera o restante descartado.
Contrato de entrega
A marca d’água avança quando o gatilho emite uma marcação, não quando o consumidor termina de agir sobre ela. Uma falha entre a emissão e a conclusão descarta essa execução em vez de repeti-la — no máximo uma vez, não pelo menos uma vez. É isso que impede que um consumidor que parou de consultar reproduza a mesma lacuna a cada reinicialização. Consumidores cujo trabalho precisa sobreviver a uma falha no meio da execução devem registrar seu próprio estado de conclusão. Se não for possível persistir uma marca d’água configurada, o fluxo cron para antes de emitir o evento afetado, em vez de enfraquecer silenciosamente essa garantia.
Gatilhos de webhook
Um webhook acessível é um ponto de entrada remoto na lógica da aplicação, portanto WebhookTrigger usa loopback por padrão e se recusa a atender em um endereço mais amplo sem um verificador.
Desenvolvimento local
use adk_agent::ambient::WebhookTrigger;
// Binds 127.0.0.1 — reachable only from this host.
let trigger = WebhookTrigger::new(8080, "/webhook");
Listeners expostos exigem um 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())));
Inscrever-se em um endereço que não seja de loopback sem um verificador falha com
agent.ambient.webhook_unauthenticated. A verificação ocorre no momento da inscrição, quando o
erro ainda é barato.
Importante: um verificador recebe o corpo bruto por meio de
WebhookRequest::body, porque os esquemas de assinatura são calculados sobre os bytes exatos recebidos. Verifique antes de confiar em qualquer forma analisada da solicitação.
Tratamento de solicitações
| Condição | Resposta | Evento |
|---|---|---|
Corpo acima do limite (padrão de 1 MiB, with_max_body_bytes) | 413 | nenhum |
| Verificador rejeita | 401 | nenhum |
| O corpo não é válido JSON | 400 | nenhum |
O corpo não é válido JSON, accept_non_json() definido | 200 | o payload é uma string JSON |
| O assinante se desconectou | 503 | nenhum |
| Caso contrário | 200 | o payload é o JSON analisado |
Um 401 não contém detalhes sobre qual parte de uma credencial falhou; o motivo é registrado em vez disso, portanto o endpoint não pode ser usado para sondar credenciais.
accept_non_json é desativado por padrão porque um corpo malformado encapsulado como uma string produz um evento de acionamento indistinguível de um evento deliberado.
Ciclo de vida
O listener HTTP pertence ao fluxo retornado por subscribe. Descartar o fluxo encerra o servidor normalmente e libera a porta, de modo que a mesma porta possa ser vinculada novamente:
let stream = trigger.subscribe().await?;
// ... consume events ...
drop(stream); // the listener stops and the port is free
Observação: anteriormente, o listener sobrevivia ao consumidor. Descartar o fluxo deixava o servidor vinculado, ainda aceitando solicitações que não podia entregar, e uma reinicialização na mesma porta falhava.