diff --git a/agents/dispatch_agent.py b/agents/dispatch_agent.py index a61ba5e..b8bfa93 100644 --- a/agents/dispatch_agent.py +++ b/agents/dispatch_agent.py @@ -27,6 +27,18 @@ DISPATCH_AGENT_AUTONOMOUS = os.getenv("DISPATCH_AGENT_AUTONOMOUS", "false").lowe # 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")) + +# A miler counts as available only if the Go backend's presence key says so. +# The milers:locations geo index keeps last-known positions indefinitely, so +# membership alone says nothing about whether anyone is actually on duty. +MILER_AVAILABLE_STATUSES = {"available"} + class DispatchAgent(SpecializedAgent): """ @@ -126,34 +138,60 @@ class DispatchAgent(SpecializedAgent): # Redis GEO lookup # # ------------------------------------------------------------------ # + async def _miler_is_available(self, miler_id) -> bool: + """Liveness check against the backend-owned miler_status: key + ({"userid": .., "status": "Available" | "Break" | ...}). Missing or + unparseable → not available: geo-index presence alone is not evidence.""" + try: + raw = await self._redis.get(f"miler_status:{miler_id}") + if not raw: + return False + status = json.loads(raw).get("status", "") + return str(status).lower() in MILER_AVAILABLE_STATUSES + except Exception: + return False + 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.""" + """GEORADIUS sweep on milers:locations filtered to milers whose status + key says they are available. Returns the nearest available miler (with + how many index entries were skipped as stale/off-duty), or None.""" try: results = await self._redis.georadius( "milers:locations", lon, lat, radius_km, "km", sort="ASC", - count=5, + count=10, ) 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)), - } + for candidate in results: + if await self._miler_is_available(candidate): + miler_info = await self._redis.hgetall(f"miler:{candidate}") + return { + "miler_id": candidate, + "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)), + "stale_candidates_skipped": len(results) - 1, + } + return None except Exception as e: logger.warning(f"Redis GEORADIUS error: {e}") return None + 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 # # ------------------------------------------------------------------ # @@ -217,6 +255,23 @@ class DispatchAgent(SpecializedAgent): facts, count = await self._gather_assignment_facts(zone_id, lat, lon) 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 DISPATCH_REALERT_EVERY and count % DISPATCH_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 {DISPATCH_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 + decision = await decide_assignment_failure(build_assignment_failure_context(facts)) if decision is None: @@ -283,11 +338,17 @@ class DispatchAgent(SpecializedAgent): break if not has_coords: - facts["nearest_miler_within_km"] = "unknown (no coordinates)" + facts["nearest_available_miler_within_km"] = "unknown (no coordinates)" elif nearest_km is None: - facts["nearest_miler_within_km"] = "none within 30km" + facts["nearest_available_miler_within_km"] = "none within 30km" else: - facts["nearest_miler_within_km"] = nearest_km + facts["nearest_available_miler_within_km"] = nearest_km + + # 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: @@ -298,25 +359,57 @@ class DispatchAgent(SpecializedAgent): 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 - async def _ops_alert(self, zone_id, booking_id, reasoning): + 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="CUSTOMER_AGENT", - message_type=MessageType.AGENT_TASK, - payload={"task_type": "ops_alert", "zone_id": zone_id, "reasoning": reasoning}, + 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={ - "task_type": "customer_notify", - "booking_id": booking_id, - "message": "We're finding the right rider for your delivery — thanks for your patience.", + # 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), ) @@ -328,18 +421,20 @@ class DispatchAgent(SpecializedAgent): ) await self.send_message( recipient="JARVIS", - message_type=MessageType.AGENT_TASK, + message_type=MessageType.EXCEPTION_DETECTED, payload={ - "task_type": "human_review", - "reason": "assignment_failed", + "exception_type": "assignment_failed", + "severity": "high", "zone_id": zone_id, - "booking_id": booking_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) # diff --git a/agents/exception_agent.py b/agents/exception_agent.py index 28be9bf..2eed928 100644 --- a/agents/exception_agent.py +++ b/agents/exception_agent.py @@ -204,10 +204,10 @@ class ExceptionAgent(SpecializedAgent): logger.error(f"Location handler error: {e}") finally: await msg.ack() - except nats.errors.TimeoutError: - pass + except (nats.errors.TimeoutError, asyncio.TimeoutError): + pass # empty fetch — normal except Exception as e: - logger.error(f"Location pull loop error: {e}") + logger.error(f"Location pull loop error: {type(e).__name__}: {e}") await asyncio.sleep(2) async def _pull_stalled_loop(self): @@ -224,10 +224,10 @@ class ExceptionAgent(SpecializedAgent): logger.error(f"Stalled handler error: {e}") finally: await msg.ack() - except nats.errors.TimeoutError: - pass + except (nats.errors.TimeoutError, asyncio.TimeoutError): + pass # empty fetch — normal except Exception as e: - logger.error(f"Stalled pull loop error: {e}") + logger.error(f"Stalled pull loop error: {type(e).__name__}: {e}") await asyncio.sleep(2) # ── Stall detection: per location ping ─────────────────────────────────── diff --git a/core/agent.py b/core/agent.py index a9134f0..847e698 100644 --- a/core/agent.py +++ b/core/agent.py @@ -208,6 +208,10 @@ class MasterAgent(Agent): self._sub_agents: Dict[str, Agent] = {} self._active_orders: Dict[str, Dict] = {} self._decision_log: List[Dict] = [] + # Human-review inbox: every escalation/ops alert a sub-agent raised. + # The only sink for proposals the agents aren't authorised to act on; + # surfaced by pending_escalations() for the Command Center. + self._escalations: List[Dict] = [] def register_sub_agent(self, agent: Agent): self._sub_agents[agent.agent_id] = agent @@ -218,7 +222,7 @@ class MasterAgent(Agent): must decide) are logged at WARNING and recorded so they are visible rather than silently dropped — the endpoint of the human-review path.""" if message.message_type == MessageType.EXCEPTION_DETECTED: - logger.warning(f"JARVIS: exception escalation from {message.sender}: {message.payload}") + self._record_escalation(message.sender, message.payload) else: logger.info(f"JARVIS: {message.message_type.value} from {message.sender}") self._decision_log.append({ @@ -240,9 +244,30 @@ class MasterAgent(Agent): return await self._handle_exception(task) elif task.task_type == "generate_report": return await self._generate_report(task) + elif task.task_type in ("human_review", "ops_alert"): + # Escalations sent as tasks land here instead of dead-lettering. + self._record_escalation(task.data.get("source", "unknown"), task.data) + return {"status": "recorded", "escalations_pending": len(self._escalations)} else: + logger.warning(f"JARVIS: unknown task type {task.task_type!r} — dropped") return {"status": "unknown_task", "task_type": task.task_type} + def _record_escalation(self, sender: str, payload: Dict[str, Any]): + logger.warning(f"JARVIS: escalation from {sender}: {payload}") + self._escalations.append({ + "timestamp": datetime.now(), + "from": sender, + "exception_type": payload.get("exception_type"), + "severity": payload.get("severity"), + "payload": payload, + }) + if len(self._escalations) > 500: + self._escalations = self._escalations[-500:] + + def pending_escalations(self, limit: int = 50) -> List[Dict[str, Any]]: + """Most recent escalations first — the human-review inbox.""" + return list(reversed(self._escalations[-limit:])) + async def _orchestrate_order(self, task: AgentTask) -> Dict[str, Any]: order_data = task.data.get("order", {}) order_id = order_data.get("order_id", "unknown") diff --git a/core/llm.py b/core/llm.py index 7c42a1a..60fe250 100644 --- a/core/llm.py +++ b/core/llm.py @@ -181,10 +181,14 @@ fallback failed, so no rider was assigned. Decide how to react. Choose exactly o - "ops_alert": a genuine coverage gap in this zone (repeated failures, no nearby riders). Alert operations to onboard or redirect riders here. - "escalate": ambiguous or contradictory — e.g. riders ARE nearby yet assignment keeps failing, which suggests a systemic issue rather than a coverage gap. A human dispatcher should look. -Weigh how many times assignment has failed in this zone today, how far the nearest available rider is, -the time of day, and whether coordinates were even available. A single failure with a rider nearby is -usually transient; repeated failures with no nearby rider is a coverage gap; repeated failures *despite* -nearby riders is not a coverage gap and warrants a human. Report an honest confidence in [0, 1].""" +Weigh how many times assignment has failed in this zone today, how far the nearest *available* rider is, +the time of day, and whether coordinates were even available. "nearest_available_miler_within_km" counts +only riders whose live status is Available; "milers_in_geo_index_within_30km" is the raw count of riders +ever seen nearby (it includes off-duty and stale entries), so a large index count with no available rider +means riders exist here but nobody is on duty — a staffing gap, not a systemic fault. A single failure with +an available rider nearby is usually transient; repeated failures with no available rider is a coverage +gap; repeated failures *despite* available riders nearby is not a coverage gap and warrants a human. +Report an honest confidence in [0, 1].""" _ASSIGNMENT_SCHEMA = { "type": "object", diff --git a/evals/assignment_cases.jsonl b/evals/assignment_cases.jsonl index 7cbbb17..e37043f 100644 --- a/evals/assignment_cases.jsonl +++ b/evals/assignment_cases.jsonl @@ -1,10 +1,11 @@ -{"id": "first-failure-close-rider", "hour_of_day_utc": 14, "facts": {"zone_id": "hyderabad", "failures_today": 1, "nearest_miler_within_km": 10, "has_coordinates": true}, "acceptable": ["monitor", "notify_customer"], "ideal": "monitor", "rationale": "Single failure with a rider only 10km away — almost certainly transient; let the backend retry."} -{"id": "repeated-zone-gap", "hour_of_day_utc": 12, "facts": {"zone_id": "pune", "failures_today": 4, "nearest_miler_within_km": "none within 30km", "has_coordinates": true}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Repeated failures and no rider within 30km — a genuine coverage gap; alert ops."} -{"id": "second-failure-far-rider", "hour_of_day_utc": 18, "facts": {"zone_id": "bangalore", "failures_today": 2, "nearest_miler_within_km": 30, "has_coordinates": true}, "acceptable": ["notify_customer", "monitor"], "ideal": "notify_customer", "rationale": "A couple of failures with the nearest rider 30km out during rush — real delay; keep the customer informed."} -{"id": "no-coordinates-first", "hour_of_day_utc": 10, "facts": {"zone_id": "unknown", "failures_today": 1, "nearest_miler_within_km": "unknown (no coordinates)", "has_coordinates": false}, "acceptable": ["monitor", "escalate"], "ideal": "monitor", "rationale": "Coverage can't be assessed without coordinates, but a single failure is most likely transient."} -{"id": "many-failures-riders-near", "hour_of_day_utc": 15, "facts": {"zone_id": "mumbai_west", "failures_today": 5, "nearest_miler_within_km": 10, "has_coordinates": true}, "acceptable": ["escalate", "ops_alert"], "ideal": "escalate", "rationale": "Riders ARE nearby yet assignment keeps failing — not a coverage gap; a systemic issue for a human."} -{"id": "night-coverage-gap", "hour_of_day_utc": 2, "facts": {"zone_id": "kolkata", "failures_today": 3, "nearest_miler_within_km": "none within 30km", "has_coordinates": true}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Repeated overnight failures with no nearby rider — a coverage gap worth flagging to ops."} -{"id": "single-far-rider", "hour_of_day_utc": 9, "facts": {"zone_id": "hyderabad", "failures_today": 1, "nearest_miler_within_km": 30, "has_coordinates": true}, "acceptable": ["monitor", "notify_customer"], "ideal": "monitor", "rationale": "One failure, distant rider — likely resolves on retry."} -{"id": "persistent-far-rush", "hour_of_day_utc": 18, "facts": {"zone_id": "mumbai_east", "failures_today": 3, "nearest_miler_within_km": 30, "has_coordinates": true}, "acceptable": ["notify_customer", "ops_alert"], "ideal": "notify_customer", "rationale": "Rush-hour, repeated, distant rider — the customer should hear about the delay."} -{"id": "borderline-count", "hour_of_day_utc": 13, "facts": {"zone_id": "north_delhi", "failures_today": 3, "nearest_miler_within_km": 20, "has_coordinates": true}, "acceptable": ["notify_customer", "ops_alert", "monitor"], "ideal": "notify_customer", "rationale": "Genuinely borderline — a few failures with a moderately-close rider; any of these is defensible."} -{"id": "severe-gap", "hour_of_day_utc": 11, "facts": {"zone_id": "pune", "failures_today": 8, "nearest_miler_within_km": "none within 30km", "has_coordinates": true}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Many failures and no rider anywhere near — a clear, severe coverage gap."} +{"id": "first-failure-close-rider", "hour_of_day_utc": 14, "facts": {"zone_id": "hyderabad", "failures_today": 1, "has_coordinates": true, "nearest_available_miler_within_km": 10, "milers_in_geo_index_within_30km": 3, "alert_already_sent_today": false}, "acceptable": ["monitor", "notify_customer"], "ideal": "monitor", "rationale": "Single failure with a rider only 10km away \u2014 almost certainly transient; let the backend retry."} +{"id": "repeated-zone-gap", "hour_of_day_utc": 12, "facts": {"zone_id": "pune", "failures_today": 4, "has_coordinates": true, "nearest_available_miler_within_km": "none within 30km", "milers_in_geo_index_within_30km": 0, "alert_already_sent_today": false}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Repeated failures and no rider within 30km \u2014 a genuine coverage gap; alert ops."} +{"id": "second-failure-far-rider", "hour_of_day_utc": 18, "facts": {"zone_id": "bangalore", "failures_today": 2, "has_coordinates": true, "nearest_available_miler_within_km": 30, "milers_in_geo_index_within_30km": 3, "alert_already_sent_today": false}, "acceptable": ["notify_customer", "monitor"], "ideal": "notify_customer", "rationale": "A couple of failures with the nearest rider 30km out during rush \u2014 real delay; keep the customer informed."} +{"id": "no-coordinates-first", "hour_of_day_utc": 10, "facts": {"zone_id": "unknown", "failures_today": 1, "has_coordinates": false, "nearest_available_miler_within_km": "unknown (no coordinates)", "alert_already_sent_today": false}, "acceptable": ["monitor", "escalate"], "ideal": "monitor", "rationale": "Coverage can't be assessed without coordinates, but a single failure is most likely transient."} +{"id": "many-failures-riders-near", "hour_of_day_utc": 15, "facts": {"zone_id": "mumbai_west", "failures_today": 5, "has_coordinates": true, "nearest_available_miler_within_km": 10, "milers_in_geo_index_within_30km": 3, "alert_already_sent_today": false}, "acceptable": ["escalate", "ops_alert"], "ideal": "escalate", "rationale": "Riders ARE nearby yet assignment keeps failing \u2014 not a coverage gap; a systemic issue for a human."} +{"id": "night-coverage-gap", "hour_of_day_utc": 2, "facts": {"zone_id": "kolkata", "failures_today": 3, "has_coordinates": true, "nearest_available_miler_within_km": "none within 30km", "milers_in_geo_index_within_30km": 0, "alert_already_sent_today": false}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Repeated overnight failures with no nearby rider \u2014 a coverage gap worth flagging to ops."} +{"id": "single-far-rider", "hour_of_day_utc": 9, "facts": {"zone_id": "hyderabad", "failures_today": 1, "has_coordinates": true, "nearest_available_miler_within_km": 30, "milers_in_geo_index_within_30km": 3, "alert_already_sent_today": false}, "acceptable": ["monitor", "notify_customer"], "ideal": "monitor", "rationale": "One failure, distant rider \u2014 likely resolves on retry."} +{"id": "persistent-far-rush", "hour_of_day_utc": 18, "facts": {"zone_id": "mumbai_east", "failures_today": 3, "has_coordinates": true, "nearest_available_miler_within_km": 30, "milers_in_geo_index_within_30km": 3, "alert_already_sent_today": false}, "acceptable": ["notify_customer", "ops_alert"], "ideal": "notify_customer", "rationale": "Rush-hour, repeated, distant rider \u2014 the customer should hear about the delay."} +{"id": "borderline-count", "hour_of_day_utc": 13, "facts": {"zone_id": "north_delhi", "failures_today": 3, "has_coordinates": true, "nearest_available_miler_within_km": 20, "milers_in_geo_index_within_30km": 3, "alert_already_sent_today": false}, "acceptable": ["notify_customer", "ops_alert", "monitor"], "ideal": "notify_customer", "rationale": "Genuinely borderline \u2014 a few failures with a moderately-close rider; any of these is defensible."} +{"id": "severe-gap", "hour_of_day_utc": 11, "facts": {"zone_id": "pune", "failures_today": 8, "has_coordinates": true, "nearest_available_miler_within_km": "none within 30km", "milers_in_geo_index_within_30km": 0, "alert_already_sent_today": false}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Many failures and no rider anywhere near \u2014 a clear, severe coverage gap."} +{"id": "index-full-nobody-on-duty", "hour_of_day_utc": 7, "facts": {"zone_id": "coimbatore", "failures_today": 6, "nearest_available_miler_within_km": "none within 30km", "milers_in_geo_index_within_30km": 26, "has_coordinates": true, "alert_already_sent_today": false}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "26 riders have been seen here but none is Available \u2014 a staffing/on-duty gap; ops needs to get riders online, not a systemic fault."} diff --git a/evals/assignment_eval.py b/evals/assignment_eval.py index dbdc702..e0f16f6 100644 --- a/evals/assignment_eval.py +++ b/evals/assignment_eval.py @@ -4,7 +4,8 @@ Offline eval for the DispatchAgent assignment-failure decision Same scoring as the stall eval (see evals/_harness.py): acceptable-rate is the headline, exact-rate secondary. The cases use the real gatherer keys -(zone_id, failures_today, nearest_miler_within_km, has_coordinates) so the eval +(zone_id, failures_today, nearest_available_miler_within_km, +milers_in_geo_index_within_30km, has_coordinates) so the eval reflects what DispatchAgent actually sends. Usage: diff --git a/tests/test_dispatch_agent.py b/tests/test_dispatch_agent.py index e6b08e9..ba784ae 100644 --- a/tests/test_dispatch_agent.py +++ b/tests/test_dispatch_agent.py @@ -17,13 +17,42 @@ from agents.dispatch_agent import DispatchAgent class FakeRedis: - """decode_responses=True style. `geo` is the GEORADIUS result for every call.""" - def __init__(self, geo=None, counter_start=0, raise_on=()): + """decode_responses=True style. `geo` is the GEORADIUS result for every call. + `available` is the set of miler ids whose miler_status: says Available; + every other geo member is treated as present-but-off-duty (the stale-index case).""" + def __init__(self, geo=None, counter_start=0, raise_on=(), available=None, alerted=False): self._geo = list(geo) if geo is not None else [] self._counter = counter_start self._raise_on = set(raise_on) + self._available = set(available) if available is not None else set(self._geo) + self._kv = {} + if alerted: + self._kv["__alerted__"] = "1" self.incr_calls = [] self.expire_calls = [] + self.set_calls = [] + + async def get(self, key): + if key.startswith("miler_status:"): + mid = key.split(":", 1)[1] + if mid in self._available: + return '{"userid": %s, "status": "Available"}' % (mid if mid.isdigit() else 0) + if mid in self._geo: + return '{"userid": 0, "status": "Break"}' + return None + return self._kv.get(key) + + async def exists(self, key): + if key.startswith("assignment_alert_sent:"): + return 1 if "__alerted__" in self._kv or key in self._kv else 0 + return 1 if key in self._kv else 0 + + async def set(self, key, value, ex=None, nx=False): + self.set_calls.append((key, ex, nx)) + if nx and key in self._kv: + return None + self._kv[key] = value + return True async def georadius(self, key, lon, lat, radius, unit, sort=None, count=None): if "georadius" in self._raise_on: @@ -88,7 +117,9 @@ class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase): facts, count = await make_agent(redis)._gather_assignment_facts("hyderabad", 17.4, 78.4) self.assertEqual(facts["zone_id"], "hyderabad") self.assertTrue(facts["has_coordinates"]) - self.assertEqual(facts["nearest_miler_within_km"], 10) + self.assertEqual(facts["nearest_available_miler_within_km"], 10) + self.assertEqual(facts["milers_in_geo_index_within_30km"], 1) + self.assertFalse(facts["alert_already_sent_today"]) self.assertEqual(facts["failures_today"], 1) self.assertEqual(count, 1) self.assertEqual(redis.expire_calls[0][1], 172800) @@ -96,31 +127,147 @@ class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase): async def test_no_rider_within_30km(self): redis = FakeRedis(geo=[], counter_start=2) # every sweep empty facts, count = await make_agent(redis)._gather_assignment_facts("pune", 18.5, 73.8) - self.assertEqual(facts["nearest_miler_within_km"], "none within 30km") + self.assertEqual(facts["nearest_available_miler_within_km"], "none within 30km") + self.assertEqual(facts["milers_in_geo_index_within_30km"], 0) self.assertEqual(count, 3) async def test_no_coordinates(self): redis = FakeRedis(geo=["m1"]) facts, count = await make_agent(redis)._gather_assignment_facts("unknown", None, None) self.assertFalse(facts["has_coordinates"]) - self.assertEqual(facts["nearest_miler_within_km"], "unknown (no coordinates)") + self.assertEqual(facts["nearest_available_miler_within_km"], "unknown (no coordinates)") + self.assertNotIn("milers_in_geo_index_within_30km", facts) self.assertEqual(count, 1) # counter still incremented self.assertEqual(redis.incr_calls and 1, 1) async def test_georadius_error_is_swallowed(self): redis = FakeRedis(raise_on={"georadius"}) facts, count = await make_agent(redis)._gather_assignment_facts("z", 1.0, 2.0) # must not raise - self.assertEqual(facts["nearest_miler_within_km"], "none within 30km") + self.assertEqual(facts["nearest_available_miler_within_km"], "none within 30km") self.assertEqual(facts["failures_today"], 1) async def test_counter_error_is_swallowed(self): redis = FakeRedis(geo=["m9"], raise_on={"incr"}) facts, count = await make_agent(redis)._gather_assignment_facts("z", 1.0, 2.0) # must not raise - self.assertEqual(facts["nearest_miler_within_km"], 10) # coverage still gathered + self.assertEqual(facts["nearest_available_miler_within_km"], 10) # coverage still gathered self.assertEqual(facts["failures_today"], 0) # counter degraded to 0 self.assertEqual(count, 0) +class TestLiveness(unittest.IsolatedAsyncioTestCase): + """The prod incident: 26 riders in the geo index, none on duty. The old + check said 'rider within 20 km' and called it a coverage gap.""" + + async def test_stale_index_entries_are_not_available_riders(self): + redis = FakeRedis(geo=["4", "7", "12"], available=[]) # all present, all off-duty + facts, _ = await make_agent(redis)._gather_assignment_facts("coimbatore", 11.0, 76.9) + self.assertEqual(facts["nearest_available_miler_within_km"], "none within 30km") + self.assertEqual(facts["milers_in_geo_index_within_30km"], 3) # but the index knows them + + async def test_first_available_candidate_wins(self): + redis = FakeRedis(geo=["4", "7", "12"], available=["12"]) + found = await make_agent(redis)._find_zone(11.0, 76.9, radius_km=10) + self.assertEqual(found["miler_id"], "12") + self.assertEqual(found["stale_candidates_skipped"], 2) + + async def test_missing_status_key_means_unavailable(self): + redis = FakeRedis(geo=["99"], available=[]) + redis._geo = [] # nothing in geo → get() returns None for status + self.assertFalse(await make_agent(redis)._miler_is_available("99")) + + +class FakeMessages: + def __init__(self): + self.sent = [] + + async def send_message(self, recipient, message_type, payload, correlation_id=None): + self.sent.append((recipient, message_type, payload)) + + +def make_handler_agent(redis, decision): + """Agent with send_message captured and the LLM stubbed to a fixed decision (or None).""" + from agents import dispatch_agent as mod + agent = DispatchAgent.__new__(DispatchAgent) + agent._redis = redis + agent.agent_id = "DISPATCH_AGENT" + box = FakeMessages() + agent.send_message = box.send_message + calls = [] + async def fake_decide(context): + calls.append(context) + return decision + mod.decide_assignment_failure = fake_decide + return agent, box, calls + + +class FakeMsg: + def __init__(self, payload): + import json + self.data = json.dumps(payload).encode() + self.acked = False + + async def ack(self): + self.acked = True + + +class TestRateLimitAndSinks(unittest.IsolatedAsyncioTestCase): + def setUp(self): + from agents import dispatch_agent as mod + self._orig = mod.decide_assignment_failure + + def tearDown(self): + from agents import dispatch_agent as mod + mod.decide_assignment_failure = self._orig + + async def test_ops_alert_goes_to_jarvis_as_exception_and_marks_zone(self): + from core.types import MessageType + from core.llm import AssignmentDecision + redis = FakeRedis(geo=[], counter_start=2) + agent, box, calls = make_handler_agent(redis, AssignmentDecision("ops_alert", "gap", 0.9)) + msg = FakeMsg({"booking_id": 1, "lat": 11.0, "lon": 76.9}) + await agent._on_nats_booking_assignment_failed(msg) + self.assertTrue(msg.acked) + self.assertEqual(len(calls), 1) # LLM consulted once + recipient, mtype, payload = box.sent[0] + self.assertEqual(recipient, "JARVIS") # not CUSTOMER_AGENT + self.assertEqual(mtype, MessageType.EXCEPTION_DETECTED) # the path JARVIS handles + self.assertEqual(payload["exception_type"], "coverage_gap") + self.assertTrue(any(k.startswith("assignment_alert_sent:") for k, _, _ in redis.set_calls)) + + async def test_burst_after_alert_skips_llm_and_realerts_every_n(self): + from agents import dispatch_agent as mod + from core.llm import AssignmentDecision + redis = FakeRedis(geo=[], counter_start=0, alerted=True) # zone already alerted today + agent, box, calls = make_handler_agent(redis, AssignmentDecision("ops_alert", "gap", 0.9)) + old = mod.DISPATCH_REALERT_EVERY + mod.DISPATCH_REALERT_EVERY = 100 + try: + for i in range(1, 251): + await agent._on_nats_booking_assignment_failed(FakeMsg({"booking_id": i, "lat": 11.0, "lon": 76.9})) + finally: + mod.DISPATCH_REALERT_EVERY = old + self.assertEqual(calls, []) # zero LLM calls during the burst + self.assertEqual(len(box.sent), 2) # re-alerts at 100 and 200 only + self.assertIn("200 failed assignments", box.sent[1][2]["reasoning"]) + + async def test_llm_unavailable_falls_back_to_heuristic(self): + redis = FakeRedis(geo=[], counter_start=2) + agent, box, calls = make_handler_agent(redis, None) # LLM down + await agent._on_nats_booking_assignment_failed(FakeMsg({"booking_id": 5, "lat": 1.0, "lon": 2.0})) + self.assertEqual(len(box.sent), 1) # count reached 3 → heuristic alert + self.assertEqual(box.sent[0][0], "JARVIS") + + async def test_escalate_goes_to_jarvis_as_exception(self): + from core.types import MessageType + from core.llm import AssignmentDecision + redis = FakeRedis(geo=["4"], available=["4"], counter_start=4) + agent, box, calls = make_handler_agent(redis, AssignmentDecision("escalate", "riders near, still failing", 0.8)) + await agent._on_nats_booking_assignment_failed(FakeMsg({"booking_id": 9, "lat": 1.0, "lon": 2.0})) + recipient, mtype, payload = box.sent[0] + self.assertEqual((recipient, mtype), ("JARVIS", MessageType.EXCEPTION_DETECTED)) + self.assertEqual(payload["proposed_action"], "escalate") + + class TestBindConsumer(unittest.IsolatedAsyncioTestCase): async def test_binds_to_discovered_stream_without_creating(self): js = FakeJS(stream_name="TRACKING") diff --git a/tests/test_jarvis_inbox.py b/tests/test_jarvis_inbox.py new file mode 100644 index 0000000..d131858 --- /dev/null +++ b/tests/test_jarvis_inbox.py @@ -0,0 +1,49 @@ +"""JARVIS is the only sink for escalations; nothing sent there may dead-letter.""" +import pytest + +from datetime import datetime + +from core.agent import MasterAgent +from core.types import AgentMessage, AgentTask, MessageType, Priority + + +def make_jarvis(): + j = MasterAgent.__new__(MasterAgent) + j.agent_id = "JARVIS" + j._escalations = [] + j._decision_log = [] + return j + + +@pytest.mark.asyncio +async def test_exception_detected_message_lands_in_inbox(): + j = make_jarvis() + await j.handle_message(AgentMessage( + message_id="m1", timestamp=datetime.now(), + sender="DISPATCH_AGENT", recipient="JARVIS", + message_type=MessageType.EXCEPTION_DETECTED, + payload={"exception_type": "coverage_gap", "severity": "medium", "zone_id": "z"}, + )) + inbox = j.pending_escalations() + assert len(inbox) == 1 + assert inbox[0]["from"] == "DISPATCH_AGENT" + assert inbox[0]["exception_type"] == "coverage_gap" + + +@pytest.mark.asyncio +async def test_human_review_task_is_recorded_not_dropped(): + j = make_jarvis() + task = AgentTask(task_id="t1", agent_type="master", task_type="human_review", priority=Priority.HIGH, + data={"source": "EXCEPTION_AGENT", "exception_type": "miler_stalled", "severity": "high"}) + result = await j.handle_task(task) + assert result["status"] == "recorded" + assert j.pending_escalations()[0]["from"] == "EXCEPTION_AGENT" + + +@pytest.mark.asyncio +async def test_inbox_is_newest_first_and_capped(): + j = make_jarvis() + for i in range(600): + j._record_escalation("X", {"exception_type": "t", "severity": "low", "n": i}) + assert len(j._escalations) == 500 + assert j.pending_escalations(limit=3)[0]["payload"]["n"] == 599