diff --git a/PORTFOLIO.md b/PORTFOLIO.md new file mode 100644 index 0000000..31c07c5 --- /dev/null +++ b/PORTFOLIO.md @@ -0,0 +1,191 @@ +# LogiFlow AI — What Was Actually Built Here + +**One-line summary:** a Python multi-agent "operations brain" that sits beside a production last-mile delivery platform (a Go backend called Doormile), watches real infrastructure events (rider GPS, booking assignments), uses an LLM to make judgment calls about operational exceptions, and acts on them through the backend's API — with careful autonomy gating, deterministic fallbacks, an offline eval suite, and a live mission-control dashboard. + +> ⚠️ **Commercially sensitive — fix before publishing this repo or portfolio:** `config/system_config.py` contains **hardcoded fallback credentials** (NATS, Redis, and Postgres passwords, an internal API key) and **real public server IPs**; `main.py`'s help text also prints the NATS host/user, and a `.env` file is present in the working tree. None of that is reproduced in this document, but it must be scrubbed and the credentials rotated before the repo is shared. The `docs/` handoff note and code comments also reveal internals of the Doormile backend (package names, NATS topology) — fine for an internal doc, worth reviewing before public posting. + +The codebase has two distinct layers, and an honest portfolio should distinguish them: + +- A **production integration layer** (the real engineering): the message bus, the ExceptionAgent stall pipeline, the DispatchAgent coverage-gap watcher, the LLM decision layer, the eval harness, the HTTP client, and the Command Center. These talk to real NATS/Redis/Postgres and are tested. +- A **demo/simulation layer**: FleetAgent, HubAgent, RouteOptimizerAgent, and the two Streamlit UIs operate on hardcoded in-memory data (a fictional fleet across Indian cities). They demonstrate the agent framework but don't touch real infrastructure. + +--- + +## 1. The Agent Runtime and Message Bus (`core/agent.py`, `core/message_bus.py`) + +**What it is.** The framework every agent runs on. Each agent is an async worker with its own task queue; agents talk to each other through a central message bus rather than calling each other directly. The problem this solves: in a logistics operation, many concerns (orders, dispatch, exceptions, customer comms) need to react to the same events without being tangled together. Decoupling them behind a bus means any agent can be restarted, replaced, or run in a different process without the others knowing. + +**How it's built.** `Agent` (abstract base) owns an `asyncio.Queue` of tasks and a processing loop: pull a task, dispatch to a handler, record success/failure, emit telemetry. `MasterAgent` ("JARVIS") is the orchestrator — it registers sub-agents, fans out order work, and is the sink for human-review escalations. `MessageBus` is a singleton with **two transport modes**: when connected, it publishes over **NATS JetStream** (durable, persisted streams); when not, it silently degrades to in-process dispatch so the whole system still runs on a laptop with no infrastructure. + +``` +publish(msg) + ├─ NATS connected? ── JetStream publish + │ direct msg → subject logistics.direct. (durable consumer per agent) + │ broadcast → subject logistics. (durable consumer per type) + └─ not connected ── local dispatch straight into the recipient's queue +``` + +**The approach.** Directed messages carrying a `task_type` are converted into `AgentTask`s and enqueued — so a message from another agent is processed by exactly the same code path as a locally submitted task. Messages without a `task_type` go to `handle_message()` as notifications. A comment in `deliver()` records the bug this design fixed: previously, directed messages "sat unread in the bus queue forever." If a recipient isn't present in the process, the message is retained in a pull queue rather than dropped. Task IDs are derived from the message's correlation ID, which threads one booking's journey across every agent that touches it. + +**Trade-offs.** Message history is an in-memory list capped at 1,000 entries (observability, not durability — durability is JetStream's job). The local-fallback mode has weaker delivery guarantees than JetStream, which is the price of a zero-infrastructure dev mode. + +--- + +## 2. Miler Stall Detection and LLM-Gated Exception Handling (`agents/exception_agent.py`) — the flagship + +**What it is.** "Milers" are delivery riders. Sometimes a rider stops moving mid-delivery — traffic, a handoff, a breakdown, a dead phone. The customer is waiting either way. This agent detects the stall from GPS data, figures out *why it probably happened*, and chooses among four responses: wait, reassure the customer, reassign the booking to another rider, or escalate to a human dispatcher. The hard part isn't detection — it's that the right response is a **judgment call** (a rider parked at the delivery address is fine; a rider frozen on a highway for 50 minutes is not), and one of the possible actions (reassignment) is customer-visible and irreversible. + +**How it's built.** Detection runs on two independent paths that feed one decision pipeline: + +``` +Path A (event-driven): Path B (sweep, every 60s): +NATS TRACKING stream Postgres: all active bookings + miler.location.updated (pull-sub) ─┐ └─ for each, read Redis movement + compare to last position in │ hash; if last GPS ping is + Redis hash miler::movement │ stale ≥10 min ──────────────┐ + if unmoved ≥10 min ─────────────┤ │ + ▼ ▼ + Redis SET NX "stall_notified:" ← dedup claim (6h TTL) + │ winner only + ▼ + publish miler.stalled (JetStream) + ▼ + pull-consumer of miler.stalled + Redis SET NX "stall_handled:" ← redelivery guard + ▼ + gather facts (Postgres booking phase/age, + Redis GPS freshness, per-miler stall count today) + ▼ + Claude decides: wait │ notify_only │ reassign │ escalate + ▼ + wait → log notify → POST /internal/notify (Go API) + reassign → gated: only if EXCEPTION_AGENT_AUTONOMOUS=true + AND confidence ≥ 0.7, else escalate to JARVIS + escalate → message to JARVIS (human review endpoint) +``` + +**The approach, step by step:** + +1. **Why two detection paths?** The event path catches a rider whose phone keeps reporting the *same* position. The sweep path catches the opposite failure: a phone that stops reporting entirely (no events → the event path is blind). Together they cover both "frozen position" and "gone dark." +2. **Why the Redis `SET NX` claims?** Both paths can detect the same stall, and every ping after minute 10 re-detects it. The old design used an in-memory set (unbounded, lost on restart, racy). The replacement is an atomic Redis `SET NX` with a 6-hour TTL — bounded memory, restart-safe, and race-free between the two paths. There are *two* claims: one dedupes the *alert* (producer side), one guards *handling* (consumer side), because JetStream is at-least-once and a redelivered event must never re-run an irreversible reassign. Both claims deliberately **fail open** on Redis errors — a genuine stall must never be silently muted; a duplicate alert is the cheaper mistake. +3. **Context gathering before deciding.** The agent pulls real signals — booking phase (assigned vs. pickup scheduled), booking age, minutes since last GPS ping, position-unchanged duration, how many times this rider has stalled today — all read-only and best-effort. A thin-context decision still runs, and the prompt is written so low information skews toward escalation. +4. **The LLM decides, but with a blast-radius hierarchy.** The four actions are ordered by reversibility. The prompt explicitly instructs "prefer the least disruptive action that fits the evidence" and to report honest confidence. The only irreversible action is double-gated in *code*, not in the prompt: an env-var autonomy switch (default **off**) AND a confidence threshold. Below the bar, the model's proposal is packaged with its reasoning and confidence and escalated to a human via JARVIS — the AI proposes, the human disposes. +5. **Deterministic fallback.** If the LLM call fails entirely (network, refusal, truncation, bad JSON), the agent falls back to the pre-LLM behavior — reassign + notify — so detection never silently stops acting just because the model is down. +6. **Timezone discipline.** All stall math goes through `_utcnow()` (tz-aware UTC) and a parser that assumes naive timestamps are UTC — explicitly so pre-existing Redis entries stay comparable during rollout and container timezone can't corrupt the arithmetic. There are dedicated unit tests for exactly this. + +**Trade-offs.** The dedup TTL (6h) means a rider who stalls, recovers, and stalls again within the window won't re-alert — accepted to avoid alert storms. Exception records live in an in-memory dict (observability, not system of record — Postgres owns booking truth). + +--- + +## 3. The Assignment-Failure Watcher (`agents/dispatch_agent.py` + `docs/handoff-assignment-failed-event.md`) + +**What it is.** When a new booking comes in, the Go backend tries to assign a rider (AI assignment plus a fallback). Sometimes both fail — nobody gets assigned and, previously, nothing happened except a synchronous flag. This agent consumes assignment events and answers a subtle operational question: **is this failure noise, a real delay, a coverage gap (not enough riders in this zone), or something systemically broken?** Each answer warrants a different response. + +**How it's built.** The agent binds durable JetStream consumers to the backend-owned `ASSIGNMENTS` stream. It is deliberately a **pure consumer** — a code comment documents the bug this fixed: creating its own stream over the backend's subjects raised a silent overlap error that left the consumer unbound. Instead it *discovers* which existing stream carries the subject (`find_stream_name_by_subject`) and binds to it, with a clear actionable log if no stream carries it. + +On `booking.assigned`, it forwards hub-preparation work. On `booking.assignment_failed`, the interesting path runs: + +``` +booking.assignment_failed (from Go backend) + ▼ +gather coverage facts: + • Redis GEORADIUS on milers:locations at 10 / 20 / 30 km rings + → "nearest rider within N km" or "none within 30km" + • Redis INCR failed_assignment:: → failures-today counter (48h TTL) + ▼ +Claude decides: monitor │ notify_customer │ ops_alert │ escalate + ▼ +notify_customer is gated behind DISPATCH_AGENT_AUTONOMOUS (default off: + it becomes an internal ops proposal instead of a customer message) +escalate → JARVIS with the model's proposed action + confidence +LLM unavailable → deterministic fallback: ops alert only if failures ≥ 3 today +``` + +**The approach.** The clever part is the *fact design*. The expanding GEORADIUS rings turn a geo query into a coarse feature the model can reason about ("nearest rider within 20km"), and the per-zone daily counter distinguishes one-off from chronic. The prompt encodes the key diagnostic asymmetry: repeated failures with **no** nearby rider = coverage gap → alert ops; repeated failures **despite** nearby riders = something is broken in the assignment system itself → escalate to a human. Missing coordinates degrade gracefully to "coverage unknown," which the prompt treats as a reason for caution. + +**A notable engineering artifact:** the consumer, decision logic, evals, and tests are all built, but the feature is currently **inert** — the Go backend doesn't publish the failure event yet. `docs/handoff-assignment-failed-event.md` is a precise cross-team handoff spec (exact subject, payload schema, stream change, verification steps, discovered by inspecting the running Go binary's symbols) telling the backend team the one thing they need to publish to light it up. That doc is good evidence of working across a language/team boundary. + +--- + +## 4. The LLM Decision Layer (`core/llm.py`) + +**What it is.** A single, thin module that owns every model interaction, so agent code never touches the Anthropic SDK directly. It solves the problem of making free-form model output safe to act on in an automated pipeline. + +**The approach.** + +- **Constrained decisions, not open generation.** Every call uses structured output against a strict JSON Schema — the model must return `{action ∈ fixed enum, reasoning, confidence}`. The action is re-validated in Python even after schema enforcement. +- **Fail-to-None contract.** Any failure — network error (one retry), refusal, `max_tokens` truncation, unparseable JSON, out-of-enum action — collapses to `None`, and every caller has a documented deterministic fallback. The model can therefore *only ever* improve behavior over the pre-LLM baseline, never break it. +- **Async client, lazily constructed** so a missing API key can't break agent startup and a slow model call can't freeze the event loop. +- **Prompt-building functions are shared with the evals** (`build_stall_context`, `build_assignment_failure_context`) so the offline eval scores *exactly* the prompt that runs in production — no drift between what's tested and what ships. +- Model, effort, token budget, and timeout are all env-configurable; a comment records the operational lesson that adaptive thinking eats into `max_tokens`, so the budget must be generous or the JSON gets truncated. + +--- + +## 5. The Offline Eval Harness (`evals/`) + +**What it is.** Before letting the model act autonomously, you need a number that says how good its judgment is. This is a small purpose-built eval framework: labeled scenario datasets (JSONL), a shared scoring harness, and one runner per decision type. + +**The approach.** + +- **Cases encode that judgment problems have multiple right answers.** Each case has an `ideal` action and an `acceptable` set (e.g., a 25-minute stall with no GPS ping for 25 minutes: ideal `escalate`, but `reassign` also acceptable). The headline metric is *acceptable-rate* — "was the decision defensible?" — with exact-match as secondary. This is a more honest metric design than forcing a single gold label onto ambiguous scenarios. +- **Majority voting across N samples** per case (`--runs 3`) measures decision *stability*, not just a lucky single sample. Model errors are counted separately from wrong answers. +- **Reproducibility:** cases pin a fixed date so only hour-of-day varies between runs; `--dry-run` prints every prompt without API calls; `--model` swaps models for comparison; `--min-pass-rate` turns the eval into a CI gate. +- The 25 cases (15 stall, 10 assignment) cover the genuinely tricky territory: rush-hour ambiguity, a rider at the delivery door, thin overnight coverage, a repeat-staller pattern, a dead phone vs. a frozen position, and the riders-nearby-yet-failing systemic case. Several are labeled as using the *real* fact keys the production gatherer emits. +- **The intended workflow is explicit in the docstrings:** collect real production cases into the JSONL, get the acceptable-rate trustworthy, and only then flip the autonomy env vars. The eval is the graduation exam for autonomy. + +--- + +## 6. The JARVIS Command Center (`dashboard/command_center.py` + one static HTML page) + +**What it is.** A real-time mission-control web dashboard (port 8600) showing every agent's live status, every message crossing the bus, task throughput/failure counts, per-agent task-duration stats, and message rate — the answer to "what is this distributed swarm doing *right now*?" + +**How it's built.** + +``` +Agents ──(JetStream)── logistics.> ┐ +Agents ──(plain NATS)── telemetry.> ┤→ Command Center process (FastAPI) + │ in-memory aggregation (deques, counters) + │ REST: /api/state, /api/agent/{id} + └──→ WebSocket fan-out → browser UI +``` + +**The approach.** The key design decision is that the dashboard is a **read-only tap that cannot perturb the system it observes**: it uses plain (non-JetStream) subscriptions with no durable consumers and no acks, so it never steals or holds messages the real agents need. Symmetrically, on the producer side, telemetry is published on plain NATS — deliberately *not* JetStream — and is a no-op when disconnected, so observability can never accumulate in the persistent stream or affect agent behavior (the base `Agent` class emits state every 5 seconds and a lifecycle event per task). The server reconnects to NATS forever, marks agents offline after 15 seconds of telemetry silence, computes a rolling messages-per-minute rate from a timestamp deque, and every buffer is a bounded `deque` so it can run indefinitely. New browsers get a full state snapshot, then deltas. + +--- + +## 7. Supporting Cast (briefer, and honestly labeled) + +- **HTTP client (`core/http_client.py`)** — one shared `aiohttp` session with connection pooling, and retry logic that encodes real HTTP semantics: 4xx returns `None` immediately (retrying a client error is pointless), 5xx retries with exponential backoff, and only 2xx/3xx yields a body. The docstring nails why: callers gate side effects on a non-None result, so "server rejected it" must never look like "it worked." +- **OrderAgent** — order intake/validation/categorization (required fields, phone/pincode regex, weight limits), zone typing via pincode prefix (same prefix → last-mile; same region → hub-to-spoke; else hub-to-hub), persisting through the Go backend's CRM booking API and adopting the backend's booking ID as the canonical order ID. +- **FleetAgent / HubAgent / RouteOptimizerAgent** — the simulation layer: an in-memory fleet with hub capacity accounting (with visible bug-fix comments about counter drift on double-release), and route estimation using **haversine distance with time-of-day traffic multipliers and a TTL route cache**. These demonstrate the framework, not production logic. +- **Tests (`tests/`, ~40 cases)** — genuinely targeted at the failure modes, not happy paths: dedup claim semantics including fail-open on Redis errors, timezone parsing edge cases, HTTP retry matrices (5xx→2xx recovery, 4xx no-retry), LLM refusal/truncation/invalid-action handling, stream binding without creation, and message-bus delivery routing. All use hand-rolled async fakes (FakeRedis with real `SET NX` semantics, fake PG pools, fake Anthropic clients) rather than mocking frameworks — the fakes model the *behavior* being relied on. +- **Deployment** — Dockerfile + docker-compose with resource limits, all secrets via env (the compose file does this correctly; it's the config file's *fallback defaults* that leak). + +--- + +## The Stack + +- **Language:** Python 3.12, `asyncio` throughout +- **Messaging:** NATS JetStream (durable streams, push + pull consumers, at-least-once delivery) +- **State/coordination:** Redis (GEO radius queries, hashes for GPS state, atomic `SET NX` claims, daily counters with TTL) +- **Database:** PostgreSQL via `asyncpg` (read-only against the backend's booking tables) +- **AI:** Anthropic Claude via the async SDK — structured output (JSON Schema), adaptive thinking, configurable effort +- **Web:** FastAPI + uvicorn + WebSockets (command center); Streamlit (admin dashboard & customer portal demos) +- **HTTP:** aiohttp (pooled shared session) +- **Integration target:** a Go backend ("Doormile") via internal REST API +- **Testing:** `unittest` + `IsolatedAsyncioTestCase`; custom JSONL eval harness +- **Deployment:** Docker, docker-compose + +--- + +## What's Genuinely Hard Here + +1. **Exactly-once *effects* on an at-least-once bus.** JetStream redelivers; two detection paths race; every GPS ping re-detects a stall. Getting from that to "the customer is reassigned at most once" required the two-tier atomic Redis claim design (alert claim + handled claim), with the deliberate fail-open asymmetry — duplicate alerts are cheaper than muted stalls. This is the classic distributed-systems dedup problem, solved correctly and unit-tested. +2. **Making an LLM safe to put in an actuation loop.** The whole shape of the safety story — actions as a closed enum with schema-enforced output, a reversibility hierarchy where only the irreversible action is code-gated behind autonomy-flag + confidence-threshold, human escalation carrying the model's proposal and reasoning, and a deterministic fallback so model downtime degrades to the old behavior instead of inaction. The gates live in code, not in the prompt — the prompt is advice, the gate is enforcement. +3. **Evaluating judgment, not correctness.** Recognizing that stall handling has no single right answer and designing the metric around a defensible-action set with majority-vote stability — then wiring the eval to consume the *identical* prompt-building code as production, and making autonomy contingent on the eval passing on real collected cases. That's a disciplined promote-to-autonomy pipeline in miniature. +4. **Being a good citizen in someone else's infrastructure.** The agents consume streams owned by a Go backend they don't control. The stream-discovery-instead-of-creation fix (with the documented overlap-error war story), the read-only no-ack dashboard tap, telemetry deliberately kept off JetStream, and the reverse-engineered handoff spec for the missing event all show real cross-system integration work — the unglamorous kind that actually breaks projects when done wrong. +5. **Dual detection with complementary blind spots.** The frozen-position event path and the gone-dark database sweep each catch what the other structurally cannot. Recognizing that "no data" is itself a signal requiring a separate mechanism is a nice piece of failure-mode thinking. + +--- + +*Notes for portfolio framing:* the git history and comments show honest iteration (bug-fix comments explaining *why* the old design was wrong — dead code after a `return`, unbound consumers, the in-memory dedup set), which reads well in an engineering narrative. And once more so it isn't missed: **rotate and remove the credentials in `config/system_config.py` and the committed `.env` before this repo goes anywhere public.** diff --git a/README.md b/README.md index 8f0b3cc..c16e161 100644 --- a/README.md +++ b/README.md @@ -239,7 +239,3 @@ Contributions welcome! Please: - -Resume this session with: -claude --resume d1f8a432-2298-4ff5-bee4-15e64d04bfa4 -PS C:\Users\Admin\Downloads\logiflow-ai-logistics-agent-system\logistics-ai> diff --git a/agents/dispatch_agent.py b/agents/dispatch_agent.py index a1160fa..a61ba5e 100644 --- a/agents/dispatch_agent.py +++ b/agents/dispatch_agent.py @@ -1,7 +1,8 @@ """Dispatch Agent — watches NATS assignment events and escalates coverage gaps.""" import asyncio import json -from datetime import datetime +import os +from datetime import datetime, timezone from typing import Dict, List, Any, Optional import nats @@ -11,18 +12,28 @@ 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 config.system_config import ( NATS_HOST, NATS_PORT, NATS_USER, NATS_PASSWORD, REDIS_HOST, REDIS_PORT, REDIS_PASSWORD, ) +# Notifying a customer is outward-facing, so it is gated: without autonomy the +# agent proposes it as an internal ops alert instead of messaging the customer. +DISPATCH_AGENT_AUTONOMOUS = os.getenv("DISPATCH_AGENT_AUTONOMOUS", "false").lower() == "true" + +# The Go backend owns the JetStream stream carrying booking.* events; this agent +# is a pure consumer and must not create it. Leave empty to auto-discover the +# stream by subject, or pin it to a stream name if discovery isn't desired. +DISPATCH_STREAM = os.getenv("DISPATCH_STREAM", "") + class DispatchAgent(SpecializedAgent): """ - Watches NATS ASSIGNMENTS stream for booking.assigned and - booking.assignment_failed events published by the Go backend - after routemate AI or fallback assignment. Does not trigger - assignment itself. + Consumes booking.assigned and booking.assignment_failed events published by + the Go backend (after routemate AI or fallback assignment) by binding to the + backend-owned JetStream stream that carries them. It is a pure consumer — it + does not create the stream, and does not trigger assignment itself. """ def __init__(self): @@ -71,34 +82,46 @@ class DispatchAgent(SpecializedAgent): ) self._nats_js = self._nats_nc.jetstream() - try: - await self._nats_js.add_stream( - name="ASSIGNMENTS", - subjects=["booking.>"], - ) - logger.info("NATS stream 'ASSIGNMENTS' created") - except nats.js.errors.BadRequestError: - pass # stream already exists - - sub_assigned = await self._nats_js.subscribe( - "booking.assigned", - durable="dispatch-booking-assigned", - cb=self._on_nats_booking_assigned, + # Bind consumers to the *existing* backend-owned stream. Do NOT create + # a stream here: booking.* is published by the Go backend, so creating + # a second stream over the same subjects raises an overlap error that + # (previously swallowed) left the consumer unbound. Mirrors the + # explicit stream= binding ExceptionAgent uses for the TRACKING stream. + sub_assigned = await self._bind_consumer( + "booking.assigned", "dispatch-booking-assigned", + self._on_nats_booking_assigned, ) - sub_failed = await self._nats_js.subscribe( - "booking.assignment_failed", - durable="dispatch-booking-assignment-failed", - cb=self._on_nats_booking_assignment_failed, - ) - self._nats_subs.extend([sub_assigned, sub_failed]) - logger.info( - "DispatchAgent subscribed to booking.assigned " - "and booking.assignment_failed" + sub_failed = await self._bind_consumer( + "booking.assignment_failed", "dispatch-booking-assignment-failed", + self._on_nats_booking_assignment_failed, ) + self._nats_subs.extend(s for s in (sub_assigned, sub_failed) if s is not None) + if self._nats_subs: + logger.info(f"DispatchAgent bound {len(self._nats_subs)} booking-event consumer(s)") except Exception as e: logger.error(f"DispatchAgent NATS connect error: {e}") + async def _bind_consumer(self, subject, durable, cb): + """Bind a durable push consumer to whichever existing stream carries + `subject`. Returns the subscription, or None (with a clear, actionable + log) if no stream carries it — never a swallowed or opaque error.""" + try: + stream = DISPATCH_STREAM or await self._nats_js.find_stream_name_by_subject(subject) + sub = await self._nats_js.subscribe(subject, durable=durable, stream=stream, cb=cb) + logger.info(f"DispatchAgent bound '{subject}' on stream '{stream}' (durable={durable})") + return sub + except nats.js.errors.NotFoundError: + logger.error( + f"DispatchAgent: no JetStream stream carries '{subject}'. The publisher " + f"(Go backend) must create a stream covering it, or set DISPATCH_STREAM to " + f"the stream name. Not subscribing to '{subject}'." + ) + return None + except Exception as e: + logger.error(f"DispatchAgent failed to bind '{subject}': {e}") + return None + # ------------------------------------------------------------------ # # Redis GEO lookup # # ------------------------------------------------------------------ # @@ -177,9 +200,10 @@ class DispatchAgent(SpecializedAgent): booking.assignment_failed — Go publishes this when decide-assignment returns escalate=true or both AI and fallback find nobody. - Widens the GEORADIUS search to gauge coverage depth, tracks daily - failures per zone in Redis, and escalates to CUSTOMER_AGENT as an - ops_alert when a zone hits 3 failures in a day. + Gathers coverage context (nearest-rider sweep + daily failure counter), + then asks Claude to choose: monitor / notify_customer / ops_alert / + escalate. Customer notification is gated behind autonomy. Falls back to + the previous count>=3 heuristic if the LLM is unavailable. """ try: data = json.loads(msg.data.decode()) @@ -188,64 +212,135 @@ class DispatchAgent(SpecializedAgent): lat = data.get("lat") lon = data.get("lon") - logger.warning( - f"[DISPATCH] booking.assignment_failed — " - f"booking={booking_id} zone={zone_id}" + logger.warning(f"[DISPATCH] booking.assignment_failed — booking={booking_id} zone={zone_id}") + + facts, count = await self._gather_assignment_facts(zone_id, lat, lon) + logger.info(f"[DISPATCH] booking={booking_id} context facts: {facts}") + + decision = await decide_assignment_failure(build_assignment_failure_context(facts)) + + if decision is None: + logger.warning(f"No LLM decision for booking {booking_id}; falling back to count>=3 heuristic") + if count >= 3: + await self._ops_alert( + zone_id, booking_id, + f"Zone {zone_id} has had {count} failed assignments today — recommend onboarding milers here.", + ) + else: + logger.info(f"[DISPATCH] zone {zone_id} failure count {count}/3 — transient, no escalation") + return + + logger.info( + f"[DISPATCH] booking={booking_id} zone={zone_id} decision={decision.action} " + f"confidence={decision.confidence:.2f} reasoning={decision.reasoning!r}" ) - # Wider sweeps to gauge how far the nearest miler actually is. - has_coords = lat is not None and lon is not None - found_20km = await self._find_zone(lat, lon, radius_km=20) if has_coords else None - found_30km = await self._find_zone(lat, lon, radius_km=30) if has_coords else None + if decision.action == "monitor": + pass # transient — the backend will retry - if found_20km: - logger.info( - f"[DISPATCH] Nearest miler within 20 km: {found_20km['miler_id']}" - ) - elif found_30km: - logger.info( - f"[DISPATCH] Nearest miler within 30 km: {found_30km['miler_id']}" - ) - else: - logger.warning( - f"[DISPATCH] No miler found within 30 km of zone {zone_id}" - ) + elif decision.action == "notify_customer": + if DISPATCH_AGENT_AUTONOMOUS: + await self._notify_customer_delay(booking_id) + else: + # Outward-facing action not authorised — raise an internal proposal instead. + await self._ops_alert( + zone_id, booking_id, + f"[proposed: notify_customer, conf={decision.confidence:.2f}] {decision.reasoning}", + ) - # Increment daily failure counter (48 h TTL covers day rollover). - today = datetime.utcnow().strftime("%Y-%m-%d") - counter_key = f"failed_assignment:{zone_id}:{today}" - count = await self._redis.incr(counter_key) - await self._redis.expire(counter_key, 172800) + elif decision.action == "ops_alert": + await self._ops_alert(zone_id, booking_id, decision.reasoning) - if count >= 3: - logger.warning( - f"[DISPATCH] Coverage gap detected: zone {zone_id} has " - f"{count} failed assignments today — escalating" - ) - await self.send_message( - recipient="CUSTOMER_AGENT", - message_type=MessageType.AGENT_TASK, - payload={ - "task_type": "ops_alert", - "zone_id": zone_id, - "reasoning": ( - f"Zone {zone_id} has had {count} failed assignments today " - f"— recommend onboarding milers here." - ), - }, - correlation_id=str(booking_id), - ) - else: - logger.info( - f"[DISPATCH] zone {zone_id} failure count {count}/3 — " - f"transient, no escalation" - ) + else: # "escalate" + await self._escalate_dispatch(zone_id, booking_id, decision) except Exception as e: logger.error(f"DispatchAgent booking.assignment_failed handler error: {e}") finally: await msg.ack() + # ------------------------------------------------------------------ # + # Context gathering + actions # + # ------------------------------------------------------------------ # + + async def _gather_assignment_facts(self, zone_id, lat, lon): + """Read-only coverage context plus the daily failure counter for this + zone. Returns (facts, failures_today). Best-effort — any failure just + yields fewer facts, never raises.""" + facts: Dict[str, Any] = {"zone_id": zone_id} + has_coords = lat is not None and lon is not None + facts["has_coordinates"] = has_coords + + nearest_km = None + if has_coords: + for radius in (10, 20, 30): + try: + if await self._find_zone(lat, lon, radius_km=radius): + nearest_km = radius + break + except Exception as e: + logger.warning(f"assignment facts: GEORADIUS {radius}km failed: {e}") + break + + if not has_coords: + facts["nearest_miler_within_km"] = "unknown (no coordinates)" + elif nearest_km is None: + facts["nearest_miler_within_km"] = "none within 30km" + else: + facts["nearest_miler_within_km"] = nearest_km + + count = 0 + try: + today = datetime.now(timezone.utc).strftime("%Y-%m-%d") + key = f"failed_assignment:{zone_id}:{today}" + count = await self._redis.incr(key) + await self._redis.expire(key, 172800) # 48h covers day rollover + except Exception as e: + logger.warning(f"assignment facts: failure counter failed for zone {zone_id}: {e}") + facts["failures_today"] = count + + return facts, count + + async def _ops_alert(self, zone_id, booking_id, reasoning): + await self.send_message( + recipient="CUSTOMER_AGENT", + message_type=MessageType.AGENT_TASK, + payload={"task_type": "ops_alert", "zone_id": zone_id, "reasoning": reasoning}, + correlation_id=str(booking_id), + ) + + async def _notify_customer_delay(self, booking_id): + await self.send_message( + recipient="CUSTOMER_AGENT", + message_type=MessageType.AGENT_TASK, + payload={ + "task_type": "customer_notify", + "booking_id": booking_id, + "message": "We're finding the right rider for your delivery — thanks for your patience.", + }, + correlation_id=str(booking_id), + ) + + async def _escalate_dispatch(self, zone_id, booking_id, decision): + logger.warning( + f"[DISPATCH] Escalating booking {booking_id} (zone {zone_id}) to a human " + f"(proposed={decision.action}, confidence={decision.confidence:.2f}): {decision.reasoning}" + ) + await self.send_message( + recipient="JARVIS", + message_type=MessageType.AGENT_TASK, + payload={ + "task_type": "human_review", + "reason": "assignment_failed", + "zone_id": zone_id, + "booking_id": booking_id, + "proposed_action": decision.action, + "confidence": decision.confidence, + "reasoning": decision.reasoning, + }, + correlation_id=str(booking_id), + ) + # ------------------------------------------------------------------ # # Task handler (no inbound task types remain) # # ------------------------------------------------------------------ # diff --git a/agents/exception_agent.py b/agents/exception_agent.py index 7a54d32..28be9bf 100644 --- a/agents/exception_agent.py +++ b/agents/exception_agent.py @@ -10,8 +10,9 @@ StallDetector sweep (every 60 s): """ import asyncio import json +import os import uuid -from datetime import datetime, timedelta +from datetime import datetime, timedelta, timezone from typing import Dict, List, Any, Optional, Set from dataclasses import dataclass, field from enum import Enum @@ -25,6 +26,7 @@ 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 config.system_config import ( GO_API_BASE_URL, INTERNAL_API_KEY, DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD, @@ -36,6 +38,35 @@ STALL_MINUTES = 10 ACTIVE_STATUSES = ["Miler_Assigned", "Pickup_Scheduled"] TRACKING_STREAM = "TRACKING" +# Reassignment is customer-visible and hard to reverse, so it is gated: the LLM +# only *proposes* a reassign unless the agent is explicitly put in autonomous +# mode AND the model is confident enough. Otherwise the proposal is escalated to +# a human via JARVIS. +AUTONOMOUS_REASSIGN = os.getenv("EXCEPTION_AGENT_AUTONOMOUS", "false").lower() == "true" +REASSIGN_MIN_CONFIDENCE = float(os.getenv("EXCEPTION_AGENT_REASSIGN_CONFIDENCE", "0.7")) +# How long a booking stays "already alerted" so we don't re-alert the same stall. +# Self-expiring so dedup state is bounded and survives restarts (unlike a set). +STALL_NOTIFIED_TTL = int(os.getenv("STALL_NOTIFIED_TTL_SEC", "21600")) # 6h + + +def _utcnow() -> datetime: + """Timezone-aware current time in UTC. All stall math uses this so it stays + correct regardless of the container's local timezone, and never mixes naive + with tz-aware timestamps.""" + return datetime.now(timezone.utc) + + +def _parse_ts(value: Optional[str]) -> Optional[datetime]: + """Parse an ISO timestamp as tz-aware UTC. Naive values are assumed UTC so + pre-existing (naive) Redis entries remain comparable during rollout.""" + if not value: + return None + try: + dt = datetime.fromisoformat(value) + except (ValueError, TypeError): + return None + return dt if dt.tzinfo else dt.replace(tzinfo=timezone.utc) + class ExceptionType(str, Enum): DELAY = "delay" @@ -95,7 +126,6 @@ class ExceptionAgent(SpecializedAgent): self._location_sub = None self._stalled_sub = None - self._notified_stalls: Set[str] = set() self._exceptions: Dict[str, ExceptionRecord] = {} self._strategies = self._init_strategies() self._escalation_rules = self._init_escalation_rules() @@ -211,7 +241,7 @@ class ExceptionAgent(SpecializedAgent): miler_id = str(data.get("miler_id", "")) lat = str(round(float(data.get("lat", 0)), 6)) lon = str(round(float(data.get("lon", 0)), 6)) - now_ts = datetime.now().isoformat() + now_ts = _utcnow().isoformat() redis_key = f"miler:{miler_id}:movement" prev = await self._redis.hgetall(redis_key) @@ -227,20 +257,13 @@ class ExceptionAgent(SpecializedAgent): else: await self._redis.hset(redis_key, mapping={"lat": lat, "lon": lon, "updated_at": now_ts}) - unchanged_since_str = prev.get("position_unchanged_since", now_ts) - try: - unchanged_since = datetime.fromisoformat(unchanged_since_str) - except ValueError: - unchanged_since = datetime.now() - - minutes_stalled = (datetime.now() - unchanged_since).total_seconds() / 60 + unchanged_since = _parse_ts(prev.get("position_unchanged_since")) or _utcnow() + minutes_stalled = (_utcnow() - unchanged_since).total_seconds() / 60 if minutes_stalled >= STALL_MINUTES: booking = await self._get_active_booking(miler_id) if booking: - booking_id = booking["booking_id"] - if booking_id not in self._notified_stalls: - await self._publish_stall(miler_id, booking_id, minutes_stalled) + await self._publish_stall(miler_id, booking["booking_id"], minutes_stalled) # ── Stall detection: background sweep ──────────────────────────────────── @@ -274,23 +297,19 @@ class ExceptionAgent(SpecializedAgent): logger.error(f"StallDetector Postgres error: {e}") return - now = datetime.now() + now = _utcnow() for row in rows: miler_id = row["miler_id"] booking_id = row["booking_id"] - if booking_id in self._notified_stalls: - continue - movement = await self._redis.hgetall(f"miler:{miler_id}:movement") if not movement or "updated_at" not in movement: continue - try: - updated_at = datetime.fromisoformat(movement["updated_at"]) - minutes_stale = (now - updated_at).total_seconds() / 60 - except ValueError: + updated_at = _parse_ts(movement.get("updated_at")) + if updated_at is None: continue + minutes_stale = (now - updated_at).total_seconds() / 60 if minutes_stale >= STALL_MINUTES: logger.warning(f"StallDetector: miler {miler_id} stale {minutes_stale:.1f} min (booking {booking_id})") @@ -298,8 +317,57 @@ class ExceptionAgent(SpecializedAgent): # ── Stall: publish + handle ─────────────────────────────────────────────── + def _stall_counter_key(self, miler_id) -> str: + return f"miler:{miler_id}:stalls:{_utcnow().strftime('%Y-%m-%d')}" + + def _stall_notified_key(self, booking_id) -> str: + return f"stall_notified:{booking_id}" + + def _stall_handled_key(self, booking_id) -> str: + return f"stall_handled:{booking_id}" + + async def _claim_stall(self, booking_id) -> bool: + """Atomically claim the first stall ALERT for a booking (Redis SET NX with + TTL). True = publish this alert; False = already alerted within the window. + Replaces the old unbounded in-memory set: self-expiring (bounded memory), + durable across restarts, and race-free between the ping and sweep paths. + Fail-open on Redis error so a genuine stall is never silently muted.""" + try: + got = await self._redis.set( + self._stall_notified_key(booking_id), "1", ex=STALL_NOTIFIED_TTL, nx=True, + ) + return bool(got) + except Exception as e: + logger.warning(f"stall claim failed for {booking_id}: {e}") + return True + + async def _claim_stall_handled(self, booking_id) -> bool: + """Atomically claim HANDLING of a stall event (consumer side). Guards + against JetStream at-least-once redelivery re-running the decision and + repeating an irreversible reassign. Fail-open on Redis error.""" + try: + got = await self._redis.set( + self._stall_handled_key(booking_id), "1", ex=STALL_NOTIFIED_TTL, nx=True, + ) + return bool(got) + except Exception as e: + logger.warning(f"stall handled-claim failed for {booking_id}: {e}") + return True + + async def _incr_stall_counter(self, miler_id): + """Bump this miler's per-day stall count (Redis, best-effort). 48h TTL + covers the day rollover; failures never block stall handling.""" + try: + key = self._stall_counter_key(miler_id) + await self._redis.incr(key) + await self._redis.expire(key, 172800) + except Exception as e: + logger.warning(f"stall counter incr failed for miler {miler_id}: {e}") + async def _publish_stall(self, miler_id: str, booking_id: str, minutes_stalled: float): - self._notified_stalls.add(booking_id) + if not await self._claim_stall(booking_id): + return # already alerted for this booking within the dedup window + await self._incr_stall_counter(miler_id) payload = json.dumps({ "miler_id": miler_id, "booking_id": booking_id, @@ -312,7 +380,15 @@ class ExceptionAgent(SpecializedAgent): logger.error(f"Failed to publish miler.stalled: {e}") async def _on_miler_stalled(self, msg): - """Reassign booking and notify customer via Go API.""" + """Decide (via Claude) how to handle a stalled miler, then act. + + 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. + """ try: data = json.loads(msg.data.decode()) except Exception: @@ -324,36 +400,175 @@ class ExceptionAgent(SpecializedAgent): logger.warning(f"EXCEPTION_AGENT stall: miler={miler_id} booking={booking_id} ({minutes_stalled} min)") - headers = {"X-Internal-Key": INTERNAL_API_KEY} + if not await self._claim_stall_handled(booking_id): + logger.info(f"Stall for booking {booking_id} already handled; skipping redelivery") + return - reassign_result = await api_post( + facts = await self._gather_stall_facts(miler_id, booking_id) + if facts: + 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, + ) + decision = await decide_stall_response(context) + + 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"], + ) + return + + logger.info( + f"[EXCEPTION] booking={booking_id} decision={decision.action} " + f"confidence={decision.confidence:.2f} reasoning={decision.reasoning!r}" + ) + + actions: List[str] = [] + + if decision.action == "wait": + actions.append("wait") + + elif decision.action == "notify_only": + await self._notify_customer( + booking_id, "We're keeping an eye on your delivery — thanks for your patience." + ) + actions.append("notify_customer") + + elif decision.action == "reassign": + if AUTONOMOUS_REASSIGN and decision.confidence >= REASSIGN_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"] + else: + # Not authorised to auto-reassign (mode off or low confidence): + # hand the proposal to a human rather than taking the action. + await self._escalate_to_human(miler_id, booking_id, minutes_stalled, decision) + actions.append("escalate_reassign_for_review") + + else: # "escalate" + await self._escalate_to_human(miler_id, booking_id, minutes_stalled, decision) + actions.append("escalate") + + self._record_stall_exception( + miler_id, booking_id, minutes_stalled, + resolution=f"[{decision.action}] {decision.reasoning}", + actions=actions, + ) + + # ── Stall: action helpers ───────────────────────────────────────────────── + + async def _reassign(self, booking_id: str, reason: str) -> bool: + result = await api_post( f"{GO_API_BASE_URL}/api/v1/internal/bookings/{booking_id}/reassign", - json={"reason": "miler_stalled"}, - headers=headers, + json={"reason": reason}, + headers={"X-Internal-Key": INTERNAL_API_KEY}, ) - logger.info(f"Reassign {'OK' if reassign_result is not None else 'FAILED'} for booking {booking_id}") + ok = result is not None + logger.info(f"Reassign {'OK' if ok else 'FAILED'} for booking {booking_id}") + return ok - notify_result = await api_post( + async def _notify_customer(self, booking_id: str, message: str) -> bool: + result = await api_post( f"{GO_API_BASE_URL}/api/v1/internal/notify", - json={ - "booking_id": booking_id, - "message": "We detected a delay, finding you a new miler", - "target": "customer", - }, - headers=headers, + json={"booking_id": booking_id, "message": message, "target": "customer"}, + headers={"X-Internal-Key": INTERNAL_API_KEY}, ) - logger.info(f"Notify {'OK' if notify_result is not None else 'FAILED'} for booking {booking_id}") + ok = result is not None + logger.info(f"Notify {'OK' if ok else 'FAILED'} for booking {booking_id}") + return ok - exc_id = f"EXC-{datetime.now().strftime('%Y%m%d')}-{uuid.uuid4().hex[:6].upper()}" + async def _escalate_to_human(self, miler_id, booking_id, minutes_stalled, decision: StallDecision): + logger.warning( + f"[EXCEPTION] Escalating booking {booking_id} to human dispatcher " + f"(proposed={decision.action}, confidence={decision.confidence:.2f}): {decision.reasoning}" + ) + await self.send_message( + recipient="JARVIS", + message_type=MessageType.EXCEPTION_DETECTED, + payload={ + "order_id": booking_id, + "miler_id": miler_id, + "exception_type": ExceptionType.MILER_STALLED.value, + "severity": ExceptionSeverity.HIGH.value, + "proposed_action": decision.action, + "confidence": decision.confidence, + "reasoning": decision.reasoning, + "minutes_stalled": minutes_stalled, + "urgency": "high", + }, + ) + + async def _gather_stall_facts(self, miler_id, booking_id) -> Dict[str, Any]: + """Read-only context gathering: pull real signals from Postgres + Redis + to inform the decision. Best-effort — any failure is logged and skipped, + never raised, so a thin-context decision still runs (which skews safe).""" + facts: Dict[str, Any] = {} + now = _utcnow() + + # Booking phase + age (Postgres, read-only). + if self._pg: + try: + async with self._pg.acquire() as conn: + row = await conn.fetchrow( + "SELECT status, createdat FROM pickupbookings WHERE bookingid::text = $1", + str(booking_id), + ) + if row: + facts["booking_status"] = row["status"] + created = row["createdat"] + if created is not None: + if created.tzinfo is None: + created = created.replace(tzinfo=timezone.utc) + facts["minutes_since_booking_created"] = round((now - created).total_seconds() / 60, 1) + except Exception as e: + logger.warning(f"stall facts: booking query failed for {booking_id}: {e}") + + # Movement freshness (Redis, read-only). + try: + movement = await self._redis.hgetall(f"miler:{miler_id}:movement") + except Exception as e: + logger.warning(f"stall facts: redis read failed for miler {miler_id}: {e}") + movement = {} + + if movement: + last_ping = _parse_ts(movement.get("updated_at")) + if last_ping: + facts["last_gps_ping_minutes_ago"] = round((now - last_ping).total_seconds() / 60, 1) + unchanged = _parse_ts(movement.get("position_unchanged_since")) + if unchanged: + facts["position_unchanged_minutes"] = round((now - unchanged).total_seconds() / 60, 1) + + # Per-miler stall count today (Redis, read-only). Already includes the + # current stall, since _publish_stall increments before the consumer runs. + try: + raw = await self._redis.get(self._stall_counter_key(miler_id)) + if raw is not None: + facts["stalls_today"] = int(raw) + except Exception as e: + logger.warning(f"stall facts: stall counter read failed for miler {miler_id}: {e}") + + return facts + + def _record_stall_exception(self, miler_id, booking_id, minutes_stalled, resolution, actions): + escalated = any(a.startswith("escalate") for a in actions) + exc_id = f"EXC-{_utcnow().strftime('%Y%m%d')}-{uuid.uuid4().hex[:6].upper()}" self._exceptions[exc_id] = ExceptionRecord( exception_id=exc_id, order_id=booking_id, exception_type=ExceptionType.MILER_STALLED, severity=ExceptionSeverity.HIGH, description=f"Miler {miler_id} stalled {minutes_stalled} min", - detected_at=datetime.now(), resolved_at=datetime.now(), - resolution="Reassignment triggered + customer notified", - assigned_to=None, status="resolved", - actions_taken=["reassign", "notify_customer"], + detected_at=_utcnow(), + resolved_at=None if escalated else _utcnow(), + resolution=resolution, + assigned_to=None, + status="escalated" if escalated else "resolved", + actions_taken=list(actions), ) # ── Postgres helper ─────────────────────────────────────────────────────── diff --git a/agents/fleet_agent.py b/agents/fleet_agent.py index 38b9a3c..80e4069 100644 --- a/agents/fleet_agent.py +++ b/agents/fleet_agent.py @@ -117,9 +117,12 @@ class FleetAgent(SpecializedAgent): vehicle_id = suitable_vehicles[0] vehicle = self._vehicles[vehicle_id] + was_available = vehicle["status"] == "available" vehicle["status"] = "assigned" - if vehicle["hub"] in self._hub_capacity: + # Only move the counter when the availability boundary is actually + # crossed, so it can't drift out of [0, total]. + if was_available and vehicle["hub"] in self._hub_capacity: self._hub_capacity[vehicle["hub"]]["available"] -= 1 assignment_id = f"ASN-{uuid.uuid4().hex[:8].upper()}" @@ -160,11 +163,20 @@ class FleetAgent(SpecializedAgent): return {"status": "error", "message": f"Vehicle {vehicle_id} not found"} vehicle = self._vehicles[vehicle_id] + was_available = vehicle["status"] == "available" vehicle["status"] = "available" - if vehicle["hub"] in self._hub_capacity: + # Only increment when the vehicle was actually unavailable, so a + # double-release (or releasing an available vehicle) can't exceed total. + if not was_available and vehicle["hub"] in self._hub_capacity: self._hub_capacity[vehicle["hub"]]["available"] += 1 + # Close any active assignment for this vehicle so it isn't reported as + # still on the road (and stale 'active' rows don't accumulate). + for assignment in self._assignments.values(): + if assignment.vehicle_id == vehicle_id and assignment.status == "active": + assignment.status = "returned" + logger.info(f"Vehicle {vehicle_id} released") return {"status": "released", "vehicle_id": vehicle_id, "hub": vehicle["hub"]} @@ -239,7 +251,12 @@ class FleetAgent(SpecializedAgent): else: self._maintenance_schedule[vehicle_id] = datetime.now() + timedelta(days=3) - self._vehicles[vehicle_id]["status"] = "maintenance" + # Taking a vehicle out of service reduces availability (only if it was + # counted as available), keeping hub capacity accurate. + vehicle = self._vehicles[vehicle_id] + if vehicle["status"] == "available" and vehicle["hub"] in self._hub_capacity: + self._hub_capacity[vehicle["hub"]]["available"] -= 1 + vehicle["status"] = "maintenance" return { "status": "scheduled", diff --git a/agents/order_agent.py b/agents/order_agent.py index 708c651..ecd3144 100644 --- a/agents/order_agent.py +++ b/agents/order_agent.py @@ -254,11 +254,11 @@ class OrderAgent(SpecializedAgent): if email and not self._is_valid_email(email): warnings.append("Invalid email format") - pickup_pincode = order_data.get("pickup_address", {}).get("pincode", "") + pickup_pincode = self._addr_pincode(order_data, "pickup_address") if not self._is_valid_pincode(pickup_pincode): errors.append("Invalid pickup pincode") - delivery_pincode = order_data.get("delivery_address", {}).get("pincode", "") + delivery_pincode = self._addr_pincode(order_data, "delivery_address") if not self._is_valid_pincode(delivery_pincode): errors.append("Invalid delivery pincode") @@ -280,8 +280,8 @@ class OrderAgent(SpecializedAgent): } def _categorize_order_data(self, order_data: Dict) -> Dict[str, Any]: - pickup_pincode = order_data.get("pickup_address", {}).get("pincode", "")[:3] - delivery_pincode = order_data.get("delivery_address", {}).get("pincode", "")[:3] + pickup_pincode = self._addr_pincode(order_data, "pickup_address")[:3] + delivery_pincode = self._addr_pincode(order_data, "delivery_address")[:3] if pickup_pincode == delivery_pincode: zone_type = ZoneType.LAST_MILE.value @@ -330,7 +330,18 @@ class OrderAgent(SpecializedAgent): import re return bool(re.match(self._validation_rules["email_pattern"], email)) - def _is_valid_pincode(self, pincode: str) -> bool: + @staticmethod + def _addr_pincode(order_data: Dict, key: str) -> str: + """Pincode as a string, tolerant of a missing/None address block or a + numeric pincode (valid JSON may send either). Never raises.""" + addr = order_data.get(key) or {} + if not isinstance(addr, dict): + return "" + pincode = addr.get("pincode") + return str(pincode) if pincode is not None else "" + + def _is_valid_pincode(self, pincode) -> bool: + pincode = str(pincode) if pincode is not None else "" return len(pincode) == self._validation_rules["pincode_length"] and pincode.isdigit() async def think(self, context: str, options: List[str] = None) -> str: diff --git a/core/agent.py b/core/agent.py index 9b2bbed..a9134f0 100644 --- a/core/agent.py +++ b/core/agent.py @@ -1,5 +1,6 @@ """Base Agent class - All agents inherit from this.""" import asyncio +import time import uuid from datetime import datetime from typing import Dict, List, Optional, Any @@ -28,6 +29,7 @@ class Agent(ABC): self._task_queue: asyncio.Queue = asyncio.Queue() self._running = False self._task_handlers: Dict[str, callable] = {} + self._last_state_emit = 0.0 # Register with message bus message_bus.register_agent(self) @@ -49,6 +51,9 @@ class Agent(ABC): except asyncio.TimeoutError: await self._heartbeat() + if time.monotonic() - self._last_state_emit >= 5.0: + await self._emit_state() + except Exception as e: logger.error(f"Error in agent {self.agent_id}: {e}") self.state.status = "error" @@ -64,6 +69,8 @@ class Agent(ABC): """Process a task from the queue.""" self.state.status = "working" self.state.current_task = task.task_id + started = time.monotonic() + await self._emit_state() try: if task.task_type in self._task_handlers: @@ -85,9 +92,19 @@ class Agent(ABC): self.state.tasks_failed += 1 finally: + duration_ms = int((time.monotonic() - started) * 1000) + await message_bus.publish_telemetry("task", { + "agent_id": self.agent_id, + "task_id": task.task_id, + "task_type": task.task_type, + "status": task.status, + "error": task.error, + "duration_ms": duration_ms, + }) self.state.current_task = None self.state.last_active = datetime.now() self.state.status = "idle" + await self._emit_state() @abstractmethod async def handle_task(self, task: AgentTask) -> Dict[str, Any]: @@ -98,6 +115,18 @@ class Agent(ABC): """Called periodically when idle. Override for custom behavior.""" pass + async def _emit_state(self): + """Publish this agent's live state as telemetry (for the command center).""" + self._last_state_emit = time.monotonic() + await message_bus.publish_telemetry("agent", { + "agent_id": self.agent_id, + "agent_type": self.agent_type, + "status": self.state.status, + "current_task": self.state.current_task, + "tasks_completed": self.state.tasks_completed, + "tasks_failed": self.state.tasks_failed, + }) + def register_task_handler(self, task_type: str, handler: callable): self._task_handlers[task_type] = handler @@ -124,6 +153,39 @@ class Agent(ABC): async def receive_messages(self) -> List[AgentMessage]: return await message_bus.get_messages(self.agent_id) + async def deliver(self, message: AgentMessage): + """Consume an inbound directed message. + + This is what makes agent-to-agent messaging actually work: the message + bus calls it for every directed message addressed to this agent. Messages + that carry a ``task_type`` in their payload are enqueued as tasks so + ``handle_task`` processes them exactly like a submitted task; everything + else is handed to ``handle_message`` so the agent can react. Without this, + directed messages sat unread in the bus queue forever. + """ + payload = message.payload if isinstance(message.payload, dict) else {} + task_type = payload.get("task_type") + if task_type: + await self._task_queue.put(AgentTask( + task_id=f"{message.correlation_id or message.message_id}:{task_type}", + agent_type=self.agent_type, + task_type=task_type, + data=payload, + )) + else: + try: + await self.handle_message(message) + except Exception as e: + logger.error(f"{self.agent_id} handle_message error: {e}") + + async def handle_message(self, message: AgentMessage): + """React to a non-task directed message (a notification). Default just + logs it; agents override this to act on events like EXCEPTION_DETECTED.""" + logger.debug( + f"{self.agent_id} received {message.message_type.value} " + f"from {message.sender} (no task_type; not handled)" + ) + def subscribe_to(self, message_type: MessageType, callback: callable): message_bus.subscribe(message_type, callback) @@ -151,6 +213,24 @@ class MasterAgent(Agent): self._sub_agents[agent.agent_id] = agent logger.info(f"JARVIS registered sub-agent: {agent.agent_id}") + async def handle_message(self, message: AgentMessage): + """Surface notifications from sub-agents. Escalations (proposals a human + must decide) are logged at WARNING and recorded so they are visible + rather than silently dropped — the endpoint of the human-review path.""" + if message.message_type == MessageType.EXCEPTION_DETECTED: + logger.warning(f"JARVIS: exception escalation from {message.sender}: {message.payload}") + else: + logger.info(f"JARVIS: {message.message_type.value} from {message.sender}") + self._decision_log.append({ + "timestamp": datetime.now(), + "action": "received_notification", + "from": message.sender, + "type": message.message_type.value, + "payload": message.payload, + }) + if len(self._decision_log) > 500: + self._decision_log = self._decision_log[-500:] + async def handle_task(self, task: AgentTask) -> Dict[str, Any]: if task.task_type == "orchestrate_order": return await self._orchestrate_order(task) diff --git a/core/http_client.py b/core/http_client.py index f83e58d..9327c98 100644 --- a/core/http_client.py +++ b/core/http_client.py @@ -43,8 +43,22 @@ async def api_patch(url: str, **kwargs) -> Optional[Dict[str, Any]]: return await _request("patch", url, **kwargs) +async def _safe_text(resp) -> str: + try: + return (await resp.text())[:300] + except Exception: + return "" + + async def _request(method: str, url: str, **kwargs) -> Optional[Dict[str, Any]]: - """Execute an HTTP request with up to 3 attempts and exponential backoff.""" + """Execute an HTTP request with up to 3 attempts and exponential backoff. + + Returns the parsed body ONLY on a 2xx/3xx response. A 4xx is a client error + and returns None immediately (no retry — it won't fix itself). A 5xx is + retried like a network error, then returns None. This is the difference + between "the action succeeded" and "the server rejected it": callers gate on + a non-None result, so an error response must not look like success. + """ max_retries = 3 backoff = 1.0 session = await get_session() @@ -52,6 +66,25 @@ async def _request(method: str, url: str, **kwargs) -> Optional[Dict[str, Any]]: for attempt in range(max_retries): try: async with getattr(session, method)(url, **kwargs) as resp: + if resp.status >= 500: + body = await _safe_text(resp) + if attempt < max_retries - 1: + logger.warning( + f"{method.upper()} {url} -> {resp.status} (server error), " + f"retrying in {backoff:.0f}s: {body}" + ) + await asyncio.sleep(backoff) + backoff *= 2 + continue + logger.error( + f"{method.upper()} {url} -> {resp.status} after {max_retries} attempts: {body}" + ) + return None + if resp.status >= 400: + logger.error( + f"{method.upper()} {url} -> {resp.status} (client error): {await _safe_text(resp)}" + ) + return None if resp.content_type == "application/json": return await resp.json() return {"status_code": resp.status} diff --git a/core/llm.py b/core/llm.py new file mode 100644 index 0000000..7c42a1a --- /dev/null +++ b/core/llm.py @@ -0,0 +1,220 @@ +""" +LLM-backed decision-making for agents. + +Currently backed by Claude (Anthropic). It is deliberately kept behind a thin +function boundary (`decide_stall_response`) and a `dataclass` result so the +model — or the whole provider — can be swapped later (e.g. a self-hosted model) +without touching any agent logic. + +Configuration (env): + ANTHROPIC_API_KEY required — the Anthropic API key + LLM_MODEL default "claude-opus-4-8" + LLM_EFFORT default "high" (low | medium | high | xhigh | max) +""" +import json +import os +from dataclasses import dataclass +from datetime import datetime, timezone +from typing import Optional + +from core.logger import logger + +LLM_MODEL = os.getenv("LLM_MODEL", "claude-opus-4-8") +LLM_EFFORT = os.getenv("LLM_EFFORT", "high") +# Thinking tokens count against max_tokens; adaptive thinking at high effort can +# use a lot, so the budget must be generous or the JSON output gets truncated. +LLM_MAX_TOKENS = int(os.getenv("LLM_MAX_TOKENS", "8192")) +LLM_TIMEOUT_S = float(os.getenv("LLM_TIMEOUT_S", "30")) + +_client = None + + +def _get_client(): + """Lazily construct the *async* Anthropic client so a missing dependency or + key never breaks agent startup — only the actual decision call fails, and the + caller falls back to deterministic behaviour. Async so the decision never + blocks the event loop (a sync client would freeze every other coroutine for + the full duration of the call).""" + global _client + if _client is None: + import anthropic # lazy import + _client = anthropic.AsyncAnthropic() # reads ANTHROPIC_API_KEY from env + return _client + + +@dataclass +class StallDecision: + action: str # "reassign" | "notify_only" | "wait" | "escalate" + reasoning: str + confidence: float # 0.0 .. 1.0 + + +DEFAULT_STALL_THRESHOLD_MIN = 10 + + +def build_stall_context(minutes_stalled, *, stall_threshold_min=DEFAULT_STALL_THRESHOLD_MIN, + now=None, facts=None) -> str: + """Assemble the prompt context for a stall decision. + + Shared by the ExceptionAgent and the eval harness so the eval tests exactly + what runs in production. `facts` is an optional dict of extra situational + detail (location, notes, prior-stall count, ...) — today the agent passes + none, but context tools will populate it later. + """ + now = now or datetime.now(timezone.utc) + lines = [ + "A miler has stalled during an active delivery.", + f"- minutes_stalled: {round(float(minutes_stalled), 1)}", + f"- stall_threshold_minutes: {stall_threshold_min}", + f"- current_time_utc: {now.isoformat()}", + f"- hour_of_day_utc: {now.hour}", + ] + for key, value in (facts or {}).items(): + lines.append(f"- {key}: {value}") + lines.append("Decide the best response.") + return "\n".join(lines) + + +_VALID_ACTIONS = {"reassign", "notify_only", "wait", "escalate"} + +_STALL_SYSTEM = """You are the exception controller for a last-mile delivery network. +A "miler" (delivery rider) assigned to an active booking has stopped moving for a while. +Decide the single best response. Choose exactly one action: + +- "wait": the stall is plausibly benign (traffic, a short pickup wait, a customer handoff). Take no action yet. +- "notify_only": the miler can likely recover, but the customer deserves reassurance. Notify the customer; keep the miler. +- "reassign": the miler is genuinely stuck and unlikely to recover soon. Hand the booking to another miler. This is customer-visible and hard to reverse — choose it only when the evidence clearly supports it. +- "escalate": the situation is ambiguous or high-stakes; a human dispatcher should decide. + +Weigh how long it has been stalled relative to the threshold, the time of day, and how far into the job it is. +Prefer the least disruptive action that fits the evidence. Report an honest confidence in [0, 1] — +a low confidence is a signal to escalate rather than act.""" + +_STALL_SCHEMA = { + "type": "object", + "properties": { + "action": {"type": "string", "enum": sorted(_VALID_ACTIONS)}, + "reasoning": {"type": "string"}, + "confidence": {"type": "number"}, + }, + "required": ["action", "reasoning", "confidence"], + "additionalProperties": False, +} + + +async def _decide(system: str, schema: dict, valid_actions, context: str): + """Run one structured-output decision against Claude. + + Returns (action, reasoning, confidence), or None on any failure (network, + timeout, refusal, truncation, bad parse) so callers can fall back to + deterministic behaviour. Retries once on a transient error before giving up. + """ + resp = None + last_err = None + 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, + ) + break + except Exception as e: + last_err = e + logger.warning(f"LLM decision request failed (attempt {attempt + 1}/2): {e}") + if resp is None: + logger.error(f"LLM decision failed after retries: {last_err}") + return None + + stop = getattr(resp, "stop_reason", None) + if stop == "refusal": + logger.warning("LLM refused the decision request") + return None + if stop == "max_tokens": + logger.error("LLM decision truncated (max_tokens); raise LLM_MAX_TOKENS. Treating as no decision.") + return None + + try: + text = next(b.text for b in resp.content if b.type == "text") + data = json.loads(text) + action = data["action"] + if action not in valid_actions: + raise ValueError(f"unexpected action {action!r}") + return action, str(data.get("reasoning", "")), float(data.get("confidence", 0.0)) + except (StopIteration, KeyError, ValueError, json.JSONDecodeError) as e: + logger.error(f"LLM decision parse failed: {e}") + return None + + +async def decide_stall_response(context: str) -> 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) + return StallDecision(*r) if r else None + + +# ── Assignment-failure decision (DispatchAgent) ────────────────────────────── + +@dataclass +class AssignmentDecision: + action: str # "monitor" | "notify_customer" | "ops_alert" | "escalate" + reasoning: str + confidence: float + + +_ASSIGNMENT_ACTIONS = {"monitor", "notify_customer", "ops_alert", "escalate"} + +_ASSIGNMENT_SYSTEM = """You are the dispatch controller for a last-mile delivery network. +The backend tried to assign a booking to a miler (rider) — both its AI assignment and the +fallback failed, so no rider was assigned. Decide how to react. Choose exactly one action: + +- "monitor": likely a transient gap (a momentary lack of free riders); the backend will retry. Take no action. +- "notify_customer": a real but recoverable delay; tell the customer we're finding a rider so they aren't left guessing. +- "ops_alert": a genuine coverage gap in this zone (repeated failures, no nearby riders). Alert operations to onboard or redirect riders here. +- "escalate": ambiguous or contradictory — e.g. riders ARE nearby yet assignment keeps failing, which suggests a systemic issue rather than a coverage gap. A human dispatcher should look. + +Weigh how many times assignment has failed in this zone today, how far the nearest available rider is, +the time of day, and whether coordinates were even available. A single failure with a rider nearby is +usually transient; repeated failures with no nearby rider is a coverage gap; repeated failures *despite* +nearby riders is not a coverage gap and warrants a human. Report an honest confidence in [0, 1].""" + +_ASSIGNMENT_SCHEMA = { + "type": "object", + "properties": { + "action": {"type": "string", "enum": sorted(_ASSIGNMENT_ACTIONS)}, + "reasoning": {"type": "string"}, + "confidence": {"type": "number"}, + }, + "required": ["action", "reasoning", "confidence"], + "additionalProperties": False, +} + + +def build_assignment_failure_context(facts, *, now=None) -> str: + """Assemble the prompt context for an assignment-failure decision. Shared by + DispatchAgent and the eval harness so the eval tests exactly what runs.""" + now = now or datetime.now(timezone.utc) + lines = [ + "A booking could not be assigned to a miler — the backend's AI and fallback assignment both failed.", + f"- current_time_utc: {now.isoformat()}", + f"- hour_of_day_utc: {now.hour}", + ] + for key, value in (facts or {}).items(): + lines.append(f"- {key}: {value}") + lines.append("Decide the best response.") + return "\n".join(lines) + + +async def decide_assignment_failure(context: str) -> 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) + return AssignmentDecision(*r) if r else None diff --git a/core/message_bus.py b/core/message_bus.py index 6c34545..a79ac9a 100644 --- a/core/message_bus.py +++ b/core/message_bus.py @@ -171,6 +171,22 @@ class MessageBus: ) -> str: return await self.send_to_agent(sender, "ALL", message_type, payload, correlation_id) + async def _deliver_to_agent(self, agent_id: str, message: AgentMessage): + """Hand a directed message to a locally-registered agent so it is + actually consumed (via its ``deliver``). Only if no such agent exists in + this process does it fall back to the pull queue — so a message is never + silently lost, but is delivered live whenever the recipient is present.""" + agent = self._agents.get(agent_id) + if agent is not None and hasattr(agent, "deliver"): + try: + await agent.deliver(message) + return + except Exception as e: + logger.error(f"Delivery to {agent_id} failed: {e}") + return + async with self._lock: + self._queues[agent_id].append(QueuedMessage(message)) + async def get_messages(self, agent_id: str) -> List[AgentMessage]: async with self._lock: queued = self._queues.pop(agent_id, []) @@ -179,6 +195,26 @@ class MessageBus: async def peek_messages(self, agent_id: str) -> List[AgentMessage]: return [q.message for q in self._queues.get(agent_id, [])] + # ------------------------------------------------------------------ # + # Telemetry # + # ------------------------------------------------------------------ # + + async def publish_telemetry(self, kind: str, payload: Dict[str, Any]): + """ + Fire-and-forget observability event on plain NATS (subject `telemetry.`). + + Deliberately NOT JetStream: telemetry is ephemeral fan-out for dashboards + and must never accumulate in the persistent `logistics` stream. No-op + when NATS is not connected — telemetry must never affect agent behavior. + """ + if self._nc is None or not self._nc.is_connected: + return + try: + body = {"kind": kind, "ts": datetime.now().isoformat(), **payload} + await self._nc.publish(f"telemetry.{kind}", json.dumps(body).encode()) + except Exception as e: + logger.debug(f"Telemetry publish failed [{kind}]: {e}") + # ------------------------------------------------------------------ # # Hooks # # ------------------------------------------------------------------ # @@ -266,8 +302,7 @@ class MessageBus: async def handler(msg): try: agent_msg = AgentMessage.from_json(msg.data.decode()) - async with self._lock: - self._queues[agent_id].append(QueuedMessage(agent_msg)) + await self._deliver_to_agent(agent_id, agent_msg) except Exception as e: logger.error(f"NATS agent-sub decode error [{agent_id}]: {e}") finally: @@ -282,8 +317,7 @@ class MessageBus: async def _local_dispatch(self, message: AgentMessage): """In-process dispatch used when NATS is not connected.""" if message.recipient != "ALL": - async with self._lock: - self._queues[message.recipient].append(QueuedMessage(message)) + await self._deliver_to_agent(message.recipient, message) else: for callback in list(self._subscribers.get(message.message_type, [])): try: diff --git a/dashboard/command_center.py b/dashboard/command_center.py new file mode 100644 index 0000000..e5f5a03 --- /dev/null +++ b/dashboard/command_center.py @@ -0,0 +1,300 @@ +""" +LogiFlow Command Center — real-time JARVIS-style mission control. + +A read-only tap on the live infrastructure: + - subscribes to `logistics.>` (every real agent-to-agent message on the bus) + - subscribes to `telemetry.>` (agent state + task lifecycle events) +and streams everything to the browser over WebSocket. + +Runs as a separate process — it never interferes with the agents: +plain (core) NATS subscriptions only, no durable consumers, no JetStream acks. + +Run: + python main.py --command-center # http://localhost:8600 + python dashboard/command_center.py # same, standalone +""" +import asyncio +import json +import sys +import time +from collections import deque +from datetime import datetime +from pathlib import Path +from typing import Any, Dict, List, Set + +# Allow running as `python dashboard/command_center.py` +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) + +from dotenv import load_dotenv +load_dotenv() + +import nats +import uvicorn +from fastapi import FastAPI, WebSocket, WebSocketDisconnect +from fastapi.responses import FileResponse + +from config.system_config import NATS_HOST, NATS_PORT, NATS_USER, NATS_PASSWORD, AGENT_CONFIG +from core.logger import logger + +STATIC_DIR = Path(__file__).resolve().parent / "static" +PORT = 8600 + +KNOWN_AGENTS = [ + "JARVIS", "ORDER_AGENT", "DISPATCH_AGENT", "FLEET_AGENT", + "HUB_AGENT", "CUSTOMER_AGENT", "EXCEPTION_AGENT", "ROUTE_OPTIMIZER", +] + +AGENT_DESCRIPTIONS = {cfg["name"]: cfg["description"] for cfg in AGENT_CONFIG.values()} + + +class CommandCenterState: + """Aggregated live view of the agent system, fed by NATS events.""" + + def __init__(self): + self.bus_connected = False + self.agents: Dict[str, Dict[str, Any]] = { + a: {"agent_id": a, "status": "offline", "current_task": None, + "tasks_completed": 0, "tasks_failed": 0, "last_seen": None, + "msgs_sent": 0, "msgs_received": 0, + "description": AGENT_DESCRIPTIONS.get(a, "")} + for a in KNOWN_AGENTS + } + self.agent_recent: Dict[str, deque] = {a: deque(maxlen=40) for a in KNOWN_AGENTS} + self.agent_durations: Dict[str, deque] = {a: deque(maxlen=50) for a in KNOWN_AGENTS} + self.events: deque = deque(maxlen=300) + self.totals = {"messages": 0, "tasks_completed": 0, "tasks_failed": 0, "exceptions": 0} + self.msg_times: deque = deque(maxlen=5000) # unix ts of bus messages, for rate + + def rate_per_min(self) -> int: + cutoff = time.time() - 60 + return sum(1 for t in self.msg_times if t >= cutoff) + + def snapshot(self) -> Dict[str, Any]: + return { + "type": "snapshot", + "bus_connected": self.bus_connected, + "agents": self.agents, + "events": list(self.events), + "totals": self.totals, + "rate_per_min": self.rate_per_min(), + } + + +state = CommandCenterState() +websockets: Set[WebSocket] = set() +app = FastAPI(title="LogiFlow Command Center") + + +async def broadcast(event: Dict[str, Any]): + """Push one event to every connected browser.""" + dead = [] + for ws in websockets: + try: + await ws.send_text(json.dumps(event)) + except Exception: + dead.append(ws) + for ws in dead: + websockets.discard(ws) + + +def _compact(payload: Any, limit: int = 240) -> str: + try: + s = json.dumps(payload, default=str) + except Exception: + s = str(payload) + return s[:limit] + ("…" if len(s) > limit else "") + + +async def _on_bus_message(msg): + """Any real agent message on `logistics.>` (direct or broadcast).""" + now = time.time() + state.totals["messages"] += 1 + state.msg_times.append(now) + try: + data = json.loads(msg.data.decode()) + except Exception: + data = {"raw": msg.data.decode(errors="replace")} + + msg_type = data.get("message_type", msg.subject.split(".")[-1]) + if msg_type == "EXCEPTION_DETECTED": + state.totals["exceptions"] += 1 + + event = { + "type": "bus", + "ts": data.get("timestamp", datetime.now().isoformat()), + "sender": data.get("sender", "?"), + "recipient": data.get("recipient", msg.subject.split(".")[-1] if ".direct." in msg.subject else "ALL"), + "msg_type": msg_type, + "summary": _compact(data.get("payload", data)), + "correlation_id": data.get("correlation_id"), + } + for role, key in (("sender", "msgs_sent"), ("recipient", "msgs_received")): + aid = event[role] + if aid in state.agents: + state.agents[aid][key] = state.agents[aid].get(key, 0) + 1 + state.agent_recent.setdefault(aid, deque(maxlen=40)).append(event) + state.events.append(event) + await broadcast(event) + + +async def _on_telemetry(msg): + """Agent state / task lifecycle events on `telemetry.>`.""" + try: + data = json.loads(msg.data.decode()) + except Exception: + return + + kind = data.get("kind") + if kind == "agent": + agent_id = data.get("agent_id") + if not agent_id: + return + entry = state.agents.setdefault(agent_id, { + "msgs_sent": 0, "msgs_received": 0, + "description": AGENT_DESCRIPTIONS.get(agent_id, ""), + }) + entry.update({ + "agent_id": agent_id, + "status": data.get("status", "idle"), + "current_task": data.get("current_task"), + "tasks_completed": data.get("tasks_completed", 0), + "tasks_failed": data.get("tasks_failed", 0), + "last_seen": data.get("ts"), + }) + await broadcast({"type": "agent", **entry}) + + elif kind == "task": + if data.get("status") == "completed": + state.totals["tasks_completed"] += 1 + elif data.get("status") == "failed": + state.totals["tasks_failed"] += 1 + event = { + "type": "task", + "ts": data.get("ts", datetime.now().isoformat()), + "agent_id": data.get("agent_id"), + "task_id": data.get("task_id"), + "task_type": data.get("task_type"), + "status": data.get("status"), + "error": data.get("error"), + "duration_ms": data.get("duration_ms"), + } + aid = event["agent_id"] + if aid: + state.agent_recent.setdefault(aid, deque(maxlen=40)).append(event) + if event["duration_ms"] is not None: + state.agent_durations.setdefault(aid, deque(maxlen=50)).append(event["duration_ms"]) + state.events.append(event) + await broadcast(event) + + +async def nats_tap(): + """Connect to NATS (with retry forever) and hold plain subscriptions.""" + while True: + try: + nc = await nats.connect( + servers=[f"nats://{NATS_HOST}:{NATS_PORT}"], + user=NATS_USER, + password=NATS_PASSWORD, + max_reconnect_attempts=-1, + disconnected_cb=_bus_down, + reconnected_cb=_bus_up, + ) + await nc.subscribe("logistics.>", cb=_on_bus_message) + await nc.subscribe("telemetry.>", cb=_on_telemetry) + await _bus_up() + logger.info(f"Command Center tapped into NATS at {NATS_HOST}:{NATS_PORT}") + while nc.is_connected or nc.is_reconnecting: + await asyncio.sleep(2) + await _bus_down() + except Exception as e: + logger.warning(f"Command Center NATS connect failed: {e} — retrying in 5s") + await _bus_down() + await asyncio.sleep(5) + + +async def _bus_up(): + state.bus_connected = True + await broadcast({"type": "sys", "bus_connected": True}) + + +async def _bus_down(): + if state.bus_connected: + logger.warning("Command Center lost NATS connection") + state.bus_connected = False + await broadcast({"type": "sys", "bus_connected": False}) + + +async def rate_ticker(): + """Push the message rate + mark stale agents offline every 2s.""" + while True: + await asyncio.sleep(2) + now = datetime.now() + for entry in state.agents.values(): + last = entry.get("last_seen") + if last and entry.get("status") != "offline": + try: + age = (now - datetime.fromisoformat(last)).total_seconds() + if age > 15: + entry["status"] = "offline" + await broadcast({"type": "agent", **entry}) + except ValueError: + pass + if websockets: + await broadcast({ + "type": "rate", + "rate_per_min": state.rate_per_min(), + "totals": state.totals, + }) + + +@app.on_event("startup") +async def startup(): + asyncio.create_task(nats_tap()) + asyncio.create_task(rate_ticker()) + + +@app.get("/") +async def index(): + return FileResponse(STATIC_DIR / "command_center.html") + + +@app.get("/api/state") +async def api_state(): + return state.snapshot() + + +@app.get("/api/agent/{agent_id}") +async def api_agent(agent_id: str): + """Full dossier for one agent — used by the UI detail panel.""" + entry = state.agents.get(agent_id) + if not entry: + return {"error": "unknown agent"} + durations = list(state.agent_durations.get(agent_id, [])) + return { + **entry, + "avg_task_ms": round(sum(durations) / len(durations)) if durations else None, + "recent": list(state.agent_recent.get(agent_id, []))[-30:], + } + + +@app.websocket("/ws") +async def ws_endpoint(ws: WebSocket): + await ws.accept() + websockets.add(ws) + await ws.send_text(json.dumps(state.snapshot())) + try: + while True: + await ws.receive_text() # keepalive pings from the browser + except WebSocketDisconnect: + pass + finally: + websockets.discard(ws) + + +def run(port: int = PORT): + logger.info(f"LogiFlow Command Center → http://localhost:{port}") + uvicorn.run(app, host="0.0.0.0", port=port, log_level="warning") + + +if __name__ == "__main__": + run() diff --git a/dashboard/static/command_center.html b/dashboard/static/command_center.html new file mode 100644 index 0000000..abd08e0 --- /dev/null +++ b/dashboard/static/command_center.html @@ -0,0 +1,848 @@ + + + + + +LogiFlow — Command Center + + + +
+ + + +
+
+
+

Agent operations

+
LogiFlow workspace · autonomous logistics grid
+
+
Live
+
+ +
+
+
Agents online
+
0 / 8
+
+
+
+
Messages / min
+
0
+
0 total on bus
+
+
+
Tasks completed
+
0
+
across all agents
+
+
+
Tasks failed
+
0
+
requires attention
+
+
+
Exceptions
+
0
+
detected on bus
+
+
+ +
+

Agent constellation

· click an agent for its dossier
+ +
+ ⚠ BUS LINK OFFLINE + No NATS connection — waiting to re-establish uplink… +
+
+ +
+
🤖
+

AGENT

+
+
OFFLINE +
+
+
Tasks done
0
+
Tasks failed
0
+
Msgs out / in
0 / 0
+
Avg task time
—
+
+
Current task
+
—
+

Recent activity — this agent

+
+
+
+ +
+

Live activity

0 events
+
Waiting for agent activity…
Tasks, bus messages and exceptions appear here in real time.
+
+ +
+
Bus throughput · messages per 2s · last 3 min
+ +
+
+
+ +
+ + + + diff --git a/docs/handoff-assignment-failed-event.md b/docs/handoff-assignment-failed-event.md new file mode 100644 index 0000000..df3d9cb --- /dev/null +++ b/docs/handoff-assignment-failed-event.md @@ -0,0 +1,129 @@ +# Handoff: publish a `booking.assignment_failed` NATS event + +**To:** whoever owns the Go backend (`assignment` package) +**From:** logistics-ai (Python agents) side +**Why:** the DispatchAgent's assignment-failure decision path is built and tested but **inert** — it consumes an event the backend does not currently publish. + +--- + +## TL;DR + +When miler assignment fails, the backend currently handles it **synchronously** (an `escalate` flag in the decision response) and publishes **nothing** to NATS. Confirmed by inspecting the running binary: `booking.assigned` is published, but there is no `booking.assignment_failed` (or any failure) subject anywhere in the binary. + +To activate the DispatchAgent coverage-gap / escalation logic, the backend needs to **publish one event** on assignment failure and **add its subject to the existing `ASSIGNMENTS` stream**. No further Python changes are required — the consumer, its decision logic, its eval set, and its tests are already in place. + +--- + +## 1. The event to publish + +- **Subject:** `booking.assignment_failed` +- **Transport:** JetStream (the code already uses `PublishAsync` — reuse it; publishing must not block the request path). +- **When:** in the failure branch of assignment — i.e. where the flow currently decides it cannot assign a miler. Based on the binary's symbols, the natural call sites are in the `assignment` package: + - `collectEligibleCandidates` returns empty (no eligible milers), and/or + - `callDecisionEngine` / `selectMilerWithAI` returns escalate, and/or + - `AutoAssignResult` / `tryAssign` / `tryCustomerAssign` resolves to "not assigned". + + Publish alongside the existing success publish (`publishAssignment` / `publishCustomerAssignment`) so success and failure are emitted from symmetric places. + +### Payload (JSON) + +The consumer reads these fields (extra fields are ignored, so it is forward-compatible): + +| Field | Type | Required | Notes | +|---|---|---|---| +| `booking_id` | string or int | **yes** | The booking that could not be assigned. | +| `zone_id` | string | recommended | Used for the per-zone daily failure counter and coverage-gap detection. Defaults to `"unknown"` if omitted. | +| `lat` | float | recommended | Pickup latitude. Enables the nearest-rider coverage sweep (GEORADIUS). | +| `lon` | float | recommended | Pickup longitude. | +| `reason` | string | optional | e.g. `"no_candidates"`, `"ai_escalate"`, `"provider_unavailable"`. Not required today; useful context for the decision and for future logic. | + +**Without `lat`/`lon`** the decision still runs but degrades to the safe side (it can't measure coverage, so it leans toward "monitor"/"escalate"). Include them whenever available. + +### Example + +```json +{ + "booking_id": 24, + "zone_id": "hyderabad", + "lat": 17.4486, + "lon": 78.3908, + "reason": "no_candidates" +} +``` + +This mirrors the shape of the existing `booking.assigned` payload (`booking_id`, `miler_id`, `hub_id`, `confidence`, `reasoning`), just for the failure case. + +--- + +## 2. Add the subject to the `ASSIGNMENTS` stream + +The stream already exists (Go-owned): + +``` +STREAM: ASSIGNMENTS | subjects: ['booking.assigned', 'booking.reassigned'] +``` + +Add `booking.assignment_failed` so JetStream captures it. **Do this before/at the same time as the first publish** — messages published to a subject no stream covers are dropped. + +Wherever the stream is provisioned in code (the same place that created `ASSIGNMENTS`), update its subject list to: + +``` +booking.assigned, booking.reassigned, booking.assignment_failed +``` + +Or, as a one-off from the NATS host: + +```sh +nats stream edit ASSIGNMENTS --subjects="booking.assigned,booking.reassigned,booking.assignment_failed" +``` + +(Prefer updating the provisioning code so it stays correct on redeploys / fresh environments.) + +--- + +## 3. Consumer side — already done, nothing to change + +The DispatchAgent (`agents/dispatch_agent.py`) already: + +- **auto-discovers** the stream for `booking.assignment_failed` (`find_stream_name_by_subject`), so once the subject is on `ASSIGNMENTS` it binds automatically on next restart — no config needed. (An optional `DISPATCH_STREAM` env var can pin the stream name if you ever want to bypass discovery.) +- binds a durable push consumer named **`dispatch-booking-assignment-failed`**. +- on each event: gathers coverage context (nearest-rider sweep + per-zone daily failure counter), asks the model to choose `monitor` / `notify_customer` / `ops_alert` / `escalate`, and acts (customer notification is gated behind `DISPATCH_AGENT_AUTONOMOUS`; default off = proposes internally). + +Delivery is at-least-once (JetStream), and the handler is safe to re-run, but keep publishes reasonable — one event per failed assignment attempt. + +--- + +## 4. How to verify end to end + +1. Deploy the DispatchAgent stream fix (the binding change) and the backend change together. +2. Confirm the subject is on the stream: + ```sh + nats stream info ASSIGNMENTS # subjects should include booking.assignment_failed + ``` +3. Trigger a real (or forced) assignment failure, or publish a test event: + ```sh + nats pub booking.assignment_failed '{"booking_id":9999,"zone_id":"testzone","lat":17.44,"lon":78.39,"reason":"no_candidates"}' + ``` +4. In the DispatchAgent logs, you should see: + ``` + [DISPATCH] booking.assignment_failed — booking=9999 zone=testzone + [DISPATCH] booking=9999 context facts: {...} + [DISPATCH] booking=9999 zone=testzone decision= confidence= reasoning='...' + ``` +5. `nats consumer info ASSIGNMENTS dispatch-booking-assignment-failed` should show delivered/ack'd counts incrementing. + +Once you see that decision line on a real failure, the assignment-failure path is live — and it can be graduated to autonomy the same way as the stall path (turn on `DISPATCH_AGENT_AUTONOMOUS` only after the eval passes on real cases). + +--- + +## Reference — current NATS topology (for context) + +``` +ASSIGNMENTS booking.assigned, booking.reassigned (add: booking.assignment_failed) +BOOKINGS api.v1.bookings.create/update/cancel +CHAT chat.room.created/closed +NOTIFICATIONS notification.send +STATUS booking.status.updated +TRACKING miler.location.updated, miler.stalled (ExceptionAgent consumes) +logistics logistics.> (internal agent-to-agent bus) +``` diff --git a/evals/__init__.py b/evals/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/evals/_harness.py b/evals/_harness.py new file mode 100644 index 0000000..0531f48 --- /dev/null +++ b/evals/_harness.py @@ -0,0 +1,131 @@ +""" +Shared scoring / reporting for agent decision evals. + +A concrete eval (stall, assignment, ...) supplies two callables: + - to_context(case) -> str : turn a dataset case into the model prompt context + - decide(context) -> obj : the decision function (returns an object with + .action / .reasoning / .confidence, or None) + +and calls run_cli(...). Scoring treats a decision as correct if the majority +action across --runs samples is in the case's ``acceptable`` set (headline +"acceptable-rate"); matching the single ``ideal`` is the secondary "exact-rate". +""" +import argparse +import asyncio +import json +import sys +from collections import Counter +from datetime import datetime, timezone +from pathlib import Path +from statistics import mean + + +def load_cases(path: Path): + cases = [] + for lineno, raw in enumerate(path.read_text(encoding="utf-8").splitlines(), 1): + line = raw.strip() + if not line or line.startswith("#"): + continue + try: + cases.append(json.loads(line)) + except json.JSONDecodeError as e: + print(f"skipping malformed case on line {lineno}: {e}", file=sys.stderr) + return cases + + +def case_now(case): + """Fixed date so only hour-of-day varies — keeps contexts reproducible.""" + return datetime(2024, 1, 1, int(case.get("hour_of_day_utc", 12)), 0, tzinfo=timezone.utc) + + +async def evaluate(cases, to_context, decide, runs: int): + rows = [] + for case in cases: + ctx = to_context(case) + decisions = [] + for _ in range(runs): + d = await decide(ctx) + if d is not None: + decisions.append(d) + errors = runs - len(decisions) + + if decisions: + counts = Counter(d.action for d in decisions) + majority, agree = counts.most_common(1)[0] + avg_conf = mean(d.confidence for d in decisions) + reasoning = next((d.reasoning for d in decisions if d.action == majority), "") + else: + majority, agree, avg_conf, reasoning = None, 0, 0.0, "" + + rows.append({ + "id": case["id"], "ideal": case.get("ideal"), "acceptable_set": case["acceptable"], + "got": majority, "agree": agree, "runs": runs, "errors": errors, + "avg_conf": avg_conf, + "acceptable": majority in case["acceptable"] if majority else False, + "exact": majority == case.get("ideal") if majority else False, + "reasoning": reasoning, + }) + return rows + + +def print_report(rows): + print(f"\n{'case':<28} {'ideal':<14} {'got':<14} {'agree':<7} {'conf':<6} result") + print("-" * 82) + for r in rows: + if r["got"] is None: + result = "ERROR (no valid decision)" + elif r["exact"]: + result = "EXACT" + elif r["acceptable"]: + result = "ok (acceptable)" + else: + result = f"MISS (allowed: {', '.join(r['acceptable_set'])})" + print(f"{r['id']:<28} {str(r['ideal']):<14} {str(r['got']):<14} " + f"{r['agree']}/{r['runs']:<5} {r['avg_conf']:<6.2f} {result}") + if r["got"] is not None and not r["acceptable"]: + print(f"{'':<28} └─ why: {r['reasoning'][:110]}") + + scored = [r for r in rows if r["got"] is not None] + n_scored = len(scored) + acc = sum(r["acceptable"] for r in scored) + exact = sum(r["exact"] for r in scored) + print("-" * 82) + print(f"cases: {len(rows)} scored: {n_scored} model-errors: {sum(r['errors'] for r in rows)}") + if n_scored: + print(f"acceptable-rate: {acc}/{n_scored} = {acc / n_scored:.0%}") + print(f"exact-rate: {exact}/{n_scored} = {exact / n_scored:.0%}") + return (acc / n_scored) if n_scored else 0.0 + + +def run_cli(default_cases: Path, to_context, decide, llm_module): + """argparse + dry-run + evaluate + report, shared across eval scripts. + ``llm_module`` is core.llm so --model can override LLM_MODEL for the run.""" + ap = argparse.ArgumentParser(description="Eval an agent decision against labelled cases.") + ap.add_argument("--cases", type=Path, default=default_cases, help="JSONL dataset") + ap.add_argument("--runs", type=int, default=1, help="samples per case (majority vote)") + ap.add_argument("--model", help="override LLM_MODEL for this run") + ap.add_argument("--dry-run", action="store_true", help="print contexts and labels; no API calls") + ap.add_argument("--min-pass-rate", type=float, default=0.0, help="exit non-zero if acceptable-rate below this") + args = ap.parse_args() + + if args.model: + llm_module.LLM_MODEL = args.model + + cases = load_cases(args.cases) + if not cases: + print(f"no cases found in {args.cases}", file=sys.stderr) + sys.exit(2) + + if args.dry_run: + for case in cases: + print(f"\n=== {case['id']} (ideal={case.get('ideal')}, allowed={case['acceptable']}) ===") + print(to_context(case)) + print(f"\n[dry-run] {len(cases)} cases, no API calls made.") + return + + print(f"model: {llm_module.LLM_MODEL} cases: {len(cases)} runs/case: {args.runs}") + rows = asyncio.run(evaluate(cases, to_context, decide, args.runs)) + pass_rate = print_report(rows) + if pass_rate < args.min_pass_rate: + print(f"\nFAIL: acceptable-rate {pass_rate:.0%} < required {args.min_pass_rate:.0%}", file=sys.stderr) + sys.exit(1) diff --git a/evals/assignment_cases.jsonl b/evals/assignment_cases.jsonl new file mode 100644 index 0000000..7cbbb17 --- /dev/null +++ b/evals/assignment_cases.jsonl @@ -0,0 +1,10 @@ +{"id": "first-failure-close-rider", "hour_of_day_utc": 14, "facts": {"zone_id": "hyderabad", "failures_today": 1, "nearest_miler_within_km": 10, "has_coordinates": true}, "acceptable": ["monitor", "notify_customer"], "ideal": "monitor", "rationale": "Single failure with a rider only 10km away — almost certainly transient; let the backend retry."} +{"id": "repeated-zone-gap", "hour_of_day_utc": 12, "facts": {"zone_id": "pune", "failures_today": 4, "nearest_miler_within_km": "none within 30km", "has_coordinates": true}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Repeated failures and no rider within 30km — a genuine coverage gap; alert ops."} +{"id": "second-failure-far-rider", "hour_of_day_utc": 18, "facts": {"zone_id": "bangalore", "failures_today": 2, "nearest_miler_within_km": 30, "has_coordinates": true}, "acceptable": ["notify_customer", "monitor"], "ideal": "notify_customer", "rationale": "A couple of failures with the nearest rider 30km out during rush — real delay; keep the customer informed."} +{"id": "no-coordinates-first", "hour_of_day_utc": 10, "facts": {"zone_id": "unknown", "failures_today": 1, "nearest_miler_within_km": "unknown (no coordinates)", "has_coordinates": false}, "acceptable": ["monitor", "escalate"], "ideal": "monitor", "rationale": "Coverage can't be assessed without coordinates, but a single failure is most likely transient."} +{"id": "many-failures-riders-near", "hour_of_day_utc": 15, "facts": {"zone_id": "mumbai_west", "failures_today": 5, "nearest_miler_within_km": 10, "has_coordinates": true}, "acceptable": ["escalate", "ops_alert"], "ideal": "escalate", "rationale": "Riders ARE nearby yet assignment keeps failing — not a coverage gap; a systemic issue for a human."} +{"id": "night-coverage-gap", "hour_of_day_utc": 2, "facts": {"zone_id": "kolkata", "failures_today": 3, "nearest_miler_within_km": "none within 30km", "has_coordinates": true}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Repeated overnight failures with no nearby rider — a coverage gap worth flagging to ops."} +{"id": "single-far-rider", "hour_of_day_utc": 9, "facts": {"zone_id": "hyderabad", "failures_today": 1, "nearest_miler_within_km": 30, "has_coordinates": true}, "acceptable": ["monitor", "notify_customer"], "ideal": "monitor", "rationale": "One failure, distant rider — likely resolves on retry."} +{"id": "persistent-far-rush", "hour_of_day_utc": 18, "facts": {"zone_id": "mumbai_east", "failures_today": 3, "nearest_miler_within_km": 30, "has_coordinates": true}, "acceptable": ["notify_customer", "ops_alert"], "ideal": "notify_customer", "rationale": "Rush-hour, repeated, distant rider — the customer should hear about the delay."} +{"id": "borderline-count", "hour_of_day_utc": 13, "facts": {"zone_id": "north_delhi", "failures_today": 3, "nearest_miler_within_km": 20, "has_coordinates": true}, "acceptable": ["notify_customer", "ops_alert", "monitor"], "ideal": "notify_customer", "rationale": "Genuinely borderline — a few failures with a moderately-close rider; any of these is defensible."} +{"id": "severe-gap", "hour_of_day_utc": 11, "facts": {"zone_id": "pune", "failures_today": 8, "nearest_miler_within_km": "none within 30km", "has_coordinates": true}, "acceptable": ["ops_alert", "escalate"], "ideal": "ops_alert", "rationale": "Many failures and no rider anywhere near — a clear, severe coverage gap."} diff --git a/evals/assignment_eval.py b/evals/assignment_eval.py new file mode 100644 index 0000000..dbdc702 --- /dev/null +++ b/evals/assignment_eval.py @@ -0,0 +1,30 @@ +""" +Offline eval for the DispatchAgent assignment-failure decision +(core.llm.decide_assignment_failure). + +Same scoring as the stall eval (see evals/_harness.py): acceptable-rate is the +headline, exact-rate secondary. The cases use the real gatherer keys +(zone_id, failures_today, nearest_miler_within_km, has_coordinates) so the eval +reflects what DispatchAgent actually sends. + +Usage: + export ANTHROPIC_API_KEY=... + python -m evals.assignment_eval --runs 3 + python -m evals.assignment_eval --dry-run + python -m evals.assignment_eval --model claude-haiku-4-5 --min-pass-rate 0.9 +""" +from pathlib import Path + +import core.llm as llm +from core.llm import build_assignment_failure_context, decide_assignment_failure +from evals._harness import case_now, run_cli + +CASES_PATH = Path(__file__).with_name("assignment_cases.jsonl") + + +def to_context(case): + return build_assignment_failure_context(case.get("facts", {}), now=case_now(case)) + + +if __name__ == "__main__": + run_cli(CASES_PATH, to_context, decide_assignment_failure, llm) diff --git a/evals/stall_cases.jsonl b/evals/stall_cases.jsonl new file mode 100644 index 0000000..ae65b7c --- /dev/null +++ b/evals/stall_cases.jsonl @@ -0,0 +1,15 @@ +{"id": "just-over-traffic", "minutes_stalled": 11, "hour_of_day_utc": 18, "facts": {"location": "major intersection", "note": "stop-go traffic reported in the area"}, "acceptable": ["wait", "notify_only"], "ideal": "wait", "rationale": "Just over threshold during evening rush; most likely ordinary traffic."} +{"id": "moderate-at-delivery", "minutes_stalled": 18, "hour_of_day_utc": 13, "facts": {"location": "residential delivery address", "note": "arrived at the delivery point"}, "acceptable": ["notify_only", "wait"], "ideal": "notify_only", "rationale": "Stopped at the delivery point — likely completing handoff; reassure the customer, keep the miler."} +{"id": "long-highway-no-stop", "minutes_stalled": 42, "hour_of_day_utc": 15, "facts": {"location": "on a highway stretch", "note": "no position change, not at any known stop"}, "acceptable": ["reassign", "escalate"], "ideal": "reassign", "rationale": "Long stall mid-highway, not at a stop — unlikely to recover on its own."} +{"id": "late-night-thin-coverage", "minutes_stalled": 25, "hour_of_day_utc": 2, "facts": {"note": "very few active milers in this zone at this hour"}, "acceptable": ["escalate", "notify_only"], "ideal": "escalate", "rationale": "Reassignment is risky with thin overnight coverage; a human should decide."} +{"id": "very-long-no-context", "minutes_stalled": 60, "hour_of_day_utc": 11, "facts": {}, "acceptable": ["reassign", "escalate"], "ideal": "escalate", "rationale": "Very long stall but zero situational detail — low information favors a human."} +{"id": "repeat-staller", "minutes_stalled": 15, "hour_of_day_utc": 10, "facts": {"stalls_today": 3, "position_unchanged_minutes": 15}, "acceptable": ["reassign", "escalate"], "ideal": "reassign", "rationale": "Real gatherer signal: 3rd stall today for this miler — a pattern that says this miler is unreliable right now."} +{"id": "near-completion", "minutes_stalled": 20, "hour_of_day_utc": 16, "facts": {"note": "route is ~90% complete, miler is at the final delivery cluster"}, "acceptable": ["wait", "notify_only"], "ideal": "notify_only", "rationale": "Almost done — reassigning would waste a near-complete trip."} +{"id": "vehicle-breakdown", "minutes_stalled": 35, "hour_of_day_utc": 12, "facts": {"location": "roadside", "note": "miler reported a vehicle breakdown"}, "acceptable": ["reassign"], "ideal": "reassign", "rationale": "Explicit breakdown — the miler cannot recover; must reassign."} +{"id": "ambiguous-mid", "minutes_stalled": 22, "hour_of_day_utc": 14, "facts": {"location": "commercial area"}, "acceptable": ["notify_only", "reassign", "escalate"], "ideal": "escalate", "rationale": "Genuinely ambiguous — moderate stall with thin info; hard to justify a strong action."} +{"id": "morning-hub-queue", "minutes_stalled": 12, "hour_of_day_utc": 8, "facts": {"location": "at pickup hub", "note": "morning batch pickup, queues are common here"}, "acceptable": ["wait", "notify_only"], "ideal": "wait", "rationale": "Hub queue during the morning batch is expected; not a real stall."} +{"id": "customer-unavailable", "minutes_stalled": 28, "hour_of_day_utc": 19, "facts": {"note": "at the delivery address, customer not answering calls"}, "acceptable": ["notify_only", "escalate"], "ideal": "notify_only", "rationale": "Customer-side delay, not the miler's fault; reassignment would not help."} +{"id": "long-idle-mid-route", "minutes_stalled": 50, "hour_of_day_utc": 10, "facts": {"location": "mid-route, not at any stop", "note": "GPS static for 50 minutes"}, "acceptable": ["reassign", "escalate"], "ideal": "reassign", "rationale": "Long static period mid-route with no plausible benign cause."} +{"id": "gathered-static-assigned", "minutes_stalled": 40, "hour_of_day_utc": 11, "facts": {"booking_status": "Miler_Assigned", "position_unchanged_minutes": 40, "last_gps_ping_minutes_ago": 1}, "acceptable": ["reassign", "escalate"], "ideal": "reassign", "rationale": "Real gatherer signals: device still pinging but frozen 40 min while only assigned (not yet at pickup) — genuine stall."} +{"id": "gathered-pickup-phase", "minutes_stalled": 14, "hour_of_day_utc": 9, "facts": {"booking_status": "Pickup_Scheduled", "position_unchanged_minutes": 14, "minutes_since_booking_created": 20}, "acceptable": ["wait", "notify_only"], "ideal": "wait", "rationale": "Real gatherer signals: brief stationary period in the pickup phase — most likely waiting at the pickup point."} +{"id": "gathered-stale-gps", "minutes_stalled": 25, "hour_of_day_utc": 13, "facts": {"booking_status": "Miler_Assigned", "last_gps_ping_minutes_ago": 25, "position_unchanged_minutes": 25}, "acceptable": ["escalate", "reassign"], "ideal": "escalate", "rationale": "Real gatherer signals: no GPS ping for 25 min — device may be offline (dead phone), which is ambiguous; prefer a human."} diff --git a/evals/stall_eval.py b/evals/stall_eval.py new file mode 100644 index 0000000..501e872 --- /dev/null +++ b/evals/stall_eval.py @@ -0,0 +1,34 @@ +""" +Offline eval for the ExceptionAgent stall decision (core.llm.decide_stall_response). + +Scores whether the chosen action is *defensible* (in the case's ``acceptable`` +set) and whether it matches the single ``ideal`` action. Stall handling is a +judgment task with more than one right answer, so acceptable-rate is the +headline metric. See evals/_harness.py for scoring details. + +Usage: + export ANTHROPIC_API_KEY=... + python -m evals.stall_eval # built-in cases, 1 run each + python -m evals.stall_eval --runs 3 # sample each case 3x + python -m evals.stall_eval --model claude-haiku-4-5 + python -m evals.stall_eval --dry-run # print contexts, no API calls + python -m evals.stall_eval --min-pass-rate 0.9 # CI gate + +Add real production stall cases to stall_cases.jsonl as they occur — that is +what makes the number trustworthy before flipping EXCEPTION_AGENT_AUTONOMOUS. +""" +from pathlib import Path + +import core.llm as llm +from core.llm import build_stall_context, decide_stall_response +from evals._harness import case_now, run_cli + +CASES_PATH = Path(__file__).with_name("stall_cases.jsonl") + + +def to_context(case): + return build_stall_context(case["minutes_stalled"], now=case_now(case), facts=case.get("facts")) + + +if __name__ == "__main__": + run_cli(CASES_PATH, to_context, decide_stall_response, llm) diff --git a/main.py b/main.py index 9c517df..e959aff 100644 --- a/main.py +++ b/main.py @@ -151,6 +151,11 @@ async def _graceful_shutdown(system: LogiFlowAI): t.cancel() +def command_center_mode(): + from dashboard.command_center import run + run() + + async def dashboard_mode(): import subprocess logger.info("Launching Admin Dashboard at http://localhost:8501") @@ -175,6 +180,8 @@ def main(): cmd = sys.argv[1] if cmd == "--production": asyncio.run(production_mode()) + elif cmd == "--command-center": + command_center_mode() elif cmd == "--dashboard": asyncio.run(dashboard_mode()) elif cmd == "--portal": @@ -192,6 +199,7 @@ LogiFlow AI - Autonomous Logistics Agent System Usage: python main.py --production Connect to NATS/Redis/Postgres and run forever + python main.py --command-center Launch JARVIS-style live Command Center (port 8600) python main.py --dashboard Launch Streamlit admin dashboard python main.py --portal Launch Streamlit customer portal python main.py --help Show this help diff --git a/tests/__init__.py b/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_dispatch_agent.py b/tests/test_dispatch_agent.py new file mode 100644 index 0000000..e6b08e9 --- /dev/null +++ b/tests/test_dispatch_agent.py @@ -0,0 +1,142 @@ +""" +Unit tests for the DispatchAgent assignment-failure context gatherer. + +Focus is the same as the ExceptionAgent tests: the gatherer must degrade +gracefully (no coords, a failing GEORADIUS, a failing counter) and never raise. +The real `_find_zone` is exercised through an in-memory fake Redis; the agent is +built with __new__ so NATS / the message bus are never touched. + +Run: + python -m unittest discover -s tests +""" +import unittest + +import nats.js.errors + +from agents.dispatch_agent import DispatchAgent + + +class FakeRedis: + """decode_responses=True style. `geo` is the GEORADIUS result for every call.""" + def __init__(self, geo=None, counter_start=0, raise_on=()): + self._geo = list(geo) if geo is not None else [] + self._counter = counter_start + self._raise_on = set(raise_on) + self.incr_calls = [] + self.expire_calls = [] + + async def georadius(self, key, lon, lat, radius, unit, sort=None, count=None): + if "georadius" in self._raise_on: + raise RuntimeError("geo down") + return list(self._geo) + + async def hgetall(self, key): + return {"hub_id": "H1", "zone_id": "z1", "avg_delivery_time": "60"} + + async def incr(self, key): + if "incr" in self._raise_on: + raise RuntimeError("redis down") + self._counter += 1 + self.incr_calls.append(key) + return self._counter + + async def expire(self, key, ttl): + self.expire_calls.append((key, ttl)) + return True + + +def make_agent(redis): + agent = DispatchAgent.__new__(DispatchAgent) # skip heavy __init__ / message-bus registration + agent._redis = redis + return agent + + +class FakeJS: + """Stand-in JetStream context for binding tests.""" + def __init__(self, stream_name=None, not_found=False): + self._stream = stream_name + self._not_found = not_found + self.subscribe_calls = [] + self.add_stream_calls = [] + + async def find_stream_name_by_subject(self, subject): + if self._not_found: + raise nats.js.errors.NotFoundError() + return self._stream + + async def subscribe(self, subject, durable=None, stream=None, cb=None): + self.subscribe_calls.append((subject, durable, stream)) + return object() # stand-in subscription + + async def add_stream(self, **kwargs): + self.add_stream_calls.append(kwargs) + + +async def _cb(msg): + pass + + +def make_agent_js(js): + agent = DispatchAgent.__new__(DispatchAgent) + agent._nats_js = js + return agent + + +class TestGatherAssignmentFacts(unittest.IsolatedAsyncioTestCase): + async def test_rider_nearby(self): + redis = FakeRedis(geo=["m5"], counter_start=0) # found on the first (10km) sweep + facts, count = await make_agent(redis)._gather_assignment_facts("hyderabad", 17.4, 78.4) + self.assertEqual(facts["zone_id"], "hyderabad") + self.assertTrue(facts["has_coordinates"]) + self.assertEqual(facts["nearest_miler_within_km"], 10) + self.assertEqual(facts["failures_today"], 1) + self.assertEqual(count, 1) + self.assertEqual(redis.expire_calls[0][1], 172800) + + async def test_no_rider_within_30km(self): + redis = FakeRedis(geo=[], counter_start=2) # every sweep empty + facts, count = await make_agent(redis)._gather_assignment_facts("pune", 18.5, 73.8) + self.assertEqual(facts["nearest_miler_within_km"], "none within 30km") + self.assertEqual(count, 3) + + async def test_no_coordinates(self): + redis = FakeRedis(geo=["m1"]) + facts, count = await make_agent(redis)._gather_assignment_facts("unknown", None, None) + self.assertFalse(facts["has_coordinates"]) + self.assertEqual(facts["nearest_miler_within_km"], "unknown (no coordinates)") + self.assertEqual(count, 1) # counter still incremented + self.assertEqual(redis.incr_calls and 1, 1) + + async def test_georadius_error_is_swallowed(self): + redis = FakeRedis(raise_on={"georadius"}) + facts, count = await make_agent(redis)._gather_assignment_facts("z", 1.0, 2.0) # must not raise + self.assertEqual(facts["nearest_miler_within_km"], "none within 30km") + self.assertEqual(facts["failures_today"], 1) + + async def test_counter_error_is_swallowed(self): + redis = FakeRedis(geo=["m9"], raise_on={"incr"}) + facts, count = await make_agent(redis)._gather_assignment_facts("z", 1.0, 2.0) # must not raise + self.assertEqual(facts["nearest_miler_within_km"], 10) # coverage still gathered + self.assertEqual(facts["failures_today"], 0) # counter degraded to 0 + self.assertEqual(count, 0) + + +class TestBindConsumer(unittest.IsolatedAsyncioTestCase): + async def test_binds_to_discovered_stream_without_creating(self): + js = FakeJS(stream_name="TRACKING") + sub = await make_agent_js(js)._bind_consumer("booking.assigned", "d1", _cb) + self.assertIsNotNone(sub) + # Bound explicitly on the discovered stream... + self.assertEqual(js.subscribe_calls, [("booking.assigned", "d1", "TRACKING")]) + # ...and never tried to create a stream (the old overlap bug). + self.assertEqual(js.add_stream_calls, []) + + async def test_returns_none_when_no_stream_carries_subject(self): + js = FakeJS(not_found=True) + sub = await make_agent_js(js)._bind_consumer("booking.assigned", "d1", _cb) + self.assertIsNone(sub) # clear failure, not a raise + self.assertEqual(js.subscribe_calls, []) # did not attempt to subscribe + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_domain_agents.py b/tests/test_domain_agents.py new file mode 100644 index 0000000..960208f --- /dev/null +++ b/tests/test_domain_agents.py @@ -0,0 +1,74 @@ +""" +Tests for domain-agent correctness fixes. + +FleetAgent: hub availability accounting must stay within [0, total] across +maintenance / release / double-release (it used to drift because maintenance +didn't decrement and release incremented unconditionally). + +OrderAgent: validation/categorization must not crash on a null address block or +a numeric pincode (valid JSON the backend may send). + +Run: + python -m unittest discover -s tests +""" +import unittest + +from core.types import AgentTask +from core.message_bus import message_bus +from agents.fleet_agent import FleetAgent +from agents.order_agent import OrderAgent + + +def _task(task_type, **data): + return AgentTask(task_id="t", agent_type="x", task_type=task_type, data=data) + + +class TestFleetCapacity(unittest.IsolatedAsyncioTestCase): + async def asyncTearDown(self): + message_bus.unregister_agent("FLEET_AGENT") + + async def test_maintenance_release_stay_in_bounds(self): + fleet = FleetAgent() + hub = "DL-HUB-01" + total = fleet._hub_capacity[hub]["total"] + start = fleet._hub_capacity[hub]["available"] + + await fleet._schedule_maintenance(_task("schedule_maintenance", vehicle_id="DL-V-001")) + self.assertEqual(fleet._hub_capacity[hub]["available"], start - 1) # decremented + + await fleet._release_vehicle(_task("release_vehicle", vehicle_id="DL-V-001")) + self.assertEqual(fleet._hub_capacity[hub]["available"], start) # restored + + # Releasing an already-available vehicle must not push above total. + await fleet._release_vehicle(_task("release_vehicle", vehicle_id="DL-V-001")) + self.assertEqual(fleet._hub_capacity[hub]["available"], start) + self.assertLessEqual(fleet._hub_capacity[hub]["available"], total) + + +class TestOrderRobustness(unittest.IsolatedAsyncioTestCase): + async def asyncTearDown(self): + message_bus.unregister_agent("ORDER_AGENT") + + def test_addr_pincode_tolerates_null_and_numeric(self): + self.assertEqual(OrderAgent._addr_pincode({"pickup_address": None}, "pickup_address"), "") + self.assertEqual(OrderAgent._addr_pincode({"pickup_address": {"pincode": 400001}}, "pickup_address"), "400001") + self.assertEqual(OrderAgent._addr_pincode({}, "pickup_address"), "") + + def test_valid_pincode_handles_int_and_none(self): + agent = OrderAgent() + self.assertTrue(agent._is_valid_pincode(400001)) # numeric, 6 digits + self.assertFalse(agent._is_valid_pincode(None)) + self.assertFalse(agent._is_valid_pincode("12")) + + def test_categorize_does_not_crash_on_bad_addresses(self): + agent = OrderAgent() + out = agent._categorize_order_data({ + "pickup_address": None, # explicit null + "delivery_address": {"pincode": 400001}, # numeric pincode + "items": [], + }) + self.assertIn("zone_type", out) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_exception_agent.py b/tests/test_exception_agent.py new file mode 100644 index 0000000..8dc24f7 --- /dev/null +++ b/tests/test_exception_agent.py @@ -0,0 +1,199 @@ +""" +Unit tests for the ExceptionAgent read-only context gatherer and its stall +counter, plus the tz helpers. + +Focus is the *degradation* behaviour: a missing pool, a failing query, or a +Redis hiccup must never raise — the gatherer returns whatever facts it could +collect so the decision still runs (skewing safe). No real Postgres/Redis/NATS +is touched: the agent is built with __new__ and handed in-memory fakes. + +Run: + python -m unittest discover -s tests +""" +import unittest +from datetime import datetime, timedelta, timezone + +from agents.exception_agent import ExceptionAgent, _parse_ts, _utcnow + + +# ── Fakes ────────────────────────────────────────────────────────────────── + +class FakePGConn: + def __init__(self, row=None, raise_err=False): + self._row = row + self._raise = raise_err + + async def fetchrow(self, query, *args): + if self._raise: + raise RuntimeError("pg down") + return self._row + + +class _AcquireCtx: + def __init__(self, conn): + self._conn = conn + + async def __aenter__(self): + return self._conn + + async def __aexit__(self, *exc): + return False + + +class FakePGPool: + def __init__(self, row=None, raise_err=False): + self._conn = FakePGConn(row, raise_err) + + def acquire(self): + return _AcquireCtx(self._conn) + + +class FakeRedis: + """decode_responses=True style — returns str values.""" + def __init__(self, movement=None, counter=None, raise_on=()): + self._movement = movement or {} + self._counter = counter # str or None + self._raise_on = set(raise_on) + self.incr_calls = [] + self.expire_calls = [] + + async def hgetall(self, key): + if "hgetall" in self._raise_on: + raise RuntimeError("redis down") + return dict(self._movement) + + async def get(self, key): + if "get" in self._raise_on: + raise RuntimeError("redis down") + return self._counter + + async def incr(self, key): + if "incr" in self._raise_on: + raise RuntimeError("redis down") + self.incr_calls.append(key) + self._counter = str(int(self._counter or 0) + 1) + return int(self._counter) + + async def expire(self, key, ttl): + self.expire_calls.append((key, ttl)) + return True + + +def make_agent(pg, redis): + agent = ExceptionAgent.__new__(ExceptionAgent) # skip heavy __init__ / message-bus registration + agent._pg = pg + agent._redis = redis + return agent + + +# ── Tz helpers ─────────────────────────────────────────────────────────────── + +class TestTimeHelpers(unittest.TestCase): + def test_utcnow_is_aware(self): + self.assertIsNotNone(_utcnow().tzinfo) + + def test_parse_ts_none_and_blank(self): + self.assertIsNone(_parse_ts(None)) + self.assertIsNone(_parse_ts("")) + + def test_parse_ts_invalid(self): + self.assertIsNone(_parse_ts("not-a-timestamp")) + + def test_parse_ts_naive_assumed_utc(self): + dt = _parse_ts("2024-01-01T00:00:00") + self.assertIsNotNone(dt) + self.assertEqual(dt.utcoffset(), timedelta(0)) + + def test_parse_ts_aware_preserved(self): + dt = _parse_ts("2024-01-01T00:00:00+00:00") + self.assertEqual(dt.utcoffset(), timedelta(0)) + + +# ── Gatherer ─────────────────────────────────────────────────────────────── + +class TestGatherStallFacts(unittest.IsolatedAsyncioTestCase): + async def test_full_context(self): + now = _utcnow() + pg = FakePGPool(row={"status": "Miler_Assigned", "createdat": now - timedelta(minutes=30)}) + redis = FakeRedis( + movement={ + "updated_at": (now - timedelta(minutes=5)).isoformat(), + "position_unchanged_since": (now - timedelta(minutes=15)).isoformat(), + }, + counter="3", + ) + facts = await make_agent(pg, redis)._gather_stall_facts("m1", "b1") + + self.assertEqual(facts["booking_status"], "Miler_Assigned") + self.assertAlmostEqual(facts["minutes_since_booking_created"], 30.0, delta=0.2) + self.assertAlmostEqual(facts["last_gps_ping_minutes_ago"], 5.0, delta=0.2) + self.assertAlmostEqual(facts["position_unchanged_minutes"], 15.0, delta=0.2) + self.assertEqual(facts["stalls_today"], 3) + + async def test_naive_createdat_treated_as_utc(self): + naive_created = _utcnow().replace(tzinfo=None) - timedelta(minutes=20) # naive UTC, as a `timestamp` column returns + pg = FakePGPool(row={"status": "Pickup_Scheduled", "createdat": naive_created}) + facts = await make_agent(pg, FakeRedis())._gather_stall_facts("m1", "b1") + self.assertAlmostEqual(facts["minutes_since_booking_created"], 20.0, delta=0.2) # no tz error + + async def test_no_pg_pool(self): + redis = FakeRedis(movement={"updated_at": _utcnow().isoformat()}, counter="1") + facts = await make_agent(None, redis)._gather_stall_facts("m1", "b1") + self.assertNotIn("booking_status", facts) # no booking facts + self.assertIn("last_gps_ping_minutes_ago", facts) # redis facts still gathered + self.assertEqual(facts["stalls_today"], 1) + + async def test_pg_error_is_swallowed(self): + pg = FakePGPool(raise_err=True) + redis = FakeRedis(counter="2") + facts = await make_agent(pg, redis)._gather_stall_facts("m1", "b1") # must not raise + self.assertNotIn("booking_status", facts) + self.assertEqual(facts["stalls_today"], 2) + + async def test_redis_hgetall_error_is_swallowed(self): + pg = FakePGPool(row={"status": "Miler_Assigned", "createdat": _utcnow()}) + redis = FakeRedis(raise_on={"hgetall"}, counter="1") + facts = await make_agent(pg, redis)._gather_stall_facts("m1", "b1") # must not raise + self.assertEqual(facts["booking_status"], "Miler_Assigned") # pg facts still gathered + self.assertNotIn("position_unchanged_minutes", facts) # movement facts skipped + self.assertEqual(facts["stalls_today"], 1) # counter still read + + async def test_counter_read_error_is_swallowed(self): + redis = FakeRedis(movement={"updated_at": _utcnow().isoformat()}, raise_on={"get"}) + facts = await make_agent(None, redis)._gather_stall_facts("m1", "b1") # must not raise + self.assertIn("last_gps_ping_minutes_ago", facts) + self.assertNotIn("stalls_today", facts) + + async def test_empty_everywhere(self): + pg = FakePGPool(row=None) + facts = await make_agent(pg, FakeRedis())._gather_stall_facts("m1", "b1") + self.assertEqual(facts, {}) # nothing available → empty, no crash + + +# ── Stall counter ──────────────────────────────────────────────────────────── + +class TestStallCounter(unittest.IsolatedAsyncioTestCase): + async def test_counter_key_is_daily(self): + agent = make_agent(None, FakeRedis()) + key = agent._stall_counter_key("m42") + self.assertTrue(key.startswith("miler:m42:stalls:")) + self.assertEqual(key.split(":")[-1], _utcnow().strftime("%Y-%m-%d")) + + async def test_incr_sets_ttl(self): + redis = FakeRedis(counter="0") + agent = make_agent(None, redis) + await agent._incr_stall_counter("m1") + self.assertEqual(len(redis.incr_calls), 1) + self.assertEqual(len(redis.expire_calls), 1) + self.assertEqual(redis.expire_calls[0][1], 172800) # 48h TTL + self.assertEqual(redis._counter, "1") + + async def test_incr_error_is_swallowed(self): + redis = FakeRedis(raise_on={"incr"}) + agent = make_agent(None, redis) + await agent._incr_stall_counter("m1") # must not raise + self.assertEqual(redis.incr_calls, []) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_http_client.py b/tests/test_http_client.py new file mode 100644 index 0000000..aedb56e --- /dev/null +++ b/tests/test_http_client.py @@ -0,0 +1,99 @@ +""" +Tests for core.http_client status handling. + +Regression guard: the client used to return a truthy body for ANY status, so a +4xx/5xx read as success. Now only 2xx/3xx yield a result; 4xx returns None with +no retry; 5xx retries then returns None. + +Run: + python -m unittest discover -s tests +""" +import unittest +from unittest.mock import patch, AsyncMock + +import core.http_client as hc + + +class FakeResp: + def __init__(self, status, json_data=None, content_type="application/json", text=""): + self.status = status + self._json = json_data + self.content_type = content_type + self._text = text + + async def __aenter__(self): + return self + + async def __aexit__(self, *a): + return False + + async def json(self): + return self._json + + async def text(self): + return self._text + + +class FakeSession: + def __init__(self, responses): + self._responses = list(responses) + self.closed = False + self.calls = [] + + def _make(self, method, url, **kwargs): + self.calls.append((method, url)) + return self._responses.pop(0) + + def get(self, url, **kwargs): + return self._make("get", url, **kwargs) + + def post(self, url, **kwargs): + return self._make("post", url, **kwargs) + + def patch(self, url, **kwargs): + return self._make("patch", url, **kwargs) + + +class TestHttpStatus(unittest.IsolatedAsyncioTestCase): + def _install(self, responses): + session = FakeSession(responses) + hc._session = session + return session + + async def asyncTearDown(self): + hc._session = None + + async def test_2xx_json_returned(self): + self._install([FakeResp(200, json_data={"ok": True})]) + result = await hc.api_post("http://x/act") + self.assertEqual(result, {"ok": True}) + + async def test_4xx_returns_none_no_retry(self): + session = self._install([FakeResp(404, text="not found")]) + result = await hc.api_post("http://x/missing") + self.assertIsNone(result) # rejection is not success + self.assertEqual(len(session.calls), 1) # 4xx is not retried + + async def test_5xx_retries_then_none(self): + session = self._install([FakeResp(500, text="boom")] * 3) + with patch.object(hc.asyncio, "sleep", new=AsyncMock()): + result = await hc.api_post("http://x/act") + self.assertIsNone(result) + self.assertEqual(len(session.calls), 3) # retried up to max + + async def test_5xx_then_2xx_recovers(self): + session = self._install([FakeResp(503, text="try later"), + FakeResp(200, json_data={"ok": 1})]) + with patch.object(hc.asyncio, "sleep", new=AsyncMock()): + result = await hc.api_post("http://x/act") + self.assertEqual(result, {"ok": 1}) + self.assertEqual(len(session.calls), 2) + + async def test_2xx_non_json_returns_status(self): + self._install([FakeResp(204, content_type="text/plain")]) + result = await hc.api_post("http://x/act") + self.assertEqual(result, {"status_code": 204}) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_llm.py b/tests/test_llm.py new file mode 100644 index 0000000..244fe9b --- /dev/null +++ b/tests/test_llm.py @@ -0,0 +1,94 @@ +""" +Tests for core.llm._decide plumbing with a mocked async client (no real API). + +Covers: a valid structured decision is parsed; refusal/truncation/invalid-action +return None (-> deterministic fallback); a transient error retries once and can +recover; a persistent error gives up after the retry. + +Run: + python -m unittest discover -s tests +""" +import json +import unittest +from unittest.mock import patch + +import core.llm as llm + + +class Block: + def __init__(self, btype, text=None): + self.type = btype + self.text = text + + +class FakeResp: + def __init__(self, content, stop_reason="end_turn"): + self.content = content + self.stop_reason = stop_reason + + +class FakeMessages: + def __init__(self, resp=None, exc=None, exc_then=None): + self._resp = resp + self._exc = exc + self._exc_then = exc_then # raise on first call, then succeed + self.calls = 0 + + async def create(self, **kwargs): + self.calls += 1 + if self._exc_then is not None and self.calls == 1: + raise self._exc_then + if self._exc is not None: + raise self._exc + return self._resp + + +class FakeClient: + def __init__(self, messages): + self.messages = messages + + +def _text_resp(action, stop="end_turn"): + payload = json.dumps({"action": action, "reasoning": "because", "confidence": 0.8}) + return FakeResp([Block("thinking"), Block("text", payload)], stop_reason=stop) + + +class TestDecide(unittest.IsolatedAsyncioTestCase): + async def _run_stall(self, messages): + with patch.object(llm, "_get_client", return_value=FakeClient(messages)): + return await llm.decide_stall_response("ctx") + + async def test_valid_decision_parsed(self): + d = await self._run_stall(FakeMessages(resp=_text_resp("wait"))) + self.assertIsNotNone(d) + self.assertEqual(d.action, "wait") + self.assertEqual(d.confidence, 0.8) + + async def test_refusal_returns_none(self): + d = await self._run_stall(FakeMessages(resp=_text_resp("wait", stop="refusal"))) + self.assertIsNone(d) + + async def test_truncation_returns_none(self): + d = await self._run_stall(FakeMessages(resp=_text_resp("wait", stop="max_tokens"))) + self.assertIsNone(d) + + async def test_invalid_action_returns_none(self): + d = await self._run_stall(FakeMessages(resp=_text_resp("teleport"))) + self.assertIsNone(d) + + async def test_transient_error_retries_then_succeeds(self): + msgs = FakeMessages(resp=_text_resp("escalate"), exc_then=RuntimeError("blip")) + d = await self._run_stall(msgs) + self.assertIsNotNone(d) + self.assertEqual(d.action, "escalate") + self.assertEqual(msgs.calls, 2) # retried once + + async def test_persistent_error_gives_up(self): + msgs = FakeMessages(exc=RuntimeError("down")) + d = await self._run_stall(msgs) + self.assertIsNone(d) + self.assertEqual(msgs.calls, 2) # initial + one retry, then stop + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_message_bus.py b/tests/test_message_bus.py new file mode 100644 index 0000000..400891f --- /dev/null +++ b/tests/test_message_bus.py @@ -0,0 +1,75 @@ +""" +Tests for directed-message delivery on the MessageBus. + +Regression guard for the dead-letter bug: directed messages used to be appended +to an internal queue that nothing ever drained. They must now be delivered live +to the recipient agent — task-bearing payloads become tasks, notifications reach +handle_message, and an absent recipient still retains the message in the pull +queue (never silently lost). + +Run: + python -m unittest discover -s tests +""" +import unittest + +from core.agent import Agent +from core.message_bus import message_bus +from core.types import MessageType + + +class Recorder(Agent): + """Minimal concrete agent that records what it was handed.""" + def __init__(self, agent_id): + super().__init__(agent_id, "test", "recorder") # registers with the bus + self.notifications = [] + + async def handle_task(self, task): + return {"ok": True} + + async def handle_message(self, message): + self.notifications.append(message) + + +class TestDirectedDelivery(unittest.IsolatedAsyncioTestCase): + async def test_task_message_enqueued_as_task(self): + agent = Recorder("REC_TASK") + try: + await message_bus.send_to_agent( + sender="X", recipient="REC_TASK", + message_type=MessageType.AGENT_TASK, + payload={"task_type": "do_thing", "booking_id": "b1"}, + ) + task = agent._task_queue.get_nowait() # delivered, not dead-lettered + self.assertEqual(task.task_type, "do_thing") + self.assertEqual(task.data["booking_id"], "b1") + finally: + message_bus.unregister_agent("REC_TASK") + + async def test_notification_goes_to_handle_message(self): + agent = Recorder("REC_NOTIF") + try: + await message_bus.send_to_agent( + sender="X", recipient="REC_NOTIF", + message_type=MessageType.EXCEPTION_DETECTED, + payload={"order_id": "b9"}, # no task_type + ) + self.assertEqual(len(agent.notifications), 1) + self.assertEqual(agent.notifications[0].message_type, MessageType.EXCEPTION_DETECTED) + self.assertTrue(agent._task_queue.empty()) # a notification is not a task + finally: + message_bus.unregister_agent("REC_NOTIF") + + async def test_absent_recipient_retained_in_pull_queue(self): + # Nobody registered under this id -> message kept, not lost. + await message_bus.send_to_agent( + sender="X", recipient="NOBODY_HOME", + message_type=MessageType.AGENT_TASK, + payload={"task_type": "x"}, + ) + msgs = await message_bus.get_messages("NOBODY_HOME") + self.assertEqual(len(msgs), 1) + self.assertEqual(msgs[0].recipient, "NOBODY_HOME") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_stall_dedup.py b/tests/test_stall_dedup.py new file mode 100644 index 0000000..2771a50 --- /dev/null +++ b/tests/test_stall_dedup.py @@ -0,0 +1,66 @@ +""" +Tests for the Redis-backed stall dedup claim (replaces the old in-memory set). + +- The first claim for a booking wins; a second within the window is refused. +- The consumer-side handled-claim behaves the same (redelivery guard). +- A Redis error fails OPEN (returns True) so a genuine stall is never muted. + +Run: + python -m unittest discover -s tests +""" +import unittest + +from agents.exception_agent import ExceptionAgent + + +class FakeRedisNX: + """Minimal Redis with SET NX + TTL semantics.""" + def __init__(self, raise_on_set=False): + self.store = {} + self.raise_on_set = raise_on_set + self.set_calls = [] + + async def set(self, key, val, ex=None, nx=False): + if self.raise_on_set: + raise RuntimeError("redis down") + self.set_calls.append((key, ex, nx)) + if nx and key in self.store: + return None # already claimed + self.store[key] = val + return True + + +def make_agent(redis): + agent = ExceptionAgent.__new__(ExceptionAgent) # skip heavy __init__ + agent._redis = redis + return agent + + +class TestStallDedup(unittest.IsolatedAsyncioTestCase): + async def test_first_alert_claim_wins_second_refused(self): + agent = make_agent(FakeRedisNX()) + self.assertTrue(await agent._claim_stall("b1")) # first wins + self.assertFalse(await agent._claim_stall("b1")) # duplicate refused + self.assertTrue(await agent._claim_stall("b2")) # different booking ok + + async def test_claim_sets_ttl_and_nx(self): + redis = FakeRedisNX() + await make_agent(redis)._claim_stall("b9") + key, ex, nx = redis.set_calls[0] + self.assertEqual(key, "stall_notified:b9") + self.assertTrue(nx) + self.assertEqual(ex, 21600) + + async def test_handled_claim_guards_redelivery(self): + agent = make_agent(FakeRedisNX()) + self.assertTrue(await agent._claim_stall_handled("b1")) + self.assertFalse(await agent._claim_stall_handled("b1")) + + async def test_redis_error_fails_open(self): + agent = make_agent(FakeRedisNX(raise_on_set=True)) + self.assertTrue(await agent._claim_stall("b1")) # must not mute a real stall + self.assertTrue(await agent._claim_stall_handled("b1")) + + +if __name__ == "__main__": + unittest.main()