diff --git a/controllers/cxCatalogueController.go b/controllers/cxCatalogueController.go index 170d131..3c295d3 100644 --- a/controllers/cxCatalogueController.go +++ b/controllers/cxCatalogueController.go @@ -11,11 +11,11 @@ import ( "doormile/constants" "doormile/db" + "doormile/internal/milergeo" "doormile/models" "doormile/utils" "github.com/gofiber/fiber/v2" - "github.com/redis/go-redis/v9" ) // Catalogue and configuration — §5 of the customer contract. @@ -440,16 +440,7 @@ func milersWithin(lat, lng, radiusKM float64) int { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() - locs, err := db.Rdb.GeoSearchLocation(ctx, "milers:locations", &redis.GeoSearchLocationQuery{ - GeoSearchQuery: redis.GeoSearchQuery{ - Longitude: lng, - Latitude: lat, - Radius: radiusKM, - RadiusUnit: "km", - Sort: "ASC", - Count: 50, - }, - }).Result() + locs, err := milergeo.Search(ctx, db.Rdb, lat, lng, radiusKM, 50) if err != nil { utils.Warn("milersWithin: geo search failed", "error", err) return 0 diff --git a/controllers/milerController.go b/controllers/milerController.go index f08fb7a..3a4ae9c 100644 --- a/controllers/milerController.go +++ b/controllers/milerController.go @@ -17,6 +17,7 @@ import ( "doormile/internal/assignment" "doormile/internal/cxstage" "doormile/internal/legs" + "doormile/internal/milergeo" "doormile/internal/notify" "doormile/internal/routing" "doormile/models" @@ -371,11 +372,11 @@ func indexMilerLocation(milerUserID int, lat, lon float64) { val := fmt.Sprintf("%f,%f", lat, lon) db.Rdb.Set(ctx, redisKey, val, 30*time.Minute) - db.Rdb.GeoAdd(ctx, "milers:locations", &redis.GeoLocation{ - Name: strconv.Itoa(milerUserID), - Latitude: lat, - Longitude: lon, - }) + // A failed write is a rider auto-assignment can never find; say so + // instead of dropping the error. + if err := milergeo.Index(ctx, db.Rdb, milerUserID, lat, lon); err != nil { + utils.Warn("rider position not indexed for assignment", "miler_id", milerUserID, "error", err) + } } func UpdateMilerLocation(c *fiber.Ctx) error { diff --git a/internal/ai/playground/tools.go b/internal/ai/playground/tools.go index fcfd86c..a44a6e9 100644 --- a/internal/ai/playground/tools.go +++ b/internal/ai/playground/tools.go @@ -9,6 +9,8 @@ import ( "strings" "time" + "doormile/internal/milergeo" + "github.com/redis/go-redis/v9" "gorm.io/gorm" ) @@ -159,13 +161,7 @@ func nearbyMilers(rdb *redis.Client) Executor { if in.RadiusKm <= 0 || in.RadiusKm > nearbyMaxKm { in.RadiusKm = 5 } - locs, err := rdb.GeoSearchLocation(ctx, "milers:locations", &redis.GeoSearchLocationQuery{ - GeoSearchQuery: redis.GeoSearchQuery{ - Longitude: *in.Lon, Latitude: *in.Lat, - Radius: in.RadiusKm, RadiusUnit: "km", Sort: "ASC", Count: nearbyMaxCount, - }, - WithDist: true, - }).Result() + locs, err := milergeo.Search(ctx, rdb, *in.Lat, *in.Lon, in.RadiusKm, nearbyMaxCount) if err != nil { return nil, errors.New("live rider positions are unavailable") } diff --git a/internal/assignment/ai_layer.go b/internal/assignment/ai_layer.go index c054611..7314511 100644 --- a/internal/assignment/ai_layer.go +++ b/internal/assignment/ai_layer.go @@ -149,12 +149,16 @@ func collectEligibleCandidates(nearby []redis.GeoLocation) ([]*milerCandidate, [ continue } + // Only today's open stops count towards the cap. An order assigned + // yesterday or last month and never closed (rider app not updated, + // order abandoned) used to hold the rider at the cap forever, so the + // one rider actually working was skipped for every new order. var activeCount int64 db.DB.Model(&models.BookingAssignment{}). - Where("mileruserid = ? AND assignmentstatus IN ?", milerUserID, []string{ + Where("mileruserid = ? AND assignmentstatus IN ? AND assignedat >= ?", milerUserID, []string{ constants.AssignmentAssigned, constants.AssignmentAccepted, - }). + }, startOfISTDay(time.Now())). Count(&activeCount) if activeCount >= maxActiveBookings() { @@ -185,6 +189,14 @@ func collectEligibleCandidates(nearby []redis.GeoLocation) ([]*milerCandidate, [ return candidates, aiCandidates } +// startOfISTDay is midnight today in India, the boundary for "today's" +// open stops. An absolute instant, so it compares correctly with assignedat +// whether the column stores IST wall-clock or a real timestamp. +func startOfISTDay(now time.Time) time.Time { + n := now.In(utils.ISTLocation()) + return time.Date(n.Year(), n.Month(), n.Day(), 0, 0, 0, 0, utils.ISTLocation()) +} + // ─── Per-miler stats ───────────────────────────────────────────────────────── func fetchMilerStats(milerUserID int) (onTimeRate float64, completedToday int64) { diff --git a/internal/assignment/crm_assignment.go b/internal/assignment/crm_assignment.go index 8617f78..947aba8 100644 --- a/internal/assignment/crm_assignment.go +++ b/internal/assignment/crm_assignment.go @@ -12,6 +12,7 @@ import ( "doormile/constants" "doormile/db" "doormile/internal/cxstage" + "doormile/internal/milergeo" "doormile/internal/routing" "doormile/models" "doormile/utils" @@ -196,22 +197,9 @@ func queryNearbyMilers(lat, lon float64) ([]redis.GeoLocation, error) { ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) defer cancel() - locs, err := db.Rdb.GeoSearchLocation(ctx, "milers:locations", &redis.GeoSearchLocationQuery{ - GeoSearchQuery: redis.GeoSearchQuery{ - Longitude: lon, - Latitude: lat, - Radius: geoRadiusKm, - RadiusUnit: "km", - Sort: "ASC", - Count: geoMaxCount, - }, - WithDist: true, - }).Result() - if err != nil { - return nil, err - } - - return locs, nil + // milergeo falls back to GEORADIUS on Redis older than 6.2, where + // GEOSEARCH does not exist and every order would find "no riders". + return milergeo.Search(ctx, db.Rdb, lat, lon, geoRadiusKm, geoMaxCount) } // commitAssignment writes the BookingAssignment row, updates the booking and the diff --git a/internal/assignment/maxactive_test.go b/internal/assignment/maxactive_test.go index 8a452f0..a0cc38b 100644 --- a/internal/assignment/maxactive_test.go +++ b/internal/assignment/maxactive_test.go @@ -1,6 +1,9 @@ package assignment -import "testing" +import ( + "testing" + "time" +) // The per-miler concurrent-stop cap. // @@ -30,3 +33,22 @@ func TestMaxActiveBookings(t *testing.T) { } } } + +// The cap counts only stops assigned since midnight India time, so an order +// left open yesterday or last month no longer blocks a working rider. +func TestStartOfISTDay(t *testing.T) { + ist := time.FixedZone("IST", 5*3600+30*60) + cases := []struct{ now, want time.Time }{ + // 11:43 IST on 6 Oct -> 00:00 IST 6 Oct + {time.Date(2026, 10, 6, 11, 43, 0, 0, ist), time.Date(2026, 10, 6, 0, 0, 0, 0, ist)}, + // 20:00 UTC on 5 Oct is already 01:30 on 6 Oct in India + {time.Date(2026, 10, 5, 20, 0, 0, 0, time.UTC), time.Date(2026, 10, 6, 0, 0, 0, 0, ist)}, + // 18:29 UTC on 5 Oct is still 23:59 on 5 Oct in India + {time.Date(2026, 10, 5, 18, 29, 0, 0, time.UTC), time.Date(2026, 10, 5, 0, 0, 0, 0, ist)}, + } + for _, c := range cases { + if got := startOfISTDay(c.now); !got.Equal(c.want) { + t.Errorf("startOfISTDay(%v) = %v, want %v", c.now, got, c.want) + } + } +} diff --git a/internal/milergeo/milergeo.go b/internal/milergeo/milergeo.go new file mode 100644 index 0000000..04a86e9 --- /dev/null +++ b/internal/milergeo/milergeo.go @@ -0,0 +1,131 @@ +// Package milergeo is the one place that reads and writes the live rider +// positions in Redis (the `milers:locations` GEO set). Auto-assignment, the +// customer pickup-slot check and the agent playground all search it; the rider +// location ping and duty start write it. +// +// Why it exists: every caller used GEOSEARCH, which Redis only has from 6.2. +// On an older server the writes (GEOADD) succeed but every search fails, and +// the callers treated the failure as "no riders nearby" — so orders sat in +// Pending with riders standing next to the pickup and nothing in the console +// said why. Search now falls back to GEORADIUS (Redis 3.2+), which answers the +// same question, and Probe reports at boot whether the search works at all. +package milergeo + +import ( + "context" + "fmt" + "reflect" + "strings" + "sync/atomic" + + "github.com/redis/go-redis/v9" +) + +// Key is the GEO set holding each rider's last reported position, member = +// the rider's userid as a decimal string. +const Key = "milers:locations" + +// Client is the slice of the Redis client this package needs; *redis.Client +// satisfies it. An interface so the fallback can be tested without a server. +type Client interface { + GeoAdd(ctx context.Context, key string, geoLocation ...*redis.GeoLocation) *redis.IntCmd + GeoSearchLocation(ctx context.Context, key string, q *redis.GeoSearchLocationQuery) *redis.GeoSearchLocationCmd + GeoRadius(ctx context.Context, key string, longitude, latitude float64, query *redis.GeoRadiusQuery) *redis.GeoLocationCmd +} + +// legacyOnly flips to true the first time the server rejects GEOSEARCH as an +// unknown command, so later searches go straight to GEORADIUS instead of +// paying for a failed round trip every time. +var legacyOnly atomic.Bool + +// Search returns the riders within radiusKm of (lat, lon), nearest first, at +// most count of them, each with Dist (km) set. It uses GEOSEARCH and falls back +// to GEORADIUS when the server is older than Redis 6.2. +// +// Any other error (timeout, WRONGTYPE on the key, ...) is returned as is: the +// caller must not mistake a broken search for an empty street. +func Search(ctx context.Context, rdb Client, lat, lon, radiusKm float64, count int) ([]redis.GeoLocation, error) { + if isNil(rdb) { + return nil, fmt.Errorf("redis not available") + } + if !legacyOnly.Load() { + locs, err := rdb.GeoSearchLocation(ctx, Key, &redis.GeoSearchLocationQuery{ + GeoSearchQuery: redis.GeoSearchQuery{ + Longitude: lon, + Latitude: lat, + Radius: radiusKm, + RadiusUnit: "km", + Sort: "ASC", + Count: count, + }, + WithDist: true, + }).Result() + if err == nil { + return locs, nil + } + if !isUnknownCommand(err) { + return nil, fmt.Errorf("GEOSEARCH %s: %w", Key, err) + } + legacyOnly.Store(true) + } + locs, err := rdb.GeoRadius(ctx, Key, lon, lat, &redis.GeoRadiusQuery{ + Radius: radiusKm, + Unit: "km", + WithDist: true, + Count: count, + Sort: "ASC", + }).Result() + if err != nil { + return nil, fmt.Errorf("GEORADIUS %s: %w", Key, err) + } + return locs, nil +} + +// Index records a rider's position. The error is returned so callers can log +// it — a failed write here is a rider auto-assignment can never find. +func Index(ctx context.Context, rdb Client, milerUserID int, lat, lon float64) error { + if isNil(rdb) { + return fmt.Errorf("redis not available") + } + if err := rdb.GeoAdd(ctx, Key, &redis.GeoLocation{ + Name: fmt.Sprint(milerUserID), + Latitude: lat, + Longitude: lon, + }).Err(); err != nil { + return fmt.Errorf("GEOADD %s: %w", Key, err) + } + return nil +} + +// Probe runs one search to report, at boot and on /ready, whether rider +// search works on this Redis: "ok", "ok (GEORADIUS fallback: Redis older than +// 6.2)", or "error: ...". It never fails the caller. +func Probe(ctx context.Context, rdb Client) string { + if isNil(rdb) { + return "error: redis not available" + } + if _, err := Search(ctx, rdb, 11.0168, 76.9558, 1, 1); err != nil { + return "error: " + err.Error() + } + if legacyOnly.Load() { + return "ok (GEORADIUS fallback: Redis older than 6.2)" + } + return "ok" +} + +// isNil also catches a nil *redis.Client inside the interface (db.Rdb before +// InitRedis), which a plain == nil does not. +func isNil(rdb Client) bool { + if rdb == nil { + return true + } + v := reflect.ValueOf(rdb) + return v.Kind() == reflect.Ptr && v.IsNil() +} + +// isUnknownCommand: how Redis < 6.2 (and proxies that filter commands) answer +// a command they do not have, e.g. "ERR unknown command 'GEOSEARCH'". +func isUnknownCommand(err error) bool { + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "unknown command") +} diff --git a/internal/milergeo/milergeo_test.go b/internal/milergeo/milergeo_test.go new file mode 100644 index 0000000..8f9063c --- /dev/null +++ b/internal/milergeo/milergeo_test.go @@ -0,0 +1,120 @@ +package milergeo + +import ( + "context" + "errors" + "strings" + "testing" + + "github.com/redis/go-redis/v9" +) + +// fakeRedis answers the three geo commands the way a given Redis would. +type fakeRedis struct { + searchErr error // GEOSEARCH result; nil = supported + radiusErr error + addErr error + locs []redis.GeoLocation + searchHits int + radiusHits int +} + +func (f *fakeRedis) GeoSearchLocation(ctx context.Context, _ string, q *redis.GeoSearchLocationQuery) *redis.GeoSearchLocationCmd { + f.searchHits++ + cmd := redis.NewGeoSearchLocationCmd(ctx, q) + if f.searchErr != nil { + cmd.SetErr(f.searchErr) + } else { + cmd.SetVal(f.locs) + } + return cmd +} + +func (f *fakeRedis) GeoRadius(_ context.Context, _ string, _, _ float64, _ *redis.GeoRadiusQuery) *redis.GeoLocationCmd { + f.radiusHits++ + return redis.NewGeoLocationCmdResult(f.locs, f.radiusErr) +} + +func (f *fakeRedis) GeoAdd(ctx context.Context, _ string, _ ...*redis.GeoLocation) *redis.IntCmd { + cmd := redis.NewIntCmd(ctx) + if f.addErr != nil { + cmd.SetErr(f.addErr) + } else { + cmd.SetVal(1) + } + return cmd +} + +var nearby = []redis.GeoLocation{{Name: "38", Dist: 0.03}, {Name: "23", Dist: 0.4}} + +func TestSearchUsesGeosearchOnRedis62(t *testing.T) { + legacyOnly.Store(false) + f := &fakeRedis{locs: nearby} + got, err := Search(context.Background(), f, 11.0053, 76.9511, 10, 10) + if err != nil || len(got) != 2 || got[0].Name != "38" { + t.Fatalf("got %v, %v", got, err) + } + if f.searchHits != 1 || f.radiusHits != 0 { + t.Fatalf("GEOSEARCH must be used when supported: search=%d radius=%d", f.searchHits, f.radiusHits) + } +} + +// The production failure: Redis < 6.2 has no GEOSEARCH, so every order found +// "no riders". The fallback must find them, and stop retrying GEOSEARCH. +func TestSearchFallsBackOnOldRedis(t *testing.T) { + legacyOnly.Store(false) + t.Cleanup(func() { legacyOnly.Store(false) }) + f := &fakeRedis{searchErr: errors.New("ERR unknown command 'GEOSEARCH', with args beginning with: 'milers:locations'"), locs: nearby} + for i := 0; i < 3; i++ { + got, err := Search(context.Background(), f, 11.0053, 76.9511, 10, 10) + if err != nil || len(got) != 2 { + t.Fatalf("call %d: got %v, %v", i, got, err) + } + } + if f.searchHits != 1 || f.radiusHits != 3 { + t.Fatalf("after the first rejection only GEORADIUS should run: search=%d radius=%d", f.searchHits, f.radiusHits) + } + if p := Probe(context.Background(), f); !strings.Contains(p, "GEORADIUS fallback") { + t.Fatalf("probe = %q", p) + } +} + +// Any other failure must surface, not read as an empty street. +func TestSearchReturnsOtherErrors(t *testing.T) { + legacyOnly.Store(false) + f := &fakeRedis{searchErr: errors.New("WRONGTYPE Operation against a key holding the wrong kind of value")} + if _, err := Search(context.Background(), f, 11, 76, 10, 10); err == nil || !strings.Contains(err.Error(), "WRONGTYPE") { + t.Fatalf("want the WRONGTYPE error, got %v", err) + } + if f.radiusHits != 0 || legacyOnly.Load() { + t.Fatal("a non-'unknown command' error must not switch to the fallback") + } + if p := Probe(context.Background(), f); !strings.HasPrefix(p, "error: ") { + t.Fatalf("probe = %q", p) + } +} + +func TestIndexReportsWriteFailure(t *testing.T) { + if err := Index(context.Background(), &fakeRedis{}, 38, 11.0, 76.9); err != nil { + t.Fatalf("ok write: %v", err) + } + err := Index(context.Background(), &fakeRedis{addErr: errors.New("WRONGTYPE Operation")}, 38, 11.0, 76.9) + if err == nil || !strings.Contains(err.Error(), "WRONGTYPE") { + t.Fatalf("want the write error, got %v", err) + } +} + +// db.Rdb is a *redis.Client; before InitRedis it is a nil pointer, which in +// an interface is not == nil. It must be refused, not dereferenced. +func TestNilClientIsRefused(t *testing.T) { + var rdb *redis.Client + if _, err := Search(context.Background(), rdb, 11, 76, 10, 10); err == nil { + t.Fatal("nil client must be an error") + } + if err := Index(context.Background(), rdb, 1, 11, 76); err == nil { + t.Fatal("nil client must be an error") + } + if p := Probe(context.Background(), rdb); !strings.HasPrefix(p, "error") { + t.Fatalf("probe = %q", p) + } +} diff --git a/main.go b/main.go index 8f14a7c..5af2200 100644 --- a/main.go +++ b/main.go @@ -1,6 +1,7 @@ package main import ( + "context" "errors" "net/url" "os" @@ -15,6 +16,7 @@ import ( "doormile/internal/ai/playground" "doormile/internal/ai/telemetry" "doormile/internal/assignment" + "doormile/internal/milergeo" "doormile/internal/notify" "doormile/internal/routing" "doormile/internal/sms" @@ -100,6 +102,17 @@ func main() { // 2. Connect to Postgres, Redis & NATS db.Connect(cfg) db.InitRedis(cfg) + // Auto-assignment finds riders with a Redis GEO search. When that search + // fails every order silently finds "no riders", so say so at boot. + if db.Rdb != nil { + geoCtx, geoCancel := context.WithTimeout(context.Background(), 3*time.Second) + if geo := milergeo.Probe(geoCtx, db.Rdb); strings.HasPrefix(geo, "error") { + utils.Error("❌ Rider location search is broken: auto-assignment will find no riders", "redis_geo", geo) + } else { + utils.Info("✅ Rider location search ready", "redis_geo", geo) + } + geoCancel() + } db.InitNATS(cfg) notify.InitFCM() diff --git a/routes/routes.go b/routes/routes.go index ff3dfc5..651a424 100644 --- a/routes/routes.go +++ b/routes/routes.go @@ -7,6 +7,7 @@ import ( "doormile/config" "doormile/controllers" "doormile/db" + "doormile/internal/milergeo" "doormile/internal/sms" "doormile/internal/ws" "doormile/middlewares" @@ -69,11 +70,21 @@ func RegisterRoutes(app *fiber.App, cfg *config.Config) { status = fiber.StatusServiceUnavailable } + // Whether rider search works on this Redis. Reported, not gating: + // a broken search stops auto-assignment, not the API. + redisGeo := "error: redis not available" + if db.Rdb != nil { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + redisGeo = milergeo.Probe(ctx, db.Rdb) + } + return c.Status(status).JSON(fiber.Map{ "status": "ready", "checks": fiber.Map{ - "postgres": dbStatus, - "redis": redisStatus, + "postgres": dbStatus, + "redis": redisStatus, + "redis_geo": redisGeo, // Reported, but deliberately NOT gating readiness: the miler and // console surfaces work perfectly without SMS. It is here because // 'the OTP never arrived' was answerable only by reading code, and