# Watches

A **Watch** is a reactive trigger primitive. You arm it on a [Stream](https://docs.ag2.ai/0.13.2/docs/beta/advanced/stream/), and it fires a callback when its condition is met — where "condition" can be an event match, a count, a time window, a schedule, or a composition of other watches.

Watches are the mechanism behind trigger-driven [Observers](https://docs.ag2.ai/0.13.2/docs/beta/advanced/observers/), but they can be used standalone for custom reactive logic on any Stream.

## When to use a Watch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#when-to-use-a-watch "Permanent link")

Use a Watch when a simple `stream.subscribe(...)` is not enough — i.e. when you need buffering, timing, ordering, or composition:

| You need                                       | Use                                                           |
|------------------------------------------------|---------------------------------------------------------------|
| Run a callback on every matching event          | `stream.subscribe(fn, condition=...)` (no Watch needed)      |
| Fire every N events                             | `CadenceWatch(n=N, condition=...)`                            |
| Collect events over a time window               | `CadenceWatch(max_wait=seconds, condition=...)`               |
| Fire on "N events OR T seconds, whichever first" | `CadenceWatch(n=N, max_wait=seconds, condition=...)`        |
| Fire once after a delay                         | `DelayWatch(seconds)`                                         |
| Fire on a schedule                              | `IntervalWatch(seconds)` or `CronWatch(expr)`                |
| Wait for two separate events in any order       | `AllOf(w1, w2)`                                              |
| Wait for an ordered sequence                    | `Sequence(w1, w2, ...)`                                      |
| Fire on the earliest of several triggers        | `AnyOf(w1, w2, ...)`                                         |

## Anatomy of a Watch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#anatomy-of-a-watch "Permanent link")

Every Watch implements the same `Watch` protocol (importable from `autogen.beta`):

|     |     |
|-----|-----|
| ```<br> 1<br> 2<br> 3<br> 4<br> 5<br> 6<br> 7<br> 8<br> 9<br>10<br>``` | ```<br>from typing import Protocol<br>class Watch(Protocol):<br>    @property<br>    def id(self) -> str: ...<br>    @property<br>    def is_armed(self) -> bool: ...<br>    def arm(self, stream, callback) -> None: ...<br>    def disarm(self) -> None: ...<br>``` |

The callback signature is uniform across all Watch kinds:

|     |     |
|-----|-----|
| ```<br>1<br>2<br>3<br>4<br>``` | ```<br>from autogen.beta import Context<br>from autogen.beta.events import BaseEvent<br>async def callback(events: list[BaseEvent], ctx: Context) -> None: ...<br>``` |

For event-driven watches (`EventWatch`, `CadenceWatch`, `Sequence`), `events` contains the matched events. For time-driven watches (`DelayWatch`, `IntervalWatch`, `CronWatch`), `events` is empty — the trigger is the timer itself.

## Event-driven Watches [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#event-driven-watches "Permanent link")

### EventWatch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#eventwatch "Permanent link")

Fires immediately on each matching event.

|     |     |
|-----|-----|
| ```<br> 1<br> 2<br> 3<br> 4<br> 5<br> 6<br> 7<br> 8<br> 9<br>10<br>11<br>``` | ```<br>from autogen.beta import MemoryStream<br>from autogen.beta.watch import EventWatch<br>from autogen.beta.events import ModelResponse<br>stream = MemoryStream()<br>watch = EventWatch(ModelResponse)<br>async def on_response(events, ctx):<br>    print(f"Model responded: {events[0].content}")<br>watch.arm(stream, on_response)<br>``` |

`EventWatch` also supports field conditions and negation:

|     |     |
|-----|-----|
| ```<br>1<br>2<br>3<br>4<br>5<br>6<br>7<br>``` | ```<br>from autogen.beta.watch import EventWatch<br>from autogen.beta.events import ToolCallEvent<br>watch = EventWatch(ToolCallEvent.name == "search")<br># Fire on every event that is NOT a ToolCallEvent<br>watch = EventWatch(~ToolCallEvent)<br>``` |

### CadenceWatch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#cadencewatch "Permanent link")

Buffers matching events and fires once the buffer reaches **size**`n`, once **time**`max_wait` seconds have elapsed since the first buffered event, or whichever comes first when both are set. At least one of `n` and `max_wait` is required.

|     |     |
|-----|-----|
| ```<br> 1<br> 2<br> 3<br> 4<br> 5<br> 6<br> 7<br> 8<br> 9<br>10<br>11<br>12<br>13<br>14<br>``` | ```<br>from autogen.beta.watch import CadenceWatch<br>from autogen.beta.events import ModelResponse<br># Count-only: fire every 5 model responses<br>batch = CadenceWatch(n=5, condition=ModelResponse)<br># Time-only: flush the buffer once a minute, whenever there's something in it<br>window = CadenceWatch(max_wait=60.0, condition=ModelResponse)<br># Size OR time: fire at 5 events OR 60 seconds after the first event, whichever first<br>hybrid = CadenceWatch(n=5, max_wait=60.0, condition=ModelResponse)<br>async def summarize(events, ctx):<br>    print(f"Batched {len(events)} responses")<br>``` |

The timer starts on the **first** buffered event in a cadence — a quiet stream produces no firings. Any unfilled buffer at disarm-time is discarded.

## Time-driven Watches [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#time-driven-watches "Permanent link")

### DelayWatch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#delaywatch "Permanent link")

Fires exactly once after a delay, then auto-disarms.

|     |     |
|-----|-----|
| ```<br>1<br>2<br>3<br>4<br>5<br>6<br>7<br>8<br>9<br>``` | ```<br>from autogen.beta.watch import DelayWatch<br>watch = DelayWatch(30.0)<br>async def timeout_guard(events, ctx):<br>    # events is always [] for time-driven watches<br>    print("30 seconds elapsed")<br>watch.arm(stream, timeout_guard)<br>``` |

### IntervalWatch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#intervalwatch "Permanent link")

Fires periodically at a fixed interval until disarmed.

|     |     |
|-----|-----|
| ```<br>1<br>2<br>3<br>4<br>5<br>6<br>``` | ```<br>from autogen.beta.watch import IntervalWatch<br>watch = IntervalWatch(60.0)<br>async def heartbeat(events, ctx):<br>    print("tick")<br>``` |

### CronWatch [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#cronwatch "Permanent link")

Fires on a standard 5-field cron expression.

|     |     |
|-----|-----|
| ```<br>1<br>2<br>3<br>``` | ```<br>from autogen.beta.watch import CronWatch<br>watch = CronWatch("0 9 * * MON")  # every Monday at 9am<br>``` |

Supports `*`, `*/n`, `a-b`, comma lists, and `SUN`–`SAT` day-of-week names.

## Composite Watches [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#composite-watches "Permanent link")

### AllOf [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#allof "Permanent link")

Fires once when **every** sub-watch has fired at least once.

|     |     |
|-----|-----|
| ```<br> 1<br> 2<br> 3<br> 4<br> 5<br> 6<br> 7<br> 8<br> 9<br>10<br>11<br>``` | ```<br>from autogen.beta.watch import AllOf, EventWatch<br>from autogen.beta.events import ModelResponse, ToolCallEvent<br>watch = AllOf(<br>    EventWatch(ModelResponse),<br>    EventWatch(ToolCallEvent),<br>)<br>async def both_seen(events, ctx):<br>    # Gets the combined events from all sub-watches<br>    print(f"Saw both types, got {len(events)} events total")<br>``` |

After firing, the gate resets — both sub-watches must fire again for the next firing.

### AnyOf [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#anyof "Permanent link")

Fires on **any** sub-watch, every time.

|     |     |
|-----|-----|
| ```<br>1<br>2<br>3<br>4<br>5<br>6<br>7<br>``` | ```<br>from autogen.beta.watch import AnyOf, EventWatch<br>from autogen.beta.events import ObserverAlert<br>watch = AnyOf(<br>    EventWatch(ObserverAlert.severity == "critical"),<br>    EventWatch(ObserverAlert.severity == "fatal"),<br>)<br>``` |

### Sequence [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#sequence "Permanent link")

Fires when sub-watches trigger **in order**. Each sub-watch is armed only after the previous has fired.

|     |     |
|-----|-----|
| ```<br> 1<br> 2<br> 3<br> 4<br> 5<br> 6<br> 7<br> 8<br> 9<br>10<br>11<br>``` | ```<br>from autogen.beta.watch import Sequence, EventWatch<br>from autogen.beta.events import ModelRequest, ModelResponse<br>watch = Sequence(<br>    EventWatch(ModelRequest),<br>    EventWatch(ModelResponse),<br>)<br>async def round_trip(events, ctx):<br>    # events contains one ModelRequest followed by one ModelResponse<br>    print("round-trip complete")<br>``` |

After the last sub-watch fires, the sequence resets.

## Emitting events from a callback [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#emitting-events-from-a-callback "Permanent link")

A Watch callback can send events back onto the stream — this is how [trigger-driven Observers](https://docs.ag2.ai/0.13.2/docs/beta/advanced/observers/#trigger-driven-observers-baseobserver) emit alerts:

|     |     |
|-----|-----|
| ```<br> 1<br> 2<br> 3<br> 4<br> 5<br> 6<br> 7<br> 8<br> 9<br>10<br>11<br>12<br>``` | ```<br>from autogen.beta import MemoryStream<br>from autogen.beta.watch import DelayWatch<br>from autogen.beta.events import ObserverAlert, Severity<br>watch = DelayWatch(30.0)<br>async def timeout_alert(events, ctx):<br>    await ctx.send(ObserverAlert(<br>        source="timeout-watch",<br>        severity=Severity.WARNING,<br>        message="Agent has been running for >30s.",<br>    ))<br>``` |

Tip

Emitting events from within a callback means other subscribers and watches can react. This is the foundation for composing reactive workflows.

## Arming and lifecycle [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#arming-and-lifecycle "Permanent link")

A Watch is a stateful object. Calling `arm()` on an already-armed Watch will first disarm it, so re-arming is safe. `disarm()` cleans up any subscriptions and cancels any timers.

Most of the time you won't call `arm()` / `disarm()` directly — you'll hand the Watch to a [BaseObserver](https://docs.ag2.ai/0.13.2/docs/beta/advanced/observers/#trigger-driven-observers-baseobserver), which manages the lifecycle against the agent's stream.

## Next steps [\#](https://docs.ag2.ai/0.13.2/docs/beta/advanced/watches/#next-steps "Permanent link")

- Use Watches inside a [BaseObserver](https://docs.ag2.ai/0.13.2/docs/beta/advanced/observers/#trigger-driven-observers-baseobserver) for agent-scoped monitoring.
- See [Stream](https://docs.ag2.ai/0.13.2/docs/beta/advanced/stream/) for event publishing and subscription mechanics.
