feat: ExpressDispatchAgent — tenant-scoped batch assignment + sequencing

New agent that owns the DoormileExpress dispatch flow. The normal B2C booking
flow is untouched; this handles express orders that a console operator has let
accumulate batch by batch, then dispatches in one go.

Triggered by express.dispatch_requested (the console's manual dispatch action).
For the tenant it: reads the pending bookings and the tenant's available riders
(Go internal API), distributes them across riders by proximity + load with a
radius guard, asks routes.workolik.com (Valhalla) for each rider's stop order,
and writes the assignments back through the Go internal API — which stays the
single writer of assignment state; the agent only decides who and what order.

- agents/express_dispatch_agent.py: the agent (mirrors DispatchAgent's NATS
  binding; pure consumer of the backend-owned EXPRESS stream).
- config/system_config.py: ROUTE_OPTIMIZER_URL.
- registered under JARVIS in main.py / agents package.
- tests: distribution (radius, load, cap, no-GPS fallback) and sequencing
  (single-stop, remap, optimizer-down, unknown-id) — deterministic, no network.

Gated off by default on the backend (EXPRESS_AGENT_ENABLED); EXPRESS_AGENT_
AUTONOMOUS=false runs it observe-only (logs the plan, writes nothing).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Suriyakumarvijayanayagam
2026-08-17 19:33:36 +05:30
parent ad753be7a5
commit aa4d4c6549
5 changed files with 542 additions and 1 deletions

View File

@@ -6,6 +6,7 @@ from agents.hub_agent import HubAgent
from agents.customer_agent import CustomerAgent
from agents.exception_agent import ExceptionAgent
from agents.route_optimizer_agent import RouteOptimizerAgent
from agents.express_dispatch_agent import ExpressDispatchAgent
__all__ = [
'OrderAgent',
@@ -14,5 +15,6 @@ __all__ = [
'HubAgent',
'CustomerAgent',
'ExceptionAgent',
'RouteOptimizerAgent'
'RouteOptimizerAgent',
'ExpressDispatchAgent'
]

View File

@@ -0,0 +1,391 @@
"""
Express Dispatch Agent — owns the DoormileExpress (console) batch-booking flow.
The normal B2C booking flow is untouched: each booking is auto-assigned by the Go
backend's per-booking engine (routemate AI + fallback). Express *batch* bookings
are different — a console operator drops N bookings for one tenant at once (e.g. a
DailyGrubs lunch rush across several kitchens), and those should be assigned as a
set to that tenant's own riders, then road-sequenced, not each picked off
independently by nearest-rider.
This agent is triggered by `express.dispatch_requested` — the console operator's
manual "dispatch" action, fired once orders have accumulated batch by batch (the
bulk endpoint only creates them, unassigned). It then:
1. reads the batch's bookings and the tenant's AVAILABLE riders (Go internal API),
2. distributes bookings across those riders (proximity + load, radius-guarded),
3. asks the Route Optimization API (routes.workolik.com, Valhalla road routing)
for each rider's stop order,
4. writes the assignments back through the Go internal API — which stays the
single writer of assignment state; this agent decides, Go persists.
Division of labour, on purpose: the agent chooses *who* and *what order*; every DB
invariant (already-assigned guard, transactional assign, notify) lives in Go.
"""
import asyncio
import json
import os
from math import radians, cos, sin, asin, sqrt
from typing import Dict, List, Any, Optional
import nats
import nats.js.errors
from core.agent import SpecializedAgent
from core.types import AgentTask
from core.logger import logger
from core.http_client import api_get, api_post
from config.system_config import (
NATS_HOST, NATS_PORT, NATS_USER, NATS_PASSWORD,
GO_API_BASE_URL, INTERNAL_API_KEY, ROUTE_OPTIMIZER_URL,
)
# When off, the agent computes and logs the full assignment plan but does NOT
# write it — an observation/dry-run mode to watch its decisions before it acts.
# Defaults on, because assigning an unassigned booking is the intended happy path
# (unlike the ExceptionAgent's irreversible reassignment, which defaults off).
EXPRESS_AGENT_AUTONOMOUS = os.getenv("EXPRESS_AGENT_AUTONOMOUS", "true").lower() == "true"
# The Go backend owns the EXPRESS stream carrying express.dispatch_requested;
# this agent is a pure consumer and must not create it. Empty = auto-discover.
EXPRESS_STREAM = os.getenv("EXPRESS_STREAM", "")
# A booking is only handed to a rider within this radius of it. Guards against a
# Coimbatore order landing on a Nagercoil rider without needing a zone-id mapping:
# they are >100 km apart, so the radius excludes it naturally.
MAX_RADIUS_KM = float(os.getenv("EXPRESS_MAX_RADIUS_KM", "30"))
# Ceiling on stops per rider in one batch, so the load spreads instead of piling
# onto whichever rider happens to be nearest every kitchen.
MAX_PER_RIDER = int(os.getenv("EXPRESS_MAX_PER_RIDER", "5"))
# Each stop already on a rider adds this much "virtual distance" when scoring, so
# an equally-near but emptier rider wins. Purely a load-balancing lever.
LOAD_PENALTY_KM = float(os.getenv("EXPRESS_LOAD_PENALTY_KM", "3"))
_INTERNAL_HEADERS = {"X-Internal-Key": INTERNAL_API_KEY}
_SEQUENCE_PATH = "/api/v1/optimization/doormile/sequence"
def _haversine_km(lat1: float, lon1: float, lat2: float, lon2: float) -> float:
"""Great-circle distance in km. Same formula as the Go side's haversineKM."""
lat1, lon1, lat2, lon2 = map(radians, (lat1, lon1, lat2, lon2))
dlat = lat2 - lat1
dlon = lon2 - lon1
a = sin(dlat / 2) ** 2 + cos(lat1) * cos(lat2) * sin(dlon / 2) ** 2
return 2 * asin(sqrt(a)) * 6371.0
class ExpressDispatchAgent(SpecializedAgent):
"""Consumes express.dispatch_requested and runs tenant-scoped batch
assignment + route sequencing for DoormileExpress bookings."""
def __init__(self):
super().__init__(
agent_id="EXPRESS_DISPATCH_AGENT",
domain="express_dispatch",
description="Assigns and sequences DoormileExpress batch bookings per tenant",
)
self._nats_nc = None
self._nats_js = None
self._nats_subs: list = []
# ------------------------------------------------------------------ #
# Lifecycle #
# ------------------------------------------------------------------ #
async def start(self):
await self._connect_nats()
await super().start()
async def stop(self):
await super().stop()
for sub in self._nats_subs:
try:
await sub.unsubscribe()
except Exception:
pass
if self._nats_nc:
await self._nats_nc.drain()
async def _connect_nats(self):
try:
self._nats_nc = await nats.connect(
servers=[f"nats://{NATS_HOST}:{NATS_PORT}"],
user=NATS_USER,
password=NATS_PASSWORD,
max_reconnect_attempts=10,
)
self._nats_js = self._nats_nc.jetstream()
sub = await self._bind_consumer(
"express.dispatch_requested", "express-dispatch-requested",
self._on_nats_express_dispatch,
)
if sub is not None:
self._nats_subs.append(sub)
logger.info("ExpressDispatchAgent bound express.dispatch_requested consumer")
except Exception as e:
logger.error(f"ExpressDispatchAgent NATS connect error: {e}")
async def _bind_consumer(self, subject, durable, cb):
"""Bind a durable push consumer to whichever existing stream carries
`subject`. Returns None (with an actionable log) if no stream carries it,
never a swallowed error — mirrors DispatchAgent."""
try:
stream = EXPRESS_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"ExpressDispatchAgent bound '{subject}' on stream '{stream}' (durable={durable})")
return sub
except nats.js.errors.NotFoundError:
logger.error(
f"ExpressDispatchAgent: no JetStream stream carries '{subject}'. The Go "
f"backend must declare a stream covering it (db/streams.go EXPRESS), or set "
f"EXPRESS_STREAM. Not subscribing."
)
return None
except Exception as e:
logger.error(f"ExpressDispatchAgent failed to bind '{subject}': {e}")
return None
# ------------------------------------------------------------------ #
# Event handler #
# ------------------------------------------------------------------ #
async def _on_nats_express_dispatch(self, msg):
# Ack first-thing safety: the handler is idempotent (the Go writeback
# rejects already-assigned bookings), so a redelivery cannot double-assign.
try:
data = json.loads(msg.data.decode())
tenant_id = data.get("tenantid")
booking_ids = data.get("booking_ids") or []
logger.info(
f"[EXPRESS] dispatch received — tenant={tenant_id} bookings={len(booking_ids)}"
)
if tenant_id and booking_ids:
await self._handle_batch(tenant_id, booking_ids)
except Exception as e:
logger.error(f"ExpressDispatchAgent batch handler error: {e}")
finally:
await msg.ack()
async def _handle_batch(self, tenant_id: int, booking_ids: List[int]):
bookings = await self._fetch_bookings(booking_ids)
# Only assign bookings that are genuinely open — a booking already picked
# up by the single-booking path, or cancelled, must be left alone.
pending = [
b for b in bookings
if not b.get("assignedmileruserid")
and b.get("status") not in ("Cancelled", "Converted_To_Consignment")
and _valid_coords(b)
]
if not pending:
logger.info(f"[EXPRESS] tenant={tenant_id}: no assignable bookings in batch")
return
riders = await self._fetch_riders(tenant_id)
if not riders:
logger.warning(
f"[EXPRESS] tenant={tenant_id}: {len(pending)} bookings but no available "
f"riders — leaving them for manual assignment."
)
return
# 1. Distribute bookings across riders (who carries what).
rider_stops = self._distribute(pending, riders)
assigned_count = sum(len(v) for v in rider_stops.values())
unassigned = len(pending) - assigned_count
logger.info(
f"[EXPRESS] tenant={tenant_id}: distributed {assigned_count}/{len(pending)} "
f"across {len(rider_stops)} rider(s), {unassigned} left unassigned"
)
# 2. Sequence each rider's stops (what order) and build the writeback.
assignments: List[Dict[str, Any]] = []
for miler_user_id, stops in rider_stops.items():
rider = next(r for r in riders if r["miler_user_id"] == miler_user_id)
sequenced = await self._sequence(tenant_id, rider, stops)
assignments.extend(sequenced)
if not assignments:
logger.warning(f"[EXPRESS] tenant={tenant_id}: nothing to write back")
return
# 3. Write back (or, in observation mode, only log the plan).
if not EXPRESS_AGENT_AUTONOMOUS:
logger.info(
f"[EXPRESS] tenant={tenant_id} OBSERVE-ONLY plan "
f"(EXPRESS_AGENT_AUTONOMOUS=false): {json.dumps(assignments)}"
)
return
result = await api_post(
f"{GO_API_BASE_URL}/api/v1/internal/express/assign",
json={"assignments": assignments},
headers=_INTERNAL_HEADERS,
)
if result is None:
logger.error(
f"[EXPRESS] tenant={tenant_id}: writeback failed — bookings remain "
f"unassigned and assignable by hand from the console."
)
return
logger.info(
f"[EXPRESS] tenant={tenant_id}: wrote {result.get('assigned')}/"
f"{result.get('total')} assignments"
)
# ------------------------------------------------------------------ #
# Distribution — proximity + load, radius-guarded #
# ------------------------------------------------------------------ #
def _distribute(self, bookings: List[Dict], riders: List[Dict]) -> Dict[int, List[Dict]]:
"""Greedy: each booking goes to the reachable, under-capacity rider with
the lowest (distance-to-pickup + load penalty). A booking with no rider
within MAX_RADIUS_KM is left unassigned rather than forced onto a distant
rider. Deterministic and explainable — the road optimizer does the clever
part (ordering); this only decides who."""
load: Dict[int, int] = {r["miler_user_id"]: 0 for r in riders}
rider_stops: Dict[int, List[Dict]] = {}
# Assign larger-pickup-cluster bookings first is unnecessary; simple
# stable order keeps it predictable and testable.
for b in bookings:
plat, plon = b["pickuplatitude"], b["pickuplongitude"]
best_rider = None
best_score = None
for r in riders:
if load[r["miler_user_id"]] >= MAX_PER_RIDER:
continue
rlat, rlon = _rider_location(r, plat, plon)
dist = _haversine_km(rlat, rlon, plat, plon)
if dist > MAX_RADIUS_KM:
continue
score = dist + LOAD_PENALTY_KM * load[r["miler_user_id"]]
if best_score is None or score < best_score:
best_score = score
best_rider = r
if best_rider is None:
continue # unreachable / all full — leave for manual assignment
mid = best_rider["miler_user_id"]
load[mid] += 1
rider_stops.setdefault(mid, []).append(b)
return rider_stops
# ------------------------------------------------------------------ #
# Sequencing — routes.workolik.com (Valhalla road routing) #
# ------------------------------------------------------------------ #
async def _sequence(self, tenant_id: int, rider: Dict, stops: List[Dict]) -> List[Dict[str, Any]]:
"""Return the writeback rows for one rider's stops, road-sequenced. A
single stop needs no solve — it is step 1. Two or more go to the optimizer;
if that call fails, the stops are still assigned, just at step 0
(unsequenced) so a slow/broken optimizer never blocks the assignment."""
mid = rider["miler_user_id"]
if len(stops) == 1:
b = stops[0]
return [_writeback_row(b["booking_id"], mid, 1, 0, 0, 0, 0)]
# Rider start = current GPS, or the first pickup if they have no live fix.
start_lat, start_lon = _rider_location(rider, stops[0]["pickuplatitude"], stops[0]["pickuplongitude"])
payload = {
"tenantid": tenant_id,
"mileruserid": mid,
"startlatitude": start_lat,
"startlongitude": start_lon,
"bookings": [
{
"bookingid": b["booking_id"],
"bookingno": b.get("booking_no"),
"pickuplatitude": b["pickuplatitude"],
"pickuplongitude": b["pickuplongitude"],
"deliverylatitude": b["deliverylatitude"],
"deliverylongitude": b["deliverylongitude"],
}
for b in stops
],
}
resp = await api_post(f"{ROUTE_OPTIMIZER_URL}{_SEQUENCE_PATH}", json=payload)
if not resp or not resp.get("success") or not resp.get("stops"):
logger.warning(
f"[EXPRESS] rider={mid}: optimizer returned no sequence — assigning "
f"{len(stops)} stops unsequenced (step 0)"
)
return [_writeback_row(b["booking_id"], mid, 0, 0, 0, 0, 0) for b in stops]
by_id = {b["booking_id"]: b for b in stops}
rows: List[Dict[str, Any]] = []
for s in resp["stops"]:
bid = s.get("bookingid")
if bid not in by_id:
continue # never write a step for a booking we did not send
rows.append(_writeback_row(
bid, mid,
int(s.get("step") or 0),
float(s.get("previouskms") or 0),
float(s.get("cumulativekms") or 0),
int(s.get("etaminutes") or 0),
int(s.get("cumulativeeta") or 0),
))
# Any stop the optimizer dropped is still assigned, unsequenced.
seen = {r["booking_id"] for r in rows}
for b in stops:
if b["booking_id"] not in seen:
rows.append(_writeback_row(b["booking_id"], mid, 0, 0, 0, 0, 0))
return rows
# ------------------------------------------------------------------ #
# Go internal API reads #
# ------------------------------------------------------------------ #
async def _fetch_bookings(self, booking_ids: List[int]) -> List[Dict]:
ids = ",".join(str(i) for i in booking_ids)
resp = await api_get(
f"{GO_API_BASE_URL}/api/v1/internal/express/bookings",
params={"ids": ids},
headers=_INTERNAL_HEADERS,
)
return (resp or {}).get("bookings", []) if resp else []
async def _fetch_riders(self, tenant_id: int) -> List[Dict]:
resp = await api_get(
f"{GO_API_BASE_URL}/api/v1/internal/express/riders",
params={"tenantid": tenant_id},
headers=_INTERNAL_HEADERS,
)
return (resp or {}).get("riders", []) if resp else []
# SpecializedAgent requires a task handler; this agent is event-driven, so
# there is no queued-task path.
async def handle_task(self, task: AgentTask) -> Dict[str, Any]:
return {"status": "ignored", "reason": "event-driven agent"}
def _valid_coords(b: Dict) -> bool:
return bool(
b.get("pickuplatitude") and b.get("pickuplongitude")
and b.get("deliverylatitude") and b.get("deliverylongitude")
)
def _rider_location(r: Dict, fallback_lat: float, fallback_lon: float):
"""A rider freshly flipped Available may have no GPS fix yet (0,0). Assume
they are at the pickup they'd serve rather than off the coast of Africa, so
they stay assignable instead of being excluded by the radius guard."""
lat, lon = r.get("latitude") or 0, r.get("longitude") or 0
if lat == 0 or lon == 0:
return fallback_lat, fallback_lon
return lat, lon
def _writeback_row(booking_id, miler_user_id, step, previouskms, cumulativekms, etaminutes, cumulativeeta):
return {
"booking_id": booking_id,
"miler_user_id": miler_user_id,
"step": step,
"previouskms": previouskms,
"cumulativekms": cumulativekms,
"etaminutes": etaminutes,
"cumulativeeta": cumulativeeta,
}