Compare commits
6 Commits
58bfa07385
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| 9cbb6f9704 | |||
| 5a7d32cc04 | |||
| aa3bb35733 | |||
| 008df49605 | |||
| 5159c8d1a7 | |||
| 4283c602f6 |
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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) #
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 #
|
||||
|
||||
@@ -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
66
core/decisions.py
Normal 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
|
||||
60
core/llm.py
60
core/llm.py
@@ -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
122
core/registry.py
Normal 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()
|
||||
@@ -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"
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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."}
|
||||
|
||||
@@ -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:
|
||||
|
||||
7
main.py
7
main.py
@@ -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
6
requirements-dev.txt
Normal 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
|
||||
117
tests/test_customer_notify.py
Normal file
117
tests/test_customer_notify.py
Normal 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'])}"
|
||||
)
|
||||
@@ -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")
|
||||
|
||||
49
tests/test_jarvis_inbox.py
Normal file
49
tests/test_jarvis_inbox.py
Normal 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
|
||||
413
tests/test_registry_phase5.py
Normal file
413
tests/test_registry_phase5.py
Normal 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()
|
||||
Reference in New Issue
Block a user