From a5fecde311962f226ec30d749c7a995c9bc5bd3e Mon Sep 17 00:00:00 2001 From: Daniel Sutton Date: Fri, 21 Aug 2026 19:04:48 +0100 Subject: [PATCH] feat(run-engine): waitpoint shard scripts for create, register and complete --- .../engine/waitpointCoordinator/scripts.ts | 336 +++++++++++++++++ .../storeCoordinator.test.ts | 340 ++++++++++++++++++ .../waitpointCoordinator/storeCoordinator.ts | 239 ++++++++++++ 3 files changed, 915 insertions(+) create mode 100644 internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts create mode 100644 internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts create mode 100644 internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts diff --git a/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts b/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts new file mode 100644 index 000000000..1e60f6516 --- /dev/null +++ b/internal-packages/run-engine/src/engine/waitpointCoordinator/scripts.ts @@ -0,0 +1,336 @@ +import type { Callback, Redis, Result } from "@internal/redis"; + +/** + * Lua for the waitpoint coordination protocol. Three rules hold throughout: + * + * 1. Every key a script touches is declared in KEYS. No script builds a key name inside + * Lua. ioredis prefixes only the KEYS array, so a key minted in Lua would be + * unprefixed while the client wrote a prefixed one — and a script with a single + * declared key gives the caller's single-slot assertion nothing to compare. + * 2. Lua never parses JSON. Each script branches only on a short status string and moves + * opaque blobs, so every encoding decision stays in TypeScript. + * 3. Every returned slot is coerced with `or ''`. A Lua false or nil TRUNCATES the reply + * array at that position, silently shortening it. + * + * STORED_COMPLETED is the value written into the record's `status` field and is + * UPPERCASE. The outcome tokens below are lowercase and are a separate vocabulary: they + * name what a script DID, not what a record IS. Sharing one constant between the two + * makes an already-completed record invisible to every script. + */ + +const STORED_COMPLETED = "COMPLETED"; + +const MISSING = "missing"; +const CREATED = "created"; +const EXISTS = "exists"; +const REGISTERED = "registered"; +const DID_COMPLETE = "completed"; +const ALREADY = "already"; +const RESERVED = "reserved"; +const CLEARED = "cleared"; +const DRAINED = "drained"; + +export function registerWaitpointCommands(redis: Redis): void { + // KEYS: record. ARGV: recordJson, status ('PENDING'|'COMPLETED'), completionJson (''). + redis.defineCommand("wpCreateIfAbsent", { + numberOfKeys: 1, + lua: ` + local record = KEYS[1] + + -- EXISTS-then-HSET inside one script, rather than a field-by-field HSETNX: the + -- record and its status must appear together or not at all. + if redis.call('EXISTS', record) == 1 then + local vals = redis.call('HMGET', record, 'r', 'status', 'c') + return { '${EXISTS}', vals[1] or '', vals[2] or '', vals[3] or '' } + end + + redis.call('HSET', record, 'r', ARGV[1], 'status', ARGV[2]) + if ARGV[3] ~= '' then + redis.call('HSET', record, 'c', ARGV[3]) + end + + return { '${CREATED}' } + `, + }); + + // KEYS: record, watchers. ARGV: watcherField, watcherJson. + redis.defineCommand("wpRegisterOrReport", { + numberOfKeys: 2, + lua: ` + local record, watchers = KEYS[1], KEYS[2] + + -- A missing waitpoint is never a silent no-op: the caller throws. Defaulting to + -- "not blocked" here would resume a run whose waitpoint never completed. + if redis.call('EXISTS', record) == 0 then + return { '${MISSING}' } + end + + if redis.call('HGET', record, 'status') == '${STORED_COMPLETED}' then + return { '${DID_COMPLETE}', redis.call('HGET', record, 'c') or '' } + end + + -- The watcher lands before any flip can read the watcher hash, because this script + -- and wpComplete are both atomic on this same shard. So a register either appears + -- in the flip's watcher list, or it observes COMPLETED above. + redis.call('HSET', watchers, ARGV[1], ARGV[2]) + return { '${REGISTERED}' } + `, + }); + + // KEYS: record, watchers. ARGV: completionJson. + redis.defineCommand("wpComplete", { + numberOfKeys: 2, + lua: ` + local record, watchers = KEYS[1], KEYS[2] + + if redis.call('EXISTS', record) == 0 then + return { '${MISSING}' } + end + + local outcome = '${DID_COMPLETE}' + if redis.call('HGET', record, 'status') == '${STORED_COMPLETED}' then + -- Double completion is not an error, and the FIRST completion wins. This is the + -- guard a conditional UPDATE ... WHERE status = 'PENDING' used to provide. + outcome = '${ALREADY}' + else + redis.call('HSET', record, 'status', '${STORED_COMPLETED}', 'c', ARGV[1]) + end + + -- Returning the watchers here is what removes the reverse fan-out query. The + -- envelope comes back too, because delivery runs on each watcher's own shard and + -- cannot read this key. + local out = { outcome, redis.call('HGET', record, 'c') or '' } + local entries = redis.call('HVALS', watchers) + for i = 1, #entries do + out[#out + 1] = entries[i] + end + + return out + `, + }); + + // KEYS: idempotency key. ARGV: waitpointId, expiresAtMs ('' for no expiry). + redis.defineCommand("wpIdemReserve", { + numberOfKeys: 1, + lua: ` + local key = KEYS[1] + + -- SET NX returns a status reply on success and false on conflict. + if redis.call('SET', key, ARGV[1], 'NX') then + -- Expiry only when the caller has one. A reservation with no expiry is the common + -- case and must never grow one here. + if ARGV[2] ~= '' then + redis.call('PEXPIREAT', key, tonumber(ARGV[2])) + end + return { '${RESERVED}', ARGV[1] } + end + + return { '${EXISTS}', redis.call('GET', key) or '' } + `, + }); + + // KEYS: pend, done, edge. + // ARGV: n, then n groups of 4 — waitpointId, edgeField, edgeJson, reportedJson (''). + redis.defineCommand("runAbsorbBlockers", { + numberOfKeys: 3, + lua: ` + local pend, done, edge = KEYS[1], KEYS[2], KEYS[3] + local n = tonumber(ARGV[1]) + + -- countedPending and seenDelivered make both outputs DISTINCT BY ID. The count this + -- replaces was a COUNT(*) over waitpoint rows, so two edges for one waitpoint + -- contributed one. Counting per edge would inflate it. + local countedPending = {} + local seenDelivered = {} + local pendingOfRequested = 0 + local out = { '0', '0' } + + for i = 0, n - 1 do + local id = ARGV[2 + i * 4] + local field = ARGV[3 + i * 4] + local edgeJson = ARGV[4 + i * 4] + local reported = ARGV[5 + i * 4] + + -- HSETNX is the ON CONFLICT DO NOTHING of the edge write: a retry must not + -- overwrite the first attempt's metadata. + redis.call('HSETNX', edge, field, edgeJson) + + if reported ~= '' then + -- Already COMPLETED when the watcher registered. It never becomes pending. + redis.call('HSET', done, id, reported) + redis.call('SREM', pend, id) + if not seenDelivered[id] then + seenDelivered[id] = true + out[#out + 1] = id + out[#out + 1] = reported + end + else + -- Check the delivered set FIRST. A completion that landed between register and + -- absorb has already delivered here, and that delivery wins. + local delivered = redis.call('HGET', done, id) + if delivered then + if not seenDelivered[id] then + seenDelivered[id] = true + out[#out + 1] = id + out[#out + 1] = delivered + end + else + redis.call('SADD', pend, id) + if not countedPending[id] then + countedPending[id] = true + pendingOfRequested = pendingOfRequested + 1 + end + end + end + end + + out[1] = tostring(pendingOfRequested) + out[2] = tostring(redis.call('SCARD', pend)) + return out + `, + }); + + // KEYS: pend, done. ARGV: waitpointId, completionJson. + redis.defineCommand("runDeliverCompletion", { + numberOfKeys: 2, + lua: ` + local pend, done = KEYS[1], KEYS[2] + + redis.call('HSET', done, ARGV[1], ARGV[2]) + redis.call('SREM', pend, ARGV[1]) + + -- The caller treats this as a wakeup trigger, not as the resume decision: the + -- resume is decided under the run lock, and this count covers store-resident + -- blockers only. + return { tostring(redis.call('SCARD', pend)) } + `, + }); + + // KEYS: pend, done, edge. + redis.defineCommand("runReadBlockState", { + numberOfKeys: 3, + lua: ` + local pend, done, edge = KEYS[1], KEYS[2], KEYS[3] + + local pendIds = redis.call('SMEMBERS', pend) + -- HKEYS, never HGETALL: the delivered set's values are completion envelopes with + -- inline outputs, and materializing those inside a single-threaded script would + -- block the shard. + local doneIds = redis.call('HKEYS', done) + local edges = redis.call('HGETALL', edge) + + local out = { tostring(#pendIds), tostring(#doneIds), tostring(#edges) } + for i = 1, #pendIds do out[#out + 1] = pendIds[i] end + for i = 1, #doneIds do out[#out + 1] = doneIds[i] end + for i = 1, #edges do out[#out + 1] = edges[i] end + return out + `, + }); + + // KEYS: pend, done, edge. ARGV: n, then n edge fields. n = 0 clears everything. + redis.defineCommand("runClear", { + numberOfKeys: 3, + lua: ` + local pend, done, edge = KEYS[1], KEYS[2], KEYS[3] + local n = tonumber(ARGV[1]) + + if n == 0 then + redis.call('DEL', pend, done, edge) + return { '${CLEARED}' } + end + + for i = 1, n do + redis.call('HDEL', edge, ARGV[1 + i]) + end + + -- Reconcile rather than delete by name. The edge set is the authority: after the + -- drain, pend and done may only hold ids that some surviving edge still references. + -- + -- Two reasons this is a superset of "remove the drained ids". First, one waitpoint + -- can hold several edges at different batch indexes, so a drained field must not + -- evict a delivery another edge still needs. Second, runDeliverCompletion writes + -- done[id] unconditionally, so a crash between register and absorb can leave a + -- delivered entry with no edge at all, which no name-derived drain could reach. + local remaining = {} + local fields = redis.call('HKEYS', edge) + for i = 1, #fields do + local sep = string.find(fields[i], '#[^#]*$') + if sep then + remaining[string.sub(fields[i], 1, sep - 1)] = true + end + end + + local doneIds = redis.call('HKEYS', done) + for i = 1, #doneIds do + if not remaining[doneIds[i]] then + redis.call('HDEL', done, doneIds[i]) + end + end + + local pendIds = redis.call('SMEMBERS', pend) + for i = 1, #pendIds do + if not remaining[pendIds[i]] then + redis.call('SREM', pend, pendIds[i]) + end + end + + return { '${DRAINED}' } + `, + }); +} + +declare module "@internal/redis" { + interface RedisCommander { + wpCreateIfAbsent( + recordKey: string, + recordJson: string, + status: string, + completionJson: string, + callback?: Callback + ): Result; + wpRegisterOrReport( + recordKey: string, + watchersKey: string, + watcherField: string, + watcherJson: string, + callback?: Callback + ): Result; + wpComplete( + recordKey: string, + watchersKey: string, + completionJson: string, + callback?: Callback + ): Result; + wpIdemReserve( + key: string, + waitpointId: string, + expiresAtMs: string, + callback?: Callback + ): Result; + runAbsorbBlockers( + pendKey: string, + doneKey: string, + edgeKey: string, + ...args: Array> + ): Result; + runDeliverCompletion( + pendKey: string, + doneKey: string, + waitpointId: string, + completionJson: string, + callback?: Callback + ): Result; + runReadBlockState( + pendKey: string, + doneKey: string, + edgeKey: string, + callback?: Callback + ): Result; + runClear( + pendKey: string, + doneKey: string, + edgeKey: string, + ...args: Array> + ): Result; + } +} diff --git a/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts new file mode 100644 index 000000000..9ec5e73c7 --- /dev/null +++ b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.test.ts @@ -0,0 +1,340 @@ +// Redis-only suite: the coordinator holds no Prisma reference, so no Postgres container +// is needed. redisTest FLUSHALLs before every test, so ids may be reused across describes. +import { createRedisClient, type RedisOptions } from "@internal/redis"; +import { redisTest } from "@internal/testcontainers"; +import { describe, expect } from "vitest"; +import { WaitpointKeyTagError } from "./keys.js"; +import { + WaitpointNotFoundError, + WaitpointStoreCoordinator, + type WaitpointCompletion, + type WaitpointRecordInput, +} from "./storeCoordinator.js"; + +const ENV_ID = "env_1"; +const PROJECT_ID = "proj_1"; +const NOW = "2026-08-21T12:00:00.000Z"; + +function coordinator(redisOptions: RedisOptions) { + return new WaitpointStoreCoordinator({ redisOptions }); +} + +function record(id: string, overrides: Partial = {}): WaitpointRecordInput { + return { + id, + friendlyId: `waitpoint_${id}`, + type: "MANUAL", + environmentId: ENV_ID, + projectId: PROJECT_ID, + createdAt: NOW, + updatedAt: NOW, + userProvidedIdempotencyKey: false, + tags: [], + ...overrides, + }; +} + +function completion(overrides: Partial = {}): WaitpointCompletion { + return { + completedAt: NOW, + outputType: "application/json", + outputIsError: false, + output: { inline: '{"ok":true}' }, + ...overrides, + }; +} + +describe("createIfAbsent", () => { + redisTest("creates a PENDING record and reports created", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + const result = await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + expect(result.outcome).toBe("created"); + } finally { + await store.quit(); + } + }); + + redisTest("returns the existing record on a second call", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + const second = await store.createIfAbsent({ + record: record("w_a", { friendlyId: "waitpoint_DIFFERENT" }), + status: "PENDING", + }); + + expect(second.outcome).toBe("exists"); + if (second.outcome !== "exists") throw new Error("unreachable"); + // The first write wins: a retry must not overwrite the stored record. + expect(second.record.friendlyId).toBe("waitpoint_w_a"); + expect(second.status).toBe("PENDING"); + expect(second.completion).toBeUndefined(); + } finally { + await store.quit(); + } + }); + + redisTest("preserves every record field through a round trip", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + const full = record("w_a", { + type: "RUN", + idempotencyKey: "key-1", + userProvidedIdempotencyKey: true, + idempotencyKeyExpiresAt: NOW, + completedAfter: NOW, + completedByTaskRunId: "run_child", + completedByBatchId: "batch_1", + tags: ["one", "two"], + }); + + await store.createIfAbsent({ record: full, status: "PENDING" }); + const read = await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + + expect(read.outcome).toBe("exists"); + if (read.outcome !== "exists") throw new Error("unreachable"); + // Every field the frozen return shapes need must survive the blob round trip. + expect(read.record).toEqual(full); + } finally { + await store.quit(); + } + }); + + redisTest( + "can create an already-COMPLETED record with no completion envelope", + async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + // This is the shape that catches a status-casing mismatch: the record is stored + // COMPLETED, and a register must see it as completed rather than pending. + await store.createIfAbsent({ record: record("w_a"), status: "COMPLETED" }); + + const reported = await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + createdAt: NOW, + }); + + expect(reported.outcome).toBe("completed"); + if (reported.outcome !== "completed") throw new Error("unreachable"); + expect(reported.completion).toBeUndefined(); + } finally { + await store.quit(); + } + } + ); + + redisTest( + "can create an already-COMPLETED record with a completion", + async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ + record: record("w_a", { type: "RUN" }), + status: "COMPLETED", + completion: completion(), + }); + + const reported = await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + createdAt: NOW, + }); + + expect(reported.outcome).toBe("completed"); + if (reported.outcome !== "completed") throw new Error("unreachable"); + expect(reported.completion?.output).toEqual({ inline: '{"ok":true}' }); + } finally { + await store.quit(); + } + } + ); +}); + +describe("registerOrReport", () => { + redisTest("registers a watcher against a PENDING waitpoint", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + const result = await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + createdAt: NOW, + }); + expect(result.outcome).toBe("registered"); + } finally { + await store.quit(); + } + }); + + redisTest("reports the completion inline for a COMPLETED waitpoint", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.complete({ waitpointId: "w_a", completion: completion() }); + + const result = await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + createdAt: NOW, + }); + + expect(result.outcome).toBe("completed"); + if (result.outcome !== "completed") throw new Error("unreachable"); + expect(result.completion?.output).toEqual({ inline: '{"ok":true}' }); + } finally { + await store.quit(); + } + }); + + redisTest("throws for a waitpoint that does not exist", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await expect( + store.registerOrReport({ waitpointId: "w_missing", runId: "run_1", createdAt: NOW }) + ).rejects.toThrow(WaitpointNotFoundError); + } finally { + await store.quit(); + } + }); + + redisTest("keeps one watcher entry per batch index", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + batchIndex: 0, + createdAt: NOW, + }); + await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + batchIndex: 2, + createdAt: NOW, + }); + + const completed = await store.complete({ waitpointId: "w_a", completion: completion() }); + expect(completed.watchers).toHaveLength(2); + expect(completed.watchers.map((w) => w.batchIndex).sort()).toEqual([0, 2]); + } finally { + await store.quit(); + } + }); + + redisTest("carries spanIdToComplete through to the watcher entry", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.registerOrReport({ + waitpointId: "w_a", + runId: "run_1", + spanIdToComplete: "span_abc", + createdAt: NOW, + }); + + const completed = await store.complete({ waitpointId: "w_a", completion: completion() }); + expect(completed.watchers[0]!.spanIdToComplete).toBe("span_abc"); + expect(completed.watchers[0]!.runId).toBe("run_1"); + expect(completed.watchers[0]!.createdAt).toBe(NOW); + } finally { + await store.quit(); + } + }); +}); + +describe("complete", () => { + redisTest("flips PENDING to COMPLETED and returns the watchers", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.registerOrReport({ waitpointId: "w_a", runId: "run_1", createdAt: NOW }); + await store.registerOrReport({ waitpointId: "w_a", runId: "run_2", createdAt: NOW }); + + const result = await store.complete({ waitpointId: "w_a", completion: completion() }); + + expect(result.outcome).toBe("completed"); + expect(result.watchers.map((w) => w.runId).sort()).toEqual(["run_1", "run_2"]); + } finally { + await store.quit(); + } + }); + + redisTest("is idempotent and returns the watchers again", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.registerOrReport({ waitpointId: "w_a", runId: "run_1", createdAt: NOW }); + + const first = await store.complete({ waitpointId: "w_a", completion: completion() }); + const second = await store.complete({ + waitpointId: "w_a", + completion: completion({ output: { inline: '{"second":true}' } }), + }); + + expect(first.outcome).toBe("completed"); + expect(second.outcome).toBe("already"); + // The FIRST completion wins, matching the guard on status = PENDING. + expect(second.completion?.output).toEqual({ inline: '{"ok":true}' }); + expect(second.watchers.map((w) => w.runId)).toEqual(["run_1"]); + } finally { + await store.quit(); + } + }); + + redisTest("throws for a waitpoint that does not exist", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await expect( + store.complete({ waitpointId: "w_missing", completion: completion() }) + ).rejects.toThrow(WaitpointNotFoundError); + } finally { + await store.quit(); + } + }); + + redisTest("returns an empty watcher list when nobody is blocked", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + const result = await store.complete({ waitpointId: "w_a", completion: completion() }); + expect(result.watchers).toEqual([]); + } finally { + await store.quit(); + } + }); + + redisTest("sets no TTL on the record or the watcher key", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + const probe = createRedisClient(redisOptions); + try { + await store.createIfAbsent({ record: record("w_a"), status: "PENDING" }); + await store.registerOrReport({ waitpointId: "w_a", runId: "run_1", createdAt: NOW }); + await store.complete({ waitpointId: "w_a", completion: completion() }); + + // -1 means the key exists with no expiry. Anything >= 0 breaks the retention rule. + expect(await probe.pttl("wp:{w_a}")).toBe(-1); + expect(await probe.pttl("wp:{w_a}:w")).toBe(-1); + } finally { + probe.disconnect(); + await store.quit(); + } + }); +}); + +describe("the single-slot guard", () => { + redisTest("rejects an invocation whose keys span two tags", async ({ redisOptions }) => { + const store = coordinator(redisOptions); + try { + // Reaches the same wrapper every operation goes through, so this proves the guard + // is live at the call path and not only in the pure unit test. + expect(() => + store.assertKeysForTest("wpComplete", ["wp:{w_a}", "wp:run:{run_1}:pend"]) + ).toThrow(WaitpointKeyTagError); + } finally { + await store.quit(); + } + }); +}); diff --git a/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts new file mode 100644 index 000000000..f98436e6d --- /dev/null +++ b/internal-packages/run-engine/src/engine/waitpointCoordinator/storeCoordinator.ts @@ -0,0 +1,239 @@ +import { createRedisClient, type Redis, type RedisOptions } from "@internal/redis"; +import { Logger } from "@trigger.dev/core/logger"; +import { assertSingleSlot, waitpointKeys, watcherField } from "./keys.js"; +import { registerWaitpointCommands } from "./scripts.js"; + +/** The values written into a record's `status` field. Uppercase, and never a token. */ +export type WaitpointStatus = "PENDING" | "COMPLETED"; + +/** Every script this coordinator may invoke. The wrapper below is the only entry point. */ +type ScriptName = + | "wpCreateIfAbsent" + | "wpRegisterOrReport" + | "wpComplete" + | "wpIdemReserve" + | "runAbsorbBlockers" + | "runDeliverCompletion" + | "runReadBlockState" + | "runClear"; + +/** + * The immutable half of a waitpoint, written once at creation. Carries every field the + * legacy-shaped return types need, including the two that gate the executor-visible + * idempotency key and the token surface. + */ +export type WaitpointRecordInput = { + id: string; + friendlyId: string; + type: "RUN" | "BATCH" | "DATETIME" | "MANUAL"; + environmentId: string; + projectId: string; + createdAt: string; + updatedAt: string; + userProvidedIdempotencyKey: boolean; + tags: string[]; + idempotencyKey?: string; + idempotencyKeyExpiresAt?: string; + completedAfter?: string; + completedByTaskRunId?: string; + completedByBatchId?: string; +}; + +/** + * A stored output: a small inline value, an already-offloaded reference, or null when the + * value is re-derivable from a business fact and is therefore never copied forward. + */ +export type WaitpointCompletionOutput = { inline: string } | { ref: string } | null; + +/** + * The completion half of a waitpoint, written at the flip. + * + * This is the coordinator's OWN type, deliberately not a projection of any frozen record + * type. The store treats a completion as an opaque blob: it writes it, returns it, and + * never inspects a field. Whoever owns the read-time resolver maps between this and the + * frozen record shape, so the two can evolve without a type dependency in either + * direction. + */ +export type WaitpointCompletion = { + /** ISO 8601. */ + completedAt: string; + outputType: string; + outputIsError: boolean; + output: WaitpointCompletionOutput; +}; + +export type WatcherEntry = { + runId: string; + batchIndex?: number; + spanIdToComplete?: string; + createdAt: string; +}; + +export type CreateIfAbsentResult = + | { outcome: "created" } + | { + outcome: "exists"; + record: WaitpointRecordInput; + status: WaitpointStatus; + completion?: WaitpointCompletion; + }; + +export type RegisterOrReportResult = + | { outcome: "registered" } + | { outcome: "completed"; completion?: WaitpointCompletion }; + +export type CompleteResult = { + outcome: "completed" | "already"; + completion?: WaitpointCompletion; + watchers: WatcherEntry[]; +}; + +export class WaitpointNotFoundError extends Error { + constructor(waitpointId: string) { + super(`Waitpoint ${waitpointId} is not present in the store`); + this.name = "WaitpointNotFoundError"; + } +} + +export type WaitpointStoreCoordinatorOptions = { + redisOptions: RedisOptions; + logger?: Logger; +}; + +// Lua returns '' for an absent value, never nil, because every reply slot is coerced to +// keep the array from truncating. So a nullish check would not fire and JSON.parse('') +// throws. One helper, used at every decode site. +function parseJson(raw: string | undefined): T | undefined { + return raw ? (JSON.parse(raw) as T) : undefined; +} + +export class WaitpointStoreCoordinator { + private readonly redis: Redis; + private readonly logger: Logger; + #quit?: Promise; + + constructor(options: WaitpointStoreCoordinatorOptions) { + this.logger = options.logger ?? new Logger("WaitpointStoreCoordinator", "debug"); + this.redis = createRedisClient(options.redisOptions, { + onError: (error) => + this.logger.error("WaitpointStoreCoordinator redis client error", { error }), + }); + registerWaitpointCommands(this.redis); + } + + // Idempotent and error-swallowing: every test calls this in a finally, and a double quit + // must never mask the real assertion failure. + async quit(): Promise { + if (!this.#quit) { + this.#quit = this.redis.quit().then( + () => undefined, + () => undefined + ); + } + await this.#quit; + } + + /** + * The ONLY way this class invokes a script. Routing every call through one place is what + * makes the single-slot guard un-forgettable: a method added later cannot reach a script + * without passing its keys through this assertion. + * + * Every script's signature is (...keys, ...argv) => string[], so one cast covers them + * all. The typed RedisCommander augmentation in scripts.ts documents each shape. + */ + #call(script: ScriptName, keys: string[], ...argv: string[]): Promise { + assertSingleSlot(script, keys); + const command = this.redis[script] as (...args: string[]) => Promise; + return command.call(this.redis, ...keys, ...argv); + } + + /** Exposed for the guard's own test. Asserts and returns; never invokes a script. */ + assertKeysForTest(operation: string, keys: string[]): void { + assertSingleSlot(operation, keys); + } + + async createIfAbsent(args: { + record: WaitpointRecordInput; + status: WaitpointStatus; + completion?: WaitpointCompletion; + }): Promise { + const keys = waitpointKeys(args.record.id); + + const reply = await this.#call( + "wpCreateIfAbsent", + [keys.record], + JSON.stringify(args.record), + args.status, + args.completion ? JSON.stringify(args.completion) : "" + ); + + if (reply[0] === "created") { + return { outcome: "created" }; + } + + return { + outcome: "exists", + record: JSON.parse(reply[1] ?? "{}") as WaitpointRecordInput, + status: (reply[2] ?? "PENDING") as WaitpointStatus, + completion: parseJson(reply[3]), + }; + } + + async registerOrReport(args: { + waitpointId: string; + runId: string; + batchIndex?: number | null; + spanIdToComplete?: string; + createdAt: string; + }): Promise { + const keys = waitpointKeys(args.waitpointId); + + // batchIndex is nullable at the boundary (matching the column) and undefined inside, + // because JSON.stringify drops an undefined field but keeps a null one. + const watcher: WatcherEntry = { + runId: args.runId, + batchIndex: args.batchIndex ?? undefined, + spanIdToComplete: args.spanIdToComplete, + createdAt: args.createdAt, + }; + + const reply = await this.#call( + "wpRegisterOrReport", + [keys.record, keys.watchers], + watcherField(args.runId, args.batchIndex), + JSON.stringify(watcher) + ); + + if (reply[0] === "missing") { + throw new WaitpointNotFoundError(args.waitpointId); + } + if (reply[0] === "completed") { + return { outcome: "completed", completion: parseJson(reply[1]) }; + } + + return { outcome: "registered" }; + } + + async complete(args: { + waitpointId: string; + completion: WaitpointCompletion; + }): Promise { + const keys = waitpointKeys(args.waitpointId); + + const reply = await this.#call( + "wpComplete", + [keys.record, keys.watchers], + JSON.stringify(args.completion) + ); + + if (reply[0] === "missing") { + throw new WaitpointNotFoundError(args.waitpointId); + } + + return { + outcome: reply[0] as "completed" | "already", + completion: parseJson(reply[1]), + watchers: reply.slice(2).map((entry) => JSON.parse(entry) as WatcherEntry), + }; + } +}