From fdd8ad4d3f4f467380f01fbd1893796e67fbc41e Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 27 Aug 2026 11:40:44 -0400 Subject: [PATCH] refactor(core): centralize full-prefix KV scans --- packages/core/src/job.ts | 9 +----- packages/core/src/kv.ts | 53 ++++++++++++++++++++++------------- packages/core/test/kv.test.ts | 11 ++++++++ 3 files changed, 45 insertions(+), 28 deletions(-) diff --git a/packages/core/src/job.ts b/packages/core/src/job.ts index 1137ec0c026b..5364d9bae893 100644 --- a/packages/core/src/job.ts +++ b/packages/core/src/job.ts @@ -402,14 +402,7 @@ export const make = Effect.gen(function* () { }) const pendingBackground: Interface["pendingBackground"] = Effect.gen(function* () { - const recovered: Background[] = [] - let after: string | undefined - do { - const page = yield* kv.scan({ prefix: backgroundPrefix, after }) - recovered.push(...Array.filterMap(page.entries, (entry) => decodeBackground(entry.value))) - after = page.next - } while (after) - return recovered + return Array.filterMap(yield* kv.scanAll(backgroundPrefix), (entry) => decodeBackground(entry.value)) }).pipe(Effect.withSpan("Job.pendingBackground")) const completeBackground: Interface["completeBackground"] = Effect.fn("Job.completeBackground")((notificationID) => diff --git a/packages/core/src/kv.ts b/packages/core/src/kv.ts index d4cf15896af8..11b81067d9f7 100644 --- a/packages/core/src/kv.ts +++ b/packages/core/src/kv.ts @@ -29,6 +29,7 @@ export interface Interface { readonly set: (key: string, value: Value) => Effect.Effect readonly remove: (key: string) => Effect.Effect readonly scan: (options: ScanOptions) => Effect.Effect + readonly scanAll: (prefix: string) => Effect.Effect } export class Service extends Context.Service()("@opencode/KV") {} @@ -37,6 +38,28 @@ const layer = Layer.effect( Service, Effect.gen(function* () { const db = (yield* Database.Service).db + const scan: Interface["scan"] = Effect.fn("KV.scan")(function* (options) { + const limit = Number.isNaN(options.limit) ? 100 : Math.min(Math.max(Math.floor(options.limit ?? 100), 1), 1000) + const end = prefixEnd(options.prefix) + const rows = yield* db + .select({ key: KVTable.key, value: KVTable.value }) + .from(KVTable) + .where( + and( + options.prefix === "" ? undefined : gte(KVTable.key, options.prefix), + end === undefined ? undefined : lt(KVTable.key, end), + options.after === undefined ? undefined : gt(KVTable.key, options.after), + ), + ) + .orderBy(asc(KVTable.key)) + .limit(limit + 1) + .all() + .pipe(Effect.orDie) + const entries = rows.slice(0, limit) + if (rows.length <= limit) return { entries } + return { entries, next: entries[entries.length - 1].key } + }) + return Service.of({ get: Effect.fn("KV.get")(function* (key) { return (yield* db @@ -57,26 +80,16 @@ const layer = Layer.effect( remove: Effect.fn("KV.remove")(function* (key) { yield* db.delete(KVTable).where(eq(KVTable.key, key)).run().pipe(Effect.orDie) }), - scan: Effect.fn("KV.scan")(function* (options) { - const limit = Number.isNaN(options.limit) ? 100 : Math.min(Math.max(Math.floor(options.limit ?? 100), 1), 1000) - const end = prefixEnd(options.prefix) - const rows = yield* db - .select({ key: KVTable.key, value: KVTable.value }) - .from(KVTable) - .where( - and( - options.prefix === "" ? undefined : gte(KVTable.key, options.prefix), - end === undefined ? undefined : lt(KVTable.key, end), - options.after === undefined ? undefined : gt(KVTable.key, options.after), - ), - ) - .orderBy(asc(KVTable.key)) - .limit(limit + 1) - .all() - .pipe(Effect.orDie) - const entries = rows.slice(0, limit) - if (rows.length <= limit) return { entries } - return { entries, next: entries[entries.length - 1].key } + scan, + scanAll: Effect.fn("KV.scanAll")(function* (prefix) { + const entries: Entry[] = [] + let after: string | undefined + do { + const page = yield* scan({ prefix, after, limit: 1000 }) + entries.push(...page.entries) + after = page.next + } while (after !== undefined) + return entries }), }) }), diff --git a/packages/core/test/kv.test.ts b/packages/core/test/kv.test.ts index 5eba7e75c963..620a63c1eb3f 100644 --- a/packages/core/test/kv.test.ts +++ b/packages/core/test/kv.test.ts @@ -51,6 +51,8 @@ describe("KV", () => { entries: [{ key: `${prefix}éclair`, value: { order: 3 } }], }) expect(yield* kv.scan({ prefix: `${prefix}%_` })).toEqual({ entries: [] }) + expect(yield* kv.scanAll(`${prefix}%_`)).toEqual([]) + expect(yield* kv.scanAll(prefix)).toEqual([...first.entries, { key: `${prefix}éclair`, value: { order: 3 } }]) }), ) @@ -76,6 +78,15 @@ describe("KV", () => { expect((yield* kv.scan({ prefix, limit: 0 })).entries).toHaveLength(1) expect((yield* kv.scan({ prefix, limit: -10 })).entries).toHaveLength(1) expect((yield* kv.scan({ prefix, limit: Number.NaN })).entries).toHaveLength(100) + + const all = kv.scanAll(prefix) + const entries = Array.from({ length: 1001 }, (_, index) => { + const key = `${prefix}${index.toString().padStart(4, "0")}` + return { key, value: key } + }) + expect(yield* all).toEqual(entries) + yield* kv.remove(`${prefix}0000`) + expect(yield* all).toEqual(entries.slice(1)) }), ) })