Loading...
Loading...
Monitor an AG2 agent's stream — log events, detect repeated tool calls, track token spend, build trigger-driven observers, route observer alerts to the model, and halt on FATAL conditions. Covers `@observer(...)` (stateless), `BaseObserver` (stateful), built-ins (`TokenMonitor`, `LoopDetector`), `Watch` primitives (`EventWatch`, `CadenceWatch`, `DelayWatch`, `IntervalWatch`, `CronWatch`, `AllOf`, `AnyOf`, `Sequence`), `ObserverAlert` (`Severity.INFO/WARNING/CRITICAL/FATAL`), `AlertPolicy`, and `HaltEvent`. Use when the user wants observability, runtime safety guards, alerts, or batch/time-based reactive logic.
npx skill4agent add ag2ai/ag2-skills ag2-observers-and-alerts| Shape | When | Use |
|---|---|---|
| Stateless function | One-off event hook (logging, metrics) | |
| Stateful class | Counters / windows / thresholds / composed triggers | Subclass |
@observerfrom 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.observer(...)agent.ask("...", observers=[...])ContextInjectVariableDependsModelRequest | ModelResponseToolCallEvent.name == "search"interrupt=Truefrom 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),
],
)TokenMonitorModelResponseTaskCompletedWARNINGCRITICALObserverAlertmonitor.total_tokensLoopDetectorWARNINGrepeat_thresholdBaseObserverBaseObserverWatchprocess()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}",
)process()ObserverAlertawait ctx.send(...)| You need | Use |
|---|---|
| Every matching event | |
| Every N matching events | |
| Every T seconds (buffered events) | |
| Either threshold | |
| Once after delay | |
| Periodic timer | |
| Cron schedule | |
| All sub-watches must fire | |
| Any sub-watch fires | |
| In order | |
ag2.watchasync def cb(events: list[BaseEvent], ctx: Context) -> Noneevents=[]ObserverAlertfrom ag2.events.alert import ObserverAlert, Severity
ObserverAlert(
source="my-observer",
severity=Severity.WARNING, # INFO, WARNING, CRITICAL, FATAL
message="What happened",
)ObserverAlertAlertPolicy()assembly=[...]from ag2.policies import AlertPolicy
agent = Agent("assistant", config=config, assembly=[AlertPolicy()])HaltEventAlertPolicySeverity.FATALHaltEventassembly=[...]_HaltCheckMiddlewareHaltEventHALTED: ...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
)assets/safety_guard.pyfrom 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)| Feature | Observer | Middleware | Stream subscriber |
|---|---|---|---|
| Registered on | Agent | Agent | Stream |
| Lifecycle | Scoped to execution | Scoped to execution | Manual |
| Boilerplate | Function (or | | Function |
| Can modify events | | Yes (wraps execution) | |
| DI support | Yes | Yes | Yes |
| Use case | Monitoring, metrics, alerts | Cross-cutting (retry, auth, rate limit) | Low-level event wiring |
assets/token_watchdog.pyTokenMonitorLoopDetectorAlertConsolecode_examples/04assets/safety_guard.pyPathGuardianAlertPolicyHaltEventcode_examples/08website/docs/user-guide/advanced/observers.mdx@observerBaseObserverObserverAlertwebsite/docs/user-guide/advanced/watches.mdxwebsite/docs/user-guide/advanced/stream.mdxwheresubscribeRedisStreamwebsite/docs/user-guide/advanced/assembly.mdxAlertPolicyObserverAlertAlertPolicy()assembly=[...]AlertPolicyHaltEventassembly=[..., AlertPolicy(), ...]_HaltCheckMiddlewareAlertPolicy()eventsDelayWatchIntervalWatchCronWatchevents[]process()BaseObserver.processasync defsubscribe(fn)subscribe()stream.subscribe(fn)@stream.subscribe()CadenceWatchnmax_waitValueError