sync: capture the production server's code, which was never committed

/root/Routes-api on 31.97.228.132 is not a git repository. Work had been
done directly on the box and existed nowhere else -- a single rm -rf from
being lost, and impossible to review or roll back.

Deploying the previous HEAD over it would have silently reverted all of
this. Most visibly the Valhalla road-backend probe in main.py, whose own
comment explains why it exists: road sequencing degrades to aerial
silently by design, so an unreachable backend stays invisible, "which is
exactly how the expired Google key went unnoticed". Overwriting it would
have reintroduced precisely the failure it was written to catch, and the
service would have kept answering 200 throughout.

The server had also moved from Google Maps to Valhalla for road distance
(VALHALLA_URL, road_backend_status, +190 lines in route_optimizer),
extended docker-compose from 44 to 95 lines, and changed rider fetching,
health, dynamic config and the cache layer.

Only 10 files differ in substance. The other 27 that appeared to differ
were CRLF-vs-LF noise -- the server writes CRLF -- and are normalised to
LF here rather than committed as spurious whole-file rewrites.

Committed as-is, before any change of mine, so the diff that follows is
reviewable against what is actually running.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Suriya
2026-08-11 15:56:32 +05:30
parent f804791864
commit 9ef2f61870
10 changed files with 327 additions and 75 deletions

View File

@@ -15,6 +15,31 @@ except Exception: # pragma: no cover
logger = logging.getLogger(__name__)
def _redis_url_from_parts() -> Optional[str]:
"""
Build a redis:// URL from discrete env vars, percent-encoding the credentials.
A hand-written REDIS_URL breaks the moment the password contains a reserved
character: an '@' splits the authority section early, so the client silently
dials a garbage host instead of failing loudly. Assembling the URL here with
quote() removes that whole class of bug - callers set REDIS_HOST/REDIS_PASSWORD
and never have to think about escaping.
Returns None when REDIS_HOST is unset, so an explicit REDIS_URL still wins.
"""
host = os.getenv("REDIS_HOST")
if not host:
return None
from urllib.parse import quote
user = os.getenv("REDIS_USERNAME", "")
pwd = os.getenv("REDIS_PASSWORD", "")
port = os.getenv("REDIS_PORT", "6379")
db = os.getenv("REDIS_DB", "0")
auth = f"{quote(user, safe='')}:{quote(pwd, safe='')}@" if pwd else ""
return f"redis://{auth}{host}:{port}/{db}"
class RedisCache:
"""Lightweight Redis cache wrapper with graceful in-memory fallback."""
@@ -33,9 +58,10 @@ class RedisCache:
self._client = None
self._stats = {"hits": 0, "misses": 0, "sets": 0}
url = os.getenv(url_env)
url = os.getenv(url_env) or _redis_url_from_parts()
if not url or redis is None:
logger.warning("Redis not configured or client unavailable; falling back to local thread-safe in-memory cache")
_why = "redis client not installed" if redis is None else "no REDIS_URL/REDIS_HOST"
logger.warning(f"Redis not configured or client unavailable ({_why}); falling back to local thread-safe in-memory cache")
return
try:
self._client = redis.Redis.from_url(url, decode_responses=True)

View File

@@ -6,13 +6,45 @@ from typing import List, Dict, Any
logger = logging.getLogger(__name__)
def _roster_cache_ttl() -> int:
"""Roster cache TTL in seconds. 0 disables caching entirely."""
try:
from app.config.dynamic_config import get_config
return int(get_config().get("rider_roster_cache_ttl_seconds", 30))
except Exception:
return 30
async def fetch_active_riders() -> List[Dict[str, Any]]:
"""
Fetch active rider logs from the external API for the current date.
Returns a list of rider log dictionaries.
This upstream call measured ~6s on live traffic and was the single largest
component of riderassign latency, so successful responses are cached for
`rider_roster_cache_ttl_seconds`. The TTL is deliberately short: a rider
going on or off duty must become visible quickly, and the absent-rider
validation in the assignment path checks against this roster.
Only non-empty successes are cached. Caching an empty result or a failure
would turn a brief upstream blip into a full-TTL outage where every request
sees zero riders.
"""
today_str = datetime.now().strftime("%Y-%m-%d")
ttl = _roster_cache_ttl()
cache_key = f"riders:active:{today_str}"
if ttl > 0:
try:
from app.services import cache as _cache
cached = _cache.get_json(cache_key)
if isinstance(cached, list) and cached:
logger.debug(f"[Riders] roster cache hit ({len(cached)} riders)")
return cached
except Exception:
pass
try:
today_str = datetime.now().strftime("%Y-%m-%d")
url = "https://jupiter.nearle.app/live/api/v2/partners/getriderlogs/"
params = {
"applocationid": 1,
@@ -21,19 +53,26 @@ async def fetch_active_riders() -> List[Dict[str, Any]]:
"todate": today_str,
"keyword": ""
}
async with httpx.AsyncClient(timeout=30.0) as client:
response = await client.get(url, params=params)
response.raise_for_status()
data = response.json()
if data and data.get("code") == 200 and data.get("details"):
# Filter riders who are in our preferences list and are 'active' or 'idle' (assuming we want online riders)
# The user's example showed "onduty": 1. We might want to filter by that.
# For now, returning all logs, filtering can happen in assignment logic or here.
# Let's return the raw list as requested, filtering logic will be applied during assignment.
return data.get("details", [])
riders = data.get("details", [])
if riders and ttl > 0:
try:
from app.services import cache as _cache
_cache.set_json(cache_key, riders, ttl_seconds=ttl)
except Exception:
pass
return riders
logger.warning(f"Fetch active riders returned no details: {data}")
return []
@@ -60,7 +99,7 @@ async def fetch_created_orders() -> List[Dict[str, Any]]:
"pageno": 1
# "pagesize" intentionally omitted to fetch all
}
async with httpx.AsyncClient(timeout=60.0) as client:
response = await client.get(url, params=params)
response.raise_for_status()

View File

@@ -61,8 +61,8 @@ class RoadSequencingAgent:
from app.core.arrow_utils import calculate_haversine_matrix_vectorized
opt = self._optimizer()
if not opt.use_google_maps:
return {"evaluated": 0, "reason": "no_google_key", "mean_gain_pct": 0.0}
if not opt.use_road_matrix:
return {"evaluated": 0, "reason": "no_matrix_backend", "mean_gain_pct": 0.0}
batches = get_delivery_history_service().sample_batches(
days=days, limit=sample_batches

View File

@@ -8,7 +8,7 @@ ALGORITHM: TSP / VRP with Google OR-Tools
FEATURES:
- Automatic outlier detection and coordinate correction
- Hybrid distance calculation (Google Maps + Haversine fallback)
- Hybrid distance calculation (self-hosted Valhalla road matrix + Haversine fallback)
- Robust error handling for invalid inputs
"""
@@ -37,6 +37,37 @@ except ImportError:
logger = logging.getLogger(__name__)
async def road_backend_status() -> Dict[str, Any]:
"""
Probe the Valhalla matrix backend.
Module-level (not a RouteOptimizer method) so startup and health checks can
call it without paying for a full optimizer construction, which pulls in the
empirical ETA calculator and its history load.
"""
url = os.getenv("VALHALLA_URL", "").strip().rstrip("/")
if not url:
return {
"configured": False,
"reachable": False,
"detail": "VALHALLA_URL not set - road sequencing disabled, aerial only",
}
try:
async with httpx.AsyncClient(timeout=5.0) as client:
resp = await client.get(f"{url}/status")
resp.raise_for_status()
data = resp.json()
return {
"configured": True,
"reachable": True,
"url": url,
"version": data.get("version"),
"tileset_last_modified": data.get("tileset_last_modified"),
}
except Exception as e:
return {"configured": True, "reachable": False, "url": url, "detail": str(e)}
class RouteOptimizer:
"""Route optimization using Google OR-Tools (Async)."""
@@ -60,9 +91,11 @@ class RouteOptimizer:
# Road factor (haversine -> road distance multiplier, ML-tuned)
self.road_factor = float(_cfg.get("road_factor"))
# Google Maps API settings
self.google_maps_api_key = os.getenv("GOOGLE_MAPS_API_KEY", "")
self.use_google_maps = bool(self.google_maps_api_key)
# Road travel-time matrix backend: self-hosted Valhalla (no per-request
# cost, no rate limit). Unset VALHALLA_URL -> road sequencing stays off
# and every caller falls back to aerial ordering.
self.valhalla_url = os.getenv("VALHALLA_URL", "").strip().rstrip("/")
self.use_road_matrix = bool(self.valhalla_url)
# Solver time limit (ML-tuned)
self.search_time_limit_seconds = int(_cfg.get("search_time_limit_seconds"))
@@ -90,47 +123,94 @@ class RouteOptimizer:
# ROAD-AWARE VISITING ORDER (Phase 2 - opt-in, cached)
# ------------------------------------------------------------------
def _aerial_minutes(
self, a: Tuple[float, float], b: Tuple[float, float]
) -> float:
"""Haversine travel-time estimate (minutes), used to patch unroutable pairs."""
km = self.haversine_distance(a[0], a[1], b[0], b[1]) * self.road_factor
return (km / (self.avg_speed_kmh or 20.0)) * 60.0
async def _road_duration_matrix(
self, coords: _List[Tuple[float, float]]
) -> Optional[_List[_List[float]]]:
"""
Full NxN road travel-TIME matrix (minutes) via Google Distance Matrix.
Full NxN road travel-TIME matrix (minutes) via self-hosted Valhalla.
One /sources_to_targets call returns the whole matrix, so unlike the old
Google Distance Matrix path there is no per-request element cap and no
chunking. Costing defaults to `motorcycle`, which models the lane access
and one-way behaviour our riders actually have — a car matrix systematically
overstates their travel time.
Valhalla reports an unroutable pair as time=null. Those are patched with an
aerial estimate rather than left at 0, because a 0-cost edge would look
free to the TSP and pull the whole sequence through it. If too many pairs
are unroutable the tileset probably doesn't cover this region, so we bail
to aerial entirely.
Chunks destinations to respect Google's ~100-elements-per-request limit.
Returns None on any failure so the caller falls back to aerial ordering.
"""
if not self.use_google_maps:
if not self.use_road_matrix:
return None
cfg = get_config()
costing = str(cfg.get("routing_valhalla_costing", "motorcycle"))
timeout = float(cfg.get("routing_matrix_timeout_seconds", 15.0))
max_unroutable = float(cfg.get("routing_matrix_max_unroutable_pct", 20.0))
n = len(coords)
origins = "|".join(f"{la},{lo}" for la, lo in coords)
locations = [{"lat": float(la), "lon": float(lo)} for la, lo in coords]
matrix = [[0.0] * n for _ in range(n)]
dest_chunk = max(1, 100 // max(1, n))
try:
async with httpx.AsyncClient(timeout=15.0) as client:
for j0 in range(0, n, dest_chunk):
js = list(range(j0, min(j0 + dest_chunk, n)))
dests = "|".join(f"{coords[j][0]},{coords[j][1]}" for j in js)
resp = await client.get(
"https://maps.googleapis.com/maps/api/distancematrix/json",
params={"origins": origins, "destinations": dests,
"key": self.google_maps_api_key, "units": "metric"},
)
resp.raise_for_status()
data = resp.json()
if data.get("status") != "OK":
logger.debug(f"[RoadSeq] DistanceMatrix status={data.get('status')}")
return None
for i, row in enumerate(data.get("rows", [])):
for k, el in enumerate(row.get("elements", [])):
if el.get("status") == "OK":
dur = el.get("duration", {}).get("value")
if dur is not None:
matrix[i][js[k]] = dur / 60.0
async with httpx.AsyncClient(timeout=timeout) as client:
resp = await client.post(
f"{self.valhalla_url}/sources_to_targets",
json={
"sources": locations,
"targets": locations,
"costing": costing,
"units": "km",
},
)
resp.raise_for_status()
data = resp.json()
rows = data.get("sources_to_targets") or []
if len(rows) != n:
logger.debug(
f"[RoadSeq] Valhalla returned {len(rows)} rows, expected {n}"
)
return None
unroutable = 0
for i, row in enumerate(rows):
for el in row or []:
j = el.get("to_index")
if j is None or not (0 <= j < n):
continue
secs = el.get("time")
if secs is None:
unroutable += 1
matrix[i][j] = self._aerial_minutes(coords[i], coords[j])
else:
matrix[i][j] = secs / 60.0
pct = 100.0 * unroutable / max(1, n * n)
if pct > max_unroutable:
logger.warning(
f"[RoadSeq] {unroutable}/{n * n} pairs ({pct:.0f}%) unroutable via "
f"Valhalla - tileset likely missing this region; using aerial"
)
return None
if unroutable:
logger.debug(
f"[RoadSeq] patched {unroutable}/{n * n} unroutable pairs with aerial"
)
return matrix
except Exception as e:
logger.debug(f"[RoadSeq] matrix build failed: {e}")
return None
async def _road_optimal_order(
self,
start_lat: float,
@@ -140,27 +220,25 @@ class RouteOptimizer:
"""
Road-aware visiting order for `points`, starting from (start_lat, start_lon).
Builds a real road travel-TIME matrix (Google Distance Matrix) and solves
an OPEN TSP with OR-Tools (return-to-depot edge = 0), so the sequence
respects real road geometry/one-ways instead of straight-line distance.
Validated on live batches to cut real travel time ~5-13% vs aerial; note
Google's Directions optimize:true is NOT used - it optimises a closed loop
and measured *worse* than aerial for our open delivery routes.
Builds a real road travel-TIME matrix (Valhalla) and solves an OPEN TSP
with OR-Tools (return-to-depot edge = 0), so the sequence respects real
road geometry/one-ways instead of straight-line distance. Validated on
live batches to cut real travel time ~5-13% vs aerial.
Returns 0-based indices into `points` in optimal order, or None to signal
the caller to fall back to the existing aerial greedy + 2-opt.
Safe + cheap on the live path:
* disabled unless `routing_use_road_distance` AND a Google key is set
Safe on the live path:
* disabled unless `routing_use_road_distance` AND VALHALLA_URL is set
* only for 3..`routing_road_max_stops` stops (fewer is trivial; more
exceeds Google's waypoint-optimize cap)
costs more solver time than the ordering gain is worth)
* result cached in Redis (default 24h) keyed by the rounded coords in
input order, so the returned indices always map back correctly
"""
cfg = get_config()
if not cfg.get("routing_use_road_distance", False):
return None
if not self.use_google_maps:
if not self.use_road_matrix:
return None
n = len(points)
if n < 3 or n > int(cfg.get("routing_road_max_stops", 25)):
@@ -362,14 +440,17 @@ class RouteOptimizer:
search_params.local_search_metaheuristic = (
routing_enums_pb2.LocalSearchMetaheuristic.GUIDED_LOCAL_SEARCH
)
# TSP time limit hard-capped at 2 seconds per kitchen.
# search_time_limit_seconds can be tuned up to 8-10s via config, but at
# delivery scale (< 15 stops)
# OR-Tools finds a near-optimal solution in < 200ms. Waiting 8-10s
# per kitchen x 3 kitchens x 4 riders = 96s of unnecessary waiting.
# The VRP already has its own 3s cap. This cap applies to per-rider
# TSP and the per-kitchen beatmap solves.
search_params.time_limit.seconds = min(self.search_time_limit_seconds, 2)
# TSP time budget, ceiling 2 seconds per kitchen. This cap applies to
# per-rider TSP and the per-kitchen beatmap solves.
#
# GUIDED_LOCAL_SEARCH runs until its time limit EXPIRES - it does not
# return early once it has the optimum. A flat cap therefore burned the
# whole budget on trivial inputs, where PATH_CHEAPEST_ARC is done in
# under 200ms. Scale with stop count; large inputs still get the ceiling.
_cap_ms = min(self.search_time_limit_seconds, 2) * 1000
search_params.time_limit.FromMilliseconds(
max(200, min(_cap_ms, 200 * len(locations)))
)
solution = routing.SolveWithParameters(search_params)
@@ -742,13 +823,18 @@ class RouteOptimizer:
sp.local_search_metaheuristic = (
routing_enums_pb2.LocalSearchMetaheuristic.GUIDED_LOCAL_SEARCH
)
# VRP time budget: cap at 3 seconds so the API response is never
# VRP time budget: ceiling of 3 seconds so the API response is never
# blocked longer than that. The x 3 multiplier caused 15-second
# responses with the default 5-second search_time_limit.
# PATH_CHEAPEST_ARC typically finds a good solution in < 1 second
# for our typical sizes (<= 15 riders, <= 60 orders); GLS then
# improves it within the remaining budget.
sp.time_limit.seconds = min(self.search_time_limit_seconds, 3)
#
# GUIDED_LOCAL_SEARCH runs until its time limit EXPIRES rather than
# stopping once the optimum is found, so a flat 3s spent ~2.9s of
# pure waiting on a 1-order batch (measured on live traffic).
# Scale with problem size; large batches still get the full ceiling.
_cap_ms = min(self.search_time_limit_seconds, 3) * 1000
sp.time_limit.FromMilliseconds(
max(250, min(_cap_ms, 250 * len(orders)))
)
solution = routing.SolveWithParameters(sp)