From c9b2f1906a8aff7e64c51738783519e59a0cbf0b Mon Sep 17 00:00:00 2001 From: djgrant <1670902+djgrant@users.noreply.github.com> Date: Wed, 22 Jul 2026 09:55:01 +0100 Subject: [PATCH] Fix millisecond/second mismatch in SQLite task visibility timeout Fixes #32 --- .../sqlite-runtime/src/dao/task-queue-dao.ts | 20 +++++---- packages/sqlite-runtime/src/index.ts | 2 +- .../sqlite-runtime/src/sqlite-task-queue.ts | 6 +-- test/sqlite-task-queue.test.ts | 42 +++++++++++++++++++ 4 files changed, 57 insertions(+), 13 deletions(-) create mode 100644 test/sqlite-task-queue.test.ts diff --git a/packages/sqlite-runtime/src/dao/task-queue-dao.ts b/packages/sqlite-runtime/src/dao/task-queue-dao.ts index b371cb6..78e8995 100644 --- a/packages/sqlite-runtime/src/dao/task-queue-dao.ts +++ b/packages/sqlite-runtime/src/dao/task-queue-dao.ts @@ -28,15 +28,15 @@ export class TaskQueueDao { execution_id TEXT NOT NULL, params TEXT, context TEXT, - visible_from INTEGER DEFAULT (strftime('%s', 'now')) + visible_from INTEGER DEFAULT 0 ) `); } insertTask(event: WorkflowEvent) { const query = this.db.query( - `INSERT INTO task_queue (workflow_id, execution_id, params, context) - VALUES ($workflowId, $executionId, $params, $context)` + `INSERT INTO task_queue (workflow_id, execution_id, params, context, visible_from) + VALUES ($workflowId, $executionId, $params, $context, 0)` ); query.run({ @@ -49,11 +49,13 @@ export class TaskQueueDao { }); } - getNextTask() { + // `now` is a millisecond Unix timestamp, matching the values written by + // updateTaskVisibility + getNextTask(now: number) { const query = this.db.query( - `SELECT * FROM task_queue WHERE visible_from < strftime('%s', 'now') ORDER BY task_id LIMIT 1` + `SELECT * FROM task_queue WHERE visible_from <= $now ORDER BY task_id LIMIT 1` ); - return query.get(); + return query.get({ $now: now }); } updateTaskVisibility(taskId: number, visibleFrom: number) { @@ -70,10 +72,10 @@ export class TaskQueueDao { query.run({ $taskId: id }); } - getTaskCount() { + getTaskCount(now: number) { const query = this.db.query( - `SELECT COUNT(*) as count FROM task_queue WHERE visible_from < strftime('%s', 'now')` + `SELECT COUNT(*) as count FROM task_queue WHERE visible_from <= $now` ); - return query.get()!.count; + return query.get({ $now: now })!.count; } } diff --git a/packages/sqlite-runtime/src/index.ts b/packages/sqlite-runtime/src/index.ts index b3cec16..3c7656b 100644 --- a/packages/sqlite-runtime/src/index.ts +++ b/packages/sqlite-runtime/src/index.ts @@ -2,7 +2,7 @@ export { SqliteHeapClient } from "./sqlite-heap"; export { SqliteStoreClient } from "./sqlite-store"; export { SqliteSchedulerClient } from "./sqlite-scheduler"; export { SqliteEventLoop } from "./sqlite-event-loop"; -export { SqliteTaskQueueClient } from "./sqlite-task-queue"; +export { SqliteTaskQueue, SqliteTaskQueueClient } from "./sqlite-task-queue"; export { SqliteTimersClient } from "./sqlite-timers"; export type { SqliteDriver, diff --git a/packages/sqlite-runtime/src/sqlite-task-queue.ts b/packages/sqlite-runtime/src/sqlite-task-queue.ts index bd6fc8a..2a01aeb 100644 --- a/packages/sqlite-runtime/src/sqlite-task-queue.ts +++ b/packages/sqlite-runtime/src/sqlite-task-queue.ts @@ -17,11 +17,11 @@ export class SqliteTaskQueue { } process() { - const row = this.taskQueueDao.getNextTask(); + const now = Date.now(); + const row = this.taskQueueDao.getNextTask(now); if (!row) return undefined; - const now = Date.now(); const visibilityTimeout = now + VISIBILITY_WINDOW; this.taskQueueDao.updateTaskVisibility(row.task_id, visibilityTimeout); @@ -47,7 +47,7 @@ export class SqliteTaskQueue { } get isEmpty(): boolean { - return this.taskQueueDao.getTaskCount() === 0; + return this.taskQueueDao.getTaskCount(Date.now()) === 0; } } diff --git a/test/sqlite-task-queue.test.ts b/test/sqlite-task-queue.test.ts new file mode 100644 index 0000000..b2ea945 --- /dev/null +++ b/test/sqlite-task-queue.test.ts @@ -0,0 +1,42 @@ +import { afterEach, expect, test, vi } from "vitest"; +import { SqliteTaskQueue } from "@yieldstar/sqlite-runtime"; + +const isBun = "Bun" in globalThis; +const { createSqliteDb } = isBun + ? await import("@yieldstar/sqlite-runtime/bun") + : await import("@yieldstar/sqlite-runtime/node"); + +const VISIBILITY_WINDOW = 300000; + +afterEach(() => { + vi.useRealTimers(); +}); + +test("a claimed task becomes visible again after the visibility window", async () => { + vi.useFakeTimers(); + + const db = createSqliteDb({ path: ":memory:", wal: false }); + const taskQueue = new SqliteTaskQueue(db); + + taskQueue.add({ workflowId: "workflow-1", executionId: "execution-1" }); + + const claimed = taskQueue.process(); + expect(claimed?.event.workflowId).toBe("workflow-1"); + + // While claimed, the task is hidden from other workers + expect(taskQueue.process()).toBeUndefined(); + expect(taskQueue.isEmpty).toBe(true); + + // Simulate a worker crash: the task is never removed or made visible. + // Once the visibility window elapses, a fresh queue can claim it again. + vi.advanceTimersByTime(VISIBILITY_WINDOW + 1); + + const recoveredQueue = new SqliteTaskQueue(db); + expect(recoveredQueue.isEmpty).toBe(false); + + const reclaimed = recoveredQueue.process(); + expect(reclaimed?.taskId).toBe(claimed?.taskId); + expect(reclaimed?.event.workflowId).toBe("workflow-1"); + + db.close(); +});