updates on the ai agents and time series prediction and updates on the api to

This commit is contained in:
2026-10-09 13:50:30 +05:30
parent 29690c56f2
commit 153be40e5c
42 changed files with 4889 additions and 186 deletions

View File

@@ -0,0 +1,240 @@
package outcomes
import (
"encoding/json"
"time"
"doormile/constants"
"doormile/models"
"doormile/utils"
"gorm.io/gorm"
)
// Seeding the decision memory from history, so recall is not useless for weeks.
//
// ─── The cold-start problem this solves ────────────────────────────────────
//
// `/internal/agent-decisions/similar` filters on `outcome IS NOT NULL`. A row
// becomes eligible only after the outcome sweeper has judged it, and a decision
// is only judged once its window has closed. So the sequence for a freshly
// switched-on memory is:
//
// switch embeddings on -> wait for stalls to happen -> wait 48h per decision
// -> the sweeper labels them -> only now does recall return anything
//
// With autonomy gates off and two decision types, that is weeks of recall
// returning [] — which reads as "retrieval does not help here" rather than
// "retrieval has nothing to retrieve yet". The first conclusion is wrong and
// expensive to un-learn.
//
// This backfills decisions from bookings whose outcome is ALREADY known.
//
// ─── What it does and does not claim ───────────────────────────────────────
//
// A backfilled row is not a decision the engine made. It is a record of a
// situation that occurred and how it ended, shaped so the retrieval path can
// use it as precedent. That distinction is recorded honestly:
//
// decision.action = "none" — nothing was decided; nobody intervened
// decision.source = "backfill" — so these are distinguishable forever
// reasoning = states plainly that this is historical, not a decision
//
// Why that matters: as_prompt_block renders `action` to the model. Writing a
// plausible-looking action here would teach the model that an action it never
// took produced the outcome that followed — which is worse than no memory.
// "none" is honest: this is what happened when nothing was done.
//
// ─── It writes no embeddings ───────────────────────────────────────────────
//
// Embedding is the engine's job (AI_engine/core/embeddings.py) and needs the
// provider key, which this process does not have. Backfilled rows land with
// NULL context_embedding and are invisible to similarity search until
// something embeds them. That is deliberate: a backfill that silently created
// un-embedded rows AND claimed to have seeded the memory would be the worse
// failure. BackfillStats reports the count so the caller knows what is owed.
// BackfillStats is what one run produced.
type BackfillStats struct {
Scanned int `json:"scanned"`
Inserted int `json:"inserted"`
Skipped int `json:"skipped"`
// NeedsEmbedding is Inserted — every backfilled row still needs a vector
// before it can be retrieved. Surfaced separately so it cannot be missed.
NeedsEmbedding int `json:"needsembedding"`
}
// historicalRow is one past booking whose ending is known.
//
// Joins through consignment_booking, never consignments.bookingid — that column
// does not exist, and the legacy pickupbookings.consignmentid link names only
// the first order of a multi-destination pickup (hazard H4).
type historicalRow struct {
Bookingid int
Tenantid *uint64
Deliverypincode string
Status string
Createdat time.Time
Deliveredat *time.Time
Sladueat *time.Time
Attemptcount int
Chargeableweight float64
}
const historicalSQL = `
SELECT pb.bookingid,
pb.tenantid AS tenantid,
COALESCE(c.deliverypincode, '') AS deliverypincode,
COALESCE(c.status, pb.status) AS status,
c.createdat AS createdat,
del.createdat AS deliveredat,
c.sladueat AS sladueat,
COALESCE(c.attemptcount, 0) AS attemptcount,
COALESCE(c.chargeableweight, 0) AS chargeableweight
FROM consignments c
JOIN consignment_booking cb ON cb.consignmentid = c.consignmentid
JOIN pickupbookings pb ON pb.bookingid = cb.bookingid
LEFT JOIN (
SELECT DISTINCT ON (consignmentid) consignmentid, createdat
FROM consignmenthistory
WHERE eventstatus = ?
ORDER BY consignmentid, createdat ASC
) del ON del.consignmentid = c.consignmentid
WHERE c.deletedat IS NULL
-- Only parcels whose story has ended. An in-flight parcel has no outcome to
-- learn from, and guessing one is exactly what this must not do.
AND (del.createdat IS NOT NULL OR c.status IN ?)
-- Not already backfilled or decided for. The uniqueness is on
-- (decision_type, booking_id), enforced here rather than by a constraint
-- because real decisions legitimately repeat for one booking.
AND NOT EXISTS (
SELECT 1 FROM agent_decisions d
WHERE d.booking_id = pb.bookingid AND d.decision_type = ?
)
ORDER BY c.createdat DESC
LIMIT ?
`
// BackfillDecisionType is its own type, kept separate from the engine's
// `stall_response` and `assignment_failure`. Recall is per-type, so backfilled
// precedent is retrievable on purpose and never silently mixed into a type the
// engine thinks it authored.
const BackfillDecisionType = "historical_delivery"
// Backfill writes historical precedent rows. Idempotent: a booking that already
// has a decision of this type is skipped, so re-running adds only what is new.
//
// limit bounds one run — this scans delivery history, which is the largest
// table pair in the database.
func Backfill(gdb *gorm.DB, limit int) (BackfillStats, error) {
if gdb == nil {
return BackfillStats{}, nil
}
if limit <= 0 {
limit = 2000
}
terminal := []string{
constants.ConsignmentDelivered,
"Returned_to_Sender",
"Missing",
"Damaged",
}
var rows []historicalRow
if err := gdb.Raw(historicalSQL,
constants.ConsignmentDelivered, terminal, BackfillDecisionType, limit,
).Scan(&rows).Error; err != nil {
return BackfillStats{}, err
}
stats := BackfillStats{Scanned: len(rows)}
batch := make([]models.AgentDecision, 0, len(rows))
for _, r := range rows {
outcome := OutcomeFailure
switch {
case r.Deliveredat == nil:
// Terminal but never delivered: returned, lost or damaged.
outcome = OutcomeFailure
case r.Sladueat != nil && r.Deliveredat.After(*r.Sladueat):
// Delivered late. Counting this a success would teach that any
// eventual delivery is a good outcome — the same trap the
// stall_response rule avoids.
outcome = OutcomeFailure
default:
outcome = OutcomeSuccess
}
// The facts shape mirrors what AI_engine puts in context.facts, so an
// embedding of a backfilled row sits in the same space as a real one.
// If these diverge, retrieval returns neighbours that are near in
// vector space for the wrong reasons.
facts := map[string]any{
"booking_id": r.Bookingid,
"delivery_pincode": r.Deliverypincode,
"attempt_count": r.Attemptcount,
"chargeable_weight": r.Chargeableweight,
"final_status": r.Status,
}
if r.Deliveredat != nil {
facts["hours_to_deliver"] = int(r.Deliveredat.Sub(r.Createdat).Hours())
}
contextJSON, err := json.Marshal(map[string]any{
"facts": facts,
"model": "none",
"source": "backfill",
})
if err != nil {
stats.Skipped++
continue
}
decisionJSON, err := json.Marshal(map[string]any{
// Honest: nobody decided anything. See the package comment — a
// plausible-looking action here would be a fabricated lesson.
"action": "none",
"confidence": 0.0,
"source": "backfill",
})
if err != nil {
stats.Skipped++
continue
}
recordedAt := utils.DBNow()
batch = append(batch, models.AgentDecision{
DecisionType: BackfillDecisionType,
BookingID: u64(r.Bookingid),
TenantID: r.Tenantid,
Context: string(contextJSON),
Decision: string(decisionJSON),
Reasoning: "Historical outcome backfilled from delivery records. No agent decision was made for this booking; this row records what happened when nothing intervened.",
Outcome: &outcome,
OutcomeRecordedAt: &recordedAt,
CreatedAt: r.Createdat,
})
}
if len(batch) == 0 {
return stats, nil
}
if err := gdb.CreateInBatches(&batch, 200).Error; err != nil {
return stats, err
}
stats.Inserted = len(batch)
stats.NeedsEmbedding = len(batch)
utils.Info("outcomes: backfilled historical precedent",
"scanned", stats.Scanned, "inserted", stats.Inserted, "skipped", stats.Skipped,
"note", "rows have no embedding yet and are not retrievable until one is written")
return stats, nil
}
func u64(n int) *uint64 {
if n <= 0 {
return nil
}
v := uint64(n)
return &v
}

View File

@@ -0,0 +1,158 @@
// Package outcomes decides, after the fact, whether an agent's decision worked.
//
// This is the half of the decision memory that was missing, and without it the
// other half does nothing. `/internal/agent-decisions/similar` filters on
// `outcome IS NOT NULL`, so until something judges a decision it is invisible
// to retrieval. Embeddings could flow for months and every recall would still
// come back empty.
//
// A decision is not precedent because it was made. It is precedent because we
// know how it turned out.
//
// ─── What "worked" means ──────────────────────────────────────────────────
//
// Two decision types exist today, from exactly two call sites in the engine:
//
// stall_response AI_engine/agents/exception_agent.py:456
// assignment_failure AI_engine/agents/dispatch_agent.py:335
//
// Each gets its own definition below. A third type appearing without a rule
// here is left pending rather than guessed at — a wrong label is worse than no
// label, because it teaches the model confidently.
package outcomes
import (
"os"
"strconv"
"strings"
"time"
"doormile/utils"
)
// The outcome vocabulary. Stored in agent_decisions.outcome (varchar 30) and
// read back by the Insights tab, which groups by it.
const (
// OutcomeSuccess — the thing the decision was trying to achieve happened.
OutcomeSuccess = "success"
// OutcomeFailure — it did not.
OutcomeFailure = "failure"
// OutcomeUnknown — the window closed without enough evidence either way.
// Deliberately recorded rather than left pending: a pending row is
// retried every sweep forever, and an unjudgeable decision should stop
// costing a scan. It is excluded from retrieval the same as pending,
// because `outcome IS NOT NULL` is not the only filter that matters —
// see judgeable().
OutcomeUnknown = "unknown"
)
const (
defaultOutcomeWindowHours = 48
// A decision younger than this is left alone: the booking it concerns is
// probably still in flight, and judging it now would record a failure for
// something that simply has not finished.
minDecisionAge = 30 * time.Minute
// Caps one sweep's work; the rest are next sweep's.
sweepBatch = 500
)
// OutcomeWindow is how long after a decision its result is judged.
// AGENT_OUTCOME_WINDOW_HOURS, default 48.
//
// The window is a real tradeoff, not a tuning knob. Too short and a parcel
// that was always going to take three days is recorded as a failure of the
// decision rather than of the promise. Too long and the memory learns slowly.
// 48h matches the Standard service SLA (36h) with headroom.
func OutcomeWindow() time.Duration {
if v := strings.TrimSpace(os.Getenv("AGENT_OUTCOME_WINDOW_HOURS")); v != "" {
if n, err := strconv.Atoi(v); err == nil && n > 0 {
return time.Duration(n) * time.Hour
}
utils.Warn("AGENT_OUTCOME_WINDOW_HOURS is not a positive integer, using the default",
"value", v, "default_hours", defaultOutcomeWindowHours)
}
return defaultOutcomeWindowHours * time.Hour
}
// RetentionDays bounds how long decisions are kept. Longer than the 30 days
// aiagentruns keeps, because old precedent is the whole point of this table —
// but not unbounded, which is what it was.
func RetentionDays() int {
if v := strings.TrimSpace(os.Getenv("AGENT_DECISION_RETENTION_DAYS")); v != "" {
if n, err := strconv.Atoi(v); err == nil && n > 0 {
return n
}
}
return 180
}
// bookingFacts is what the sweeper reads about the booking a decision concerned.
// Timestamps come from columns written with CURRENT_TIMESTAMP defaults, so they
// are consistent with each other and their differences are correct regardless
// of the IST-digits-labelled-UTC convention (utils.DBNow) — the same reasoning
// internal/prediction's calibration relies on.
type bookingFacts struct {
Bookingid int
Status string
Assigned bool
DecidedAt time.Time
DeliveredAt *time.Time
SLADueAt *time.Time
AssignedAt *time.Time
Cancelled bool
}
// judge applies the per-type rule. Returns the outcome and whether the decision
// is judgeable at all; a false means leave it pending for now.
func judge(decisionType string, f bookingFacts, now time.Time, window time.Duration) (string, bool) {
age := now.Sub(f.DecidedAt)
if age < minDecisionAge {
return "", false // too soon; the booking is still in flight
}
switch decisionType {
case "stall_response":
// The agent intervened because a rider had stopped making progress.
// It worked if the parcel reached the customer, and reached them
// within the promise that was in force.
if f.Cancelled {
return OutcomeFailure, true
}
if f.DeliveredAt != nil {
if f.SLADueAt != nil && f.DeliveredAt.After(*f.SLADueAt) {
// Delivered, but late. Counting this as success would teach
// the model that any eventual delivery vindicates the action.
return OutcomeFailure, true
}
return OutcomeSuccess, true
}
if age >= window {
// Window closed, never delivered, not cancelled — stuck.
return OutcomeFailure, true
}
return "", false
case "assignment_failure":
// The agent reasoned about why no rider could be found. It worked if
// the booking subsequently got one.
if f.Cancelled {
return OutcomeFailure, true
}
if f.Assigned && f.AssignedAt != nil && f.AssignedAt.After(f.DecidedAt) {
return OutcomeSuccess, true
}
if age >= window {
return OutcomeFailure, true
}
return "", false
default:
// An unrecognised decision type. Judge it unknown once the window has
// closed so it stops being rescanned, but never guess success or
// failure — a wrong label is worse than no label.
if age >= window {
return OutcomeUnknown, true
}
return "", false
}
}

View File

@@ -0,0 +1,214 @@
package outcomes
import (
"testing"
"time"
)
func at(h int) time.Time {
return time.Date(2026, 10, 5, h, 0, 0, 0, time.UTC)
}
func tp(t time.Time) *time.Time { return &t }
const window = 48 * time.Hour
// ─── Too soon to judge ─────────────────────────────────────────────────────
func TestYoungDecisionIsLeftPending(t *testing.T) {
decided := at(10)
for _, dt := range []string{"stall_response", "assignment_failure", "something_new"} {
_, ok := judge(dt, bookingFacts{DecidedAt: decided}, decided.Add(5*time.Minute), window)
if ok {
t.Errorf("%s: judged a 5-minute-old decision; the booking is still in flight", dt)
}
}
}
func TestPendingWhileInsideTheWindow(t *testing.T) {
decided := at(10)
// An hour in: no delivery yet, but the window has not closed. Recording
// failure here would blame the decision for a parcel still on its way.
if _, ok := judge("stall_response", bookingFacts{DecidedAt: decided}, decided.Add(time.Hour), window); ok {
t.Error("stall_response judged before its window closed")
}
if _, ok := judge("assignment_failure", bookingFacts{DecidedAt: decided}, decided.Add(time.Hour), window); ok {
t.Error("assignment_failure judged before its window closed")
}
}
// ─── stall_response ────────────────────────────────────────────────────────
func TestStallDeliveredOnTimeIsSuccess(t *testing.T) {
decided := at(10)
f := bookingFacts{
DecidedAt: decided,
DeliveredAt: tp(decided.Add(3 * time.Hour)),
SLADueAt: tp(decided.Add(12 * time.Hour)),
}
got, ok := judge("stall_response", f, decided.Add(4*time.Hour), window)
if !ok || got != OutcomeSuccess {
t.Errorf("got %q ok=%v, want success", got, ok)
}
}
// The case that matters most: counting any eventual delivery as success would
// teach the model that the action is always vindicated, which is exactly the
// wrong lesson.
func TestStallDeliveredLateIsFailure(t *testing.T) {
decided := at(10)
f := bookingFacts{
DecidedAt: decided,
DeliveredAt: tp(decided.Add(20 * time.Hour)),
SLADueAt: tp(decided.Add(12 * time.Hour)),
}
got, ok := judge("stall_response", f, decided.Add(21*time.Hour), window)
if !ok || got != OutcomeFailure {
t.Errorf("got %q ok=%v, want failure — delivered 8h past SLA", got, ok)
}
}
func TestStallDeliveredWithNoSLAIsSuccess(t *testing.T) {
decided := at(10)
// No sladueat recorded (an older row). Delivery is the best evidence
// available and there is nothing to call it late against.
f := bookingFacts{DecidedAt: decided, DeliveredAt: tp(decided.Add(5 * time.Hour))}
got, ok := judge("stall_response", f, decided.Add(6*time.Hour), window)
if !ok || got != OutcomeSuccess {
t.Errorf("got %q ok=%v, want success", got, ok)
}
}
func TestStallNeverDeliveredAfterWindowIsFailure(t *testing.T) {
decided := at(10)
got, ok := judge("stall_response", bookingFacts{DecidedAt: decided}, decided.Add(window+time.Hour), window)
if !ok || got != OutcomeFailure {
t.Errorf("got %q ok=%v, want failure", got, ok)
}
}
func TestStallCancelledIsFailureImmediately(t *testing.T) {
decided := at(10)
f := bookingFacts{DecidedAt: decided, Cancelled: true}
// Cancellation is conclusive — no need to wait out the window.
got, ok := judge("stall_response", f, decided.Add(time.Hour), window)
if !ok || got != OutcomeFailure {
t.Errorf("got %q ok=%v, want failure on cancellation", got, ok)
}
}
// ─── assignment_failure ────────────────────────────────────────────────────
func TestAssignmentLaterAssignedIsSuccess(t *testing.T) {
decided := at(10)
f := bookingFacts{
DecidedAt: decided,
Assigned: true,
AssignedAt: tp(decided.Add(2 * time.Hour)),
}
got, ok := judge("assignment_failure", f, decided.Add(3*time.Hour), window)
if !ok || got != OutcomeSuccess {
t.Errorf("got %q ok=%v, want success", got, ok)
}
}
// A rider assigned BEFORE the decision is not evidence the decision worked —
// it is the assignment that was already there. Without the ordering check,
// every reassignment decision would read as an instant success.
func TestAssignmentPredatingTheDecisionIsNotSuccess(t *testing.T) {
decided := at(10)
f := bookingFacts{
DecidedAt: decided,
Assigned: true,
AssignedAt: tp(decided.Add(-2 * time.Hour)),
}
got, ok := judge("assignment_failure", f, decided.Add(time.Hour), window)
if ok && got == OutcomeSuccess {
t.Error("counted a pre-existing assignment as the decision's success")
}
}
func TestAssignmentNeverAssignedAfterWindowIsFailure(t *testing.T) {
decided := at(10)
got, ok := judge("assignment_failure", bookingFacts{DecidedAt: decided}, decided.Add(window+time.Hour), window)
if !ok || got != OutcomeFailure {
t.Errorf("got %q ok=%v, want failure", got, ok)
}
}
func TestAssignmentCancelledIsFailure(t *testing.T) {
decided := at(10)
f := bookingFacts{DecidedAt: decided, Cancelled: true}
got, ok := judge("assignment_failure", f, decided.Add(time.Hour), window)
if !ok || got != OutcomeFailure {
t.Errorf("got %q ok=%v, want failure", got, ok)
}
}
// ─── Unknown types are never guessed ───────────────────────────────────────
func TestUnknownTypeIsNeverSuccessOrFailure(t *testing.T) {
decided := at(10)
// Inside the window: pending.
if _, ok := judge("a_new_decision_type", bookingFacts{DecidedAt: decided}, decided.Add(time.Hour), window); ok {
t.Error("judged an unknown decision type inside its window")
}
// Past it: unknown, so it stops being rescanned — but never a label that
// would teach the model something nobody defined.
got, ok := judge("a_new_decision_type", bookingFacts{DecidedAt: decided}, decided.Add(window+time.Hour), window)
if !ok || got != OutcomeUnknown {
t.Errorf("got %q ok=%v, want unknown", got, ok)
}
}
// A delivered unknown type must still not be called a success: the rule for
// what success means for that type does not exist yet.
func TestUnknownTypeWithDeliveryIsStillUnknown(t *testing.T) {
decided := at(10)
f := bookingFacts{DecidedAt: decided, DeliveredAt: tp(decided.Add(time.Hour))}
got, ok := judge("a_new_decision_type", f, decided.Add(window+time.Hour), window)
if !ok || got != OutcomeUnknown {
t.Errorf("got %q ok=%v, want unknown", got, ok)
}
}
// ─── Configuration ─────────────────────────────────────────────────────────
func TestOutcomeWindowDefaultAndOverride(t *testing.T) {
t.Setenv("AGENT_OUTCOME_WINDOW_HOURS", "")
if got := OutcomeWindow(); got != defaultOutcomeWindowHours*time.Hour {
t.Errorf("default window = %v, want %v", got, defaultOutcomeWindowHours*time.Hour)
}
t.Setenv("AGENT_OUTCOME_WINDOW_HOURS", "12")
if got := OutcomeWindow(); got != 12*time.Hour {
t.Errorf("window = %v, want 12h", got)
}
// Garbage falls back rather than producing a zero window, which would
// judge every decision the instant it passed minDecisionAge.
t.Setenv("AGENT_OUTCOME_WINDOW_HOURS", "not-a-number")
if got := OutcomeWindow(); got != defaultOutcomeWindowHours*time.Hour {
t.Errorf("window on garbage = %v, want the default", got)
}
t.Setenv("AGENT_OUTCOME_WINDOW_HOURS", "0")
if got := OutcomeWindow(); got != defaultOutcomeWindowHours*time.Hour {
t.Errorf("window on 0 = %v, want the default", got)
}
}
func TestRetentionDaysDefaultAndOverride(t *testing.T) {
t.Setenv("AGENT_DECISION_RETENTION_DAYS", "")
if got := RetentionDays(); got != 180 {
t.Errorf("default retention = %d, want 180", got)
}
t.Setenv("AGENT_DECISION_RETENTION_DAYS", "90")
if got := RetentionDays(); got != 90 {
t.Errorf("retention = %d, want 90", got)
}
// Retention must be longer than aiagentruns' 30 days, because old
// precedent is the point of this table. Not enforced in code — asserted
// here so a future change to the default has to confront it.
t.Setenv("AGENT_DECISION_RETENTION_DAYS", "")
if RetentionDays() <= 30 {
t.Error("decision retention is no longer than telemetry retention; precedent will be purged before it is useful")
}
}

View File

@@ -0,0 +1,226 @@
package outcomes
import (
"context"
"os"
"strconv"
"strings"
"time"
"doormile/constants"
"doormile/db"
"doormile/utils"
"gorm.io/gorm"
)
// The outcome sweeper.
//
// Same shape as internal/assignment/sweeper.go and internal/prediction/sweeper.go:
// a ticker, a recover() per tick, and a Redis lock that expires before the next
// tick so one replica of three does the work. Without Redis every replica
// sweeps, which is wasteful but correct — each update is idempotent and scoped
// to rows that are still pending.
const (
defaultOutcomeSweepSeconds = 900 // 15m
minOutcomeSweepSeconds = 120
outcomeLockKey = "ai:outcome-sweep:lock"
)
func sweepInterval() time.Duration {
if v := strings.TrimSpace(os.Getenv("AGENT_OUTCOME_SWEEP_SECONDS")); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= 0 {
if n > 0 && n < minOutcomeSweepSeconds {
n = minOutcomeSweepSeconds
}
return time.Duration(n) * time.Second
}
utils.Warn("AGENT_OUTCOME_SWEEP_SECONDS is not a non-negative integer, using the default",
"value", v, "default_seconds", defaultOutcomeSweepSeconds)
}
return defaultOutcomeSweepSeconds * time.Second
}
// PruneFindings is set at boot by main.go to controllers.PruneAIFindings.
// A function variable rather than a direct call because controllers imports
// most of the codebase, and this package is imported BY controllers' siblings —
// calling it directly would be an import cycle.
var PruneFindings func(retentionDays int) (int, error)
// StartOutcomeSweeper judges pending decisions on a timer, and prunes old ones.
// Call once at boot, in a goroutine.
func StartOutcomeSweeper() {
interval := sweepInterval()
if interval == 0 {
utils.Info("OutcomeSweeper: disabled (AGENT_OUTCOME_SWEEP_SECONDS=0)")
return
}
utils.Info("OutcomeSweeper: started",
"interval", interval.String(), "window", OutcomeWindow().String(),
"retention_days", RetentionDays())
ticker := time.NewTicker(interval)
defer ticker.Stop()
for range ticker.C {
sweepOnce(interval)
}
}
func sweepOnce(interval time.Duration) {
defer func() {
if r := recover(); r != nil {
utils.Error("OutcomeSweeper: panic recovered", "error", r)
}
}()
if db.DB == nil {
return
}
if db.Rdb != nil {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
ttl := interval - 30*time.Second
if ttl <= 0 {
ttl = interval / 2
}
got, err := db.Rdb.SetNX(ctx, outcomeLockKey, "1", ttl).Result()
cancel()
if err == nil && !got {
return
}
}
judged, err := JudgePending(db.DB, time.Now())
if err != nil {
utils.Error("OutcomeSweeper: judging failed", "error", err)
}
pruned, err := Prune(db.DB)
if err != nil {
utils.Error("OutcomeSweeper: prune failed", "error", err)
}
// Skill findings share this tick rather than carrying their own timer:
// one more table to keep tidy, not one more goroutine. Injected so this
// package does not import controllers (which imports everything).
findingsPruned := 0
if PruneFindings != nil {
if n, err := PruneFindings(RetentionDays()); err != nil {
utils.Error("OutcomeSweeper: finding prune failed", "error", err)
} else {
findingsPruned = n
}
}
if judged > 0 || pruned > 0 || findingsPruned > 0 {
utils.Info("OutcomeSweeper: swept",
"judged", judged, "decisions_pruned", pruned, "findings_pruned", findingsPruned)
}
}
// pendingRow is one unjudged decision joined to the booking it concerned.
//
// The join runs through pickupbookings on the decision's own booking_id, which
// the engine supplies — not through consignments, which has no bookingid column
// (the hazard CLAUDE.md §8.5 and docs/prediction-plan.md H4 describe). The
// delivered timestamp therefore comes from consignmenthistory via the
// consignment_booking view, the one place that resolution is written down.
type pendingRow struct {
ID uint64
DecisionType string
Bookingid *int
Status string
Assignedat *time.Time
Assigned bool
Createdat time.Time
Deliveredat *time.Time
Sladueat *time.Time
}
const pendingSQL = `
SELECT d.id AS id,
d.decision_type AS decision_type,
d.booking_id AS bookingid,
COALESCE(pb.status, '') AS status,
ba.assignedat AS assignedat,
(pb.assignedmileruserid IS NOT NULL) AS assigned,
d.created_at AS createdat,
del.createdat AS deliveredat,
con.sladueat AS sladueat
FROM agent_decisions d
LEFT JOIN pickupbookings pb ON pb.bookingid = d.booking_id
LEFT JOIN (
SELECT DISTINCT ON (bookingid) bookingid, assignedat
FROM bookingassignments
ORDER BY bookingid, assignedat DESC NULLS LAST, bookingassignmentid DESC
) ba ON ba.bookingid = d.booking_id
LEFT JOIN consignment_booking cb ON cb.bookingid = d.booking_id
LEFT JOIN consignments con ON con.consignmentid = cb.consignmentid
LEFT JOIN (
SELECT DISTINCT ON (consignmentid) consignmentid, createdat
FROM consignmenthistory
WHERE eventstatus = ?
ORDER BY consignmentid, createdat ASC
) del ON del.consignmentid = cb.consignmentid
WHERE d.outcome IS NULL
ORDER BY d.created_at ASC
LIMIT ?
`
// JudgePending labels every pending decision it can and returns how many it
// wrote. A decision it cannot judge yet is left pending for the next sweep.
func JudgePending(gdb *gorm.DB, now time.Time) (int, error) {
if gdb == nil {
return 0, nil
}
var rows []pendingRow
if err := gdb.Raw(pendingSQL, constants.ConsignmentDelivered, sweepBatch).Scan(&rows).Error; err != nil {
return 0, err
}
window := OutcomeWindow()
judged := 0
for _, r := range rows {
facts := bookingFacts{
Status: r.Status,
Assigned: r.Assigned,
DecidedAt: r.Createdat,
DeliveredAt: r.Deliveredat,
SLADueAt: r.Sladueat,
AssignedAt: r.Assignedat,
Cancelled: r.Status == constants.BookingCancelled,
}
if r.Bookingid != nil {
facts.Bookingid = *r.Bookingid
}
outcome, ok := judge(r.DecisionType, facts, now, window)
if !ok {
continue
}
// Guarded on outcome IS NULL so two replicas sweeping concurrently
// cannot overwrite each other, and a decision is judged exactly once.
res := gdb.Exec(
`UPDATE agent_decisions SET outcome = ?, outcome_recorded_at = ? WHERE id = ? AND outcome IS NULL`,
outcome, utils.DBNow(), r.ID,
)
if res.Error != nil {
utils.Warn("OutcomeSweeper: could not record outcome", "id", r.ID, "error", res.Error)
continue
}
judged += int(res.RowsAffected)
}
return judged, nil
}
// Prune drops decisions past the retention window. agent_decisions had no
// retention at all while aiagentruns purged at 30 days — and this is the table
// the similarity query scans, so unbounded growth degrades every recall.
func Prune(gdb *gorm.DB) (int, error) {
if gdb == nil {
return 0, nil
}
cutoff := utils.DBNow().AddDate(0, 0, -RetentionDays())
res := gdb.Exec(`DELETE FROM agent_decisions WHERE created_at < ?`, cutoff)
if res.Error != nil {
return 0, res.Error
}
return int(res.RowsAffected), nil
}

View File

@@ -466,8 +466,17 @@ func TestSkillsWithNoDataSourceShipDisabled(t *testing.T) {
}
}
// Only notify_riders has an executor in the console. Every other console write
// must say REVIEW ONLY, or Agent Studio would advertise an action that cannot run.
// Four console verbs have executors now: notify_riders, alert_low_battery_rider,
// assign_riders and trigger_auto_dispatch (the last two share the batch-assign
// solver behind POST /admin/bookings/batch-assign). Every console write that
// still has none must say REVIEW ONLY, or Agent Studio would advertise an
// action that cannot run.
//
// The invariant reads Implementedat, not a list kept here: a tool whose
// Implementedat still points into lib/assistant/skills/ has no executor, and
// one pointing at lib/assistant/agent/actions.js does. That is why wiring an
// executor means moving Implementedat — the test cannot be satisfied by
// editing the description alone.
func TestConsoleWritesWithoutExecutorSayReviewOnly(t *testing.T) {
for _, tl := range SeedTools {
if !strings.HasPrefix(tl.Implementedat, consoleSrc+"lib/assistant/skills/") && !strings.Contains(tl.Implementedat, "(no executor)") {

View File

@@ -168,21 +168,21 @@ var SeedTools = []models.AITool{
Description: "Message the riders on a finding's orders. Runs only when an operator clicks it; partial success is reported as partial.", Target: "doormile_backend POST /admin/milers/:id/notify",
Implementedat: consoleSrc + "lib/assistant/agent/actions.js", Inputschema: open},
{Toolname: "assign_riders", Kind: KindWrite, Requiresconfirmation: true,
Description: "REVIEW ONLY. Would assign a finding's orders to riders, but POST /hub/bookings/batch-assign accepts hub-staff logins only, so the console cannot run it.", Target: "doormile_backend POST /hub/bookings/batch-assign (hub staff only)",
Implementedat: consoleSrc + "lib/assistant/agent/actions.js (no executor)", Inputschema: open},
Description: "Assign a finding's orders to riders. Greedy nearest on-duty rider, capped per rider, committed server-side — there is no preview step, so the operator's click is the gate.", Target: "doormile_backend POST /admin/bookings/batch-assign",
Implementedat: consoleSrc + "lib/assistant/agent/actions.js", Inputschema: open},
{Toolname: "enforce_otp_verification", Kind: KindWrite, Requiresconfirmation: true,
Description: "REVIEW ONLY. Flag a high-value COD order as requiring the receiver's OTP at handover. No executor.", Target: "none yet",
Implementedat: consoleSrc + "lib/assistant/skills/definitions/HighValueCodSkill.js", Inputschema: obj(map[string]map[string]string{"bookingId": integer("Booking to flag")}, "bookingId")},
{Toolname: "alert_low_battery_rider", Kind: KindNotify, Requiresconfirmation: true,
Description: "REVIEW ONLY. Tell a rider on a low battery to charge or report to the nearest hub. No executor.", Target: "doormile_backend POST /admin/milers/:id/notify",
Implementedat: consoleSrc + "lib/assistant/skills/definitions/RiderBatterySafetySkill.js", Inputschema: obj(map[string]map[string]string{"milerId": integer("Rider to alert")}, "milerId")},
Description: "Tell a rider on a low battery to charge or report to the nearest hub.", Target: "doormile_backend POST /admin/milers/:id/notify",
Implementedat: consoleSrc + "lib/assistant/agent/actions.js", Inputschema: obj(map[string]map[string]string{"milerId": integer("Rider to alert")}, "milerId")},
{Toolname: "dispatch_hub_idle_parcels", Kind: KindWrite, Requiresconfirmation: true,
Description: "REVIEW ONLY. Send an idle rider to collect parcels dwelling at a hub. No executor.", Target: "none yet",
Implementedat: consoleSrc + "lib/assistant/skills/definitions/HubCongestionSkill.js",
Inputschema: obj(map[string]map[string]string{"hubId": str("Hub where parcels are waiting"), "milerId": integer("Idle rider")}, "hubId", "milerId")},
{Toolname: "trigger_auto_dispatch", Kind: KindWrite, Requiresconfirmation: true,
Description: "REVIEW ONLY. Auto-assign orders that have waited too long for dispatch. No executor.", Target: "none yet",
Implementedat: consoleSrc + "lib/assistant/skills/definitions/LateDispatchSkill.js", Inputschema: open},
Description: "Auto-assign orders that have waited too long for dispatch. The same batch-assign action as assign_riders, under the name LateDispatchSkill proposes it by.", Target: "doormile_backend POST /admin/bookings/batch-assign",
Implementedat: consoleSrc + "lib/assistant/agent/actions.js", Inputschema: open},
{Toolname: "enforce_cash_handoff", Kind: KindWrite, Requiresconfirmation: true,
Description: "REVIEW ONLY. Route a rider carrying too much COD via the nearest hub. No executor.", Target: "none yet",
Implementedat: consoleSrc + "lib/assistant/skills/definitions/CashExposureSkill.js", Inputschema: open},

View File

@@ -26,7 +26,25 @@ const (
maxRetries = 5
retryDelay = 2 * time.Minute
geoRadiusKm = 10.0
geoMaxCount = 10
// How many GEO members to fetch per search.
//
// Raised from 10, and the reason matters. GEOSEARCH returns the N NEAREST
// members, and eligibility (GPS freshness, availability, the active cap) is
// applied AFTERWARDS. So any member that is near but not eligible consumes
// one of the N slots and pushes a usable rider out of the result entirely.
//
// Stale members were the worst case: nothing removed a rider from the set
// when they went off duty, so riders who finished weeks ago sat frozen near
// the hub where most pickups originate — the nearest members there are.
// Ten of those filled every slot, all ten were discarded, and the pool
// collapsed to one rider who then took every order in the city.
//
// MilerEndDuty now removes riders (internal/milergeo.Remove), which fixes
// it going forward. This is the defence for members ALREADY in a live
// Redis, and for any future reason a nearby rider turns out ineligible.
// 50 members is a cheap read and leaves room for the filters.
geoMaxCount = 50
// defaultMaxActive is how many open stops one miler may hold at once.
// Override with MILER_MAX_ACTIVE_BOOKINGS — how many parcels a rider can

View File

@@ -15,6 +15,7 @@ import (
"context"
"fmt"
"reflect"
"strconv"
"strings"
"sync/atomic"
@@ -31,6 +32,35 @@ 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
ZRem(ctx context.Context, key string, members ...interface{}) *redis.IntCmd
}
// Remove drops a rider from the GEO set.
//
// ─── Why this has to exist ─────────────────────────────────────────────────
//
// A GEO set member never expires. Nothing removed a rider when they went off
// duty, so `milers:locations` accumulated every rider who had ever started a
// shift, frozen at their last reported position — for good.
//
// That is not merely untidy, because assignment asks for the TEN NEAREST
// members (geoMaxCount in internal/assignment). Riders who finished weeks ago,
// parked near the hub where most pickups originate, are geographically the
// closest members there are. They fill all ten slots, every one of them is
// then discarded by the GPS-freshness check, and the candidate pool collapses
// to whoever happened to survive — often a single rider.
//
// The balancing logic below that (leastLoaded, betterChoice) is correct and
// irrelevant: it can only balance across the pool it is handed. The symptom is
// every order in a city going to the same person.
//
// Called on end-duty. A failure is logged by the caller and not fatal: a rider
// left in the set is the status quo, not a regression.
func Remove(ctx context.Context, rdb Client, milerUserID int) error {
if rdb == nil {
return nil
}
return rdb.ZRem(ctx, Key, strconv.Itoa(milerUserID)).Err()
}
// legacyOnly flips to true the first time the server rejects GEOSEARCH as an

View File

@@ -17,6 +17,19 @@ type fakeRedis struct {
locs []redis.GeoLocation
searchHits int
radiusHits int
remErr error
removed []interface{}
}
func (f *fakeRedis) ZRem(ctx context.Context, _ string, members ...interface{}) *redis.IntCmd {
f.removed = append(f.removed, members...)
cmd := redis.NewIntCmd(ctx)
if f.remErr != nil {
cmd.SetErr(f.remErr)
} else {
cmd.SetVal(int64(len(members)))
}
return cmd
}
func (f *fakeRedis) GeoSearchLocation(ctx context.Context, _ string, q *redis.GeoSearchLocationQuery) *redis.GeoSearchLocationCmd {
@@ -118,3 +131,45 @@ func TestNilClientIsRefused(t *testing.T) {
t.Fatalf("probe = %q", p)
}
}
// A GEO member never expires, and nothing removed one when a rider went off
// duty — so milers:locations kept every rider who had ever started a shift,
// frozen where they last reported. Assignment takes the N NEAREST members, so
// those stale entries crowded out riders who were actually working and the
// candidate pool collapsed to whoever survived the freshness filter.
func TestRemoveDropsTheRiderFromTheSet(t *testing.T) {
f := &fakeRedis{}
if err := Remove(context.Background(), f, 412); err != nil {
t.Fatalf("Remove: %v", err)
}
if len(f.removed) != 1 || f.removed[0] != "412" {
t.Errorf("removed = %v, want [\"412\"]", f.removed)
}
}
// The member is the rider's userid as a DECIMAL STRING — the same spelling
// GeoAdd writes and Search reads back. A mismatch here would remove nothing
// and report success.
func TestRemoveUsesTheSameMemberSpellingAsAdd(t *testing.T) {
f := &fakeRedis{}
_ = Remove(context.Background(), f, 7)
if f.removed[0] != "7" {
t.Errorf("member = %v, want \"7\" (decimal string, not an int)", f.removed[0])
}
}
// Best-effort at the call site: the rider is off duty either way, and the
// freshness check still excludes them. The error must reach the caller so it
// can be logged rather than swallowed here.
func TestRemoveReturnsTheError(t *testing.T) {
f := &fakeRedis{remErr: errors.New("redis down")}
if err := Remove(context.Background(), f, 1); err == nil {
t.Error("Remove swallowed the error")
}
}
func TestRemoveOnNilClientIsANoOp(t *testing.T) {
if err := Remove(context.Background(), nil, 1); err != nil {
t.Errorf("Remove(nil) = %v, want nil", err)
}
}

View File

@@ -0,0 +1,272 @@
package prediction
import (
"time"
"doormile/constants"
"doormile/db"
"doormile/utils"
"gorm.io/gorm"
)
// The calibration refresh.
//
// Every delivered consignment whose rider was sequenced carries two numbers:
// what the Route Optimization API predicted (bookingassignments.etaminutes,
// written by internal/routing) and what actually happened (the Delivered row in
// consignmenthistory). The ratio between them, grouped and taken at p80, is the
// factor ETAMinutes multiplies by.
//
// Why p80 and not the mean: docs/prediction-plan.md §8 decision 3. An ETA is a
// promise, and a mean is late half the time.
//
// Why a view and not a join here: consignments has no bookingid column (hazard
// H4 in the plan). The consignment_booking view resolves the two paths —
// bookingdestinations for multi-destination pickups, the legacy
// pickupbookings.consignmentid for console and pre-fan-out rows — once, in
// migrate.go. Joining through consignments directly attaches features to the
// wrong parcel for every multi-destination booking, silently.
// ETACalibration is one calibration cell as stored. Written only by Refresh,
// read at boot and after each refresh.
type ETACalibration struct {
Calibrationid int `gorm:"primaryKey;column:calibrationid;autoIncrement"`
Scope string `gorm:"column:scope;size:24;not null;index:idx_etacalibration_lookup,priority:1"`
Zone string `gorm:"column:zone;size:8;index:idx_etacalibration_lookup,priority:2"`
Hourbucket *int `gorm:"column:hourbucket;index:idx_etacalibration_lookup,priority:3"`
Weekday *int `gorm:"column:weekday;index:idx_etacalibration_lookup,priority:4"`
Factor float64 `gorm:"column:factor;not null"`
Handling float64 `gorm:"column:handling;not null;default:0"`
Samples int `gorm:"column:samples;not null"`
Refreshedat time.Time `gorm:"column:refreshedat;not null"`
}
func (ETACalibration) TableName() string { return "etacalibration" }
// row is one group from the refresh query.
type row struct {
Scope string
Zone string
Hourbucket *int
Weekday *int
Factor float64
Samples int
}
// Observed durations below this are almost certainly data errors — a Delivered
// event written in the same second as the consignment, or a backfill. Including
// them drags every factor toward zero.
const minObservedMinutes = 5
// And above this, the parcel sat for days for a reason that has nothing to do
// with road duration (held at hub, re-attempts, disputes). Calibrating a
// routed-duration multiplier on those teaches it the wrong thing.
const maxObservedMinutes = 48 * 60
// refreshSQL computes the p80 ratio of actual to routed duration, at three
// grains plus a global row, in one pass.
//
// The `actual` term uses consignmenthistory.createdat - consignments.createdat.
// Both come from CURRENT_TIMESTAMP defaults, so they share a tagging and their
// difference is correct regardless of hazard H1 — which is why this does not
// go near estimateddeliveryat, whose writes are not uniform
// (adminController.go:2738 uses time.Now(), elsewhere it is CURRENT_TIMESTAMP).
//
// EXTRACT(ISODOW) is 1..7 Monday-first, matching weekdayOf in eta.go. The hour
// bucket divides by 3, matching hourBucketSize. If either changes, both change.
// A booking can hold SEVERAL bookingassignments rows — Assigned, Rejected,
// Reassigned, Cancelled are all statuses it passes through. Joining them all
// would multiply one delivered parcel by its whole assignment history and pull
// the p80 toward whatever got reassigned most, so `assigned` picks exactly one
// row per booking: the most recently sequenced one, which is the routed ETA
// that was actually in force when the parcel was delivered.
//
// NULL::int / NULL::text are explicit because a UNION ALL resolves column types
// across its branches; an untyped NULL can be inferred as text and fail against
// the integer from the first branch — at runtime, against the real database,
// which is exactly where it is most expensive to discover.
const refreshSQL = `
WITH assigned AS (
SELECT DISTINCT ON (ba.bookingid)
ba.bookingid,
ba.etaminutes
FROM bookingassignments ba
WHERE ba.etaminutes > 0
ORDER BY ba.bookingid, ba.sequencedat DESC NULLS LAST, ba.bookingassignmentid DESC
),
observed AS (
SELECT left(c.deliverypincode, 3) AS zone,
(EXTRACT(HOUR FROM c.createdat)::int / ?) AS hourbucket,
EXTRACT(ISODOW FROM c.createdat)::int AS weekday,
a.etaminutes::double precision AS routed,
EXTRACT(EPOCH FROM (h.createdat - c.createdat)) / 60.0 AS actual
FROM consignments c
JOIN consignmenthistory h
ON h.consignmentid = c.consignmentid
AND h.eventstatus = ?
JOIN consignment_booking cb
ON cb.consignmentid = c.consignmentid
JOIN assigned a
ON a.bookingid = cb.bookingid
WHERE c.deliverypincode IS NOT NULL
AND length(c.deliverypincode) >= 3
),
clean AS (
SELECT * FROM observed
WHERE actual BETWEEN ? AND ?
AND routed > 0
)
SELECT 'zone_hour_weekday'::text AS scope, zone, hourbucket, weekday,
percentile_cont(0.8) WITHIN GROUP (ORDER BY actual / routed) AS factor,
count(*)::int AS samples
FROM clean GROUP BY zone, hourbucket, weekday
UNION ALL
SELECT 'zone_weekday'::text, zone, NULL::int, weekday,
percentile_cont(0.8) WITHIN GROUP (ORDER BY actual / routed), count(*)::int
FROM clean GROUP BY zone, weekday
UNION ALL
SELECT 'zone'::text, zone, NULL::int, NULL::int,
percentile_cont(0.8) WITHIN GROUP (ORDER BY actual / routed), count(*)::int
FROM clean GROUP BY zone
UNION ALL
SELECT 'global'::text, ''::text, NULL::int, NULL::int,
percentile_cont(0.8) WITHIN GROUP (ORDER BY actual / routed), count(*)::int
FROM clean
`
// Refresh recomputes the calibration from history, replaces the stored table,
// and swaps the in-memory snapshot.
//
// A failure leaves the previous calibration in place — a refresh that cannot
// run is not a reason to stop answering with the last good factors. If they go
// stale past staleAfter, ETAMinutes stops trusting them on its own.
func Refresh(gdb *gorm.DB) error {
if gdb == nil {
return nil
}
var rows []row
if err := gdb.Raw(refreshSQL,
hourBucketSize, constants.ConsignmentDelivered, minObservedMinutes, maxObservedMinutes,
).Scan(&rows).Error; err != nil {
return err
}
kept := make([]ETACalibration, 0, len(rows))
now := utils.DBNow()
for _, r := range rows {
if r.Samples < minSamples {
continue // a p80 over fewer than minSamples is noise
}
if r.Factor < minFactor || r.Factor > maxFactor {
// Out of bounds is a signal about the data, not a usable factor.
// Log it once here rather than discovering it per request.
utils.Warn("prediction: discarding out-of-bounds calibration factor",
"scope", r.Scope, "zone", r.Zone, "factor", r.Factor, "samples", r.Samples)
continue
}
kept = append(kept, ETACalibration{
Scope: r.Scope,
Zone: r.Zone,
Hourbucket: r.Hourbucket,
Weekday: r.Weekday,
Factor: r.Factor,
Samples: r.Samples,
Refreshedat: now,
})
}
if len(kept) == 0 {
// No cell cleared the floor. Expected until routing is on and history
// accumulates; ETAMinutes keeps returning false and the promise tables
// keep answering.
utils.Info("prediction: no calibration cells met the sample floor",
"groups_considered", len(rows), "min_samples", minSamples)
return nil
}
// Replace wholesale in one transaction: a partial table would serve a mix
// of old and new factors for the same zone.
if err := gdb.Transaction(func(tx *gorm.DB) error {
if err := tx.Exec(`DELETE FROM etacalibration`).Error; err != nil {
return err
}
return tx.CreateInBatches(&kept, 200).Error
}); err != nil {
return err
}
Load(kept)
utils.Info("prediction: calibration refreshed", "cells", len(kept))
return nil
}
// Load swaps the in-memory snapshot. Exported so boot can populate from the
// stored table without recomputing, and so tests can install a known
// calibration without a database.
//
// builtAt is the NEWEST Refreshedat among the rows, not time.Now(): a restart
// that loads month-old rows from the store must still look month-old to the
// staleness guard in ETAMinutes. Taking the load time here would silently
// re-arm a stale calibration on every deploy. A row with no Refreshedat (a test
// fixture) is treated as fresh.
// Refreshedat is stored through utils.DBNow, which writes IST wall-clock digits
// labelled UTC. Reading it back as an instant is off by 5h30m, so it goes
// through utils.IST first — the same correction utils/epoch.go exists for.
// Without it a fresh calibration reads as 5.5 hours old, which is the very
// hazard docs/prediction-plan.md calls H1.
func Load(rows []ETACalibration) {
t := &table{cells: make(map[string]cell, len(rows))}
for _, r := range rows {
if r.Refreshedat.IsZero() {
continue
}
if at := utils.IST(r.Refreshedat); at.After(t.builtAt) {
t.builtAt = at
}
}
if t.builtAt.IsZero() {
t.builtAt = time.Now() // a fixture with no Refreshedat counts as fresh
}
for _, r := range rows {
c := cell{factor: r.Factor, handling: r.Handling, samples: r.Samples}
switch r.Scope {
case "global":
t.global, t.hasGlobal = c, true
case "zone":
t.cells[keyZone(r.Zone)] = c
case "zone_weekday":
if r.Weekday != nil {
t.cells[keyZoneWeekday(r.Zone, *r.Weekday)] = c
}
case "zone_hour_weekday":
if r.Weekday != nil && r.Hourbucket != nil {
t.cells[keyZoneHourWeekday(r.Zone, *r.Hourbucket, *r.Weekday)] = c
}
}
}
current.store(t)
}
// Reset clears the in-memory calibration. Tests only — it makes ETAMinutes
// return false, which is the state a fresh process is in before boot loads.
func Reset() { current.store(nil) }
// LoadFromDB populates the snapshot from the stored table at boot, so a restart
// does not wait for the next refresh to start answering.
func LoadFromDB() {
if db.DB == nil {
return
}
var rows []ETACalibration
if err := db.DB.Find(&rows).Error; err != nil {
utils.Warn("prediction: could not load stored calibration", "error", err)
return
}
if len(rows) == 0 {
return
}
Load(rows)
utils.Info("prediction: calibration loaded from store", "cells", len(rows))
}

277
internal/prediction/eta.go Normal file
View File

@@ -0,0 +1,277 @@
// Package prediction turns Doormile's own past predictions into better ones.
//
// Doormile already promises a delivery time, in two places and both of them
// fixed tables: controllers/adminController.go uses service-type constants
// (Normal 24h, Fast 12h, Superfast 6h) and controllers/cxPickupFanout.go uses
// the destination district's `promise` column (same-day/next-day/2-day/3-day).
// Neither looks at a single delivery that actually happened.
//
// This package does. The Route Optimization API already returns a road-network
// duration per stop and internal/routing already stores it on
// bookingassignments.etaminutes. Pairing that stored prediction with the
// Delivered event in consignmenthistory gives a predicted-vs-actual series the
// system produced itself, and a grouped p80 over it is a calibration factor:
//
// eta = routed_minutes × factor[zone, hour_bucket, weekday] + handling[zone]
//
// That is a median, not a model. No training, no inference server, no new
// runtime — a table refreshed nightly and a map lookup on the booking path.
//
// ─── The floor rule ────────────────────────────────────────────────────────
//
// ETAMinutes returns (0, false) whenever it lacks the evidence to do better,
// and every caller keeps its existing rule for that case. So:
//
// - no routed duration (ROUTE_OPTIMIZER_URL unset, or a single-stop rider) → false
// - no calibration cell with enough samples → false
// - calibration never refreshed, or the refresh failed → false
// - a factor outside sane bounds → false
//
// Today, in the cluster, ROUTE_OPTIMIZER_URL is not set (see Phase 7 Track A1),
// so this returns false for every booking and the promise tables answer exactly
// as they do now. Switch routing on, let a few weeks of deliveries land, and it
// starts answering — with no code change and no deploy. Today's behaviour is
// the floor; this can only raise it.
//
// ─── p80, not the mean ─────────────────────────────────────────────────────
//
// The calibration is the 80th percentile of observed overrun, not the average.
// An ETA shown to a customer is a promise — "arrives by" — and a mean is late
// half the time. docs/prediction-plan.md §8 decision 3.
package prediction
import (
"strconv"
"strings"
"sync"
"time"
"doormile/utils"
)
const (
// minSamples is the floor for trusting one calibration cell. Below it the
// p80 is noise and the lookup falls back to a coarser key.
minSamples = 20
// A factor outside these bounds means the calibration is wrong, not that
// deliveries are 10× their routed duration. Refuse rather than serve it:
// an absurd ETA is worse than today's flat constant.
minFactor = 0.5
maxFactor = 5.0
// maxHandlingMinutes caps the additive hub term for the same reason.
maxHandlingMinutes = 240
// staleAfter is how long a calibration stays usable without a refresh. The
// sweeper runs far more often than this; exceeding it means the refresh has
// been failing silently, and a month-old factor should not keep answering.
staleAfter = 72 * time.Hour
// hourBucketSize groups the day into 8 three-hour buckets. Finer buckets
// split the samples too thin to clear minSamples on real volume.
hourBucketSize = 3
)
// Input is everything the estimate needs. Built by the caller from the booking
// it already has in hand; this package never queries on the request path.
type Input struct {
// RoutedMinutes is bookingassignments.etaminutes — the Route Optimization
// API's road-network duration. Zero means unknown, which is the common case
// today and the whole reason for the floor rule.
RoutedMinutes int
// DeliveryPincode keys the zone. Only its first three digits are used, the
// same grain the hub console filters on (pickuppincode LIKE '641%').
DeliveryPincode string
// At is when the estimate is being made. Interpreted through utils.IST,
// because this database stores IST wall-clock digits and a raw .Hour() on a
// value tagged UTC is off by 5h30m — the defect utils/epoch.go documents.
At time.Time
}
// Result is a calibrated estimate and the cell that produced it. Source is for
// logging and for the console to say where a number came from; nothing branches
// on it.
type Result struct {
Minutes int
Source string // "zone_hour_weekday", "zone_weekday", "zone", "global"
Samples int
}
// cell is one calibration row, in memory.
type cell struct {
factor float64
handling float64
samples int
}
// table is an immutable calibration snapshot. Replaced wholesale by the
// refresh; readers never see a half-updated map.
type table struct {
cells map[string]cell
global cell
hasGlobal bool
builtAt time.Time
}
var (
current atomic[*table]
)
// atomic is a tiny generic holder. sync/atomic.Pointer would do, but this keeps
// the zero value useful (an unset calibration reads as nil, which ETAMinutes
// treats as "no evidence") without an init func.
type atomic[T any] struct {
mu sync.RWMutex
v T
}
func (a *atomic[T]) load() T {
a.mu.RLock()
defer a.mu.RUnlock()
return a.v
}
func (a *atomic[T]) store(v T) {
a.mu.Lock()
a.v = v
a.mu.Unlock()
}
// zoneOf is the first three digits of a pincode — the same grain the hub
// console scopes on. An empty or short pincode has no zone, which the lookup
// treats as "global only".
func zoneOf(pincode string) string {
p := strings.TrimSpace(pincode)
if len(p) < 3 {
return ""
}
return p[:3]
}
// hourBucket is the 3-hour block of the IST day, 0..7.
func hourBucket(t time.Time) int {
return utils.IST(t).Hour() / hourBucketSize
}
// weekdayOf is the ISO weekday in IST, 1 (Monday) to 7 (Sunday) — matching
// Postgres's EXTRACT(ISODOW) so the Go lookup and the refresh SQL agree.
func weekdayOf(t time.Time) int {
d := int(utils.IST(t).Weekday())
if d == 0 {
return 7 // Go's Sunday is 0; ISO's is 7
}
return d
}
// keyZoneHourWeekday, keyZoneWeekday and keyZone are the three progressively
// coarser lookups. Distinct prefixes so a zone can never collide with a
// weekday-qualified key.
func keyZoneHourWeekday(zone string, bucket, weekday int) string {
return "zhw:" + zone + ":" + strconv.Itoa(bucket) + ":" + strconv.Itoa(weekday)
}
func keyZoneWeekday(zone string, weekday int) string {
return "zw:" + zone + ":" + strconv.Itoa(weekday)
}
func keyZone(zone string) string {
return "z:" + zone
}
// ETAMinutes returns a calibrated door-to-door estimate in minutes, and whether
// it is trustworthy.
//
// False means the caller must use whatever it does today — its service-type
// constants or the district promise table. A false is not an error and is not
// logged per call: it is the expected answer until routing is switched on and
// history accumulates.
func ETAMinutes(in Input) (Result, bool) {
if in.RoutedMinutes <= 0 {
// No road-network duration to calibrate against. The dominant case
// today: ROUTE_OPTIMIZER_URL is unset in the cluster, and even with it
// set, internal/routing skips riders with fewer than two stops.
return Result{}, false
}
t := current.load()
if t == nil || len(t.cells) == 0 && !t.hasGlobal {
return Result{}, false
}
if !t.builtAt.IsZero() && time.Since(t.builtAt) > staleAfter {
// A refresh has been failing for days. Fall back rather than serve a
// factor that predates whatever changed.
return Result{}, false
}
zone := zoneOf(in.DeliveryPincode)
bucket, weekday := hourBucket(in.At), weekdayOf(in.At)
type candidate struct {
key string
source string
}
candidates := []candidate{}
if zone != "" {
candidates = append(candidates,
candidate{keyZoneHourWeekday(zone, bucket, weekday), "zone_hour_weekday"},
candidate{keyZoneWeekday(zone, weekday), "zone_weekday"},
candidate{keyZone(zone), "zone"},
)
}
for _, c := range candidates {
if cl, ok := t.cells[c.key]; ok && cl.samples >= minSamples {
if m, ok := apply(in.RoutedMinutes, cl); ok {
return Result{Minutes: m, Source: c.source, Samples: cl.samples}, true
}
}
}
if t.hasGlobal && t.global.samples >= minSamples {
if m, ok := apply(in.RoutedMinutes, t.global); ok {
return Result{Minutes: m, Source: "global", Samples: t.global.samples}, true
}
}
return Result{}, false
}
// apply is the estimate itself, with the bounds check that keeps a bad
// calibration from producing an absurd promise.
func apply(routedMinutes int, c cell) (int, bool) {
if c.factor < minFactor || c.factor > maxFactor {
return 0, false
}
if c.handling < 0 || c.handling > maxHandlingMinutes {
return 0, false
}
m := float64(routedMinutes)*c.factor + c.handling
if m <= 0 {
return 0, false
}
return int(m + 0.5), true
}
// ETAAt is ETAMinutes as an absolute time, for the callers that store a
// timestamp rather than a duration. The returned time is in the same shape the
// caller's `from` was, so it round-trips into the database unchanged.
func ETAAt(in Input, from time.Time) (time.Time, Result, bool) {
r, ok := ETAMinutes(in)
if !ok {
return time.Time{}, Result{}, false
}
return from.Add(time.Duration(r.Minutes) * time.Minute), r, true
}
// Loaded reports whether a usable calibration is in memory. For
// GET /admin/ai/status and the readiness note; not used on the booking path.
func Loaded() (bool, time.Time, int) {
t := current.load()
if t == nil {
return false, time.Time{}, 0
}
return len(t.cells) > 0 || t.hasGlobal, t.builtAt, len(t.cells)
}

View File

@@ -0,0 +1,284 @@
package prediction
import (
"testing"
"time"
"doormile/utils"
)
func ptr(n int) *int { return &n }
// A calibration that is deliberately coarse-to-fine, so the fallback ladder is
// observable: the zone_hour_weekday cell has a different factor from the
// zone_weekday one, which differs from zone, which differs from global.
func fixture() []ETACalibration {
return []ETACalibration{
{Scope: "zone_hour_weekday", Zone: "641", Hourbucket: ptr(3), Weekday: ptr(1), Factor: 1.5, Samples: 100},
{Scope: "zone_weekday", Zone: "641", Weekday: ptr(1), Factor: 2.0, Samples: 100},
{Scope: "zone", Zone: "641", Factor: 2.5, Samples: 100},
{Scope: "global", Factor: 3.0, Samples: 100},
}
}
// 2026-10-05 is a Monday. 10:00 IST is hour bucket 3 (10/3). The database
// stores IST digits labelled UTC, so the fixture time is built that way on
// purpose — it is the shape a real column read produces.
func mondayTenAM() time.Time {
return time.Date(2026, 10, 5, 10, 0, 0, 0, time.UTC)
}
// ─── The floor rule: these are the cases that must return false ────────────
func TestNoRoutedDurationFallsBack(t *testing.T) {
Load(fixture())
defer Reset()
// The dominant case in production today: ROUTE_OPTIMIZER_URL is unset, so
// internal/routing never runs and etaminutes is 0.
if _, ok := ETAMinutes(Input{RoutedMinutes: 0, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Fatal("ETAMinutes answered with no routed duration; the promise table must keep answering")
}
if _, ok := ETAMinutes(Input{RoutedMinutes: -5, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Fatal("ETAMinutes answered on a negative routed duration")
}
}
func TestNoCalibrationFallsBack(t *testing.T) {
Reset()
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Fatal("ETAMinutes answered with no calibration loaded")
}
}
func TestBelowSampleFloorFallsBack(t *testing.T) {
Load([]ETACalibration{
{Scope: "zone", Zone: "641", Factor: 2.0, Samples: minSamples - 1},
})
defer Reset()
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Fatalf("ETAMinutes trusted a cell with %d samples (floor is %d)", minSamples-1, minSamples)
}
}
func TestOutOfBoundsFactorFallsBack(t *testing.T) {
for _, f := range []float64{minFactor - 0.1, maxFactor + 0.1, 0, -1} {
Load([]ETACalibration{{Scope: "zone", Zone: "641", Factor: f, Samples: 100}})
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Errorf("ETAMinutes served an out-of-bounds factor %v; an absurd ETA is worse than a flat constant", f)
}
Reset()
}
}
// This is also the H1 regression test, verified by mutation: Refreshedat is
// written through utils.DBNow (IST wall-clock digits labelled UTC), so reading
// it as a raw instant makes it appear 5h30m NEWER than it is. A calibration
// 73h old then measures as 67.5h and slips under the 72h staleAfter bound.
// Removing the utils.IST call in Load fails exactly this test.
//
// The margin matters: with staleAfter at 72h, the 5.5h error is a 7.6% window
// in which a stale calibration keeps answering. Narrow staleAfter and the bug
// gets proportionally worse, which is why the correction belongs in Load rather
// than in a wider bound here.
func TestStaleCalibrationFallsBack(t *testing.T) {
old := utils.DBNow().Add(-(staleAfter + time.Hour))
Load([]ETACalibration{
{Scope: "zone", Zone: "641", Factor: 2.0, Samples: 100, Refreshedat: old},
})
defer Reset()
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Fatal("ETAMinutes trusted a calibration older than staleAfter; refreshes have been failing")
}
}
// The companion case: a just-refreshed calibration must answer. Note this one
// passes with or without the IST correction (the raw read errs toward "newer",
// not "older"), so it documents intent rather than guarding the bug — the guard
// above is what fails on mutation.
func TestFreshCalibrationIsNotMisreadAsStale(t *testing.T) {
Load([]ETACalibration{
{Scope: "zone", Zone: "641", Factor: 2.0, Samples: 100, Refreshedat: utils.DBNow()},
})
defer Reset()
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); !ok {
t.Fatal("a just-refreshed calibration was treated as stale — the DBNow/IST correction in Load is wrong")
}
}
// Regression: loading month-old rows from the store at boot must stay stale.
// Stamping builtAt with time.Now() here would re-arm a dead calibration on
// every deploy.
func TestBootLoadDoesNotRearmStaleRows(t *testing.T) {
old := utils.DBNow().Add(-(staleAfter + 24*time.Hour))
Load([]ETACalibration{
{Scope: "zone", Zone: "641", Factor: 2.0, Samples: 100, Refreshedat: old},
{Scope: "global", Factor: 2.2, Samples: 100, Refreshedat: old},
})
defer Reset()
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Fatal("a restart re-armed a stale stored calibration")
}
}
// ─── The ladder: most specific cell wins, then progressively coarser ───────
func TestMostSpecificCellWins(t *testing.T) {
Load(fixture())
defer Reset()
r, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()})
if !ok {
t.Fatal("expected an estimate")
}
if r.Source != "zone_hour_weekday" {
t.Errorf("source = %q, want zone_hour_weekday", r.Source)
}
if r.Minutes != 60 { // 40 × 1.5
t.Errorf("minutes = %d, want 60 (40 × 1.5)", r.Minutes)
}
}
func TestFallsToZoneWeekdayThenZoneThenGlobal(t *testing.T) {
full := fixture()
defer Reset()
// Drop the finest cell: 13:00 IST is bucket 4, for which there is no
// zone_hour_weekday row, so zone_weekday answers.
Load(full)
r, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001",
At: time.Date(2026, 10, 5, 13, 0, 0, 0, time.UTC)})
if !ok || r.Source != "zone_weekday" || r.Minutes != 80 { // 40 × 2.0
t.Errorf("got %+v ok=%v, want zone_weekday 80", r, ok)
}
// A Tuesday has no weekday cell for this zone, so zone answers.
r, ok = ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001",
At: time.Date(2026, 10, 6, 13, 0, 0, 0, time.UTC)})
if !ok || r.Source != "zone" || r.Minutes != 100 { // 40 × 2.5
t.Errorf("got %+v ok=%v, want zone 100", r, ok)
}
// An unknown zone falls all the way to global.
r, ok = ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "500081", At: mondayTenAM()})
if !ok || r.Source != "global" || r.Minutes != 120 { // 40 × 3.0
t.Errorf("got %+v ok=%v, want global 120", r, ok)
}
}
func TestShortOrEmptyPincodeUsesGlobalOnly(t *testing.T) {
Load(fixture())
defer Reset()
for _, p := range []string{"", "6", "64", " "} {
r, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: p, At: mondayTenAM()})
if !ok || r.Source != "global" {
t.Errorf("pincode %q: got %+v ok=%v, want global", p, r, ok)
}
}
}
// ─── Bucket arithmetic must match the refresh SQL, or the lookup misses ────
func TestHourBucketMatchesRefreshSQL(t *testing.T) {
// refreshSQL does EXTRACT(HOUR FROM createdat)::int / hourBucketSize.
// hourBucket must agree exactly, or every finest-grain lookup misses.
for hour := 0; hour < 24; hour++ {
got := hourBucket(time.Date(2026, 10, 5, hour, 30, 0, 0, time.UTC))
if want := hour / hourBucketSize; got != want {
t.Errorf("hour %d: bucket = %d, want %d", hour, got, want)
}
}
}
func TestWeekdayMatchesPostgresISODOW(t *testing.T) {
// EXTRACT(ISODOW) is Monday=1 … Sunday=7. Go's time.Weekday is Sunday=0.
// 2026-10-05 is a Monday.
want := map[int]int{5: 1, 6: 2, 7: 3, 8: 4, 9: 5, 10: 6, 11: 7}
for day, w := range want {
got := weekdayOf(time.Date(2026, 10, day, 12, 0, 0, 0, time.UTC))
if got != w {
t.Errorf("2026-10-%02d: weekday = %d, want %d (ISODOW)", day, got, w)
}
}
}
func TestZoneIsPincodePrefix(t *testing.T) {
cases := map[string]string{
"641001": "641", "641 ": "641", "500081": "500",
"": "", "64": "", " ": "",
}
for in, want := range cases {
if got := zoneOf(in); got != want {
t.Errorf("zoneOf(%q) = %q, want %q", in, got, want)
}
}
}
// ─── ETAAt and handling term ───────────────────────────────────────────────
func TestETAAtAddsToTheGivenTime(t *testing.T) {
Load(fixture())
defer Reset()
from := mondayTenAM()
at, r, ok := ETAAt(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: from}, from)
if !ok {
t.Fatal("expected an estimate")
}
if want := from.Add(60 * time.Minute); !at.Equal(want) {
t.Errorf("ETAAt = %v, want %v", at, want)
}
// The returned time keeps the caller's location, so it round-trips into the
// database in the same shape it came out.
if at.Location() != from.Location() {
t.Errorf("ETAAt changed the location from %v to %v", from.Location(), at.Location())
}
if r.Minutes != 60 {
t.Errorf("result minutes = %d, want 60", r.Minutes)
}
}
func TestETAAtFallsBackWithoutCalibration(t *testing.T) {
Reset()
from := mondayTenAM()
if _, _, ok := ETAAt(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: from}, from); ok {
t.Fatal("ETAAt answered with no calibration")
}
}
func TestHandlingTermIsAddedAndBounded(t *testing.T) {
Load([]ETACalibration{
{Scope: "zone", Zone: "641", Factor: 2.0, Handling: 15, Samples: 100},
})
r, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()})
if !ok || r.Minutes != 95 { // 40 × 2.0 + 15
t.Errorf("got %+v ok=%v, want 95", r, ok)
}
Reset()
Load([]ETACalibration{
{Scope: "zone", Zone: "641", Factor: 2.0, Handling: maxHandlingMinutes + 1, Samples: 100},
})
if _, ok := ETAMinutes(Input{RoutedMinutes: 40, DeliveryPincode: "641001", At: mondayTenAM()}); ok {
t.Error("served a handling term above the cap")
}
Reset()
}
func TestLoadedReportsState(t *testing.T) {
Reset()
if ok, _, _ := Loaded(); ok {
t.Error("Loaded() true after Reset")
}
Load(fixture())
defer Reset()
ok, _, cells := Loaded()
if !ok || cells != 3 { // 3 keyed cells; global is held separately
t.Errorf("Loaded() = %v, cells %d; want true, 3", ok, cells)
}
}

View File

@@ -0,0 +1,185 @@
package prediction
import (
"time"
"doormile/constants"
"doormile/models"
"doormile/utils"
"gorm.io/gorm"
)
// Refining a promise after the customer has already been given one.
//
// ─── Why this cannot happen at booking-create time ─────────────────────────
//
// The obvious place to use a calibrated ETA is where the promise is first
// written — controllers/adminController.go:2738 and
// controllers/cxPickupFanout.go:224. It does not work there, and not because
// the calibration is cold: at that moment there is structurally no routed
// duration to calibrate. The order of events is
//
// booking created -> estimateddeliveryat written from the promise table
// -> rider assigned
// -> internal/routing sequences -> etaminutes written
//
// so RoutedMinutes is 0 on every create and ETAMinutes would return false
// forever. Wiring it there would be dead code.
//
// The routed duration first exists at sequencing. So refining the promise means
// revising a number the customer may already have seen.
//
// ─── What this does and deliberately does not do ───────────────────────────
//
// It updates consignments.estimateddeliveryat — the "arrives by" the customer
// is shown — and never touches sladueat. That split is the whole design:
//
// estimateddeliveryat our best current belief. Allowed to improve.
// sladueat the commitment made at booking. Must not move, or it
// stops being a commitment and every breach can be
// explained away by moving the target.
//
// It also only ever moves the estimate when the calibration is trustworthy
// (ETAMinutes true), so with ROUTE_OPTIMIZER_URL unset — the state today —
// this function reads some rows and writes nothing.
//
// A customer watching the tracking page will see the time change. Whether that
// should also emit a stage event is an open product question
// (docs/prediction-plan.md §3a); this writes the column and does not notify,
// because inventing a notification is the more consequential of the two
// choices to get wrong.
// RefinedETA is one estimate this pass revised, for logging.
type RefinedETA struct {
Consignmentid int
Was *time.Time
Now time.Time
Source string
}
// refineRow is a consignment awaiting delivery, joined to the routed duration
// of its booking's latest sequenced assignment.
//
// The join goes through consignment_booking because consignments has no
// bookingid column, and the legacy pickupbookings.consignmentid link names only
// the FIRST order of a multi-destination pickup — hazard H4. DISTINCT ON picks
// one assignment per booking: a booking passes through several
// (Assigned/Rejected/Reassigned) and the newest sequenced one is the routed ETA
// actually in force.
type refineRow struct {
Consignmentid int
Deliverypincode string
Createdat time.Time
Estimateddeliveryat *time.Time
Etaminutes int
Sequencedat *time.Time
}
const refineSQL = `
WITH assigned AS (
SELECT DISTINCT ON (ba.bookingid)
ba.bookingid, ba.etaminutes, ba.sequencedat
FROM bookingassignments ba
WHERE ba.etaminutes > 0
ORDER BY ba.bookingid, ba.sequencedat DESC NULLS LAST, ba.bookingassignmentid DESC
)
SELECT c.consignmentid,
COALESCE(c.deliverypincode, '') AS deliverypincode,
c.createdat,
c.estimateddeliveryat,
a.etaminutes,
a.sequencedat
FROM consignments c
JOIN consignment_booking cb ON cb.consignmentid = c.consignmentid
JOIN assigned a ON a.bookingid = cb.bookingid
WHERE c.status NOT IN ?
AND c.deletedat IS NULL
ORDER BY a.sequencedat DESC NULLS LAST
LIMIT ?
`
// refineBatch caps one pass. The sweeper runs often enough that a backlog
// clears over a few ticks, and an unbounded UPDATE loop on a hot table is the
// kind of thing that shows up as a latency spike somewhere unrelated.
const refineBatch = 300
// terminal statuses — a parcel that has arrived, come back, or been written off
// has nothing left to estimate.
func terminalStatuses() []string {
return []string{
constants.ConsignmentDelivered,
"Returned_to_Sender",
"Missing",
"Damaged",
}
}
// RefineInFlight recomputes estimateddeliveryat for consignments that are still
// moving and whose rider has a routed duration. Returns what it changed.
//
// Safe to call when nothing is calibrated: ETAMinutes returns false for every
// row and this writes nothing.
func RefineInFlight(gdb *gorm.DB, now time.Time) ([]RefinedETA, error) {
if gdb == nil {
return nil, nil
}
if ok, _, _ := Loaded(); !ok {
return nil, nil // nothing to refine with; the promise tables stand
}
var rows []refineRow
if err := gdb.Raw(refineSQL, terminalStatuses(), refineBatch).Scan(&rows).Error; err != nil {
return nil, err
}
changed := make([]RefinedETA, 0, len(rows))
for _, r := range rows {
// The estimate is anchored on when sequencing happened, not on the
// consignment's creation: the routed duration describes the journey
// from the rider's current position onward.
anchor := r.Createdat
if r.Sequencedat != nil {
anchor = *r.Sequencedat
}
eta, res, ok := ETAAt(Input{
RoutedMinutes: r.Etaminutes,
DeliveryPincode: r.Deliverypincode,
At: anchor,
}, anchor)
if !ok {
continue
}
// Skip a no-op write. Without this every pass rewrites every in-flight
// row with the same value, which churns WAL and makes the log useless
// for seeing what actually moved.
if r.Estimateddeliveryat != nil && withinAMinute(*r.Estimateddeliveryat, eta) {
continue
}
if err := gdb.Model(&models.Consignment{}).
Where("consignmentid = ?", r.Consignmentid).
Update("estimateddeliveryat", eta).Error; err != nil {
utils.Warn("prediction: could not refine ETA",
"consignmentid", r.Consignmentid, "error", err)
continue
}
changed = append(changed, RefinedETA{
Consignmentid: r.Consignmentid,
Was: r.Estimateddeliveryat,
Now: eta,
Source: res.Source,
})
}
return changed, nil
}
func withinAMinute(a, b time.Time) bool {
d := a.Sub(b)
if d < 0 {
d = -d
}
return d < time.Minute
}

View File

@@ -0,0 +1,130 @@
package prediction
import (
"context"
"os"
"strconv"
"strings"
"time"
"doormile/db"
"doormile/utils"
)
// The calibration sweeper.
//
// Deliberately the same shape as internal/assignment/sweeper.go: a ticker, a
// recover() per tick so one bad run cannot take the process down, and a Redis
// lock that expires before the next tick so only one replica of three does the
// work. Without Redis every replica refreshes, which is wasteful but correct —
// Refresh replaces the table in one transaction, so concurrent runs converge
// rather than interleave.
//
// This is a nightly job by default. The calibration is a p80 over weeks of
// deliveries; recomputing it more often costs a scan and changes nothing.
const (
defaultCalibrationSweepSeconds = 6 * 60 * 60 // 6h
calibrationLockKey = "prediction:calibration:lock"
// Below this a "sweep" is a scan loop against the whole delivery history.
minCalibrationSweepSeconds = 600
)
// calibrationInterval: PREDICTION_CALIBRATION_SECONDS, default 6h; 0 turns the
// sweeper off and leaves whatever is in the store. Read once at start, matching
// how assignment's sweepInterval behaves.
func calibrationInterval() time.Duration {
if v := strings.TrimSpace(os.Getenv("PREDICTION_CALIBRATION_SECONDS")); v != "" {
if n, err := strconv.Atoi(v); err == nil && n >= 0 {
if n > 0 && n < minCalibrationSweepSeconds {
n = minCalibrationSweepSeconds
}
return time.Duration(n) * time.Second
}
utils.Warn("PREDICTION_CALIBRATION_SECONDS is not a non-negative integer, using the default",
"value", v, "default_seconds", defaultCalibrationSweepSeconds)
}
return defaultCalibrationSweepSeconds * time.Second
}
// StartCalibrationSweeper loads whatever calibration is stored, then refreshes
// it on a timer. Call once at boot, in a goroutine.
//
// The load happens even when the sweeper is disabled: a stored calibration from
// a previous deploy is still the best available answer, and ETAMinutes will
// stop trusting it on its own once it passes staleAfter.
func StartCalibrationSweeper() {
LoadFromDB()
interval := calibrationInterval()
if interval == 0 {
utils.Info("CalibrationSweeper: disabled (PREDICTION_CALIBRATION_SECONDS=0)")
return
}
utils.Info("CalibrationSweeper: started", "interval", interval.String())
// One refresh shortly after boot rather than waiting a full interval, so a
// newly deployed replica is not serving a day-old calibration for six
// hours. Offset so three replicas do not all wake together.
time.Sleep(90 * time.Second)
refreshOnce(interval)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for range ticker.C {
refreshOnce(interval)
}
}
// refreshOnce makes one refresh attempt. Only one replica at a time (a Redis
// lock that expires before the next tick).
func refreshOnce(interval time.Duration) {
defer func() {
if r := recover(); r != nil {
utils.Error("CalibrationSweeper: panic recovered", "error", r)
}
}()
if db.DB == nil {
return
}
if db.Rdb != nil {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
ttl := interval - 60*time.Second
if ttl <= 0 {
ttl = interval / 2
}
got, err := db.Rdb.SetNX(ctx, calibrationLockKey, "1", ttl).Result()
cancel()
if err == nil && !got {
return // another replica has this refresh
}
}
started := time.Now()
if err := Refresh(db.DB); err != nil {
// A failed refresh leaves the previous calibration serving. That is the
// intended behaviour: the last good factors beat falling back to flat
// constants, and staleAfter puts a bound on how long that can last.
utils.Error("CalibrationSweeper: refresh failed", "error", err,
"took_ms", time.Since(started).Milliseconds())
return
}
ok, builtAt, cells := Loaded()
utils.Info("CalibrationSweeper: refresh complete",
"loaded", ok, "cells", cells, "built_at", builtAt,
"took_ms", time.Since(started).Milliseconds())
// Apply the fresh calibration to parcels still in flight. Writes
// estimateddeliveryat only — never sladueat, which is the commitment made
// at booking (see internal/prediction/refine.go). A no-op while nothing is
// calibrated.
refined, err := RefineInFlight(db.DB, time.Now())
if err != nil {
utils.Error("CalibrationSweeper: ETA refine failed", "error", err)
return
}
if len(refined) > 0 {
utils.Info("CalibrationSweeper: refined in-flight ETAs", "count", len(refined))
}
}