From 5a7d32cc04f3867be6a1429b0e1be16345bd0ffe Mon Sep 17 00:00:00 2001 From: dharaneesh-r Date: Wed, 30 Sep 2026 14:56:08 +0530 Subject: [PATCH] implemenation on the ai agents and registry --- agents/customer_agent.py | 30 +++ agents/dispatch_agent.py | 42 ++-- agents/exception_agent.py | 69 ++++-- agents/express_dispatch_agent.py | 19 +- agents/fleet_agent.py | 18 ++ agents/order_agent.py | 33 +-- core/agent.py | 4 +- core/decisions.py | 66 +++++ core/llm.py | 41 ++- core/registry.py | 122 +++++++++ core/types.py | 3 + main.py | 7 + requirements-dev.txt | 6 + tests/test_dispatch_agent.py | 2 +- tests/test_registry_phase5.py | 413 +++++++++++++++++++++++++++++++ 15 files changed, 803 insertions(+), 72 deletions(-) create mode 100644 core/decisions.py create mode 100644 core/registry.py create mode 100644 requirements-dev.txt create mode 100644 tests/test_registry_phase5.py diff --git a/agents/customer_agent.py b/agents/customer_agent.py index 11d722b..325d7de 100644 --- a/agents/customer_agent.py +++ b/agents/customer_agent.py @@ -11,6 +11,7 @@ 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 @@ -64,6 +65,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 +213,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", {}) diff --git a/agents/dispatch_agent.py b/agents/dispatch_agent.py index 213b955..7cad049 100644 --- a/agents/dispatch_agent.py +++ b/agents/dispatch_agent.py @@ -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, @@ -264,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}") @@ -305,17 +301,25 @@ class DispatchAgent(SpecializedAgent): 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: @@ -325,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") @@ -347,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. diff --git a/agents/exception_agent.py b/agents/exception_agent.py index 2eed928..6f7cef9 100644 --- a/agents/exception_agent.py +++ b/agents/exception_agent.py @@ -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"] diff --git a/agents/express_dispatch_agent.py b/agents/express_dispatch_agent.py index 309bd86..4d38d96 100644 --- a/agents/express_dispatch_agent.py +++ b/agents/express_dispatch_agent.py @@ -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 diff --git a/agents/fleet_agent.py b/agents/fleet_agent.py index 80e4069..09716de 100644 --- a/agents/fleet_agent.py +++ b/agents/fleet_agent.py @@ -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") diff --git a/agents/order_agent.py b/agents/order_agent.py index ecd3144..c8e3582 100644 --- a/agents/order_agent.py +++ b/agents/order_agent.py @@ -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 # diff --git a/core/agent.py b/core/agent.py index 847e698..3185d50 100644 --- a/core/agent.py +++ b/core/agent.py @@ -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") diff --git a/core/decisions.py b/core/decisions.py new file mode 100644 index 0000000..7f3a9f6 --- /dev/null +++ b/core/decisions.py @@ -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 diff --git a/core/llm.py b/core/llm.py index 93e9557..26b3ced 100644 --- a/core/llm.py +++ b/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 @@ -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 diff --git a/core/registry.py b/core/registry.py new file mode 100644 index 0000000..15b427d --- /dev/null +++ b/core/registry.py @@ -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() diff --git a/core/types.py b/core/types.py index 8ae53b0..08b6158 100644 --- a/core/types.py +++ b/core/types.py @@ -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" diff --git a/main.py b/main.py index b8873c2..5ded5d5 100644 --- a/main.py +++ b/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() diff --git a/requirements-dev.txt b/requirements-dev.txt new file mode 100644 index 0000000..58f3ace --- /dev/null +++ b/requirements-dev.txt @@ -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 diff --git a/tests/test_dispatch_agent.py b/tests/test_dispatch_agent.py index 10ae4e1..ff3d199 100644 --- a/tests/test_dispatch_agent.py +++ b/tests/test_dispatch_agent.py @@ -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 diff --git a/tests/test_registry_phase5.py b/tests/test_registry_phase5.py new file mode 100644 index 0000000..0dd7481 --- /dev/null +++ b/tests/test_registry_phase5.py @@ -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()