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::AgentInvoker — Runner 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éthode | Objectif | Valeur 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_policy | PerTrigger attribue à chaque événement sa propre session ; Shared(id) en réutilise une. | PerTrigger |
with_prompt | Transforme 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
| Comportement | Contrôle |
|---|---|
| Les événements et les erreurs produits par l’agent sont transmis à un canal | take_output(capacity) |
| Déclencheurs traités simultanément | with_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_blockingdans 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. UnEventSourceproduit desTriggerEvents 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
| Source | Se déclenche sur | Principal |
|---|---|---|
CronTrigger | Une planification cron | None — une planification n’a aucun appelant |
FileWatchTrigger | Une modification correspondante du système de fichiers | None — une modification de fichier n’a aucun appelant |
WebhookTrigger | Un 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 :
| Politique | Comportement | À utiliser pour |
|---|---|---|
Skip | Ignorer 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
| Condition | Réponse | Événement |
|---|---|---|
Corps dépassant la limite (1 MiB par défaut, with_max_body_bytes) | 413 | aucun |
| Le vérificateur rejette | 401 | aucun |
| Le corps n’est pas valide JSON | 400 | aucun |
Le corps n’est pas valide JSON, accept_non_json() défini | 200 | la charge utile est une chaîne JSON |
| L’abonné a disparu | 503 | aucun |
| Sinon | 200 | la 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.