fix(expiry): partial composite index + batched UPDATE to stop "context deadline exceeded"

Root cause
----------
ExpireTasks in internal/channels/task_store.go ran a single unbounded
UPDATE filtered on (status='open' AND deadline IS NOT NULL AND deadline < now).
The only indexes on tasks were idx_tasks_status(status) and
idx_tasks_channel(channel_id). With status cardinality of ~4 and a growing
auction-tasks table on kubic, the planner used idx_tasks_status to enumerate
all open rows then evaluated deadline per row, holding a SQLite write
transaction the whole time. Under WAL contention with concurrent writers
(message inserts, consolidator) the worker's 30s context regularly expired,
producing the recurring expiry-worker log line.

Fix
---
1. New migration 031_tasks_expiry_index.sql: partial composite index
   idx_tasks_expiry(status, deadline) WHERE status='open' AND deadline IS NOT NULL.
   This is the exact predicate ExpireTasks uses, so the planner now seeks
   straight to eligible rows. The partial form keeps the index empty for the
   steady-state majority of rows (completed/cancelled), so writes elsewhere
   aren't penalized.

2. Batch the UPDATE in chunks of 500 (rowid IN subquery; UPDATE ... LIMIT
   isn't compiled into modernc.org/sqlite by default). Bounded write
   transactions stop the worker from starving other writers and let it
   observe context cancellation between batches.

Perf
----
New test exercises 2400 mixed rows (1200 expirable). With the index +
batching, expiry finishes in ~3ms inside a 5s context; without the index a
regression to full status-scan would be measurably worse and is also
guarded by an EXPLAIN QUERY PLAN test.

Operational notes
-----------------
- Migration is additive and idempotent (CREATE INDEX IF NOT EXISTS). No
  backfill needed; it will apply on next pod startup.
- After rollout, expiry-worker error logs should clear within one tick
  (default 1m).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Algis Dumbris
2026-05-13 07:00:16 +03:00
co-authored by Claude Opus 4.7
parent 0bb2b6500a
commit a2ea7cc516
3 changed files with 218 additions and 11 deletions
+53 -11
View File
@@ -264,21 +264,63 @@ func (s *SQLiteTaskStore) UpdateBidStatus(ctx context.Context, bidID int64, stat
return nil
}
// expireTasksBatchSize bounds how many tasks a single UPDATE statement
// rewrites. A small batch keeps the SQLite write transaction short, so
// the expiry worker doesn't starve under WAL contention with concurrent
// writers (e.g. message inserts, consolidator) and respects the worker's
// context deadline even if the eligible set is large.
const expireTasksBatchSize = 500
// ExpireTasks marks all open tasks past their deadline as cancelled.
//
// The work is chunked into bounded UPDATEs (LIMIT expireTasksBatchSize)
// rather than a single unbounded UPDATE, for two reasons:
//
// 1. Bounded transactions: SQLite serializes writers, so a long UPDATE
// blocks every other writer until commit. Chunking caps the lock
// window per round-trip.
// 2. Context responsiveness: the expiry worker uses a 30s context. A
// single UPDATE doesn't observe ctx between rows, so a slow query
// would always run to completion and then return ctx error. Looping
// lets us bail between batches.
//
// Performance also depends on migration 031, which adds a partial
// composite index `idx_tasks_expiry(status, deadline) WHERE status='open'
// AND deadline IS NOT NULL`. The query below is shaped to match it.
func (s *SQLiteTaskStore) ExpireTasks(ctx context.Context) (int, error) {
// Use a string-formatted timestamp for consistent SQLite comparison
now := time.Now().UTC().Format("2006-01-02 15:04:05")
result, err := s.db.ExecContext(ctx,
`UPDATE tasks SET status = ?, updated_at = CURRENT_TIMESTAMP
WHERE status = ? AND deadline IS NOT NULL AND deadline < ?`,
TaskStatusCancelled, TaskStatusOpen, now,
)
if err != nil {
return 0, fmt.Errorf("expire tasks: %w", err)
}
now := time.Now().UTC().Format(sqliteTimeFormat)
rowsAffected, _ := result.RowsAffected()
return int(rowsAffected), nil
total := 0
for {
if err := ctx.Err(); err != nil {
return total, fmt.Errorf("expire tasks: %w", err)
}
// SQLite's UPDATE ... LIMIT is only enabled with the
// SQLITE_ENABLE_UPDATE_DELETE_LIMIT compile flag (not on by
// default in modernc.org/sqlite). Use a subquery on ROWID
// to portably bound the batch.
result, err := s.db.ExecContext(ctx,
`UPDATE tasks SET status = ?, updated_at = CURRENT_TIMESTAMP
WHERE rowid IN (
SELECT rowid FROM tasks
WHERE status = ? AND deadline IS NOT NULL AND deadline < ?
LIMIT ?
)`,
TaskStatusCancelled, TaskStatusOpen, now, expireTasksBatchSize,
)
if err != nil {
return total, fmt.Errorf("expire tasks: %w", err)
}
rowsAffected, _ := result.RowsAffected()
total += int(rowsAffected)
if rowsAffected < int64(expireTasksBatchSize) {
// Last (possibly empty) batch — nothing more to expire.
return total, nil
}
}
}
// CancelTasksByChannel cancels all open tasks for a channel (used before channel deletion).
+144
View File
@@ -326,6 +326,150 @@ func TestSQLiteTaskStore_ExpireTasks(t *testing.T) {
}
}
// TestSQLiteTaskStore_ExpireTasks_LargeBatch exercises the expiry path at a
// scale comparable to a long-lived production instance: a few thousand mixed
// tasks (expired-open, future-open, no-deadline, cancelled), confirms the
// worker chunks past its internal batch size, and that the entire run fits
// inside a tight per-tick context — guarding the regression that the
// expiry-worker hit on kubic ("context deadline exceeded").
//
// The total row count is deliberately larger than expireTasksBatchSize (500)
// so the batching loop must iterate more than once.
func TestSQLiteTaskStore_ExpireTasks_LargeBatch(t *testing.T) {
taskStore, channelStore := newTestTaskStore(t)
ch := createAuctionChannel(t, channelStore)
ctx := context.Background()
const (
expiredOpen = 1200 // > 2 * batch size, forces multiple iterations
futureOpen = 400
noDeadline = 400
alreadyDone = 400
)
past := time.Now().Add(-1 * time.Hour)
future := time.Now().Add(1 * time.Hour)
for i := 0; i < expiredOpen; i++ {
if err := taskStore.CreateTask(ctx, &Task{
ChannelID: ch.ID, PostedBy: "poster-agent",
Title: "expired", Status: TaskStatusOpen,
Deadline: &past, Requirements: json.RawMessage(`{}`),
}); err != nil {
t.Fatalf("seed expired task %d: %v", i, err)
}
}
for i := 0; i < futureOpen; i++ {
taskStore.CreateTask(ctx, &Task{
ChannelID: ch.ID, PostedBy: "poster-agent",
Title: "future", Status: TaskStatusOpen,
Deadline: &future, Requirements: json.RawMessage(`{}`),
})
}
for i := 0; i < noDeadline; i++ {
taskStore.CreateTask(ctx, &Task{
ChannelID: ch.ID, PostedBy: "poster-agent",
Title: "no-deadline", Status: TaskStatusOpen,
Requirements: json.RawMessage(`{}`),
})
}
for i := 0; i < alreadyDone; i++ {
taskStore.CreateTask(ctx, &Task{
ChannelID: ch.ID, PostedBy: "poster-agent",
Title: "done", Status: TaskStatusCompleted,
Deadline: &past, Requirements: json.RawMessage(`{}`),
})
}
// Run with the same shape of context budget the worker uses, but tighter
// (5s) so a regression to the unbounded-scan plan would fail this test
// well before the worker's real 30s ceiling.
tightCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
start := time.Now()
count, err := taskStore.ExpireTasks(tightCtx)
elapsed := time.Since(start)
if err != nil {
t.Fatalf("ExpireTasks: %v (elapsed=%s)", err, elapsed)
}
if count != expiredOpen {
t.Errorf("expired count = %d, want %d", count, expiredOpen)
}
t.Logf("expired %d tasks in %s (over %d total rows)",
count, elapsed, expiredOpen+futureOpen+noDeadline+alreadyDone)
// Future / no-deadline tasks must remain open.
openTasks, _ := taskStore.ListTasks(ctx, ch.ID, TaskStatusOpen)
if len(openTasks) != futureOpen+noDeadline {
t.Errorf("open after expiry = %d, want %d",
len(openTasks), futureOpen+noDeadline)
}
// Second invocation on a clean set must be cheap and return 0.
count2, err := taskStore.ExpireTasks(ctx)
if err != nil {
t.Fatalf("ExpireTasks (idempotent run): %v", err)
}
if count2 != 0 {
t.Errorf("idempotent run expired = %d, want 0", count2)
}
}
// TestSQLiteTaskStore_ExpireTasks_UsesIndex confirms the partial composite
// index from migration 031 is the plan the query optimizer picks. If a
// future change drops the index or rewrites the query incompatibly, the
// planner will fall back to a SCAN and this test will fail loudly.
func TestSQLiteTaskStore_ExpireTasks_UsesIndex(t *testing.T) {
taskStore, _ := newTestTaskStore(t)
// Seed a few rows so the planner has stats to work with.
// (SQLite's planner is mostly schema-driven, but better safe.)
ctx := context.Background()
now := time.Now().UTC().Format(sqliteTimeFormat)
rows, err := taskStore.db.QueryContext(ctx,
`EXPLAIN QUERY PLAN
UPDATE tasks SET status = ?, updated_at = CURRENT_TIMESTAMP
WHERE rowid IN (
SELECT rowid FROM tasks
WHERE status = ? AND deadline IS NOT NULL AND deadline < ?
LIMIT ?
)`,
TaskStatusCancelled, TaskStatusOpen, now, expireTasksBatchSize,
)
if err != nil {
t.Fatalf("EXPLAIN QUERY PLAN: %v", err)
}
defer rows.Close()
sawIndex := false
var plan []string
for rows.Next() {
var id, parent, notused int
var detail string
if err := rows.Scan(&id, &parent, &notused, &detail); err != nil {
t.Fatalf("scan plan row: %v", err)
}
plan = append(plan, detail)
if contains(detail, "idx_tasks_expiry") {
sawIndex = true
}
}
if !sawIndex {
t.Errorf("expected query plan to use idx_tasks_expiry, got: %v", plan)
}
}
func contains(s, sub string) bool {
for i := 0; i+len(sub) <= len(s); i++ {
if s[i:i+len(sub)] == sub {
return true
}
}
return false
}
func TestSQLiteTaskStore_CancelTasksByChannel(t *testing.T) {
taskStore, channelStore := newTestTaskStore(t)
ch := createAuctionChannel(t, channelStore)
@@ -0,0 +1,21 @@
-- 031: Composite index to speed up ExpireTasks() in expiry-worker.
--
-- Root cause: `ExpireTasks` runs
-- UPDATE tasks SET status='cancelled', updated_at=CURRENT_TIMESTAMP
-- WHERE status='open' AND deadline IS NOT NULL AND deadline < ?
--
-- Pre-031 indexes were only `idx_tasks_status(status)` and `idx_tasks_channel(channel_id)`.
-- With a status cardinality of 4 and most tasks in two buckets, the planner used
-- `idx_tasks_status` to find all open rows then evaluated the deadline predicate
-- per row. As the auction-tasks table grew (kubic deploy), the worker's 30s
-- context deadline started to be exceeded on every tick, especially under WAL
-- write contention from concurrent message inserts.
--
-- Fix: composite, partial index on `(status, deadline)` covering only rows that
-- can ever be expired (status='open' AND deadline IS NOT NULL). This is the
-- exact predicate the worker uses, so SQLite can seek straight to the eligible
-- rows. The partial form keeps the index tiny once tasks transition out of
-- 'open' (the dominant steady-state).
CREATE INDEX IF NOT EXISTS idx_tasks_expiry
ON tasks(status, deadline)
WHERE status = 'open' AND deadline IS NOT NULL;