From 66479aa9f0bbb2a7c2056c74bc0ae18562e110eb Mon Sep 17 00:00:00 2001 From: Daniel Sutton Date: Fri, 21 Aug 2026 13:08:39 +0100 Subject: [PATCH] fix(run-store): anchor liveness on the counter and make append idempotent --- .../run-store/src/redisSnapshotStore.test.ts | 99 +++++++++++++++++++ .../run-store/src/redisSnapshotStore.ts | 49 ++++++--- 2 files changed, 134 insertions(+), 14 deletions(-) diff --git a/internal-packages/run-store/src/redisSnapshotStore.test.ts b/internal-packages/run-store/src/redisSnapshotStore.test.ts index d1a7d8fe7..d3be34b19 100644 --- a/internal-packages/run-store/src/redisSnapshotStore.test.ts +++ b/internal-packages/run-store/src/redisSnapshotStore.test.ts @@ -2,6 +2,7 @@ // reference, so no Postgres container is needed. import { expect, describe } from "vitest"; import { redisTest } from "@internal/testcontainers"; +import { createRedisClient } from "@internal/redis"; import { snapshotKeys, deriveOrder, @@ -139,8 +140,106 @@ describe("append", () => { }); expect(r).toEqual({ outcome: "skippedNoKeyspace" }); expect(await store.getLatest("run_never")).toBeNull(); + + const k = snapshotKeys("run_never"); + const raw = createRedisClient(redisOptions); + try { + expect(await raw.exists(k.e, k.idx, k.cur, k.seq)).toBe(0); + } finally { + await raw.quit(); + } } finally { await store.quit(); } }); + + redisTest("skips a transition when only the seq key has expired", async ({ redisOptions }) => { + const store = new RedisSnapshotStore({ redisOptions, completedTtlMs: 1000 }); + try { + await store.append({ entry: entry({ id: "snap_1" }), kind: "birth", isTerminal: false }); + + const k = snapshotKeys("run_1"); + const raw = createRedisClient(redisOptions); + try { + await raw.del(k.seq); + } finally { + await raw.quit(); + } + + const r = await store.append({ + entry: entry({ id: "snap_2" }), + kind: "transition", + isTerminal: false, + }); + expect(r).toEqual({ outcome: "skippedNoKeyspace" }); + } finally { + await store.quit(); + } + }); + + redisTest( + "carries the original count forward on a carryForward append", + async ({ redisOptions }) => { + const store = new RedisSnapshotStore({ redisOptions, completedTtlMs: 1000 }); + try { + await store.append({ + entry: entry({ id: "snap_1" }), + kind: "birth", + isTerminal: false, + cycle: { + kind: "new", + completedWaitpoints: [ + { id: "w_a", index: 0 }, + { id: "w_b", index: 1 }, + ], + }, + }); + await store.append({ + entry: entry({ id: "snap_2" }), + kind: "transition", + isTerminal: false, + cycle: { kind: "carryForward", cycleSeq: 1 }, + }); + + const read = await store.getById("run_1", "snap_2"); + expect(read?.cycle).toEqual({ cycleSeq: 1, count: 2 }); + } finally { + await store.quit(); + } + } + ); + + redisTest( + "reports a duplicate id without overwriting the original entry", + async ({ redisOptions }) => { + const store = new RedisSnapshotStore({ redisOptions, completedTtlMs: 1000 }); + try { + const first = await store.append({ + entry: entry({ id: "snap_1", description: "created" }), + kind: "birth", + isTerminal: false, + }); + expect(first).toMatchObject({ outcome: "written", seq: 1 }); + + const dup = await store.append({ + entry: entry({ id: "snap_1", description: "different" }), + kind: "transition", + isTerminal: false, + }); + expect(dup).toEqual({ outcome: "duplicate", seq: 1 }); + + const read = await store.getById("run_1", "snap_1"); + expect(read?.entry.description).toBe("created"); + + const next = await store.append({ + entry: entry({ id: "snap_2" }), + kind: "transition", + isTerminal: false, + }); + expect(next).toMatchObject({ outcome: "written", seq: 2 }); + } finally { + await store.quit(); + } + } + ); }); diff --git a/internal-packages/run-store/src/redisSnapshotStore.ts b/internal-packages/run-store/src/redisSnapshotStore.ts index 72aae1c25..5aa7e8b97 100644 --- a/internal-packages/run-store/src/redisSnapshotStore.ts +++ b/internal-packages/run-store/src/redisSnapshotStore.ts @@ -76,7 +76,8 @@ export type AppendResult = cycleMismatch: boolean; } | { outcome: "skippedNoKeyspace" } - | { outcome: "forked"; actualCur: string }; + | { outcome: "forked"; actualCur: string } + | { outcome: "duplicate"; seq: number }; export type SnapshotStoreMetrics = { recordAppend(outcome: string, ttl: string): void; @@ -100,6 +101,7 @@ export type RedisSnapshotStoreOptions = { const SKIPPED = "skipped"; const FORKED = "forked"; const WRITTEN = "written"; +const DUPLICATE = "duplicate"; export class RedisSnapshotStore { private readonly redis: Redis; @@ -189,7 +191,8 @@ export class RedisSnapshotStore { orderJson, records, orderCount, - args.expectedCur ?? "" + args.expectedCur ?? "", + args.expectedCur !== undefined ? "1" : "0" )) as string[]; return this.#interpretAppend(reply, raw, orderJson); @@ -206,6 +209,10 @@ export class RedisSnapshotStore { this.metrics?.recordAppend("forked", "none"); return { outcome: "forked", actualCur: reply[1] ?? "" }; } + if (reply[0] === DUPLICATE) { + this.metrics?.recordAppend("duplicate", "none"); + return { outcome: "duplicate", seq: Number(reply[1]) }; + } const seq = Number(reply[1]); const cycleSeq = Number(reply[2]); const ttl = reply[3] as "none" | "completion" | "reapplied"; @@ -319,25 +326,33 @@ export class RedisSnapshotStore { local records = ARGV[10] local orderCount = ARGV[11] local expectedCur = ARGV[12] + local casEnabled = ARGV[13] == '1' - -- Liveness is ONE anchor. All keys get the same PEXPIRE but expire independently, so seq can - -- vanish while e and cur linger; anchoring on e treats a partly expired keyspace as gone, - -- once and consistently. A birth creates the keyspace; a transition that finds none writes - -- nothing. That state has two causes the caller must not merge: a completed run whose TTL - -- fired, and a run that predates this org's dual-write. - if kind == 'transition' and redis.call('EXISTS', eKey) == 0 then + -- Liveness is TWO anchors: e and seq. All keys get the same PEXPIRE but expire independently + -- (or seq can vanish under maxmemory eviction while e survives), so checking e alone lets a + -- late transition recreate seq with no TTL and restart it at 1 beside a surviving idx. A + -- birth always creates both in this same script, so this never rejects a live keyspace. + if kind == 'transition' and (redis.call('EXISTS', eKey) == 0 or redis.call('EXISTS', seqKey) == 0) then return { '${SKIPPED}' } end - -- Optional compare-and-set on cur, checked BEFORE any mutation. Absent by default, which - -- matches Postgres: it has no such guard either. - if expectedCur ~= '' then + -- Optional compare-and-set on cur, checked BEFORE any mutation. Gated on an explicit flag + -- (not on expectedCur ~= ''), so a caller asserting cur is unset (expectedCur = '') still + -- gets a real check instead of silently skipping it. + if casEnabled then local actual = redis.call('GET', curKey) if (actual or '') ~= expectedCur then return { '${FORKED}', actual or '' } end end + -- Append-only: a retried append (eg. ioredis reconnect-and-retry on a READONLY/UNBLOCKED + -- reply error) must not overwrite an existing entry's bytes or rescore it in idx. + local prior = redis.call('HGET', eKey, id .. '#s') + if prior then + return { '${DUPLICATE}', prior } + end + local seq = redis.call('HINCRBY', seqKey, 'e', 1) local cycleSeq = 0 @@ -346,14 +361,17 @@ export class RedisSnapshotStore { -- The STORE mints cycleSeq, so the sequence is dense by construction and the terminal -- PEXPIRE loop from 1..c is correct. cycleSeq = redis.call('HINCRBY', seqKey, 'c', 1) - redis.call('HSET', wpKey(cycleSeq), 'order', orderJson) + redis.call('HSET', wpKey(cycleSeq), 'order', orderJson, 'count', orderCount) if records ~= '' then redis.call('HSET', wpKey(cycleSeq), 'records', records) end elseif cycleMode == 'carry' then cycleSeq = cycleSeqIn - if redis.call('EXISTS', wpKey(cycleSeq)) == 0 then + local c = redis.call('HGET', wpKey(cycleSeq), 'count') + if not c then mismatch = 1 + else + orderCount = c end end @@ -363,7 +381,9 @@ export class RedisSnapshotStore { end -- idx indexes VALID entries only, which makes the since-cap exact. An invalid entry is still - -- reachable by id, and its seq is still readable from its own '#s' field. + -- reachable by id, and its seq is still readable from its own '#s' field. ZADD before SET cur + -- because Redis never rolls back a partially applied script: if a later call in this script + -- errored, having idx already written is the recoverable half of the pair. if isValid then redis.call('ZADD', idxKey, seq, id) redis.call('SET', curKey, id) @@ -444,6 +464,7 @@ declare module "@internal/redis" { records: string, orderCount: string, expectedCur: string, + casEnabled: string, callback?: Callback ): Result; readSnapshotById(