diff --git a/internal/channels/task_store.go b/internal/channels/task_store.go index 59a3710..1d7d7cf 100644 --- a/internal/channels/task_store.go +++ b/internal/channels/task_store.go @@ -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). diff --git a/internal/channels/task_store_test.go b/internal/channels/task_store_test.go index 5fe98d6..787ba85 100644 --- a/internal/channels/task_store_test.go +++ b/internal/channels/task_store_test.go @@ -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, ¬used, &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) diff --git a/internal/storage/schema/031_tasks_expiry_index.sql b/internal/storage/schema/031_tasks_expiry_index.sql new file mode 100644 index 0000000..cef8564 --- /dev/null +++ b/internal/storage/schema/031_tasks_expiry_index.sql @@ -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;