앰비언트 에이전트
앰비언트 에이전트 실행
트리거 핸들러가 필요합니다. 지원되는 경로는 with_invoker이며, Runner가 수행하는 것처럼 adk_core::AgentInvoker를 구현하는 모든 것을 받습니다.
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) | 실행이 기록되는 ID입니다. 트리거에는 대화형 사용자가 없으므로 "system" 또는 서비스 계정 이름을 사용합니다. | 필수 |
with_session_policy | PerTrigger는 각 이벤트에 자체 세션을 부여하고, Shared(id)는 하나의 세션을 재사용합니다. | PerTrigger |
with_prompt | 이벤트를 프롬프트 텍스트로 변환합니다. | 소스를 명시하고 페이로드를 직렬화합니다. |
PerTrigger은 기본값입니다. 1분마다 하나의 공유 세션에서 실행되는 일정은 해당 세션의 기록을 늘리고, 이후 모든 실행의 토큰 비용을 제한 없이 증가시키기 때문입니다.
Runner은 동일한 공유 세션을 대상으로 외부에서 호출된 턴을 각 이벤트 스트림이 완료될 때까지 직렬화하지만, 서로 다른 세션 ID는 여전히 동시에 실행될 수 있습니다.
AgentInvoker::invoke은 세션이 존재하지 않을 때 세션을 생성합니다. 반면 Runner::run는 세션을 생성하지 않습니다. 기존 세션을 확인하고, 그렇지 않으면 스트림을 통해 session.not_found을 반환하지만, 외부에서 트리거된 실행에는 세션을 미리 등록할 기회가 없습니다.
호출자가 실행 가능한 루트를 노출하는 경우(Runner가 그러함), with_invoker는 해당 에이전트를 주변 로깅 및 진단에 사용합니다. 이를 통해 실제로 트리거를 처리하는 에이전트와 텔레메트리에 표시되는 에이전트가 서로 달라지는 것을 방지합니다.
핸들러 직접 제공
with_trigger_handler은 Runner이 아닌 다른 대상을 구동하는 호출자를 위해 계속 사용할 수 있습니다. 이벤트와 에이전트를 전달받으며 이벤트 스트림을 반환해야 합니다. 따라서 세션 생성은 핸들러의 책임이 됩니다.
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()가 실제로는 실행되지 않은 에이전트를 실행 중인 것처럼 보였습니다.
출력 및 동시성
| 동작 | 제어 |
|---|---|
| 에이전트가 생성하는 이벤트와 오류가 채널로 전달됨 | take_output(capacity) |
| 트리거가 즉시 처리됨 | with_max_concurrent_triggers (기본값 4, 0은 1로 처리) |
생성된 이벤트는 이전에는 디버그 수준에서 기록된 후 삭제되었으므로, 호출자는 실행에서 어떤 일이 발생했는지 또는 실패했는지를 확인할 수 없었습니다. 트리거도 한 번에 엄격하게 하나씩만 처리되었습니다. 즉, 루프가 소스를 다시 폴링하기 전에 핸들러의 전체 이벤트 스트림을 비워야 했으므로, 하나의 느린 트리거가 이후의 모든 트리거를 차단했습니다.
참고: 이 제한은 병렬성이 아니라 동시성을 제어합니다. 핸들러는 동일한 주변 작업을 공유하므로, 스레드를 차단하는 핸들러는 여전히 루프를 멈춥니다. 이러한 경우에는
tokio::task::spawn_blocking를 사용하세요. 영속적인 트리거 오프셋, 배달 불가 처리 및 재시도는 여전히 호출자의 책임입니다.
주변 에이전트는 사용자의 턴이 아니라 이벤트에 반응합니다. EventSource이 TriggerEvent을 생성하면 에이전트가 각 이벤트마다 실행됩니다.
ambient 기능이 필요합니다:
[dependencies]
adk-agent = { version = "2.1.0", features = ["ambient"] }
이벤트 소스
| 소스 | 다음에 실행됨 | 주체 |
|---|---|---|
CronTrigger | cron 일정 | None — 일정에는 호출자가 없음 |
FileWatchTrigger | 일치하는 파일 시스템 변경 | None — 파일 변경에는 호출자가 없음 |
WebhookTrigger | 승인된 HTTP POST | 검증자의 결과 |
TriggerEvent::principal를 사용하면 핸들러가 모든 이벤트를 동일하게 신뢰하는 대신 인증된 트리거와 익명 트리거를 구분할 수 있습니다.
누락된 틱
CronTrigger::subscribe는 호출된 시점부터 다음 틱을 계산합니다. 따라서 중단 시간 후 다시 시작되는 트리거나 일시 중단된 호스트에서 실행되는 트리거는 다음 미래의 틱에서 재개되며, 그 사이에 도래한 모든 틱은 폐기됩니다.
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 스트림은 영향을 받는 이벤트를 내보내기 전에 중지되며 이 보장이 조용히 약화되지 않습니다.
웹훅 트리거
접근 가능한 웹훅은 애플리케이션 로직으로 들어가는 원격 진입점이므로, 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_bytes) | 413 | 없음 |
| 검증자가 거부함 | 401 | 없음 |
| 본문이 유효하지 않음 JSON | 400 | 없음 |
본문이 유효하지 않음 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
참고: 이전에는 리스너의 수명이 해당 소비자보다 길었습니다. 스트림을 삭제해도 서버는 바인딩된 상태로 남아 전달할 수 없는 요청을 계속 수락했으며, 동일한 포트에서 다시 시작할 수 없었습니다.