Files
AI_engine/agents/dispatch_agent.py

517 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)
"task_type": "send_notification",
"order_id": booking_id,
"notification_type": "delayed",
"template_vars": {"order_id": booking_id, "reason": "finding the right rider"},
},
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}"