milergeo added

This commit is contained in:
2026-10-06 12:26:57 +05:30
parent 220e934045
commit 9b94be1f07
10 changed files with 329 additions and 44 deletions

View File

@@ -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

View File

@@ -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 {

View File

@@ -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")
}

View File

@@ -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) {

View File

@@ -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

View File

@@ -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)
}
}
}

View File

@@ -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")
}

View File

@@ -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)
}
}

13
main.go
View File

@@ -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()

View File

@@ -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