"""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 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", "") 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 _find_zone( self, lat: float, lon: float, radius_km: int = 10 ) -> Optional[Dict[str, Any]]: """GEORADIUS sweep on milers:locations; returns nearest miler or None.""" try: results = await self._redis.georadius( "milers:locations", lon, lat, radius_km, "km", sort="ASC", count=5, ) if not results: return None nearest_miler = results[0] miler_info = await self._redis.hgetall(f"miler:{nearest_miler}") return { "miler_id": nearest_miler, "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)), } except Exception as e: logger.warning(f"Redis GEORADIUS error: {e}") 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}" ) await self.send_message( recipient="HUB_AGENT", message_type=MessageType.AGENT_TASK, payload={ "task_type": "prepare_receiving", "booking_id": booking_id, "miler_id": miler_id, "hub_id": hub_id, }, correlation_id=str(booking_id), ) 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") zone_id = data.get("zone_id", "unknown") lat = data.get("lat") lon = data.get("lon") logger.warning(f"[DISPATCH] booking.assignment_failed — booking={booking_id} zone={zone_id}") facts, count = await self._gather_assignment_facts(zone_id, lat, lon) logger.info(f"[DISPATCH] booking={booking_id} context facts: {facts}") decision = await decide_assignment_failure(build_assignment_failure_context(facts)) 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 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 if has_coords: for radius in (10, 20, 30): try: if await self._find_zone(lat, lon, radius_km=radius): nearest_km = radius break except Exception as e: logger.warning(f"assignment facts: GEORADIUS {radius}km failed: {e}") break if not has_coords: facts["nearest_miler_within_km"] = "unknown (no coordinates)" elif nearest_km is None: facts["nearest_miler_within_km"] = "none within 30km" else: facts["nearest_miler_within_km"] = nearest_km 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 return facts, count async def _ops_alert(self, zone_id, booking_id, reasoning): await self.send_message( recipient="CUSTOMER_AGENT", message_type=MessageType.AGENT_TASK, payload={"task_type": "ops_alert", "zone_id": zone_id, "reasoning": reasoning}, correlation_id=str(booking_id), ) async def _notify_customer_delay(self, booking_id): await self.send_message( recipient="CUSTOMER_AGENT", message_type=MessageType.AGENT_TASK, payload={ "task_type": "customer_notify", "booking_id": booking_id, "message": "We're finding the right rider for your delivery — thanks for your patience.", }, 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.AGENT_TASK, payload={ "task_type": "human_review", "reason": "assignment_failed", "zone_id": zone_id, "booking_id": booking_id, "proposed_action": decision.action, "confidence": decision.confidence, "reasoning": decision.reasoning, }, correlation_id=str(booking_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}"