Compare commits

...

6 Commits

Author SHA1 Message Date
9cbb6f9704 fix(customer): never send an unsubstituted placeholder or a fake customer id
With DISPATCH_AGENT autonomous=true in the agent registry, the assignment-failure
path sends a real customer message. It would have sent
  "Delay alert: Order 36 may arrive later than expected. New ETA: {new_eta}"
with customer_id "unknown": the DELAYED template names {new_eta} but DispatchAgent
supplied only order_id, and it has no customer id to give.

- _send_via_channel refuses to send when any {placeholder} survives substitution
  (or the template is missing), logs which variables the caller omitted, and
  records a failed notification. Silence beats a broken message.
- customer_id is omitted from the payload when it is absent or "unknown" rather
  than sent literally; the backend resolves the recipient from order_id.
- DispatchAgent supplies new_eta "being confirmed" — honest, since no ETA exists
  at assignment-failure time.
- Tests: the old bug is refused, a complete payload renders clean, the id is
  omitted/passed correctly, and DispatchAgent's vars satisfy the template.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin
2026-09-30 15:14:08 +05:30
5a7d32cc04 implemenation on the ai agents and registry 2026-09-30 14:56:08 +05:30
aa3bb35733 fix(deploy): pass NATS_HOST/NATS_PORT and autonomy gates through compose
The agents connect with NATS_HOST/NATS_PORT (NATS_URL is informational), but
docker-compose only forwarded NATS_URL. That was masked while system_config
carried the real host as a hardcoded default; after the secrets scrub the
default is localhost, so a deploy would have sent every NATS connection —
message bus, DispatchAgent, ExceptionAgent — to localhost.

Also forwards the autonomy gates (all default false) plus DISPATCH_REALERT_EVERY
and LLM_MODEL, so production behaviour is set in .env rather than by code
defaults. .env.example updated to match.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin
2026-09-22 16:43:51 +05:30
008df49605 fix(dispatch): bucket failures by geo cell while the backend omits zone_id
Every booking.assignment_failed observed in production carries no zone_id, so
the per-zone failure counter and the once-a-day alert limit collapsed into a
single global "unknown" bucket — one alert per day for the whole country.

zone_key() prefers the backend's zone_id and falls back to a coarse lat/lon
grid cell (1 dp, ~11 km, DISPATCH_GEO_BUCKET_DP). The failure reason is now
logged too.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin
2026-09-22 16:42:15 +05:30
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
21 changed files with 1495 additions and 133 deletions

View File

@@ -1,5 +1,7 @@
GO_API_BASE_URL=
NATS_URL=
NATS_HOST=
NATS_PORT=4222
REDIS_HOST=
REDIS_PORT=
REDIS_PASSWORD=
@@ -14,3 +16,10 @@ NATS_PASSWORD=
ANTHROPIC_API_KEY=
LLM_MODEL=claude-opus-4-8
LOG_LEVEL=INFO
# Autonomy gates — outward-facing/irreversible actions are off unless enabled
DISPATCH_AGENT_AUTONOMOUS=false
EXPRESS_AGENT_AUTONOMOUS=false
EXCEPTION_AGENT_AUTONOMOUS=false
# Re-alert cadence during a failure burst in one zone
DISPATCH_REALERT_EVERY=100

View File

@@ -1,4 +1,5 @@
"""Customer Agent - Handles notifications, tracking, and customer communication."""
import re
import uuid
from datetime import datetime, timedelta
from typing import Dict, List, Any, Optional
@@ -11,8 +12,12 @@ from core.agent import SpecializedAgent
from core.types import AgentTask, MessageType, OrderStatus
from core.logger import logger
from core.http_client import api_post, api_get
from core.registry import registry
from config.system_config import GO_API_BASE_URL, INTERNAL_API_KEY
# Any {placeholder} left after template substitution — a caller forgot a variable.
_UNFILLED_PLACEHOLDER = re.compile(r"\{([a-zA-Z_][a-zA-Z0-9_]*)\}")
class NotificationChannel(str, Enum):
SMS = "sms"
@@ -64,6 +69,32 @@ class CustomerAgent(SpecializedAgent):
self._notifications: Dict[str, List[Notification]] = {}
self._templates = self._init_templates()
self._preferences: Dict[str, Dict] = {}
self._received_events: List[Dict[str, Any]] = []
# Event notifications other agents send without a task_type. They used to
# fall through to the base handle_message and be dropped at DEBUG level.
# They come from ExceptionAgent's generic exception tasks, which run over
# in-memory, simulated records (their order ids are not real bookings), so
# they are RECORDED here and logged — deliberately not turned into a real
# customer message through the backend's /internal/notify.
_RECORDED_EVENTS = {MessageType.ORDER_CANCELLED, MessageType.NOTIFICATION_SENT}
async def handle_message(self, message) -> None:
if message.message_type not in self._RECORDED_EVENTS:
await super().handle_message(message)
return
payload = message.payload if isinstance(message.payload, dict) else {}
event = {
"type": message.message_type.value,
"from": message.sender,
"order_id": payload.get("order_id"),
"received_at": datetime.now().isoformat(),
}
self._received_events = (self._received_events + [event])[-200:]
logger.info(
f"Customer Agent: recorded {event['type']} from {event['from']} for order {event['order_id']} "
"(simulated exception flow; no customer message sent)"
)
def _init_templates(self) -> Dict[str, Dict]:
return {
@@ -186,6 +217,9 @@ class CustomerAgent(SpecializedAgent):
return {"status": "sent", "order_id": order_id, "notifications": notifications_sent}
async def _send_notification(self, task: AgentTask) -> Dict[str, Any]:
if not registry.skill_enabled("customer_notifications"):
logger.info("Customer Agent: skill customer_notifications is disabled in the agent registry; not sending")
return {"status": "skipped", "reason": "customer_notifications disabled in the agent registry"}
order_id = task.data.get("order_id")
notification_type = task.data.get("notification_type")
template_vars = task.data.get("template_vars", {})
@@ -330,6 +364,30 @@ class CustomerAgent(SpecializedAgent):
"action": "provide_options",
}
def _failed_notification(self, order_id, customer_id, channel, notification_type,
reason: str, message: str) -> Dict[str, Any]:
"""A not-sent result in the same shape as a sent one, recorded for audit."""
notification_id = f"NOTIF-{uuid.uuid4().hex[:8].upper()}"
key = customer_id or "unknown"
self._notifications.setdefault(key, []).append(Notification(
notification_id=notification_id,
order_id=order_id,
customer_id=key,
channel=channel,
notification_type=notification_type,
message=message,
sent_at=datetime.now(),
delivered_at=None,
status="failed",
))
return {
"notification_id": notification_id,
"channel": channel.value,
"status": "failed",
"reason": reason,
"message": message,
}
async def _send_via_channel(
self,
order_id: str,
@@ -338,24 +396,51 @@ class CustomerAgent(SpecializedAgent):
notification_type: NotificationType,
template_vars: Dict[str, Any],
) -> Dict[str, Any]:
"""Render template and call Go backend (POST /api/v1/internal/notify)."""
"""Render template and call Go backend (POST /api/v1/internal/notify).
Refuses to send a message that still contains an unsubstituted
{placeholder}: a caller that omits a template variable would otherwise
put the literal text in front of a customer. Silence is better than a
broken message, and the ERROR names the caller's missing variables."""
template = self._templates.get(notification_type, {}).get(channel.value, "")
if not template:
logger.error(
f"Customer Agent: no {channel.value} template for {notification_type.value}; not sending"
)
return self._failed_notification(order_id, customer_id, channel, notification_type,
"no template", "")
message = template
for key, value in template_vars.items():
message = message.replace(f"{{{key}}}", str(value))
missing = _UNFILLED_PLACEHOLDER.findall(message)
if missing:
logger.error(
f"Customer Agent: refusing to send {notification_type.value} for order {order_id} — "
f"template variables not supplied: {missing}"
)
return self._failed_notification(order_id, customer_id, channel, notification_type,
f"unfilled template variables: {missing}", message)
notification_id = f"NOTIF-{uuid.uuid4().hex[:8].upper()}"
payload = {
"notification_id": notification_id,
"order_id": order_id,
"channel": channel.value,
"notification_type": notification_type.value,
"message": message,
}
# Omit customer_id rather than sending a placeholder: callers that only
# know the booking (e.g. DispatchAgent on an assignment failure) leave the
# backend to resolve the recipient from order_id, which it can already do.
if customer_id and str(customer_id) != "unknown":
payload["customer_id"] = customer_id
result = await api_post(
f"{GO_API_BASE_URL}/api/v1/internal/notify",
json={
"notification_id": notification_id,
"order_id": order_id,
"customer_id": customer_id,
"channel": channel.value,
"notification_type": notification_type.value,
"message": message,
},
json=payload,
headers={"X-Internal-Key": INTERNAL_API_KEY},
timeout=aiohttp.ClientTimeout(total=10),
)
@@ -377,11 +462,13 @@ class CustomerAgent(SpecializedAgent):
delivered_at=None,
status=status,
)
if customer_id not in self._notifications:
self._notifications[customer_id] = []
self._notifications[customer_id].append(notification)
history_key = customer_id or "unknown"
self._notifications.setdefault(history_key, []).append(notification)
logger.info(f"[{channel.value.upper()}] {notification_type.value} -> {customer_id[:20]} [{status}]")
logger.info(
f"[{channel.value.upper()}] {notification_type.value} -> "
f"{str(history_key)[:20]} [{status}]"
)
return {
"notification_id": notification_id,

View File

@@ -12,7 +12,9 @@ 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 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,
@@ -27,6 +29,47 @@ 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"
# 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):
"""
@@ -126,34 +169,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 #
# ------------------------------------------------------------------ #
@@ -177,18 +266,12 @@ class DispatchAgent(SpecializedAgent):
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),
)
# 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}")
@@ -208,16 +291,49 @@ class DispatchAgent(SpecializedAgent):
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")
# 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}")
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}")
decision = await decide_assignment_failure(build_assignment_failure_context(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")
@@ -239,7 +355,7 @@ class DispatchAgent(SpecializedAgent):
pass # transient — the backend will retry
elif decision.action == "notify_customer":
if DISPATCH_AGENT_AUTONOMOUS:
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.
@@ -272,22 +388,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 +428,61 @@ 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). Every variable the template names must be supplied or
# the send is refused; we have no ETA at assignment-failure time,
# so say so rather than inventing one. No customer_id: the event
# carries none, and the backend resolves the recipient from order_id.
"task_type": "send_notification",
"order_id": booking_id,
"notification_type": "delayed",
"template_vars": {"order_id": booking_id, "new_eta": "being confirmed"},
},
correlation_id=str(booking_id),
)
@@ -328,18 +494,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

@@ -26,7 +26,9 @@ from core.agent import SpecializedAgent
from core.types import AgentTask, MessageType
from core.logger import logger
from core.http_client import api_post
from core.llm import decide_stall_response, build_stall_context, StallDecision
from core.llm import decide_stall_response, build_stall_context, StallDecision, LLM_MODEL
from core.registry import registry
from core.decisions import record_decision
from config.system_config import (
GO_API_BASE_URL, INTERNAL_API_KEY,
DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD,
@@ -34,7 +36,12 @@ from config.system_config import (
REDIS_HOST, REDIS_PORT, REDIS_PASSWORD,
)
STALL_MINUTES = 10
STALL_MINUTES = 10 # default; the agent registry's stall_response.stallMinutes overrides it
def stall_minutes() -> float:
"""Minutes without movement before a rider counts as stalled (registry, else 10)."""
return registry.threshold("stall_response", "stallMinutes", STALL_MINUTES)
ACTIVE_STATUSES = ["Miler_Assigned", "Pickup_Scheduled"]
TRACKING_STREAM = "TRACKING"
@@ -204,10 +211,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 +231,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 ───────────────────────────────────
@@ -260,7 +267,7 @@ class ExceptionAgent(SpecializedAgent):
unchanged_since = _parse_ts(prev.get("position_unchanged_since")) or _utcnow()
minutes_stalled = (_utcnow() - unchanged_since).total_seconds() / 60
if minutes_stalled >= STALL_MINUTES:
if minutes_stalled >= stall_minutes():
booking = await self._get_active_booking(miler_id)
if booking:
await self._publish_stall(miler_id, booking["booking_id"], minutes_stalled)
@@ -311,7 +318,7 @@ class ExceptionAgent(SpecializedAgent):
continue
minutes_stale = (now - updated_at).total_seconds() / 60
if minutes_stale >= STALL_MINUTES:
if minutes_stale >= stall_minutes():
logger.warning(f"StallDetector: miler {miler_id} stale {minutes_stale:.1f} min (booking {booking_id})")
await self._publish_stall(miler_id, booking_id, minutes_stale)
@@ -385,15 +392,24 @@ class ExceptionAgent(SpecializedAgent):
Claude chooses one of wait / notify_only / reassign / escalate. The only
irreversible, customer-visible action (reassign) is gated behind an
autonomy flag and a confidence threshold; otherwise it is escalated to a
human via JARVIS. If the LLM is unavailable, we fall back to the previous
deterministic behaviour (reassign + notify) so detection never silently
stops acting.
human via JARVIS. If the LLM is unavailable: an autonomous agent falls
back to reassign + notify so detection never silently stops acting; a
non-autonomous one escalates to a human instead (it used to reassign
regardless, which made "autonomy off" untrue during an LLM outage).
Settings come from the agent registry (core/registry.py), falling back
to env: skill `stall_response` on/off, its `stallMinutes` and
`reassignConfidence`, and this agent's autonomy and model.
"""
try:
data = json.loads(msg.data.decode())
except Exception:
return
if not registry.skill_enabled("stall_response"):
logger.info("Stall received but skill stall_response is disabled in the agent registry; not acting")
return
miler_id = data.get("miler_id", "")
booking_id = data.get("booking_id", "")
minutes_stalled = data.get("minutes_stalled", 0)
@@ -409,21 +425,37 @@ class ExceptionAgent(SpecializedAgent):
logger.info(f"[EXCEPTION] booking={booking_id} context facts: {facts}")
context = build_stall_context(
minutes_stalled, stall_threshold_min=STALL_MINUTES, now=_utcnow(), facts=facts,
minutes_stalled, stall_threshold_min=stall_minutes(), now=_utcnow(), facts=facts,
)
decision = await decide_stall_response(context)
model = registry.model("EXCEPTION_AGENT")
decision = await decide_stall_response(context, model)
autonomous = registry.autonomous("EXCEPTION_AGENT", AUTONOMOUS_REASSIGN)
if decision is None:
logger.warning(f"No LLM decision for booking {booking_id}; falling back to reassign + notify")
await self._reassign(booking_id, "miler_stalled")
await self._notify_customer(booking_id, "We detected a delay, finding you a new miler")
self._record_stall_exception(
miler_id, booking_id, minutes_stalled,
resolution="Fallback (LLM unavailable): reassignment triggered + customer notified",
actions=["reassign", "notify_customer"],
)
if autonomous:
logger.warning(f"No LLM decision for booking {booking_id}; autonomous — falling back to reassign + notify")
await self._reassign(booking_id, "miler_stalled")
await self._notify_customer(booking_id, "We detected a delay, finding you a new miler")
self._record_stall_exception(
miler_id, booking_id, minutes_stalled,
resolution="Fallback (LLM unavailable): reassignment triggered + customer notified",
actions=["reassign", "notify_customer"],
)
else:
logger.warning(f"No LLM decision for booking {booking_id}; not autonomous — escalating to a human")
await self._escalate_to_human(miler_id, booking_id, minutes_stalled, StallDecision(
action="escalate", reasoning="LLM unavailable; autonomy is off, so a human decides.", confidence=0.0,
))
self._record_stall_exception(
miler_id, booking_id, minutes_stalled,
resolution="Fallback (LLM unavailable, autonomy off): escalated to a human",
actions=["escalate"],
)
return
record_decision("stall_response", booking_id, facts, decision.action, decision.confidence,
decision.reasoning, model or LLM_MODEL)
logger.info(
f"[EXCEPTION] booking={booking_id} decision={decision.action} "
f"confidence={decision.confidence:.2f} reasoning={decision.reasoning!r}"
@@ -441,7 +473,8 @@ class ExceptionAgent(SpecializedAgent):
actions.append("notify_customer")
elif decision.action == "reassign":
if AUTONOMOUS_REASSIGN and decision.confidence >= REASSIGN_MIN_CONFIDENCE:
min_confidence = registry.threshold("stall_response", "reassignConfidence", REASSIGN_MIN_CONFIDENCE)
if autonomous and decision.confidence >= min_confidence:
await self._reassign(booking_id, "miler_stalled")
await self._notify_customer(booking_id, "We detected a delay, finding you a new miler")
actions += ["reassign", "notify_customer"]

View File

@@ -36,6 +36,7 @@ from core.agent import SpecializedAgent
from core.types import AgentTask
from core.logger import logger
from core.http_client import api_get, api_post
from core.registry import registry
from config.system_config import (
NATS_HOST, NATS_PORT, NATS_USER, NATS_PASSWORD,
GO_API_BASE_URL, INTERNAL_API_KEY, ROUTE_OPTIMIZER_URL,
@@ -162,7 +163,9 @@ class ExpressDispatchAgent(SpecializedAgent):
logger.info(
f"[EXPRESS] dispatch received — tenant={tenant_id} bookings={len(booking_ids)}"
)
if tenant_id and booking_ids:
if not registry.skill_enabled("express_batch_dispatch"):
logger.info("[EXPRESS] skill express_batch_dispatch is disabled in the agent registry; batch left for the console")
elif tenant_id and booking_ids:
await self._handle_batch(tenant_id, booking_ids)
except Exception as e:
logger.error(f"ExpressDispatchAgent batch handler error: {e}")
@@ -212,10 +215,10 @@ class ExpressDispatchAgent(SpecializedAgent):
return
# 3. Write back (or, in observation mode, only log the plan).
if not EXPRESS_AGENT_AUTONOMOUS:
if not registry.autonomous("EXPRESS_DISPATCH_AGENT", EXPRESS_AGENT_AUTONOMOUS):
logger.info(
f"[EXPRESS] tenant={tenant_id} OBSERVE-ONLY plan "
f"(EXPRESS_AGENT_AUTONOMOUS=false): {json.dumps(assignments)}"
f"(autonomy off): {json.dumps(assignments)}"
)
return
@@ -247,6 +250,10 @@ class ExpressDispatchAgent(SpecializedAgent):
part (ordering); this only decides who."""
load: Dict[int, int] = {r["miler_user_id"]: 0 for r in riders}
rider_stops: Dict[int, List[Dict]] = {}
# Tuned in the agent registry (skill express_batch_dispatch), else env.
max_per_rider = int(registry.threshold("express_batch_dispatch", "maxPerRider", MAX_PER_RIDER))
max_radius_km = registry.threshold("express_batch_dispatch", "maxRadiusKm", MAX_RADIUS_KM)
load_penalty_km = registry.threshold("express_batch_dispatch", "loadPenaltyKm", LOAD_PENALTY_KM)
# Assign larger-pickup-cluster bookings first is unnecessary; simple
# stable order keeps it predictable and testable.
@@ -255,13 +262,13 @@ class ExpressDispatchAgent(SpecializedAgent):
best_rider = None
best_score = None
for r in riders:
if load[r["miler_user_id"]] >= MAX_PER_RIDER:
if load[r["miler_user_id"]] >= max_per_rider:
continue
rlat, rlon = _rider_location(r, plat, plon)
dist = _haversine_km(rlat, rlon, plat, plon)
if dist > MAX_RADIUS_KM:
if dist > max_radius_km:
continue
score = dist + LOAD_PENALTY_KM * load[r["miler_user_id"]]
score = dist + load_penalty_km * load[r["miler_user_id"]]
if best_score is None or score < best_score:
best_score = score
best_rider = r

View File

@@ -89,6 +89,10 @@ class FleetAgent(SpecializedAgent):
handlers = {
"assign_vehicle": self._assign_vehicle,
"release_vehicle": self._release_vehicle,
# Sent by ExceptionAgent._cancel_order, which knows the order, not
# the vehicle. Had no handler, so every cancellation hit
# _unknown_task and the vehicle stayed marked in use.
"release_vehicle_for_cancel": self._release_vehicle_for_order,
"get_availability": self._get_availability,
"track_vehicle": self._track_vehicle,
"update_location": self._update_location,
@@ -181,6 +185,20 @@ class FleetAgent(SpecializedAgent):
return {"status": "released", "vehicle_id": vehicle_id, "hub": vehicle["hub"]}
async def _release_vehicle_for_order(self, task: AgentTask) -> Dict[str, Any]:
"""Release the vehicle carrying a cancelled order, found by order id."""
order_id = task.data.get("order_id")
for assignment in self._assignments.values():
if assignment.status == "active" and order_id in (assignment.order_ids or []):
return await self._release_vehicle(AgentTask(
task_id=f"{task.task_id}:release",
agent_type=task.agent_type,
task_type="release_vehicle",
data={"vehicle_id": assignment.vehicle_id},
))
return {"status": "no_assignment", "order_id": order_id,
"message": f"No active vehicle assignment carries order {order_id}"}
async def _get_availability(self, task: AgentTask) -> Dict[str, Any]:
hub_id = task.data.get("hub_id")
vehicle_type = task.data.get("vehicle_type")

View File

@@ -8,16 +8,18 @@ from core.types import (
AgentTask, MessageType, Priority, OrderStatus, ZoneType
)
from core.logger import logger
from core.http_client import api_get, api_post, api_patch
from config.system_config import GO_API_BASE_URL
class OrderAgent(SpecializedAgent):
"""
Order Agent - Manages the entire order lifecycle from intake to validation.
Writes go to the Doormile Go backend (POST /api/v1/admin/crmbooking).
Reads and status updates also go through the Go API.
NOT CONNECTED to the Doormile backend (agent registry status: broken).
It was written against /api/v1/admin/crmbooking, which was renamed to
/admin/expressbooking, and every /admin/* route needs a console login's
JWT — this engine holds only the internal key, which /admin/* refuses. So
the helpers below refuse up front instead of sending requests that can only
fail. Wiring it needs an /internal/* booking route on the backend first.
"""
def __init__(self):
@@ -70,23 +72,22 @@ class OrderAgent(SpecializedAgent):
# Go API helpers (use shared session + retry from http_client) #
# ------------------------------------------------------------------ #
# Returns None (what every caller already treats as "no backend result").
NO_BACKEND_ROUTE = ("OrderAgent has no working backend route: /admin/* needs a console JWT "
"and /admin/crmbooking no longer exists")
def _refuse(self, method: str, path: str) -> None:
logger.warning(f"OrderAgent: not calling {method} {path} — {self.NO_BACKEND_ROUTE}")
return None
async def _api_post(self, path: str, payload: Dict) -> Optional[Dict]:
result = await api_post(f"{GO_API_BASE_URL}{path}", json=payload)
if result is None:
logger.error(f"Go API POST {path} returned no response")
return result
return self._refuse("POST", path)
async def _api_get(self, path: str, params: Dict = None) -> Optional[Dict]:
result = await api_get(f"{GO_API_BASE_URL}{path}", params=params)
if result is None:
logger.error(f"Go API GET {path} returned no response")
return result
return self._refuse("GET", path)
async def _api_patch(self, path: str, payload: Dict) -> Optional[Dict]:
result = await api_patch(f"{GO_API_BASE_URL}{path}", json=payload)
if result is None:
logger.error(f"Go API PATCH {path} returned no response")
return result
return self._refuse("PATCH", path)
# ------------------------------------------------------------------ #
# Task handlers #

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")
@@ -255,7 +280,9 @@ class MasterAgent(Agent):
task_id=f"{order_id}_validate",
agent_type="order",
task_type="validate_order",
data={"order": order_data}
# OrderAgent._validate_order reads `order_id`; sending only
# {"order": ...} made every orchestrated order "not found".
data={"order_id": order_id, "order": order_data}
))
customer_agent = self._sub_agents.get("CUSTOMER_AGENT")

66
core/decisions.py Normal file
View File

@@ -0,0 +1,66 @@
"""
Records the engine's model decisions in the backend's decision log.
Until Phase 5 of the plan (krow_talent_app/docs/agent-platform-plan.md) the
stall and assignment-failure decisions existed only as log lines, so the
console's Insights tab could not show them. Each decision is now sent to
``POST /api/v1/internal/agent-decisions`` — the same log routemate's
assignment decisions go to.
Fire-and-forget: the post runs as a background task, so a slow or down
backend never delays an agent's reaction to a stall. A failure is logged by
the HTTP client and otherwise ignored — the decision was still made and acted
on; only its record is missing.
"""
import asyncio
from typing import Any, Dict, Optional, Set
from config.system_config import GO_API_BASE_URL, INTERNAL_API_KEY
from core.http_client import api_post
_pending: Set[asyncio.Task] = set() # keep references so tasks are not garbage-collected
def _as_booking_id(value: Any) -> Optional[int]:
"""The backend's booking_id is an unsigned integer; anything else is omitted."""
try:
n = int(value)
except (TypeError, ValueError):
return None
return n if n > 0 else None
def build_payload(decision_type: str, booking_id: Any, facts: Optional[Dict[str, Any]],
action: str, confidence: float, reasoning: str, model: str) -> Dict[str, Any]:
"""The request body. Pure, so its shape is tested without a network."""
return {
"decision_type": decision_type,
"booking_id": _as_booking_id(booking_id),
"context": {"facts": facts or {}, "model": model},
"decision": {"action": action, "confidence": round(float(confidence), 3)},
"reasoning": reasoning or "",
}
async def _post(payload: Dict[str, Any]) -> None:
await api_post(
f"{GO_API_BASE_URL}/api/v1/internal/agent-decisions",
json=payload,
headers={"X-Internal-Key": INTERNAL_API_KEY},
)
def record_decision(decision_type: str, booking_id: Any, facts: Optional[Dict[str, Any]],
action: str, confidence: float, reasoning: str, model: str) -> Optional[asyncio.Task]:
"""Schedule the record and return at once. Returns the task (for tests), or
None when there is no running loop or no key to authenticate with."""
if not INTERNAL_API_KEY:
return None
try:
loop = asyncio.get_running_loop()
except RuntimeError:
return None
task = loop.create_task(_post(build_payload(decision_type, booking_id, facts, action, confidence, reasoning, model)))
_pending.add(task)
task.add_done_callback(_pending.discard)
return task

View File

@@ -102,7 +102,25 @@ _STALL_SCHEMA = {
}
async def _decide(system: str, schema: dict, valid_actions, context: str):
def request_params(model: Optional[str] = None) -> dict:
"""The model-dependent part of a decision request.
`model` is a per-agent override from the agent registry; None means the
engine default (LLM_MODEL). Claude Haiku 4.5 rejects adaptive thinking and
`effort`, so both are left out for it — otherwise a Haiku pin set from the
console would make every decision fail and fall back to the heuristic.
"""
m = model or LLM_MODEL
output_config = {}
params = {"model": m}
if "haiku" not in m:
params["thinking"] = {"type": "adaptive"}
output_config["effort"] = LLM_EFFORT
params["output_config"] = output_config
return params
async def _decide(system: str, schema: dict, valid_actions, context: str, model: Optional[str] = None):
"""Run one structured-output decision against Claude.
Returns (action, reasoning, confidence), or None on any failure (network,
@@ -111,20 +129,17 @@ async def _decide(system: str, schema: dict, valid_actions, context: str):
"""
resp = None
last_err = None
params = request_params(model)
params["output_config"]["format"] = {"type": "json_schema", "schema": schema}
for attempt in range(2): # initial try + one retry on transient failure
try:
client = _get_client()
resp = await client.messages.create(
model=LLM_MODEL,
max_tokens=LLM_MAX_TOKENS,
system=system,
thinking={"type": "adaptive"},
output_config={
"effort": LLM_EFFORT,
"format": {"type": "json_schema", "schema": schema},
},
messages=[{"role": "user", "content": context}],
timeout=LLM_TIMEOUT_S,
**params,
)
break
except Exception as e:
@@ -154,10 +169,10 @@ async def _decide(system: str, schema: dict, valid_actions, context: str):
return None
async def decide_stall_response(context: str) -> Optional[StallDecision]:
async def decide_stall_response(context: str, model: Optional[str] = None) -> Optional[StallDecision]:
"""Ask Claude how to handle a stalled miler. None on failure (caller falls
back to deterministic behaviour)."""
r = await _decide(_STALL_SYSTEM, _STALL_SCHEMA, _VALID_ACTIONS, context)
back to deterministic behaviour). `model` overrides LLM_MODEL."""
r = await _decide(_STALL_SYSTEM, _STALL_SCHEMA, _VALID_ACTIONS, context, model)
return StallDecision(*r) if r else None
@@ -181,10 +196,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",
@@ -213,8 +239,8 @@ def build_assignment_failure_context(facts, *, now=None) -> str:
return "\n".join(lines)
async def decide_assignment_failure(context: str) -> Optional[AssignmentDecision]:
async def decide_assignment_failure(context: str, model: Optional[str] = None) -> Optional[AssignmentDecision]:
"""Ask Claude how to react to a failed assignment. None on failure (caller
falls back to deterministic behaviour)."""
r = await _decide(_ASSIGNMENT_SYSTEM, _ASSIGNMENT_SCHEMA, _ASSIGNMENT_ACTIONS, context)
falls back to deterministic behaviour). `model` overrides LLM_MODEL."""
r = await _decide(_ASSIGNMENT_SYSTEM, _ASSIGNMENT_SCHEMA, _ASSIGNMENT_ACTIONS, context, model)
return AssignmentDecision(*r) if r else None

122
core/registry.py Normal file
View File

@@ -0,0 +1,122 @@
"""
The agent registry, as this engine sees it.
Settings an operator changes in the console (Settings → Skills & Tools) are
stored in the Doormile backend's agent registry. This module polls
``GET /api/v1/internal/ai/registry`` and answers, at the moment an agent needs
it, three questions:
registry.autonomous(agent_id, env_default) may this agent act on its own?
registry.model(agent_id) which Claude model, if pinned?
registry.skill_enabled(skill_id) is this behaviour switched on?
registry.threshold(skill_id, key, default) the tuned value of a knob
Precedence — deliberate, and the same for every setting:
registry loaded → the registry's value (it is the operator's decision)
not yet / never → the env default this engine always had
So a backend that is down, unreachable, or not yet deployed leaves the engine
exactly as it behaved before the registry existed. Once the registry has been
read, the last good copy is kept through later failures: a flapping backend
does not flip autonomy back to an env default mid-shift.
Polling uses the ETag the backend sends; an unchanged registry costs a 304
with no body. Before Phase 5 of the plan (krow_talent_app/docs/
agent-platform-plan.md) every one of these settings was a module constant read
from the environment at import time, so a change needed a redeploy.
"""
import asyncio
import os
from typing import Any, Dict, Optional
import aiohttp
from core.logger import logger
POLL_SECONDS = float(os.getenv("REGISTRY_POLL_SECONDS", "30"))
class RegistryClient:
def __init__(self):
self._agents: Dict[str, Dict[str, Any]] = {}
self._skills: Dict[str, Dict[str, Any]] = {}
self._etag: Optional[str] = None
self.loaded = False
# ── Reading (sync, cheap, safe to call on every event) ──────────────────
def autonomous(self, agent_id: str, env_default: bool) -> bool:
if not self.loaded or agent_id not in self._agents:
return env_default
return bool(self._agents[agent_id].get("autonomous", False))
def model(self, agent_id: str) -> Optional[str]:
"""A pinned model id, or None for the engine's own LLM_MODEL."""
if not self.loaded:
return None
return self._agents.get(agent_id, {}).get("model") or None
def skill_enabled(self, skill_id: str, default: bool = True) -> bool:
if not self.loaded or skill_id not in self._skills:
return default
return bool(self._skills[skill_id].get("enabled", default))
def threshold(self, skill_id: str, key: str, default: float) -> float:
if not self.loaded:
return default
value = (self._skills.get(skill_id, {}).get("thresholds") or {}).get(key)
return value if isinstance(value, (int, float)) and not isinstance(value, bool) else default
# ── Loading ──────────────────────────────────────────────────────────────
def apply(self, snapshot: Dict[str, Any]) -> None:
"""Adopt a registry snapshot (the `data` of /internal/ai/registry)."""
agents = {a["agentid"]: a for a in snapshot.get("agents") or [] if a.get("agentid")}
skills = {s["skillid"]: s for s in snapshot.get("skills") or [] if s.get("skillid")}
self._agents, self._skills = agents, skills
if not self.loaded:
logger.info(f"Agent registry loaded: {len(agents)} agents, {len(skills)} skills")
self.loaded = True
async def fetch_once(self, session: aiohttp.ClientSession, base_url: str, api_key: str) -> str:
"""One poll. Returns 'updated', 'unchanged', 'unavailable' or 'refused'."""
headers = {"X-Internal-Key": api_key}
if self._etag:
headers["If-None-Match"] = self._etag
try:
async with session.get(f"{base_url}/api/v1/internal/ai/registry", headers=headers,
timeout=aiohttp.ClientTimeout(total=10)) as resp:
if resp.status == 304:
return "unchanged"
if resp.status in (401, 403):
return "refused"
if resp.status != 200:
return "unavailable"
body = await resp.json()
self.apply(body.get("data") or {})
self._etag = resp.headers.get("ETag")
return "updated"
except Exception as exc: # network, timeout, bad JSON — keep the last good copy
logger.debug(f"Agent registry poll failed: {exc}")
return "unavailable"
async def run(self, base_url: str, api_key: str, poll_seconds: float = POLL_SECONDS) -> None:
"""Poll forever. Never raises; the engine runs on env defaults meanwhile."""
if not api_key:
logger.warning("INTERNAL_API_KEY unset — agent registry not read; running on env defaults")
return
last = None
async with aiohttp.ClientSession() as session:
while True:
outcome = await self.fetch_once(session, base_url, api_key)
if outcome != last and outcome in ("refused", "unavailable"):
logger.warning(
f"Agent registry {outcome} at {base_url}; "
+ ("keeping the last loaded copy" if self.loaded else "running on env defaults")
)
last = outcome
await asyncio.sleep(poll_seconds)
registry = RegistryClient()

View File

@@ -16,6 +16,9 @@ class MessageType(str, Enum):
ORDER_DELIVERED = "ORDER_DELIVERED"
ORDER_CANCELLED = "ORDER_CANCELLED"
ORDER_DELAYED = "ORDER_DELAYED"
# OrderAgent._update_status publishes this; it was missing from the enum, so
# every update_status task raised AttributeError.
ORDER_STATUS_UPDATE = "ORDER_STATUS_UPDATE"
HUB_STATUS_UPDATE = "HUB_STATUS_UPDATE"
VEHICLE_ASSIGNED = "VEHICLE_ASSIGNED"
ROUTE_OPTIMIZED = "ROUTE_OPTIMIZED"

View File

@@ -13,8 +13,12 @@ services:
container_name: doormile_agents
restart: unless-stopped
environment:
# NATS — event bus
# NATS — event bus. The agents connect with NATS_HOST/NATS_PORT (NATS_URL
# is informational), so both must be passed through or every connection
# falls back to localhost.
- NATS_URL=${NATS_URL}
- NATS_HOST=${NATS_HOST}
- NATS_PORT=${NATS_PORT}
- NATS_USER=${NATS_USER}
- NATS_PASSWORD=${NATS_PASSWORD}
# Redis — GEO queries, miler presence
@@ -33,6 +37,12 @@ services:
# AI layer — decision engine
- AI_LAYER_BASE_URL=${AI_LAYER_BASE_URL}
- ANTHROPIC_API_KEY=${ANTHROPIC_API_KEY}
# Agent behaviour. Outward-facing actions stay gated unless explicitly enabled.
- DISPATCH_AGENT_AUTONOMOUS=${DISPATCH_AGENT_AUTONOMOUS:-false}
- EXPRESS_AGENT_AUTONOMOUS=${EXPRESS_AGENT_AUTONOMOUS:-false}
- EXCEPTION_AGENT_AUTONOMOUS=${EXCEPTION_AGENT_AUTONOMOUS:-false}
- DISPATCH_REALERT_EVERY=${DISPATCH_REALERT_EVERY:-100}
- LLM_MODEL=${LLM_MODEL:-claude-opus-4-8}
# Logging
- LOG_LEVEL=${LOG_LEVEL:-INFO}
volumes:

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

@@ -28,6 +28,8 @@ from core.logger import logger
from core.agent import MasterAgent, Agent
from core.message_bus import message_bus
from core.http_client import close_session
from core.registry import registry
from config.system_config import GO_API_BASE_URL, INTERNAL_API_KEY
from core.types import AgentTask, MessageType, Priority
from agents.order_agent import OrderAgent
@@ -125,6 +127,11 @@ async def production_mode():
system = LogiFlowAI()
system.initialize_agents()
# Operator settings from the backend's agent registry (autonomy, model,
# skills on/off, thresholds). Polls with ETag; until it answers, and if it
# never does, every agent runs on its env defaults exactly as before.
registry_task = asyncio.create_task(registry.run(GO_API_BASE_URL, INTERNAL_API_KEY)) # noqa: F841 — cancelled at shutdown
logger.info("Infrastructure connected. Press Ctrl+C to stop.")
loop = asyncio.get_running_loop()

6
requirements-dev.txt Normal file
View File

@@ -0,0 +1,6 @@
# Test-only dependencies. Not installed in the production image (the Dockerfile
# installs requirements.txt only). Install locally with:
# pip install -r requirements-dev.txt
-r requirements.txt
pytest>=8.0.0
pytest-asyncio>=0.23.0

View File

@@ -0,0 +1,117 @@
"""CUSTOMER_AGENT must never put an unsubstituted {placeholder} in front of a
customer, and must not send a literal "unknown" customer id.
The DispatchAgent assignment-failure path knows only a booking id, so this is
the realistic caller: it supplies order_id and new_eta and no customer.
"""
import pytest
from agents.customer_agent import (
CustomerAgent, NotificationChannel, NotificationType, _UNFILLED_PLACEHOLDER,
)
def make_agent(captured):
agent = CustomerAgent.__new__(CustomerAgent)
agent.agent_id = "CUSTOMER_AGENT"
agent._notifications = {}
agent._templates = CustomerAgent._init_templates(agent)
async def fake_post(url, json=None, headers=None, timeout=None):
captured.append(json)
return {"ok": True}
import agents.customer_agent as mod
mod.api_post = fake_post
return agent
@pytest.fixture(autouse=True)
def restore_api_post():
import agents.customer_agent as mod
original = mod.api_post
yield
mod.api_post = original
@pytest.mark.asyncio
async def test_missing_template_var_is_refused_not_sent():
sent = []
agent = make_agent(sent)
result = await agent._send_via_channel(
order_id="36", customer_id="unknown",
channel=NotificationChannel.SMS,
notification_type=NotificationType.DELAYED,
template_vars={"order_id": "36"}, # new_eta omitted — the old bug
)
assert result["status"] == "failed"
assert "new_eta" in result["reason"]
assert sent == [] # nothing reached the backend
@pytest.mark.asyncio
async def test_all_vars_supplied_sends_clean_message():
sent = []
agent = make_agent(sent)
result = await agent._send_via_channel(
order_id="36", customer_id="unknown",
channel=NotificationChannel.SMS,
notification_type=NotificationType.DELAYED,
template_vars={"order_id": "36", "new_eta": "being confirmed"},
)
assert result["status"] == "sent"
assert len(sent) == 1
assert not _UNFILLED_PLACEHOLDER.findall(sent[0]["message"])
assert sent[0]["message"] == (
"Delay alert: Order 36 may arrive later than expected. New ETA: being confirmed"
)
@pytest.mark.asyncio
async def test_unknown_customer_id_is_omitted_from_payload():
sent = []
agent = make_agent(sent)
await agent._send_via_channel(
order_id="36", customer_id="unknown",
channel=NotificationChannel.SMS,
notification_type=NotificationType.DELAYED,
template_vars={"order_id": "36", "new_eta": "being confirmed"},
)
assert "customer_id" not in sent[0] # backend resolves from order_id
assert sent[0]["order_id"] == "36"
@pytest.mark.asyncio
async def test_real_customer_id_is_sent():
sent = []
agent = make_agent(sent)
await agent._send_via_channel(
order_id="36", customer_id="CUST-9",
channel=NotificationChannel.SMS,
notification_type=NotificationType.DELAYED,
template_vars={"order_id": "36", "new_eta": "being confirmed"},
)
assert sent[0]["customer_id"] == "CUST-9"
@pytest.mark.asyncio
async def test_dispatch_supplies_every_template_var():
"""The DispatchAgent payload must satisfy the DELAYED template exactly."""
from agents.dispatch_agent import DispatchAgent
sent_msgs = []
async def capture(recipient, message_type, payload, correlation_id=None):
sent_msgs.append(payload)
agent = DispatchAgent.__new__(DispatchAgent)
agent.send_message = capture
await agent._notify_customer_delay("36")
payload = sent_msgs[0]
assert payload["task_type"] == "send_notification"
template = CustomerAgent._init_templates(CustomerAgent.__new__(CustomerAgent))[
NotificationType.DELAYED]["sms"]
required = set(_UNFILLED_PLACEHOLDER.findall(template))
assert required <= set(payload["template_vars"]), (
f"template needs {required}, dispatch supplies {set(payload['template_vars'])}"
)

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,159 @@ 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, model=None):
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 TestZoneKey(unittest.IsolatedAsyncioTestCase):
"""The backend sends no zone_id today; without a fallback every failure in
the country shares one bucket and the daily alert limit becomes global."""
async def test_real_zone_id_wins(self):
from agents.dispatch_agent import zone_key
self.assertEqual(zone_key("hyderabad", 17.4, 78.4), "hyderabad")
async def test_missing_zone_falls_back_to_geo_cell(self):
from agents.dispatch_agent import zone_key
self.assertEqual(zone_key(None, 11.0182714, 76.9677744), "geo:11.0,77.0")
self.assertEqual(zone_key("unknown", 11.0182714, 76.9677744), "geo:11.0,77.0")
async def test_distant_pickups_get_different_buckets(self):
from agents.dispatch_agent import zone_key
self.assertNotEqual(zone_key(None, 11.01, 76.96), zone_key(None, 17.44, 78.39))
async def test_no_zone_and_no_coords_is_unknown(self):
from agents.dispatch_agent import zone_key
self.assertEqual(zone_key(None, None, None), "unknown")
self.assertEqual(zone_key("unknown", "bad", "coords"), "unknown")
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

View File

@@ -0,0 +1,413 @@
"""
Tests for Phase 5 of krow_talent_app/docs/agent-platform-plan.md: the engine
reads operator settings from the backend's agent registry, records its model
decisions, and the six agent-to-agent message bugs are fixed.
No network, no NATS, no Redis — every external edge is faked.
Run:
python -m unittest discover -s tests
"""
import json
import unittest
from datetime import datetime
from unittest.mock import patch
import core.decisions as decisions
import core.llm as llm
from core.message_bus import message_bus
from core.registry import RegistryClient, registry
from core.types import AgentMessage, AgentTask, MessageType
SNAPSHOT = {
"agents": [
{"agentid": "EXCEPTION_AGENT", "autonomous": True, "model": "claude-sonnet-5-5"},
{"agentid": "DISPATCH_AGENT", "autonomous": False, "model": ""},
],
"skills": [
{"skillid": "stall_response", "enabled": True,
"thresholds": {"stallMinutes": 15, "reassignConfidence": 0.9, "bogus": "x", "flag": True}},
{"skillid": "customer_notifications", "enabled": False},
],
}
def _task(task_type, **data):
return AgentTask(task_id="t", agent_type="x", task_type=task_type, data=data)
class _RegistryIsolation:
"""Save and restore the shared registry singleton the agents read."""
def _save_registry(self):
self._saved = (registry._agents, registry._skills, registry._etag, registry.loaded)
def _restore_registry(self):
registry._agents, registry._skills, registry._etag, registry.loaded = self._saved
# ── RegistryClient: precedence ───────────────────────────────────────────────
class TestRegistryPrecedence(unittest.TestCase):
def test_env_defaults_before_load(self):
r = RegistryClient()
self.assertTrue(r.autonomous("EXCEPTION_AGENT", True))
self.assertFalse(r.autonomous("EXCEPTION_AGENT", False))
self.assertIsNone(r.model("EXCEPTION_AGENT"))
self.assertTrue(r.skill_enabled("customer_notifications"))
self.assertFalse(r.skill_enabled("anything", default=False))
self.assertEqual(r.threshold("stall_response", "stallMinutes", 10), 10)
def test_registry_wins_once_loaded(self):
r = RegistryClient()
r.apply(SNAPSHOT)
self.assertTrue(r.loaded)
self.assertTrue(r.autonomous("EXCEPTION_AGENT", False))
self.assertFalse(r.autonomous("DISPATCH_AGENT", True))
self.assertEqual(r.model("EXCEPTION_AGENT"), "claude-sonnet-5-5")
self.assertIsNone(r.model("DISPATCH_AGENT")) # "" means engine default
self.assertFalse(r.skill_enabled("customer_notifications"))
self.assertEqual(r.threshold("stall_response", "stallMinutes", 10), 15)
self.assertEqual(r.threshold("stall_response", "reassignConfidence", 0.75), 0.9)
def test_unknown_ids_fall_back(self):
r = RegistryClient()
r.apply(SNAPSHOT)
self.assertTrue(r.autonomous("HUB_AGENT", True))
self.assertTrue(r.skill_enabled("not_in_registry"))
self.assertEqual(r.threshold("not_in_registry", "k", 3), 3)
def test_non_numeric_thresholds_ignored(self):
r = RegistryClient()
r.apply(SNAPSHOT)
self.assertEqual(r.threshold("stall_response", "bogus", 7), 7)
self.assertEqual(r.threshold("stall_response", "flag", 7), 7) # bool is not a number here
def test_rows_without_ids_skipped(self):
r = RegistryClient()
r.apply({"agents": [{"autonomous": True}], "skills": [{"enabled": False}]})
self.assertTrue(r.loaded)
self.assertEqual(r._agents, {})
self.assertEqual(r._skills, {})
# ── RegistryClient: polling ──────────────────────────────────────────────────
class FakeResp:
def __init__(self, status, body=None, etag=None, raise_on_json=False):
self.status = status
self._body = body
self.headers = {"ETag": etag} if etag else {}
self._raise = raise_on_json
async def json(self):
if self._raise:
raise ValueError("bad json")
return self._body
async def __aenter__(self):
return self
async def __aexit__(self, *exc):
return False
class FakeSession:
def __init__(self, *responses, exc=None):
self._responses = list(responses)
self._exc = exc
self.calls = []
def get(self, url, headers=None, timeout=None):
self.calls.append((url, dict(headers or {})))
if self._exc:
raise self._exc
return self._responses.pop(0)
class TestRegistryFetch(unittest.IsolatedAsyncioTestCase):
async def test_200_then_304_uses_etag(self):
r = RegistryClient()
s = FakeSession(FakeResp(200, {"data": SNAPSHOT}, etag='"v1"'), FakeResp(304))
self.assertEqual(await r.fetch_once(s, "http://api", "k"), "updated")
self.assertTrue(r.loaded)
self.assertEqual(await r.fetch_once(s, "http://api", "k"), "unchanged")
url, first = s.calls[0]
self.assertEqual(url, "http://api/api/v1/internal/ai/registry")
self.assertEqual(first["X-Internal-Key"], "k")
self.assertNotIn("If-None-Match", first)
self.assertEqual(s.calls[1][1]["If-None-Match"], '"v1"')
async def test_refused_and_unavailable_keep_last_copy(self):
r = RegistryClient()
r.apply(SNAPSHOT)
for resp, want in ((FakeResp(401), "refused"), (FakeResp(403), "refused"),
(FakeResp(500), "unavailable"), (FakeResp(200, raise_on_json=True), "unavailable")):
self.assertEqual(await r.fetch_once(FakeSession(resp), "http://api", "k"), want)
self.assertEqual(await r.fetch_once(FakeSession(exc=OSError("down")), "http://api", "k"), "unavailable")
# Still the operator's settings, not the env defaults.
self.assertTrue(r.autonomous("EXCEPTION_AGENT", False))
self.assertFalse(r.skill_enabled("customer_notifications"))
async def test_run_without_key_returns_at_once(self):
r = RegistryClient()
await r.run("http://api", "", poll_seconds=0) # must not loop forever
self.assertFalse(r.loaded)
# ── Decision log payload ─────────────────────────────────────────────────────
class TestDecisions(unittest.IsolatedAsyncioTestCase):
def test_build_payload_shape(self):
p = decisions.build_payload("stall_response", "42", {"eta": 5}, "reassign", 0.87654, "slow", "m")
self.assertEqual(p, {
"decision_type": "stall_response",
"booking_id": 42,
"context": {"facts": {"eta": 5}, "model": "m"},
"decision": {"action": "reassign", "confidence": 0.877},
"reasoning": "slow",
})
def test_bad_booking_ids_omitted(self):
for bad in (None, "", "abc", 0, -3):
self.assertIsNone(decisions.build_payload("t", bad, None, "a", 0, None, "m")["booking_id"])
self.assertEqual(decisions.build_payload("t", 7, None, "a", 0, None, "m")["context"]["facts"], {})
def test_no_running_loop_returns_none(self):
with patch.object(decisions, "INTERNAL_API_KEY", "k"):
self.assertIsNone(decisions.record_decision("t", 1, {}, "a", 0.5, "r", "m"))
async def test_no_key_returns_none(self):
with patch.object(decisions, "INTERNAL_API_KEY", ""):
self.assertIsNone(decisions.record_decision("t", 1, {}, "a", 0.5, "r", "m"))
async def test_posts_to_decision_log(self):
sent = []
async def fake_post(url, json=None, headers=None, **kw):
sent.append((url, json, headers))
with patch.object(decisions, "INTERNAL_API_KEY", "k"), \
patch.object(decisions, "GO_API_BASE_URL", "http://api"), \
patch.object(decisions, "api_post", fake_post):
task = decisions.record_decision("assignment_failure", 9, {}, "notify", 0.5, "r", "m")
self.assertIsNotNone(task)
await task
url, body, headers = sent[0]
self.assertEqual(url, "http://api/api/v1/internal/agent-decisions")
self.assertEqual(body["booking_id"], 9)
self.assertEqual(headers, {"X-Internal-Key": "k"})
# ── Model-dependent request params ──────────────────────────────────────────
class TestRequestParams(unittest.TestCase):
def test_default_model_has_thinking_and_effort(self):
p = llm.request_params()
self.assertEqual(p["model"], llm.LLM_MODEL)
self.assertEqual(p["thinking"], {"type": "adaptive"})
self.assertIn("effort", p["output_config"])
def test_override_model_used(self):
self.assertEqual(llm.request_params("claude-sonnet-5-5")["model"], "claude-sonnet-5-5")
def test_haiku_omits_thinking_and_effort(self):
p = llm.request_params("claude-haiku-4-5")
self.assertNotIn("thinking", p)
self.assertNotIn("effort", p["output_config"])
def test_calls_do_not_share_output_config(self):
a, b = llm.request_params(), llm.request_params()
a["output_config"]["format"] = {}
self.assertNotIn("format", b["output_config"])
# ── ExceptionAgent: registry gates and the LLM-down fallback ────────────────
class FakeStallMsg:
def __init__(self, payload):
self.data = json.dumps(payload).encode()
class TestExceptionAgentRegistry(_RegistryIsolation, unittest.IsolatedAsyncioTestCase):
def setUp(self):
from agents import exception_agent as mod
self.mod = mod
self._save_registry()
self._orig_decide = mod.decide_stall_response
self._orig_record = mod.record_decision
def tearDown(self):
self._restore_registry()
self.mod.decide_stall_response = self._orig_decide
self.mod.record_decision = self._orig_record
def _agent(self, decision=None):
from agents.exception_agent import ExceptionAgent
agent = ExceptionAgent.__new__(ExceptionAgent)
agent.acts = []
async def claim(_b):
return True
async def facts(_m, _b):
return {}
async def reassign(b, reason):
agent.acts.append("reassign")
return True
async def notify(b, m):
agent.acts.append("notify")
return True
async def escalate(m, b, mins, d):
agent.acts.append(("escalate", d.action))
agent._claim_stall_handled = claim
agent._gather_stall_facts = facts
agent._reassign = reassign
agent._notify_customer = notify
agent._escalate_to_human = escalate
agent._record_stall_exception = lambda *a, **k: None
self.models = []
self.recorded = []
async def fake_decide(context, model=None):
self.models.append(model)
return decision
self.mod.decide_stall_response = fake_decide
self.mod.record_decision = lambda *a: self.recorded.append(a)
return agent
def _load(self, autonomous, enabled=True, confidence=0.75):
registry.apply({
"agents": [{"agentid": "EXCEPTION_AGENT", "autonomous": autonomous, "model": "claude-sonnet-5-5"}],
"skills": [{"skillid": "stall_response", "enabled": enabled,
"thresholds": {"reassignConfidence": confidence}}],
})
async def _stall(self, agent):
await agent._on_miler_stalled(FakeStallMsg({"miler_id": "m1", "booking_id": "5", "minutes_stalled": 12}))
async def test_llm_down_not_autonomous_escalates(self):
self._load(autonomous=False)
agent = self._agent(decision=None)
await self._stall(agent)
self.assertEqual(agent.acts, [("escalate", "escalate")])
async def test_llm_down_autonomous_reassigns(self):
self._load(autonomous=True)
agent = self._agent(decision=None)
await self._stall(agent)
self.assertEqual(agent.acts, ["reassign", "notify"])
async def test_skill_disabled_does_nothing(self):
self._load(autonomous=True, enabled=False)
agent = self._agent(decision=None)
await self._stall(agent)
self.assertEqual(agent.acts, [])
self.assertEqual(self.models, []) # the model was never asked
async def test_model_pin_and_decision_recorded(self):
self._load(autonomous=True, confidence=0.95)
agent = self._agent(decision=llm.StallDecision("reassign", "stuck", 0.9))
await self._stall(agent)
self.assertEqual(self.models, ["claude-sonnet-5-5"])
# 0.9 is below the registry's 0.95, so the reassign goes to a human.
self.assertEqual(agent.acts, [("escalate", "reassign")])
self.assertEqual(self.recorded[0][0], "stall_response")
self.assertEqual(self.recorded[0][-1], "claude-sonnet-5-5")
async def test_stall_minutes_from_registry(self):
registry.loaded = False
self.assertEqual(self.mod.stall_minutes(), self.mod.STALL_MINUTES)
registry.apply({"skills": [{"skillid": "stall_response", "enabled": True, "thresholds": {"stallMinutes": 20}}]})
self.assertEqual(self.mod.stall_minutes(), 20)
# ── CustomerAgent: skill gate and recorded events ───────────────────────────
class TestCustomerAgent(_RegistryIsolation, unittest.IsolatedAsyncioTestCase):
def setUp(self):
self._save_registry()
async def asyncTearDown(self):
self._restore_registry()
message_bus.unregister_agent("CUSTOMER_AGENT")
async def test_notifications_disabled_skips_send(self):
from agents.customer_agent import CustomerAgent
registry.apply({"skills": [{"skillid": "customer_notifications", "enabled": False}]})
result = await CustomerAgent()._send_notification(_task("send_notification", order_id="O1"))
self.assertEqual(result["status"], "skipped")
async def test_cancel_and_notification_events_recorded(self):
from agents.customer_agent import CustomerAgent
agent = CustomerAgent()
for mt in (MessageType.ORDER_CANCELLED, MessageType.NOTIFICATION_SENT, MessageType.HEARTBEAT):
await agent.handle_message(AgentMessage("id", "EXCEPTION_AGENT", "CUSTOMER_AGENT", mt,
{"order_id": "O9"}, datetime.now()))
self.assertEqual([e["type"] for e in agent._received_events], ["ORDER_CANCELLED", "NOTIFICATION_SENT"])
self.assertEqual(agent._received_events[0]["order_id"], "O9")
# ── Message-bug fixes ────────────────────────────────────────────────────────
class TestMessageFixes(unittest.IsolatedAsyncioTestCase):
async def asyncTearDown(self):
for agent_id in ("FLEET_AGENT", "ORDER_AGENT", "JARVIS", "FAKE_ORDER", "FAKE_CUSTOMER"):
message_bus.unregister_agent(agent_id)
async def test_fleet_release_for_cancel_finds_vehicle_by_order(self):
from agents.fleet_agent import FleetAgent, VehicleAssignment
fleet = FleetAgent()
vehicle_id = next(iter(fleet._vehicles))
fleet._vehicles[vehicle_id]["status"] = "in_transit"
fleet._assignments["A1"] = VehicleAssignment(
"A1", vehicle_id, "R1", ["O1", "O2"], datetime.now(), datetime.now(), {})
missing = await fleet.handle_task(_task("release_vehicle_for_cancel", order_id="NOPE"))
self.assertEqual(missing["status"], "no_assignment")
await fleet.handle_task(_task("release_vehicle_for_cancel", order_id="O2"))
self.assertEqual(fleet._vehicles[vehicle_id]["status"], "available")
self.assertEqual(fleet._assignments["A1"].status, "returned")
async def test_order_agent_refuses_backend_calls(self):
from agents.order_agent import OrderAgent
agent = OrderAgent()
with patch("core.http_client.api_post") as post, patch("core.http_client.api_get") as get:
self.assertIsNone(await agent._api_post("/api/v1/admin/crmbooking", {}))
self.assertIsNone(await agent._api_get("/api/v1/admin/x"))
self.assertIsNone(await agent._api_patch("/api/v1/admin/x", {}))
post.assert_not_called()
get.assert_not_called()
async def test_jarvis_validate_payload_carries_order_id(self):
from core.agent import MasterAgent
class Sink:
def __init__(self, agent_id):
self.agent_id = agent_id
self.tasks = []
async def submit_task(self, task):
self.tasks.append(task)
jarvis = MasterAgent()
order, customer = Sink("ORDER_AGENT"), Sink("CUSTOMER_AGENT")
jarvis._sub_agents = {"ORDER_AGENT": order, "CUSTOMER_AGENT": customer}
await jarvis.handle_task(_task("orchestrate_order", order={"order_id": "O7"}))
self.assertEqual(order.tasks[0].task_type, "validate_order")
self.assertEqual(order.tasks[0].data["order_id"], "O7")
def test_order_status_update_message_type_exists(self):
self.assertEqual(MessageType.ORDER_STATUS_UPDATE.value, "ORDER_STATUS_UPDATE")
if __name__ == "__main__":
unittest.main()