Compare commits

..

2 Commits

Author SHA1 Message Date
5159c8d1a7 fix(dispatch): presence is three-state — absent status key is unknown, not off-duty
The previous commit treated a missing miler_status:<id> key as "not available".
In production only 2 of 34 milers have that key at all, so the agent would have
reported "no available rider" for nearly every failure — a confident wrong
answer in the opposite direction from the bug it fixed.

- _miler_presence returns available / unavailable / unknown. No key, an
  unparseable value, or an unrecognised status reads as unknown.
- _find_zone prefers a confirmed-available miler, otherwise reports the nearest
  unknown-presence one (it may well be assignable), and returns None only when
  every nearby candidate is confirmed off duty.
- Facts carry nearest_miler_within_km + nearest_miler_presence; the prompt
  states plainly that unknown presence is not evidence of a coverage gap and
  should lean to monitor/escalate rather than ops_alert.
- New eval case for riders-nearby-but-no-presence-data; tests for all three
  presence states.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin
2026-09-22 16:39:33 +05:30
4283c602f6 fix(dispatch): liveness-aware coverage check, escalation rate limit, real alert sink
Prod findings (2026-09-22): a backend retry sweep failed 1,001 bookings in
60s; the agent called each a "coverage gap" because GEORADIUS found a miler
in the geo index (last seen in June), then sent 1,001 ops_alert tasks to
CUSTOMER_AGENT, which has no such handler.

- _find_zone filters GEORADIUS candidates by the backend's miler_status:<id>
  key; only status=Available counts. Facts now carry
  nearest_available_miler_within_km plus milers_in_geo_index_within_30km so
  the decision can separate "no riders here" from "riders exist, none on duty".
- Rate limit per zone per day: after the first alert, further failures only
  bump the counter (no LLM call); a summary re-alert goes out every
  DISPATCH_REALERT_EVERY (default 100).
- _ops_alert / _escalate_dispatch send EXCEPTION_DETECTED to JARVIS (the path
  that is actually handled); customer delay notice uses CUSTOMER_AGENT's real
  send_notification contract.
- JARVIS: escalation inbox (_escalations, pending_escalations()) and
  human_review/ops_alert task types are recorded instead of dropped.
- ExceptionAgent pull loops: also catch asyncio.TimeoutError (distinct from
  nats.errors.TimeoutError on 3.11) and log the exception type — the blank
  "pull loop error:" lines.
- Prompt + eval cases updated for the renamed facts; new case for the
  observed index-full/nobody-on-duty pattern. Tests for liveness filtering,
  burst suppression, fallback heuristic, sinks, and the JARVIS inbox.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin
2026-09-22 16:26:27 +05:30
8 changed files with 434 additions and 48 deletions

View File

@@ -27,6 +27,26 @@ 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"))
# 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"
class DispatchAgent(SpecializedAgent):
"""
@@ -126,34 +146,80 @@ class DispatchAgent(SpecializedAgent):
# 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; returns nearest miler or None."""
"""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=5,
count=10,
)
if not results:
return None
nearest_miler = results[0]
miler_info = await self._redis.hgetall(f"miler:{nearest_miler}")
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
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)),
}
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 #
# ------------------------------------------------------------------ #
@@ -217,6 +283,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:
@@ -272,22 +355,36 @@ class DispatchAgent(SpecializedAgent):
facts["has_coordinates"] = has_coords
nearest_km = None
nearest_presence = 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
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:
@@ -298,25 +395,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 +457,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) #

View File

@@ -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 ───────────────────────────────────

View File

@@ -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")

View File

@@ -181,10 +181,21 @@ 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 rider is and whether
that rider is actually on duty, the time of day, and whether coordinates were even available.
Read the rider facts carefully:
- "nearest_miler_within_km": distance to the nearest rider in the location index.
- "nearest_miler_presence": "available" = the backend confirms that rider is on duty;
"unavailable" = confirmed off duty / on break; "unknown" = there is NO presence data for nearby
riders. Unknown is not evidence of a coverage gap — do not treat it as "nobody is working".
- "milers_in_geo_index_within_30km": raw count of riders ever seen nearby, including stale entries.
A single failure with a rider nearby is usually transient. Repeated failures with no rider at all within
30 km is a genuine coverage gap. Repeated failures *despite* a confirmed-available rider nearby is not a
coverage gap but a systemic issue — escalate to a human. When presence is unknown, lean toward monitoring
or escalating for a human to check rather than declaring a coverage gap.
Report an honest confidence in [0, 1]."""
_ASSIGNMENT_SCHEMA = {
"type": "object",

View File

@@ -1,10 +1,12 @@
{"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, "nearest_miler_within_km": 10, "nearest_miler_presence": "available", "milers_in_geo_index_within_30km": 3, "has_coordinates": true, "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, "nearest_miler_within_km": "none within 30km", "nearest_miler_presence": "none found", "milers_in_geo_index_within_30km": 0, "has_coordinates": true, "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, "nearest_miler_within_km": 30, "nearest_miler_presence": "available", "milers_in_geo_index_within_30km": 3, "has_coordinates": true, "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, "nearest_miler_within_km": "unknown (no coordinates)", "nearest_miler_presence": "unknown (no coordinates)", "has_coordinates": false, "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, "nearest_miler_within_km": 10, "nearest_miler_presence": "available", "milers_in_geo_index_within_30km": 3, "has_coordinates": true, "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, "nearest_miler_within_km": "none within 30km", "nearest_miler_presence": "none found", "milers_in_geo_index_within_30km": 0, "has_coordinates": true, "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, "nearest_miler_within_km": 30, "nearest_miler_presence": "available", "milers_in_geo_index_within_30km": 3, "has_coordinates": true, "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, "nearest_miler_within_km": 30, "nearest_miler_presence": "available", "milers_in_geo_index_within_30km": 3, "has_coordinates": true, "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, "nearest_miler_within_km": 20, "nearest_miler_presence": "available", "milers_in_geo_index_within_30km": 3, "has_coordinates": true, "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, "nearest_miler_within_km": "none within 30km", "nearest_miler_presence": "none found", "milers_in_geo_index_within_30km": 0, "has_coordinates": true, "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_miler_within_km": "none within 30km", "nearest_miler_presence": "none found", "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."}
{"id": "rider-near-presence-unknown", "hour_of_day_utc": 11, "facts": {"zone_id": "coimbatore", "failures_today": 4, "nearest_miler_within_km": 10, "nearest_miler_presence": "unknown", "milers_in_geo_index_within_30km": 26, "has_coordinates": true, "alert_already_sent_today": false}, "acceptable": ["escalate", "monitor", "notify_customer"], "ideal": "escalate", "rationale": "Riders are nearby but the backend reports no presence data \u2014 we cannot tell on-duty from off-duty, so this is NOT a coverage gap; a human should check why assignment fails."}

View File

@@ -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_miler_within_km, nearest_miler_presence,
milers_in_geo_index_within_30km, has_coordinates) so the eval
reflects what DispatchAgent actually sends.
Usage:

View File

@@ -17,13 +17,45 @@ 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` / `on_break` are miler ids with a miler_status:<id> key saying so.
Any geo member in neither set has NO status key at all — the common case in
production, which must read as unknown presence, not off-duty."""
def __init__(self, geo=None, counter_start=0, raise_on=(), available=None,
on_break=(), 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._on_break = set(on_break)
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._on_break:
return '{"userid": 0, "status": "Break"}'
return None # no status key — unknown, not off-duty
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:
@@ -89,6 +121,9 @@ class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase):
self.assertEqual(facts["zone_id"], "hyderabad")
self.assertTrue(facts["has_coordinates"])
self.assertEqual(facts["nearest_miler_within_km"], 10)
self.assertEqual(facts["nearest_miler_presence"], "available")
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)
@@ -97,6 +132,7 @@ class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase):
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["milers_in_geo_index_within_30km"], 0)
self.assertEqual(count, 3)
async def test_no_coordinates(self):
@@ -104,6 +140,7 @@ class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase):
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.assertNotIn("milers_in_geo_index_within_30km", facts)
self.assertEqual(count, 1) # counter still incremented
self.assertEqual(redis.incr_calls and 1, 1)
@@ -121,6 +158,136 @@ class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase):
self.assertEqual(count, 0)
class TestPresence(unittest.IsolatedAsyncioTestCase):
"""Presence is three-state. The prod incident had 26 riders in the geo index;
only 2 of 34 milers had a status key at all, so "no key" must read as unknown
— reporting it as off-duty would be a confident wrong answer."""
async def test_confirmed_available_rider_is_preferred(self):
redis = FakeRedis(geo=["4", "7", "12"], available=["12"], on_break=["4", "7"])
found = await make_agent(redis)._find_zone(11.0, 76.9, radius_km=10)
self.assertEqual(found["miler_id"], "12")
self.assertEqual(found["presence"], "available")
async def test_no_status_key_reads_as_unknown_not_unavailable(self):
redis = FakeRedis(geo=["4", "7"], available=[], on_break=[]) # no status keys at all
agent = make_agent(redis)
self.assertEqual(await agent._miler_presence("4"), "unknown")
found = await agent._find_zone(11.0, 76.9, radius_km=10)
self.assertIsNotNone(found) # still reported — may be assignable
self.assertEqual(found["presence"], "unknown")
async def test_all_confirmed_off_duty_returns_none(self):
redis = FakeRedis(geo=["4", "7"], available=[], on_break=["4", "7"])
self.assertIsNone(await make_agent(redis)._find_zone(11.0, 76.9, radius_km=10))
async def test_facts_report_presence_unknown_with_rider_nearby(self):
redis = FakeRedis(geo=["4", "7", "12"], available=[], on_break=[])
facts, _ = await make_agent(redis)._gather_assignment_facts("coimbatore", 11.0, 76.9)
self.assertEqual(facts["nearest_miler_within_km"], 10)
self.assertEqual(facts["nearest_miler_presence"], "unknown")
self.assertEqual(facts["milers_in_geo_index_within_30km"], 3)
async def test_facts_report_none_found_when_all_off_duty(self):
redis = FakeRedis(geo=["4"], available=[], on_break=["4"])
facts, _ = await make_agent(redis)._gather_assignment_facts("coimbatore", 11.0, 76.9)
self.assertEqual(facts["nearest_miler_within_km"], "none within 30km")
self.assertEqual(facts["nearest_miler_presence"], "none found")
self.assertEqual(facts["milers_in_geo_index_within_30km"], 1) # index still knows them
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")

View File

@@ -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