With DISPATCH_AGENT autonomous=true in the agent registry, the assignment-failure
path sends a real customer message. It would have sent
"Delay alert: Order 36 may arrive later than expected. New ETA: {new_eta}"
with customer_id "unknown": the DELAYED template names {new_eta} but DispatchAgent
supplied only order_id, and it has no customer id to give.
- _send_via_channel refuses to send when any {placeholder} survives substitution
(or the template is missing), logs which variables the caller omitted, and
records a failed notification. Silence beats a broken message.
- customer_id is omitted from the payload when it is absent or "unknown" rather
than sent literally; the backend resolves the recipient from order_id.
- DispatchAgent supplies new_eta "being confirmed" — honest, since no ETA exists
at assignment-failure time.
- Tests: the old bug is refused, a complete payload renders clean, the id is
omitted/passed correctly, and DispatchAgent's vars satisfy the template.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin
521 lines
23 KiB
Python
521 lines
23 KiB
Python
"""Dispatch Agent — watches NATS assignment events and escalates coverage gaps."""
|
|
import asyncio
|
|
import json
|
|
import os
|
|
from datetime import datetime, timezone
|
|
from typing import Dict, List, Any, Optional
|
|
|
|
import nats
|
|
import nats.js.errors
|
|
import redis.asyncio as aioredis
|
|
|
|
from core.agent import SpecializedAgent
|
|
from core.types import AgentTask, MessageType
|
|
from core.logger import logger
|
|
from core.llm import decide_assignment_failure, build_assignment_failure_context, LLM_MODEL
|
|
from core.registry import registry
|
|
from core.decisions import record_decision
|
|
from config.system_config import (
|
|
NATS_HOST, NATS_PORT, NATS_USER, NATS_PASSWORD,
|
|
REDIS_HOST, REDIS_PORT, REDIS_PASSWORD,
|
|
)
|
|
|
|
# Notifying a customer is outward-facing, so it is gated: without autonomy the
|
|
# agent proposes it as an internal ops alert instead of messaging the customer.
|
|
DISPATCH_AGENT_AUTONOMOUS = os.getenv("DISPATCH_AGENT_AUTONOMOUS", "false").lower() == "true"
|
|
|
|
# The Go backend owns the JetStream stream carrying booking.* events; this agent
|
|
# is a pure consumer and must not create it. Leave empty to auto-discover the
|
|
# stream by subject, or pin it to a stream name if discovery isn't desired.
|
|
DISPATCH_STREAM = os.getenv("DISPATCH_STREAM", "")
|
|
|
|
# Escalation rate limit. A backend retry sweep can fail hundreds of bookings in
|
|
# one zone within a minute (observed: 1,001 in 60 s). Once a zone has been
|
|
# alerted today, further failures only bump the counter; a summary re-alert goes
|
|
# out every DISPATCH_REALERT_EVERY failures so ops can see it's still growing.
|
|
# The LLM is only consulted before the first alert, so cost is bounded per zone.
|
|
DISPATCH_REALERT_EVERY = int(os.getenv("DISPATCH_REALERT_EVERY", "100"))
|
|
|
|
# Presence, as reported by the Go backend's miler_status:<id> key. The
|
|
# milers:locations geo index keeps last-known positions indefinitely, so
|
|
# membership alone says nothing about whether anyone is actually on duty.
|
|
#
|
|
# Three-state on purpose. Only 2 of 34 milers had a status key in production
|
|
# (2026-09-22), so treating "no key" as off-duty would report "nobody
|
|
# available" even where the backend can assign — a confident wrong answer.
|
|
# Absent presence data is reported as unknown and the decision is told so.
|
|
MILER_AVAILABLE_STATUSES = {"available"}
|
|
MILER_UNAVAILABLE_STATUSES = {"break", "offline", "busy", "onbreak", "off_duty", "unavailable"}
|
|
|
|
AVAILABLE, UNAVAILABLE, UNKNOWN = "available", "unavailable", "unknown"
|
|
|
|
# The Go backend does not yet send zone_id on booking.assignment_failed (every
|
|
# event observed carries none). Without a fallback, every failure in the country
|
|
# shares the single "unknown" bucket, so the per-zone counter and the once-a-day
|
|
# alert limit would collapse into one global bucket. Derive a coarse grid cell
|
|
# from the pickup coordinates instead: 1 decimal place is ~11 km, about the size
|
|
# of a delivery zone. Real zone_ids, once sent, take precedence.
|
|
GEO_BUCKET_DP = int(os.getenv("DISPATCH_GEO_BUCKET_DP", "1"))
|
|
|
|
|
|
def zone_key(zone_id, lat, lon) -> str:
|
|
"""Stable bucket for counting failures: the backend's zone_id when present,
|
|
else a coarse lat/lon grid cell, else 'unknown'."""
|
|
if zone_id and str(zone_id).lower() not in ("", "unknown", "none", "null"):
|
|
return str(zone_id)
|
|
if lat is not None and lon is not None:
|
|
try:
|
|
return f"geo:{round(float(lat), GEO_BUCKET_DP)},{round(float(lon), GEO_BUCKET_DP)}"
|
|
except (TypeError, ValueError):
|
|
pass
|
|
return "unknown"
|
|
|
|
|
|
class DispatchAgent(SpecializedAgent):
|
|
"""
|
|
Consumes booking.assigned and booking.assignment_failed events published by
|
|
the Go backend (after routemate AI or fallback assignment) by binding to the
|
|
backend-owned JetStream stream that carries them. It is a pure consumer — it
|
|
does not create the stream, and does not trigger assignment itself.
|
|
"""
|
|
|
|
def __init__(self):
|
|
super().__init__(
|
|
agent_id="DISPATCH_AGENT",
|
|
domain="dispatch",
|
|
description="Watches assignment events and escalates coverage gaps"
|
|
)
|
|
|
|
self._redis = aioredis.Redis(
|
|
host=REDIS_HOST,
|
|
port=REDIS_PORT,
|
|
password=REDIS_PASSWORD,
|
|
decode_responses=True,
|
|
)
|
|
|
|
self._nats_nc = None
|
|
self._nats_js = None
|
|
self._nats_subs: list = []
|
|
|
|
# ------------------------------------------------------------------ #
|
|
# Lifecycle #
|
|
# ------------------------------------------------------------------ #
|
|
|
|
async def start(self):
|
|
await self._connect_nats()
|
|
await super().start()
|
|
|
|
async def stop(self):
|
|
await super().stop()
|
|
for sub in self._nats_subs:
|
|
try:
|
|
await sub.unsubscribe()
|
|
except Exception:
|
|
pass
|
|
if self._nats_nc:
|
|
await self._nats_nc.drain()
|
|
|
|
async def _connect_nats(self):
|
|
try:
|
|
self._nats_nc = await nats.connect(
|
|
servers=[f"nats://{NATS_HOST}:{NATS_PORT}"],
|
|
user=NATS_USER,
|
|
password=NATS_PASSWORD,
|
|
max_reconnect_attempts=10,
|
|
)
|
|
self._nats_js = self._nats_nc.jetstream()
|
|
|
|
# Bind consumers to the *existing* backend-owned stream. Do NOT create
|
|
# a stream here: booking.* is published by the Go backend, so creating
|
|
# a second stream over the same subjects raises an overlap error that
|
|
# (previously swallowed) left the consumer unbound. Mirrors the
|
|
# explicit stream= binding ExceptionAgent uses for the TRACKING stream.
|
|
sub_assigned = await self._bind_consumer(
|
|
"booking.assigned", "dispatch-booking-assigned",
|
|
self._on_nats_booking_assigned,
|
|
)
|
|
sub_failed = await self._bind_consumer(
|
|
"booking.assignment_failed", "dispatch-booking-assignment-failed",
|
|
self._on_nats_booking_assignment_failed,
|
|
)
|
|
self._nats_subs.extend(s for s in (sub_assigned, sub_failed) if s is not None)
|
|
if self._nats_subs:
|
|
logger.info(f"DispatchAgent bound {len(self._nats_subs)} booking-event consumer(s)")
|
|
|
|
except Exception as e:
|
|
logger.error(f"DispatchAgent NATS connect error: {e}")
|
|
|
|
async def _bind_consumer(self, subject, durable, cb):
|
|
"""Bind a durable push consumer to whichever existing stream carries
|
|
`subject`. Returns the subscription, or None (with a clear, actionable
|
|
log) if no stream carries it — never a swallowed or opaque error."""
|
|
try:
|
|
stream = DISPATCH_STREAM or await self._nats_js.find_stream_name_by_subject(subject)
|
|
sub = await self._nats_js.subscribe(subject, durable=durable, stream=stream, cb=cb)
|
|
logger.info(f"DispatchAgent bound '{subject}' on stream '{stream}' (durable={durable})")
|
|
return sub
|
|
except nats.js.errors.NotFoundError:
|
|
logger.error(
|
|
f"DispatchAgent: no JetStream stream carries '{subject}'. The publisher "
|
|
f"(Go backend) must create a stream covering it, or set DISPATCH_STREAM to "
|
|
f"the stream name. Not subscribing to '{subject}'."
|
|
)
|
|
return None
|
|
except Exception as e:
|
|
logger.error(f"DispatchAgent failed to bind '{subject}': {e}")
|
|
return None
|
|
|
|
# ------------------------------------------------------------------ #
|
|
# Redis GEO lookup #
|
|
# ------------------------------------------------------------------ #
|
|
|
|
async def _miler_presence(self, miler_id) -> str:
|
|
"""AVAILABLE / UNAVAILABLE / UNKNOWN from the backend-owned
|
|
miler_status:<id> key ({"userid": .., "status": "Available"|"Break"|..}).
|
|
|
|
No key, an unparseable value, or an unrecognised status → UNKNOWN, never
|
|
UNAVAILABLE: most milers have no key at all, and absence of evidence is
|
|
not evidence of absence."""
|
|
try:
|
|
raw = await self._redis.get(f"miler_status:{miler_id}")
|
|
if not raw:
|
|
return UNKNOWN
|
|
status = str(json.loads(raw).get("status", "")).lower()
|
|
except Exception:
|
|
return UNKNOWN
|
|
if status in MILER_AVAILABLE_STATUSES:
|
|
return AVAILABLE
|
|
if status in MILER_UNAVAILABLE_STATUSES:
|
|
return UNAVAILABLE
|
|
return UNKNOWN
|
|
|
|
async def _find_zone(
|
|
self, lat: float, lon: float, radius_km: int = 10
|
|
) -> Optional[Dict[str, Any]]:
|
|
"""GEORADIUS sweep on milers:locations, preferring a miler the backend
|
|
reports Available. If none is confirmed Available but some have no
|
|
presence data, the nearest of those is returned with presence=UNKNOWN —
|
|
it may well be assignable. Returns None only when every candidate is
|
|
confirmed unavailable (or there are none)."""
|
|
try:
|
|
results = await self._redis.georadius(
|
|
"milers:locations",
|
|
lon, lat,
|
|
radius_km, "km",
|
|
sort="ASC",
|
|
count=10,
|
|
)
|
|
if not results:
|
|
return None
|
|
|
|
first_unknown = None
|
|
for candidate in results:
|
|
presence = await self._miler_presence(candidate)
|
|
if presence == AVAILABLE:
|
|
return await self._miler_record(candidate, AVAILABLE, results)
|
|
if presence == UNKNOWN and first_unknown is None:
|
|
first_unknown = candidate
|
|
|
|
if first_unknown is not None:
|
|
return await self._miler_record(first_unknown, UNKNOWN, results)
|
|
return None
|
|
except Exception as e:
|
|
logger.warning(f"Redis GEORADIUS error: {e}")
|
|
return None
|
|
|
|
async def _miler_record(self, miler_id, presence, candidates) -> Dict[str, Any]:
|
|
miler_info = await self._redis.hgetall(f"miler:{miler_id}")
|
|
return {
|
|
"miler_id": miler_id,
|
|
"presence": presence,
|
|
"hub_id": miler_info.get("hub_id", "UNKNOWN"),
|
|
"zone_id": miler_info.get("zone_id", "unknown"),
|
|
"avg_delivery_time": int(miler_info.get("avg_delivery_time", 60)),
|
|
"candidates_considered": len(candidates),
|
|
}
|
|
|
|
async def _count_in_geo_index(self, lat: float, lon: float, radius_km: int = 30) -> Optional[int]:
|
|
"""Raw geo-index membership (any status) — reported separately so the
|
|
decision can tell 'nobody registered here' from 'riders exist but none on duty'."""
|
|
try:
|
|
results = await self._redis.georadius("milers:locations", lon, lat, radius_km, "km", count=50)
|
|
return len(results or [])
|
|
except Exception:
|
|
return None
|
|
|
|
# ------------------------------------------------------------------ #
|
|
# NATS event handlers #
|
|
# ------------------------------------------------------------------ #
|
|
|
|
async def _on_nats_booking_assigned(self, msg):
|
|
"""
|
|
booking.assigned — Go publishes this after a successful AI or
|
|
fallback assignment. Forward a prepare_receiving task to HUB_AGENT
|
|
using the real miler_id/hub_id from the event.
|
|
"""
|
|
try:
|
|
data = json.loads(msg.data.decode())
|
|
booking_id = data.get("booking_id")
|
|
miler_id = data.get("miler_id")
|
|
hub_id = data.get("hub_id")
|
|
confidence = data.get("confidence")
|
|
reasoning = data.get("reasoning", "")
|
|
|
|
logger.info(
|
|
f"[DISPATCH] booking.assigned — booking={booking_id} "
|
|
f"miler={miler_id} hub={hub_id} confidence={confidence} "
|
|
f"reasoning={reasoning!r}"
|
|
)
|
|
# No longer forwarded to HUB_AGENT as `prepare_receiving`. HUB_AGENT
|
|
# is a simulation over eight fictional hubs keyed "XX-HUB-01"; handed a
|
|
# real backend hub id and booking id it answered "Hub not found" for
|
|
# every assignment, so the hand-off did nothing but log noise. Restore
|
|
# it only once HUB_AGENT reads real hubs (agent registry status:
|
|
# simulation).
|
|
|
|
except Exception as e:
|
|
logger.error(f"DispatchAgent booking.assigned handler error: {e}")
|
|
finally:
|
|
await msg.ack()
|
|
|
|
async def _on_nats_booking_assignment_failed(self, msg):
|
|
"""
|
|
booking.assignment_failed — Go publishes this when decide-assignment
|
|
returns escalate=true or both AI and fallback find nobody.
|
|
|
|
Gathers coverage context (nearest-rider sweep + daily failure counter),
|
|
then asks Claude to choose: monitor / notify_customer / ops_alert /
|
|
escalate. Customer notification is gated behind autonomy. Falls back to
|
|
the previous count>=3 heuristic if the LLM is unavailable.
|
|
"""
|
|
try:
|
|
data = json.loads(msg.data.decode())
|
|
booking_id = data.get("booking_id")
|
|
lat = data.get("lat")
|
|
lon = data.get("lon")
|
|
# Falls back to a coarse geo cell while the backend omits zone_id.
|
|
zone_id = zone_key(data.get("zone_id"), lat, lon)
|
|
|
|
logger.warning(
|
|
f"[DISPATCH] booking.assignment_failed — booking={booking_id} zone={zone_id} "
|
|
f"reason={data.get('reason')!r}"
|
|
)
|
|
|
|
if not registry.skill_enabled("assignment_failure_triage"):
|
|
logger.info(
|
|
f"[DISPATCH] booking={booking_id}: skill assignment_failure_triage is disabled "
|
|
"in the agent registry; not triaging"
|
|
)
|
|
return
|
|
|
|
facts, count = await self._gather_assignment_facts(zone_id, lat, lon)
|
|
realert_every = int(registry.threshold("assignment_failure_triage", "realertEvery", DISPATCH_REALERT_EVERY))
|
|
logger.info(f"[DISPATCH] booking={booking_id} context facts: {facts}")
|
|
|
|
# Rate limit: once this zone has been alerted today, don't re-decide
|
|
# (and don't pay for an LLM call) on every further failure.
|
|
if facts.get("alert_already_sent_today"):
|
|
if realert_every and count % realert_every == 0:
|
|
await self._ops_alert(
|
|
zone_id, booking_id,
|
|
f"Zone {zone_id}: {count} failed assignments today and still climbing "
|
|
f"(re-alert every {realert_every}; earlier alert already raised).",
|
|
severity="high",
|
|
)
|
|
else:
|
|
logger.info(
|
|
f"[DISPATCH] zone {zone_id} already alerted today — "
|
|
f"failure #{count} counted, no new alert"
|
|
)
|
|
return
|
|
|
|
model = registry.model("DISPATCH_AGENT")
|
|
decision = await decide_assignment_failure(build_assignment_failure_context(facts), model)
|
|
if decision is not None:
|
|
record_decision("assignment_failure", booking_id, facts, decision.action, decision.confidence,
|
|
decision.reasoning, model or LLM_MODEL)
|
|
|
|
if decision is None:
|
|
logger.warning(f"No LLM decision for booking {booking_id}; falling back to count>=3 heuristic")
|
|
if count >= 3:
|
|
await self._ops_alert(
|
|
zone_id, booking_id,
|
|
f"Zone {zone_id} has had {count} failed assignments today — recommend onboarding milers here.",
|
|
)
|
|
else:
|
|
logger.info(f"[DISPATCH] zone {zone_id} failure count {count}/3 — transient, no escalation")
|
|
return
|
|
|
|
logger.info(
|
|
f"[DISPATCH] booking={booking_id} zone={zone_id} decision={decision.action} "
|
|
f"confidence={decision.confidence:.2f} reasoning={decision.reasoning!r}"
|
|
)
|
|
|
|
if decision.action == "monitor":
|
|
pass # transient — the backend will retry
|
|
|
|
elif decision.action == "notify_customer":
|
|
if registry.autonomous("DISPATCH_AGENT", DISPATCH_AGENT_AUTONOMOUS):
|
|
await self._notify_customer_delay(booking_id)
|
|
else:
|
|
# Outward-facing action not authorised — raise an internal proposal instead.
|
|
await self._ops_alert(
|
|
zone_id, booking_id,
|
|
f"[proposed: notify_customer, conf={decision.confidence:.2f}] {decision.reasoning}",
|
|
)
|
|
|
|
elif decision.action == "ops_alert":
|
|
await self._ops_alert(zone_id, booking_id, decision.reasoning)
|
|
|
|
else: # "escalate"
|
|
await self._escalate_dispatch(zone_id, booking_id, decision)
|
|
|
|
except Exception as e:
|
|
logger.error(f"DispatchAgent booking.assignment_failed handler error: {e}")
|
|
finally:
|
|
await msg.ack()
|
|
|
|
# ------------------------------------------------------------------ #
|
|
# Context gathering + actions #
|
|
# ------------------------------------------------------------------ #
|
|
|
|
async def _gather_assignment_facts(self, zone_id, lat, lon):
|
|
"""Read-only coverage context plus the daily failure counter for this
|
|
zone. Returns (facts, failures_today). Best-effort — any failure just
|
|
yields fewer facts, never raises."""
|
|
facts: Dict[str, Any] = {"zone_id": zone_id}
|
|
has_coords = lat is not None and lon is not None
|
|
facts["has_coordinates"] = has_coords
|
|
|
|
nearest_km = None
|
|
nearest_presence = None
|
|
if has_coords:
|
|
for radius in (10, 20, 30):
|
|
try:
|
|
found = await self._find_zone(lat, lon, radius_km=radius)
|
|
except Exception as e:
|
|
logger.warning(f"assignment facts: GEORADIUS {radius}km failed: {e}")
|
|
break
|
|
if found:
|
|
nearest_km = radius
|
|
nearest_presence = found.get("presence", UNKNOWN)
|
|
break
|
|
|
|
if not has_coords:
|
|
facts["nearest_miler_within_km"] = "unknown (no coordinates)"
|
|
facts["nearest_miler_presence"] = "unknown (no coordinates)"
|
|
elif nearest_km is None:
|
|
facts["nearest_miler_within_km"] = "none within 30km"
|
|
facts["nearest_miler_presence"] = "none found"
|
|
else:
|
|
facts["nearest_miler_within_km"] = nearest_km
|
|
# "available" = backend confirms on duty; "unknown" = no presence
|
|
# data for any nearby miler, so this is not evidence of a gap.
|
|
facts["nearest_miler_presence"] = nearest_presence
|
|
|
|
# Distinguish "no riders registered here" from "riders exist, none on duty".
|
|
if has_coords:
|
|
indexed = await self._count_in_geo_index(lat, lon, radius_km=30)
|
|
if indexed is not None:
|
|
facts["milers_in_geo_index_within_30km"] = indexed
|
|
|
|
count = 0
|
|
try:
|
|
today = datetime.now(timezone.utc).strftime("%Y-%m-%d")
|
|
key = f"failed_assignment:{zone_id}:{today}"
|
|
count = await self._redis.incr(key)
|
|
await self._redis.expire(key, 172800) # 48h covers day rollover
|
|
except Exception as e:
|
|
logger.warning(f"assignment facts: failure counter failed for zone {zone_id}: {e}")
|
|
facts["failures_today"] = count
|
|
facts["alert_already_sent_today"] = await self._alert_already_sent_today(zone_id)
|
|
|
|
return facts, count
|
|
|
|
def _alert_key(self, zone_id) -> str:
|
|
today = datetime.now(timezone.utc).strftime("%Y-%m-%d")
|
|
return f"assignment_alert_sent:{zone_id}:{today}"
|
|
|
|
async def _alert_already_sent_today(self, zone_id) -> bool:
|
|
try:
|
|
return bool(await self._redis.exists(self._alert_key(zone_id)))
|
|
except Exception:
|
|
return False
|
|
|
|
async def _mark_alert_sent(self, zone_id):
|
|
try:
|
|
await self._redis.set(self._alert_key(zone_id), "1", ex=172800, nx=True)
|
|
except Exception as e:
|
|
logger.warning(f"could not mark alert sent for zone {zone_id}: {e}")
|
|
|
|
async def _ops_alert(self, zone_id, booking_id, reasoning, severity="medium"):
|
|
"""Internal ops alert. Goes to JARVIS as EXCEPTION_DETECTED — the same
|
|
path the ExceptionAgent uses — which is logged at WARNING and kept in
|
|
JARVIS's escalation inbox. (Previously sent as an 'ops_alert' task to
|
|
CUSTOMER_AGENT, which had no handler for it: 1,001 alerts went nowhere.)"""
|
|
logger.warning(f"[DISPATCH] OPS ALERT zone={zone_id} booking={booking_id}: {reasoning}")
|
|
await self.send_message(
|
|
recipient="JARVIS",
|
|
message_type=MessageType.EXCEPTION_DETECTED,
|
|
payload={
|
|
"exception_type": "coverage_gap",
|
|
"severity": severity,
|
|
"zone_id": zone_id,
|
|
"order_id": booking_id,
|
|
"reasoning": reasoning,
|
|
"source": self.agent_id,
|
|
},
|
|
correlation_id=str(booking_id),
|
|
)
|
|
await self._mark_alert_sent(zone_id)
|
|
|
|
async def _notify_customer_delay(self, booking_id):
|
|
await self.send_message(
|
|
recipient="CUSTOMER_AGENT",
|
|
message_type=MessageType.AGENT_TASK,
|
|
payload={
|
|
# CUSTOMER_AGENT's real task contract (send_notification + DELAYED
|
|
# template). Every variable the template names must be supplied or
|
|
# the send is refused; we have no ETA at assignment-failure time,
|
|
# so say so rather than inventing one. No customer_id: the event
|
|
# carries none, and the backend resolves the recipient from order_id.
|
|
"task_type": "send_notification",
|
|
"order_id": booking_id,
|
|
"notification_type": "delayed",
|
|
"template_vars": {"order_id": booking_id, "new_eta": "being confirmed"},
|
|
},
|
|
correlation_id=str(booking_id),
|
|
)
|
|
|
|
async def _escalate_dispatch(self, zone_id, booking_id, decision):
|
|
logger.warning(
|
|
f"[DISPATCH] Escalating booking {booking_id} (zone {zone_id}) to a human "
|
|
f"(proposed={decision.action}, confidence={decision.confidence:.2f}): {decision.reasoning}"
|
|
)
|
|
await self.send_message(
|
|
recipient="JARVIS",
|
|
message_type=MessageType.EXCEPTION_DETECTED,
|
|
payload={
|
|
"exception_type": "assignment_failed",
|
|
"severity": "high",
|
|
"zone_id": zone_id,
|
|
"order_id": booking_id,
|
|
"proposed_action": decision.action,
|
|
"confidence": decision.confidence,
|
|
"reasoning": decision.reasoning,
|
|
"source": self.agent_id,
|
|
},
|
|
correlation_id=str(booking_id),
|
|
)
|
|
await self._mark_alert_sent(zone_id)
|
|
|
|
# ------------------------------------------------------------------ #
|
|
# Task handler (no inbound task types remain) #
|
|
# ------------------------------------------------------------------ #
|
|
|
|
async def handle_task(self, task: AgentTask) -> Dict[str, Any]:
|
|
return {"status": "error", "message": f"Unknown task: {task.task_type}"}
|
|
|
|
async def think(self, context: str, options: List[str] = None) -> str:
|
|
return f"[DISPATCH_AGENT reasoning]: {context}"
|