From 008df49605220715c6f192d33c5c7e748513ffd4 Mon Sep 17 00:00:00 2001 From: Suriyakumarvijayanayagam Date: Tue, 22 Sep 2026 16:42:15 +0530 Subject: [PATCH] fix(dispatch): bucket failures by geo cell while the backend omits zone_id MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every booking.assignment_failed observed in production carries no zone_id, so the per-zone failure counter and the once-a-day alert limit collapsed into a single global "unknown" bucket — one alert per day for the whole country. zone_key() prefers the backend's zone_id and falls back to a coarse lat/lon grid cell (1 dp, ~11 km, DISPATCH_GEO_BUCKET_DP). The failure reason is now logged too. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_012AJLYcbTHCe45fyFnMfEin --- agents/dispatch_agent.py | 29 +++++++++++++++++++++++++++-- tests/test_dispatch_agent.py | 23 +++++++++++++++++++++++ 2 files changed, 50 insertions(+), 2 deletions(-) diff --git a/agents/dispatch_agent.py b/agents/dispatch_agent.py index 3d2a035..213b955 100644 --- a/agents/dispatch_agent.py +++ b/agents/dispatch_agent.py @@ -47,6 +47,27 @@ MILER_UNAVAILABLE_STATUSES = {"break", "offline", "busy", "onbreak", "off_duty", AVAILABLE, UNAVAILABLE, UNKNOWN = "available", "unavailable", "unknown" +# The Go backend does not yet send zone_id on booking.assignment_failed (every +# event observed carries none). Without a fallback, every failure in the country +# shares the single "unknown" bucket, so the per-zone counter and the once-a-day +# alert limit would collapse into one global bucket. Derive a coarse grid cell +# from the pickup coordinates instead: 1 decimal place is ~11 km, about the size +# of a delivery zone. Real zone_ids, once sent, take precedence. +GEO_BUCKET_DP = int(os.getenv("DISPATCH_GEO_BUCKET_DP", "1")) + + +def zone_key(zone_id, lat, lon) -> str: + """Stable bucket for counting failures: the backend's zone_id when present, + else a coarse lat/lon grid cell, else 'unknown'.""" + if zone_id and str(zone_id).lower() not in ("", "unknown", "none", "null"): + return str(zone_id) + if lat is not None and lon is not None: + try: + return f"geo:{round(float(lat), GEO_BUCKET_DP)},{round(float(lon), GEO_BUCKET_DP)}" + except (TypeError, ValueError): + pass + return "unknown" + class DispatchAgent(SpecializedAgent): """ @@ -274,11 +295,15 @@ class DispatchAgent(SpecializedAgent): try: data = json.loads(msg.data.decode()) booking_id = data.get("booking_id") - zone_id = data.get("zone_id", "unknown") lat = data.get("lat") lon = data.get("lon") + # Falls back to a coarse geo cell while the backend omits zone_id. + zone_id = zone_key(data.get("zone_id"), lat, lon) - logger.warning(f"[DISPATCH] booking.assignment_failed — booking={booking_id} zone={zone_id}") + logger.warning( + f"[DISPATCH] booking.assignment_failed — booking={booking_id} zone={zone_id} " + f"reason={data.get('reason')!r}" + ) facts, count = await self._gather_assignment_facts(zone_id, lat, lon) logger.info(f"[DISPATCH] booking={booking_id} context facts: {facts}") diff --git a/tests/test_dispatch_agent.py b/tests/test_dispatch_agent.py index 47ff5f0..10ae4e1 100644 --- a/tests/test_dispatch_agent.py +++ b/tests/test_dispatch_agent.py @@ -288,6 +288,29 @@ class TestRateLimitAndSinks(unittest.IsolatedAsyncioTestCase): self.assertEqual(payload["proposed_action"], "escalate") +class TestZoneKey(unittest.IsolatedAsyncioTestCase): + """The backend sends no zone_id today; without a fallback every failure in + the country shares one bucket and the daily alert limit becomes global.""" + + async def test_real_zone_id_wins(self): + from agents.dispatch_agent import zone_key + self.assertEqual(zone_key("hyderabad", 17.4, 78.4), "hyderabad") + + async def test_missing_zone_falls_back_to_geo_cell(self): + from agents.dispatch_agent import zone_key + self.assertEqual(zone_key(None, 11.0182714, 76.9677744), "geo:11.0,77.0") + self.assertEqual(zone_key("unknown", 11.0182714, 76.9677744), "geo:11.0,77.0") + + async def test_distant_pickups_get_different_buckets(self): + from agents.dispatch_agent import zone_key + self.assertNotEqual(zone_key(None, 11.01, 76.96), zone_key(None, 17.44, 78.39)) + + async def test_no_zone_and_no_coords_is_unknown(self): + from agents.dispatch_agent import zone_key + self.assertEqual(zone_key(None, None, None), "unknown") + self.assertEqual(zone_key("unknown", "bad", "coords"), "unknown") + + class TestBindConsumer(unittest.IsolatedAsyncioTestCase): async def test_binds_to_discovered_stream_without_creating(self): js = FakeJS(stream_name="TRACKING")