ag2-observers-and-alerts

Compare original and translation side by side

🇺🇸

Original

English
🇨🇳

Translation

Chinese

Observers, watches, and alerts

观察者、监视器与告警

When to use

适用场景

  • Observability — log model responses, tool calls, token usage.
  • Runtime safety — block dangerous tool arguments, halt the agent.
  • Reactive metrics — fire on every Nth response, or every M seconds.
  • Loop / repetition detection — catch infinite tool-call loops.
  • Stateful monitoring — anything that needs to remember prior events to decide what to do next.
  • 可观测性 — 记录模型响应、工具调用、Token使用情况。
  • 运行时安全 — 拦截危险的工具参数、终止Agent运行。
  • 响应式指标 — 每N次响应或每M秒触发一次动作。
  • 循环/重复检测 — 捕获无限工具调用循环。
  • 有状态监控 — 需要记录历史事件以决定后续操作的场景。

Two observer shapes

两种观察者类型

ShapeWhenUse
Stateless functionOne-off event hook (logging, metrics)
@observer(EventType)
Stateful classCounters / windows / thresholds / composed triggersSubclass
BaseObserver
Both are stream subscribers under the hood — registered on the agent rather than directly on the stream.
类型适用场景使用方式
无状态函数一次性事件钩子(日志、指标)
@observer(EventType)
有状态类计数器/时间窗口/阈值/组合触发器继承
BaseObserver
两者本质都是流订阅者——注册在Agent上而非直接注册在流上。

60-second recipe —
@observer

60秒快速上手——
@observer

python
from ag2 import Agent, observer
from ag2.config import OpenAIConfig
from ag2.events import ModelResponse

@observer(ModelResponse)
async def log_response(event: ModelResponse) -> None:
    print(f"Model said: {event.content}")

agent = Agent(
    "assistant",
    config=OpenAIConfig(model="gpt-4o-mini"),
    observers=[log_response],
)
Or attach after construction with
@agent.observer(...)
. Per-call observers also supported (
agent.ask("...", observers=[...])
).
Observer callbacks support full dependency injection (
Context
,
Inject
,
Variable
,
Depends
). Filter by event type, multiple types (
ModelRequest | ModelResponse
), or field value (
ToolCallEvent.name == "search"
). Use
interrupt=True
to modify or suppress events before regular subscribers see them.
python
from ag2 import Agent, observer
from ag2.config import OpenAIConfig
from ag2.events import ModelResponse

@observer(ModelResponse)
async def log_response(event: ModelResponse) -> None:
    print(f"Model said: {event.content}")

agent = Agent(
    "assistant",
    config=OpenAIConfig(model="gpt-4o-mini"),
    observers=[log_response],
)
也可在Agent创建后通过
@agent.observer(...)
附加观察者。同时支持单次调用级别的观察者(
agent.ask("...", observers=[...])
)。
观察者回调支持完整的依赖注入(
Context
Inject
Variable
Depends
)。可按事件类型、多种类型(
ModelRequest | ModelResponse
)或字段值(
ToolCallEvent.name == "search"
)过滤事件。使用
interrupt=True
可在常规订阅者看到事件前修改或抑制事件。

Built-in stateful observers

内置有状态观察者

python
from ag2 import Agent
from ag2.observers import LoopDetector, TokenMonitor

agent = Agent(
    "assistant",
    config=config,
    observers=[
        TokenMonitor(warn_threshold=50_000, alert_threshold=100_000),
        LoopDetector(window_size=10, repeat_threshold=3),
    ],
)
  • TokenMonitor
    — tracks cumulative tokens across
    ModelResponse
    and
    TaskCompleted
    . Emits
    WARNING
    /
    CRITICAL
    ObserverAlert
    s as thresholds are crossed. Read state via
    monitor.total_tokens
    .
  • LoopDetector
    — sliding window of recent tool calls. Emits a
    WARNING
    alert when
    repeat_threshold
    consecutive identical calls are seen.
python
from ag2 import Agent
from ag2.observers import LoopDetector, TokenMonitor

agent = Agent(
    "assistant",
    config=config,
    observers=[
        TokenMonitor(warn_threshold=50_000, alert_threshold=100_000),
        LoopDetector(window_size=10, repeat_threshold=3),
    ],
)
  • TokenMonitor
    — 追踪
    ModelResponse
    TaskCompleted
    事件的累计Token消耗。当超过阈值时发出
    WARNING
    /
    CRITICAL
    级别的
    ObserverAlert
    告警。可通过
    monitor.total_tokens
    读取当前状态。
  • LoopDetector
    — 维护最近工具调用的滑动窗口。当连续出现
    repeat_threshold
    次相同调用时发出
    WARNING
    告警。

Custom
BaseObserver

自定义
BaseObserver

A
BaseObserver
pairs a
Watch
(when to fire) with a
process()
method (what to do):
python
from ag2 import Context
from ag2.observers import BaseObserver
from ag2.watch import CadenceWatch
from ag2.events import BaseEvent, ModelResponse
from ag2.events.alert import ObserverAlert, Severity

class AvgCompletionObserver(BaseObserver):
    """Every N responses, emit an INFO alert with avg completion-token count."""

    def __init__(self, window: int = 5) -> None:
        super().__init__("avg-completion", watch=CadenceWatch(n=window, condition=ModelResponse))
        self._window = window

    async def process(self, events: list[BaseEvent], ctx: Context) -> ObserverAlert | None:
        tokens = [e.usage.completion_tokens for e in events if isinstance(e, ModelResponse) and e.usage]
        if not tokens:
            return None
        return ObserverAlert(
            source=self.name,
            severity=Severity.INFO,
            message=f"Avg completion tokens over last {self._window}: {sum(tokens) / len(tokens):.0f}",
        )
If
process()
returns an
ObserverAlert
, the base class emits it onto the stream. You can also send events manually via
await ctx.send(...)
.
BaseObserver
Watch
(触发时机)与
process()
方法(执行动作)结合:
python
from ag2 import Context
from ag2.observers import BaseObserver
from ag2.watch import CadenceWatch
from ag2.events import BaseEvent, ModelResponse
from ag2.events.alert import ObserverAlert, Severity

class AvgCompletionObserver(BaseObserver):
    """每N次响应,发出包含平均生成Token数的INFO级告警。"""

    def __init__(self, window: int = 5) -> None:
        super().__init__("avg-completion", watch=CadenceWatch(n=window, condition=ModelResponse))
        self._window = window

    async def process(self, events: list[BaseEvent], ctx: Context) -> ObserverAlert | None:
        tokens = [e.usage.completion_tokens for e in events if isinstance(e, ModelResponse) and e.usage]
        if not tokens:
            return None
        return ObserverAlert(
            source=self.name,
            severity=Severity.INFO,
            message=f"最近{self._window}次响应的平均生成Token数: {sum(tokens) / len(tokens):.0f}",
        )
如果
process()
返回
ObserverAlert
,基类会将其发送到流中。你也可以通过
await ctx.send(...)
手动发送事件。

Watch primitives — picking when to fire

Watch原语——选择触发时机

You needUse
Every matching event
EventWatch(EventType)
or just
stream.subscribe(fn, condition=...)
Every N matching events
CadenceWatch(n=N, condition=EventType)
Every T seconds (buffered events)
CadenceWatch(max_wait=T, condition=EventType)
Either threshold
CadenceWatch(n=N, max_wait=T, condition=EventType)
Once after delay
DelayWatch(seconds)
Periodic timer
IntervalWatch(seconds)
Cron schedule
CronWatch("0 9 * * MON")
All sub-watches must fire
AllOf(w1, w2)
Any sub-watch fires
AnyOf(w1, w2)
In order
Sequence(w1, w2)
All importable from
ag2.watch
. Callback signature is uniform:
async def cb(events: list[BaseEvent], ctx: Context) -> None
. Time-driven watches pass
events=[]
.
需求使用方式
每次匹配事件触发
EventWatch(EventType)
或直接使用
stream.subscribe(fn, condition=...)
每N次匹配事件触发
CadenceWatch(n=N, condition=EventType)
每T秒触发一次(缓冲事件)
CadenceWatch(max_wait=T, condition=EventType)
满足任一阈值触发
CadenceWatch(n=N, max_wait=T, condition=EventType)
延迟一段时间后触发一次
DelayWatch(seconds)
周期性定时触发
IntervalWatch(seconds)
Cron调度触发
CronWatch("0 9 * * MON")
所有子监视器都触发才执行
AllOf(w1, w2)
任一子监视器触发即执行
AnyOf(w1, w2)
按顺序触发
Sequence(w1, w2)
所有原语均可从
ag2.watch
导入。回调签名统一为:
async def cb(events: list[BaseEvent], ctx: Context) -> None
。时间驱动的监视器会传入
events=[]

ObserverAlert
— the alert type

ObserverAlert
——告警类型

python
from ag2.events.alert import ObserverAlert, Severity

ObserverAlert(
    source="my-observer",
    severity=Severity.WARNING,    # INFO, WARNING, CRITICAL, FATAL
    message="What happened",
)
Important:
ObserverAlert
is on the stream and persisted in history, but the default provider mappers do not render it back to the LLM. To make the agent see alerts, add
AlertPolicy()
to
assembly=[...]
:
python
from ag2.policies import AlertPolicy
agent = Agent("assistant", config=config, assembly=[AlertPolicy()])
python
from ag2.events.alert import ObserverAlert, Severity

ObserverAlert(
    source="my-observer",
    severity=Severity.WARNING,    # INFO, WARNING, CRITICAL, FATAL
    message="What happened",
)
重要提示
ObserverAlert
会被发送到流中并保存在历史记录中,但默认的提供者映射不会将其返回给LLM。要让Agent看到告警,需将
AlertPolicy()
添加到
assembly=[...]
中:
python
from ag2.policies import AlertPolicy
agent = Agent("assistant", config=config, assembly=[AlertPolicy()])

FATAL alerts →
HaltEvent
→ short-circuit

FATAL告警 →
HaltEvent
→ 短路终止

AlertPolicy
does two things on
Severity.FATAL
:
  1. Emits a
    HaltEvent
    on the stream.
  2. Appends a halt notice to the system prompt.
When
assembly=[...]
is non-empty, the harness automatically wires
_HaltCheckMiddleware
which sees the
HaltEvent
and short-circuits the next LLM call with a synthetic
HALTED: ...
response.
python
from ag2 import Context
from ag2.observers import BaseObserver
from ag2.events import BaseEvent, ToolCallEvent
from ag2.events.alert import HaltEvent, ObserverAlert, Severity
from ag2.policies import AlertPolicy
from ag2.watch import EventWatch

class PathGuardian(BaseObserver):
    def __init__(self) -> None:
        super().__init__("path-guardian", watch=EventWatch(ToolCallEvent))

    async def process(self, events: list[BaseEvent], ctx: Context) -> ObserverAlert | None:
        for event in events:
            if not isinstance(event, ToolCallEvent) or event.name != "write_file":
                continue
            if "/etc/" in event.arguments or "/usr/" in event.arguments:
                return ObserverAlert(
                    source=self.name,
                    severity=Severity.FATAL,
                    message=f"blocked dangerous write: {event.arguments}",
                )
        return None

agent = Agent(
    "safe-shell",
    prompt="...",
    config=config,
    tools=[write_file],
    observers=[PathGuardian()],
    assembly=[AlertPolicy()],   # routes FATAL → HaltEvent
)
The first dangerous tool call triggers FATAL → halt; the agent's next ask is short-circuited. Full runnable demo:
assets/safety_guard.py
.
AlertPolicy
在收到
Severity.FATAL
级告警时会执行两项操作:
  1. 向流中发送
    HaltEvent
    事件。
  2. 在系统提示中追加终止通知。
assembly=[...]
非空时,执行器会自动连接
_HaltCheckMiddleware
,该中间件会检测到
HaltEvent
并通过合成的
HALTED: ...
响应短路下一次LLM调用。
python
from ag2 import Context
from ag2.observers import BaseObserver
from ag2.events import BaseEvent, ToolCallEvent
from ag2.events.alert import HaltEvent, ObserverAlert, Severity
from ag2.policies import AlertPolicy
from ag2.watch import EventWatch

class PathGuardian(BaseObserver):
    def __init__(self) -> None:
        super().__init__("path-guardian", watch=EventWatch(ToolCallEvent))

    async def process(self, events: list[BaseEvent], ctx: Context) -> ObserverAlert | None:
        for event in events:
            if not isinstance(event, ToolCallEvent) or event.name != "write_file":
                continue
            if "/etc/" in event.arguments or "/usr/" in event.arguments:
                return ObserverAlert(
                    source=self.name,
                    severity=Severity.FATAL,
                    message=f"blocked dangerous write: {event.arguments}",
                )
        return None

agent = Agent(
    "safe-shell",
    prompt="...",
    config=config,
    tools=[write_file],
    observers=[PathGuardian()],
    assembly=[AlertPolicy()],   # 将FATAL告警路由为HaltEvent
)
首次危险工具调用会触发FATAL告警并终止运行;Agent的下一次请求会被短路。完整可运行示例:
assets/safety_guard.py

Subscribing to alerts and halts from outside

从外部订阅告警与终止事件

python
from ag2 import MemoryStream
from ag2.events.alert import HaltEvent, ObserverAlert

stream = MemoryStream()
stream.where(ObserverAlert).subscribe(lambda e: print(f"[{e.severity}] {e.source}: {e.message}"))
stream.where(HaltEvent).subscribe(lambda e: print(f"HALT: {e.reason}"))
await agent.ask("...", stream=stream)
python
from ag2 import MemoryStream
from ag2.events.alert import HaltEvent, ObserverAlert

stream = MemoryStream()
stream.where(ObserverAlert).subscribe(lambda e: print(f"[{e.severity}] {e.source}: {e.message}"))
stream.where(HaltEvent).subscribe(lambda e: print(f"HALT: {e.reason}"))
await agent.ask("...", stream=stream)

Observers vs Middleware vs Stream subscribers

Observer vs Middleware vs Stream订阅者

FeatureObserverMiddlewareStream subscriber
Registered onAgentAgentStream
LifecycleScoped to executionScoped to executionManual
BoilerplateFunction (or
BaseObserver
)
BaseMiddleware
class
Function
Can modify events
interrupt=True
Yes (wraps execution)
interrupt=True
DI supportYesYesYes
Use caseMonitoring, metrics, alertsCross-cutting (retry, auth, rate limit)Low-level event wiring
特性ObserverMiddlewareStream 订阅者
注册位置AgentAgentStream
生命周期作用于执行阶段作用于执行阶段手动控制
模板代码函数(或
BaseObserver
BaseMiddleware
函数
能否修改事件
interrupt=True
时可以
可以(包装执行流程)
interrupt=True
时可以
依赖注入支持
适用场景监控、指标、告警横切关注点(重试、认证、限流)底层事件连接

Going deeper

深入学习

  • assets/token_watchdog.py
    — three observers (
    TokenMonitor
    ,
    LoopDetector
    , custom
    AlertConsole
    ) on one agent. Mirrors
    code_examples/04
    .
  • assets/safety_guard.py
    PathGuardian
    → FATAL →
    AlertPolicy
    HaltEvent
    → short-circuit. Mirrors
    code_examples/08
    .
  • Source docs:
    • website/docs/user-guide/advanced/observers.mdx
      @observer
      ,
      BaseObserver
      , registration, built-ins,
      ObserverAlert
      .
    • website/docs/user-guide/advanced/watches.mdx
      — every Watch primitive, composition rules.
    • website/docs/user-guide/advanced/stream.mdx
      — Stream API,
      where
      ,
      subscribe
      , interrupters,
      RedisStream
      .
    • website/docs/user-guide/advanced/assembly.mdx
      AlertPolicy
      ordering and dedup.
  • assets/token_watchdog.py
    — 一个Agent上同时使用三个观察者(
    TokenMonitor
    LoopDetector
    、自定义
    AlertConsole
    )。对应
    code_examples/04
  • assets/safety_guard.py
    PathGuardian
    → FATAL告警 →
    AlertPolicy
    HaltEvent
    → 短路终止。对应
    code_examples/08
  • 源码文档:
    • website/docs/user-guide/advanced/observers.mdx
      @observer
      BaseObserver
      、注册方式、内置组件、
      ObserverAlert
    • website/docs/user-guide/advanced/watches.mdx
      — 所有Watch原语、组合规则。
    • website/docs/user-guide/advanced/stream.mdx
      — Stream API、
      where
      subscribe
      、拦截器、
      RedisStream
    • website/docs/user-guide/advanced/assembly.mdx
      AlertPolicy
      的排序与去重。

Common pitfalls

常见误区

  • Alerts not reaching the model
    ObserverAlert
    events are on the stream but invisible to the LLM by default. Add
    AlertPolicy()
    to
    assembly=[...]
    .
  • FATAL not halting
    AlertPolicy
    is what creates
    HaltEvent
    . Without
    assembly=[..., AlertPolicy(), ...]
    (or any non-empty assembly chain enabling
    _HaltCheckMiddleware
    ), nothing halts.
  • Sharing one
    AlertPolicy()
    across agents
    — dedup state lives on the instance. Give each agent its own.
  • Watch callback assumes
    events
    is non-empty
    — for time-driven watches (
    DelayWatch
    ,
    IntervalWatch
    ,
    CronWatch
    ),
    events
    is always
    []
    .
  • Forgetting
    process()
    is async
    BaseObserver.process
    must be
    async def
    .
  • Subscribing with
    subscribe(fn)
    when you wanted
    subscribe()
    decorator
    — both work; the bare-call form is
    stream.subscribe(fn)
    , the decorator form is
    @stream.subscribe()
    (with parens).
  • CadenceWatch
    with no
    n
    and no
    max_wait
    — raises
    ValueError
    ; at least one is required.
  • 告警未传递给模型
    ObserverAlert
    事件存在于流中,但默认对LLM不可见。需将
    AlertPolicy()
    添加到
    assembly=[...]
    中。
  • FATAL告警未触发终止
    AlertPolicy
    是生成
    HaltEvent
    的关键。如果没有
    assembly=[..., AlertPolicy(), ...]
    (或任何非空的执行链启用
    _HaltCheckMiddleware
    ),则不会触发终止。
  • 多个Agent共享同一个
    AlertPolicy()
    实例
    — 去重状态存储在实例中。应为每个Agent创建独立的实例。
  • Watch回调假设
    events
    非空
    — 对于时间驱动的监视器(
    DelayWatch
    IntervalWatch
    CronWatch
    ),
    events
    始终为
    []
  • 忘记
    process()
    是异步函数
    BaseObserver.process
    必须定义为
    async def
  • 使用
    subscribe(fn)
    而非装饰器形式
    subscribe()
    — 两种方式都可行;直接调用形式为
    stream.subscribe(fn)
    ,装饰器形式为
    @stream.subscribe()
    (带括号)。
  • CadenceWatch
    未设置
    n
    max_wait
    — 会抛出
    ValueError
    ;至少需要设置其中一个参数。