Compare commits

...

4 Commits

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

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

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

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

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

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

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

View File

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

View File

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

View File

@@ -12,7 +12,9 @@ import redis.asyncio as aioredis
from core.agent import SpecializedAgent
from core.types import AgentTask, MessageType
from core.logger import logger
from core.llm import decide_assignment_failure, build_assignment_failure_context
from core.llm import decide_assignment_failure, build_assignment_failure_context, LLM_MODEL
from core.registry import registry
from core.decisions import record_decision
from config.system_config import (
NATS_HOST, NATS_PORT, NATS_USER, NATS_PASSWORD,
REDIS_HOST, REDIS_PORT, REDIS_PASSWORD,
@@ -47,6 +49,27 @@ MILER_UNAVAILABLE_STATUSES = {"break", "offline", "busy", "onbreak", "off_duty",
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):
"""
@@ -243,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}")
@@ -274,23 +291,35 @@ 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}")
# Rate limit: once this zone has been alerted today, don't re-decide
# (and don't pay for an LLM call) on every further failure.
if facts.get("alert_already_sent_today"):
if DISPATCH_REALERT_EVERY and count % DISPATCH_REALERT_EVERY == 0:
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 {DISPATCH_REALERT_EVERY}; earlier alert already raised).",
f"(re-alert every {realert_every}; earlier alert already raised).",
severity="high",
)
else:
@@ -300,7 +329,11 @@ class DispatchAgent(SpecializedAgent):
)
return
decision = await decide_assignment_failure(build_assignment_failure_context(facts))
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")
@@ -322,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.
@@ -441,11 +474,15 @@ class DispatchAgent(SpecializedAgent):
recipient="CUSTOMER_AGENT",
message_type=MessageType.AGENT_TASK,
payload={
# CUSTOMER_AGENT's real task contract (send_notification + DELAYED template)
# 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, "reason": "finding the right rider"},
"template_vars": {"order_id": booking_id, "new_eta": "being confirmed"},
},
correlation_id=str(booking_id),
)

View File

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

View File

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

View File

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

View File

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

View File

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

66
core/decisions.py Normal file
View File

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

View File

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

122
core/registry.py Normal file
View File

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

View File

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

View File

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

View File

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

6
requirements-dev.txt Normal file
View File

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

View File

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

View File

@@ -213,7 +213,7 @@ def make_handler_agent(redis, decision):
box = FakeMessages()
agent.send_message = box.send_message
calls = []
async def fake_decide(context):
async def fake_decide(context, model=None):
calls.append(context)
return decision
mod.decide_assignment_failure = fake_decide
@@ -288,6 +288,29 @@ class TestRateLimitAndSinks(unittest.IsolatedAsyncioTestCase):
self.assertEqual(payload["proposed_action"], "escalate")
class TestZoneKey(unittest.IsolatedAsyncioTestCase):
"""The backend sends no zone_id today; without a fallback every failure in
the country shares one bucket and the daily alert limit becomes global."""
async def test_real_zone_id_wins(self):
from agents.dispatch_agent import zone_key
self.assertEqual(zone_key("hyderabad", 17.4, 78.4), "hyderabad")
async def test_missing_zone_falls_back_to_geo_cell(self):
from agents.dispatch_agent import zone_key
self.assertEqual(zone_key(None, 11.0182714, 76.9677744), "geo:11.0,77.0")
self.assertEqual(zone_key("unknown", 11.0182714, 76.9677744), "geo:11.0,77.0")
async def test_distant_pickups_get_different_buckets(self):
from agents.dispatch_agent import zone_key
self.assertNotEqual(zone_key(None, 11.01, 76.96), zone_key(None, 17.44, 78.39))
async def test_no_zone_and_no_coords_is_unknown(self):
from agents.dispatch_agent import zone_key
self.assertEqual(zone_key(None, None, None), "unknown")
self.assertEqual(zone_key("unknown", "bad", "coords"), "unknown")
class TestBindConsumer(unittest.IsolatedAsyncioTestCase):
async def test_binds_to_discovered_stream_without_creating(self):
js = FakeJS(stream_name="TRACKING")

View File

@@ -0,0 +1,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()