Agents ambiants

Exécuter un agent ambiant

Un gestionnaire de déclencheur est requis. Le chemin pris en charge est with_invoker, qui accepte tout ce qui implémente adk_core::AgentInvokerRunner le fait :

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 contrôle trois éléments :

MéthodeObjectifValeur par défaut
new(user_id)Identité sous laquelle les exécutions sont enregistrées. Un déclencheur n’a pas d’utilisateur interactif ; utilisez donc "system" ou le nom d’un compte de service.obligatoire
with_session_policyPerTrigger attribue à chaque événement sa propre session ; Shared(id) en réutilise une.PerTrigger
with_promptTransforme l’événement en texte d’invite.indique la source et sérialise la charge utile

PerTrigger est utilisé par défaut, car un déclenchement chaque minute dans une session partagée unique fait croître sans limite l'historique de cette session — ainsi que le coût en jetons de chaque exécution ultérieure. Runner sérialise les tours invoqués de l'extérieur ciblant la même session partagée jusqu'à la fin de chaque flux d'événements, tandis que des identifiants de session différents peuvent toujours s'exécuter simultanément.

AgentInvoker::invoke crée la session lorsqu'elle n'existe pas. Runner::run ne le fait pas : il résout une session existante et produit session.not_found via le flux dans le cas contraire, ce qui laisse à une exécution déclenchée de l'extérieur aucune possibilité de la préenregistrer.

Lorsque l'appelant expose sa racine exécutable, comme le fait Runner, with_invoker utilise cet agent pour la journalisation et les diagnostics ambiants. Cela empêche la télémétrie de nommer un agent alors qu'un autre traite effectivement le déclencheur.

Fournir directement un gestionnaire

with_trigger_handler reste disponible pour les appelants qui pilotent autre chose qu'un Runner. Il reçoit l'événement et l'agent, et doit retourner le flux d'événements ; la création de la session relève alors de la responsabilité du gestionnaire.

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 échoue sans gestionnaire :

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.

Important : démarrer sans gestionnaire réussissait auparavant, puis journalisait chaque déclencheur, de sorte que AmbientAgent::new(..).start() donnait l'impression d'exécuter un agent qui ne s'était jamais exécuté.

Sortie et concurrence

ComportementContrôle
Les événements et les erreurs produits par l’agent sont transmis à un canaltake_output(capacity)
Déclencheurs traités simultanémentwith_max_concurrent_triggers (4 par défaut, zéro est traité comme un)

Les événements produits étaient auparavant consignés au niveau de débogage, puis ignorés, de sorte qu’un appelant ne pouvait pas observer ce qu’une exécution avait fait ni déterminer si elle avait échoué. Les déclencheurs étaient également traités strictement un par un : la boucle épuisait tout le flux d’événements d’un gestionnaire avant d’interroger à nouveau la source ; ainsi, un déclencheur lent bloquait tous les suivants.

Remarque : la limite régit la concurrence, et non le parallélisme. Les gestionnaires partagent la tâche ambiante ; un gestionnaire qui bloque le thread immobilise donc toujours la boucle. Utilisez tokio::task::spawn_blocking dans ce cas. La responsabilité des décalages persistants des déclencheurs, de la gestion des lettres mortes et des nouvelles tentatives reste à la charge de l’appelant. Les agents ambiants réagissent aux événements plutôt qu’à une intervention de l’utilisateur. Un EventSource produit des TriggerEvents et l’agent s’exécute pour chacun d’eux.

Nécessite la fonctionnalité ambient :

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

Sources d’événements

SourceSe déclenche surPrincipal
CronTriggerUne planification cronNone — une planification n’a aucun appelant
FileWatchTriggerUne modification correspondante du système de fichiersNone — une modification de fichier n’a aucun appelant
WebhookTriggerUn POST HTTP autoriséLe résultat du vérificateur

TriggerEvent::principal permet à un gestionnaire de distinguer un déclencheur autorisé d’un déclencheur anonyme, plutôt que de traiter chaque événement comme étant également digne de confiance.

Tics manqués

CronTrigger::subscribe calcule le prochain tic à partir du moment où il est appelé. Un déclencheur qui redémarre après une interruption, ou qui s’exécute sur un hôte suspendu, reprend donc au prochain tic futur, et chaque tic arrivé à échéance entre-temps est ignoré.

MissedTickPolicy détermine ce qui se passe pendant cet intervalle :

PolitiqueComportementÀ utiliser pour
SkipIgnorer les intervalles écoulés et attendre le prochain intervalle planifié. Par défaut.Planifications pour lesquelles une exécution tardive n’a aucune valeur
CoalesceOneÉmettre un événement couvrant toute la période écoulée.Parcours pour lesquels seul l’état actuel importe
AllÉmet un événement par cycle écoulé, du plus ancien au plus récent.Planifications où chaque occurrence possède son propre travail

Une stratégie à elle seule ne couvre que les lacunes au sein d’un même abonnement. Détecter une lacune qui s’étend sur un redémarrage du processus nécessite un TickWatermark pour enregistrer l’endroit où la planification s’est arrêtée :

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 stocke un curseur RFC 3339. Il écrit via un fichier temporaire frère unique, le synchronise, puis remplace atomiquement la destination sous Unix et Windows. Implémentez TickWatermark pour les autres magasins de sauvegarde.

Limiter une relecture

All selon une planification fréquente peut laisser des milliers de déclenchements en attente après une longue panne. with_max_catch_up limite le nombre de relectures effectuées en un seul passage (64 par défaut) ; une fois la limite atteinte, le reste de la lacune est abandonné, le curseur persistant avance au-delà de celle-ci et le déclencheur reprend au prochain déclenchement futur, en journalisant le nombre de déclenchements supprimés. Un redémarrage avant le prochain déclenchement ordinaire ne récupère pas le reste abandonné.

Contrat de livraison

Le filigrane avance lorsque le déclencheur émet un déclenchement, et non lorsque le consommateur termine son traitement. Un plantage entre l’émission et l’achèvement supprime cette exécution au lieu de la répéter — au plus une fois, et non au moins une fois. C’est ce qui empêche un consommateur qui cesse d’interroger le système de rejouer la même lacune à chaque redémarrage. Les consommateurs dont le travail doit survivre à un plantage en cours d’exécution doivent enregistrer leur propre état d’achèvement. Si un filigrane configuré ne peut pas être persisté, le flux cron s’arrête avant d’émettre l’événement concerné, plutôt que d’affaiblir silencieusement cette garantie.

Déclencheurs webhook

Un webhook accessible constitue un point d’entrée distant vers la logique de l’application ; WebhookTrigger est donc défini par défaut sur l’interface de bouclage et refuse de servir une adresse plus large sans vérificateur.

Développement local

use adk_agent::ambient::WebhookTrigger;

// Binds 127.0.0.1 — reachable only from this host.
let trigger = WebhookTrigger::new(8080, "/webhook");

Les écouteurs exposés nécessitent un vérificateur

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())));

S’abonner à une adresse qui n’est pas de bouclage local sans vérificateur échoue avec agent.ambient.webhook_unauthenticated. La vérification a lieu au moment de l’abonnement, lorsque l’erreur est encore facile à corriger.

Important : un vérificateur reçoit le corps brut via WebhookRequest::body, car les schémas de signature sont calculés sur les octets exacts reçus. Vérifiez avant de faire confiance à toute forme analysée de la requête.

Traitement des requêtes

ConditionRéponseÉvénement
Corps dépassant la limite (1 MiB par défaut, with_max_body_bytes)413aucun
Le vérificateur rejette401aucun
Le corps n’est pas valide JSON400aucun
Le corps n’est pas valide JSON, accept_non_json() défini200la charge utile est une chaîne JSON
L’abonné a disparu503aucun
Sinon200la charge utile est le JSON analysé

Un 401 ne contient aucun détail sur la partie de l’identifiant qui a échoué ; la raison est plutôt enregistrée, de sorte que le point de terminaison ne peut pas être utilisé pour sonder les identifiants.

accept_non_json est désactivé par défaut, car un corps malformé encapsulé en tant que chaîne produit un événement de déclenchement impossible à distinguer d’un événement délibéré.

Durée de vie

L’écouteur HTTP appartient au flux renvoyé par subscribe. La suppression du flux arrête gracieusement le serveur et libère le port, de sorte que le même port peut être réaffecté :

let stream = trigger.subscribe().await?;
// ... consume events ...
drop(stream); // the listener stops and the port is free

Remarque : auparavant, l’écouteur survivait à son consommateur. La suppression du flux laissait le serveur lié, continuant d’accepter des requêtes qu’il ne pouvait pas transmettre, et un redémarrage sur le même port échouait.