From a9947cc5ddcf4021c686d722a7bb5c41188559e3 Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 16:04:49 +0100 Subject: [PATCH 1/7] fix: resolve performance issues #53-60 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - #53/#55: Replace Array.shift() FIFO with head-index FIFOQueue (O(1) amortized dequeue), fixing both per-request O(D) dequeue cost and O(E²) load() replay cost - #54: RateLimiter eviction sort uses last-element lookup (O(1) per comparison) instead of Math.max(...arr) (O(T)), reducing from O(I·T·log I) to O(I·log I) - #56: RateLimiter cleanup uses binary search on sorted timestamps with fast-path skip for fresh IPs, reducing common-case from O(I·T) to O(I) - #57: Manager.save() uses saveBatch() — one file open/lock/close for all items instead of per-item - #58: FileStore.loadState() stream-parses line-by-line, eliminating 3x peak memory from split/filter/map chain - #59: Manager gains persistEnabled flag; when persistence is disabled, saveEvent is skipped entirely, preventing unbounded MemoryStore event accumulation - #60: FileStore keeps write handle open (lazy-open) instead of open/close per saveEvent; added close() to QueueStore interface for resource cleanup Co-Authored-By: Claude --- deno.lock | 2 + main.ts | 4 +- src/manager.ts | 84 +++++++++++++++++++++++++++------- src/persist.ts | 102 ++++++++++++++++++++++++++++++++---------- src/rate_limiter.ts | 34 +++++++++++--- tests/manager_test.ts | 14 ++++++ tests/persist_test.ts | 17 +++++++ 7 files changed, 210 insertions(+), 47 deletions(-) diff --git a/deno.lock b/deno.lock index 9a883ff..34d71e6 100644 --- a/deno.lock +++ b/deno.lock @@ -2,7 +2,9 @@ "version": "5", "specifiers": { "jsr:@std/assert@*": "1.0.10", + "jsr:@std/assert@1.0": "1.0.10", "jsr:@std/cli@*": "1.0.29", + "jsr:@std/cli@1.0": "1.0.29", "jsr:@std/internal@^1.0.5": "1.0.13" }, "jsr": { diff --git a/main.ts b/main.ts index af8bb1d..b7d6a72 100644 --- a/main.ts +++ b/main.ts @@ -19,7 +19,7 @@ const PERSIST_ENGINE = CONFIG.persistEnabled PERSIST_ENGINE.dir(CONFIG.persistDir); // Set up the manager, which will handle our queues for us -const MANAGER = new QueueManager(PERSIST_ENGINE, CONFIG.queueDepthLimit, CONFIG.queueCountLimit); +const MANAGER = new QueueManager(PERSIST_ENGINE, CONFIG.queueDepthLimit, CONFIG.queueCountLimit, CONFIG.persistEnabled); // Load up any existing queue data, if we're persisting if (PERSIST_ENGINE instanceof Persistency.FileStore) { @@ -49,6 +49,8 @@ async function shutdown(signal: string): Promise { MANAGER.save(); } + PERSIST_ENGINE.close(); + writeLog("Goodbye!"); Deno.exit(0); diff --git a/src/manager.ts b/src/manager.ts index 4418ea3..2697d68 100644 --- a/src/manager.ts +++ b/src/manager.ts @@ -1,4 +1,4 @@ -import { QueueStore } from "./persist.ts" +import { QueueStore, QueueEvent } from "./persist.ts" export const MAX_QUEUE_NAME_LENGTH = 128; export class QueueNameTooLongError extends Error { @@ -8,20 +8,70 @@ export class QueueNameTooLongError extends Error { } } +/** + * FIFO queue with O(1) amortized enqueue and dequeue. + * Uses a head index instead of Array.shift() to avoid O(n) reindexing. + */ +class FIFOQueue { + private items: T[] = []; + private head = 0; + + push(item: T): void { + this.items.push(item); + } + + shift(): T | undefined { + if (this.head >= this.items.length) return undefined; + const item = this.items[this.head]; + this.items[this.head] = undefined as T; // help GC + this.head++; + // Compact when the consumed prefix exceeds the remaining items + if (this.head > 16 && this.head >= (this.items.length >> 1)) { + this.items = this.items.slice(this.head); + this.head = 0; + } + return item; + } + + peek(): T | undefined { + return this.head < this.items.length ? this.items[this.head] : undefined; + } + + get length(): number { + return this.items.length - this.head; + } + + [Symbol.iterator](): Iterator { + let index = this.head; + const items = this.items; + const end = items.length; + return { + next(): IteratorResult { + if (index < end) { + return { value: items[index++], done: false }; + } + return { value: undefined as unknown as T, done: true }; + }, + }; + } +} + export default class Manager { - private queues: Map>; + private queues: Map>; private store: QueueStore; private queueDepthLimit: number; private queueCountLimit: number; + private persistEnabled: boolean; - constructor(store: QueueStore, queueDepthLimit?: number, queueCountLimit?: number) { + constructor(store: QueueStore, queueDepthLimit?: number, queueCountLimit?: number, persistEnabled?: boolean) { this.store = store; - this.queues = new Map; + this.queues = new Map(); this.queueDepthLimit = queueDepthLimit ?? 10000; this.queueCountLimit = queueCountLimit ?? 1000; + this.persistEnabled = persistEnabled ?? true; } - private register(name: string, queue: Array): Manager { + private register(name: string, queue: FIFOQueue): Manager { this.queues.set(name, queue); return this; @@ -52,13 +102,13 @@ export default class Manager { return queue.length < this.queueDepthLimit; } - private find(name: string): Array | undefined { + private find(name: string): FIFOQueue | undefined { return this.queues.get(name); } public enqueue(name: string, payload: T): Manager { this.validateName(name); - const queue = this.find(name) || []; + const queue = this.find(name) || new FIFOQueue(); if (this.registered(name) === false) { this.register(name, queue); @@ -69,14 +119,16 @@ export default class Manager { } queue.push(payload); - this.store.saveEvent(name, payload, true); + if (this.persistEnabled) { + this.store.saveEvent(name, payload, true); + } return this; } public dequeue(name: string): T | undefined { this.validateName(name); - const queue = this.find(name) || []; + const queue = this.find(name) || new FIFOQueue(); const wasRegistered = this.registered(name); @@ -95,7 +147,7 @@ export default class Manager { this.queues.delete(name); } - if (payload !== undefined) { + if (payload !== undefined && this.persistEnabled) { this.store.saveEvent(name, payload, false); } @@ -110,12 +162,12 @@ export default class Manager { return undefined; } - return queue[0]; + return queue.peek(); } public length(name: string): number { this.validateName(name); - const queue = this.find(name) || []; + const queue = this.find(name) || new FIFOQueue(); if (this.registered(name) === false) { if (!this.canCreateQueue()) { @@ -133,18 +185,20 @@ export default class Manager { public save(): void { this.store.clear(); + const events: QueueEvent[] = []; for (const [name, queue] of this.queues) { for (const item of queue) { - this.store.saveEvent(name, item, true); + events.push({ queue: name, payload: item, enqueue: true, dequeue: false }); } } + this.store.saveBatch(events); } public load(): void { const events = this.store.loadState(); events.forEach((event) => { - const queue = this.find(event.queue) || []; + const queue = this.find(event.queue) || new FIFOQueue(); if (this.registered(event.queue) === false) { this.register(event.queue, queue); @@ -168,4 +222,4 @@ export default class Manager { this.store.clear(); } -} +} \ No newline at end of file diff --git a/src/persist.ts b/src/persist.ts index 5a686e8..63b05f0 100644 --- a/src/persist.ts +++ b/src/persist.ts @@ -7,43 +7,72 @@ export interface QueueEvent { export interface QueueStore { saveEvent(queueName: string, payload: T, isEnqueue: boolean): void; + saveBatch(events: Array>): void; loadState(): Array>; clear(): void; dir(dir: string): void; + close(): void; } export class FileStore implements QueueStore { private directory: string = ''; + private writeHandle: Deno.FsFile | null = null; + private encoder = new TextEncoder(); private get path(): string { return this.directory + "persist.dat"; } + // Lazily open the write handle so that dir() with an invalid path + // doesn't throw until an actual I/O operation is attempted. + private ensureOpen(): void { + if (this.writeHandle === null) { + this.writeHandle = Deno.openSync(this.path, { write: true, create: true, append: true }); + } + } + public saveEvent(queueName: string, payload: T, isEnqueue: boolean): void { + this.ensureOpen(); const line = JSON.stringify({ queue: queueName, payload: payload, enqueue: isEnqueue, dequeue: !isEnqueue }); - const file = Deno.openSync(this.path, { write: true, create: true, append: true }); - file.lockSync(true); + this.writeHandle!.lockSync(true); + try { + this.writeHandle!.writeSync(this.encoder.encode(line + "\n")); + } finally { + this.writeHandle!.unlockSync(); + } + } + + public saveBatch(events: Array>): void { + if (events.length === 0) return; + this.ensureOpen(); + this.writeHandle!.lockSync(true); try { - file.writeSync(new TextEncoder().encode(line + "\n")); + for (const event of events) { + const line = JSON.stringify({ + queue: event.queue, + payload: event.payload, + enqueue: event.enqueue, + dequeue: event.dequeue + }); + this.writeHandle!.writeSync(this.encoder.encode(line + "\n")); + } } finally { - file.unlockSync(); - file.close(); + this.writeHandle!.unlockSync(); } } public clear(): void { - const file = Deno.openSync(this.path, { write: true, create: true }); - file.lockSync(true); + this.ensureOpen(); + this.writeHandle!.lockSync(true); try { - file.truncateSync(0); + this.writeHandle!.truncateSync(0); } finally { - file.unlockSync(); - file.close(); + this.writeHandle!.unlockSync(); } } @@ -52,27 +81,33 @@ export class FileStore implements QueueStore { const file = Deno.openSync(this.path, { read: true }); file.lockSync(false); try { - const chunks: Uint8Array[] = []; + // Stream-parse line by line to avoid 3x peak memory from split/filter/map + const events: QueueEvent[] = []; + const decoder = new TextDecoder(); const chunk = new Uint8Array(4096); - let totalRead = 0; + let leftover = ""; while (true) { const read = file.readSync(chunk); if (read === null || read <= 0) { break; } - chunks.push(chunk.slice(0, read)); - totalRead += read; + leftover += decoder.decode(chunk.subarray(0, read), { stream: true }); + let idx = leftover.indexOf("\n"); + while (idx >= 0) { + const line = leftover.slice(0, idx); + leftover = leftover.slice(idx + 1); + if (line.length > 0) { + events.push(JSON.parse(line)); + } + idx = leftover.indexOf("\n"); + } } - const buf = new Uint8Array(totalRead); - let offset = 0; - for (const c of chunks) { - buf.set(c, offset); - offset += c.length; + // Flush decoder and process any remaining line (no trailing newline) + leftover += decoder.decode(); + if (leftover.length > 0) { + events.push(JSON.parse(leftover)); } - const content = new TextDecoder().decode(buf); - return content.split("\n") - .filter((line: string) => line.length > 0) - .map((line: string) => JSON.parse(line)); + return events; } finally { file.unlockSync(); file.close(); @@ -86,8 +121,19 @@ export class FileStore implements QueueStore { } public dir(dir: string): void { + if (this.writeHandle !== null) { + this.writeHandle.close(); + this.writeHandle = null; + } this.directory = dir.replace(/\/$/, '') + "/"; } + + public close(): void { + if (this.writeHandle !== null) { + this.writeHandle.close(); + this.writeHandle = null; + } + } } export class MemoryStore implements QueueStore { @@ -102,6 +148,12 @@ export class MemoryStore implements QueueStore { }); } + public saveBatch(events: Array>): void { + for (const event of events) { + this.events.push(event); + } + } + public clear(): void { this.events = []; } @@ -111,4 +163,6 @@ export class MemoryStore implements QueueStore { } public dir(): void {} -} + + public close(): void {} +} \ No newline at end of file diff --git a/src/rate_limiter.ts b/src/rate_limiter.ts index a01eea4..f62dc74 100644 --- a/src/rate_limiter.ts +++ b/src/rate_limiter.ts @@ -29,19 +29,39 @@ export class RateLimiter { private cleanupStaleEntries(now: number): void { const cutoff = now - this.windowMs; for (const [ip, timestamps] of this.requestTimestamps) { - const fresh = timestamps.filter(ts => ts > cutoff); - if (fresh.length === 0) { + // Fast path: timestamps are sorted ascending, so if the oldest + // (first) is still fresh, all are fresh — skip without scanning. + if (timestamps.length === 0 || timestamps[0] > cutoff) { + continue; + } + // Binary search for the first fresh timestamp (array is sorted ascending) + let lo = 0, hi = timestamps.length; + while (lo < hi) { + const mid = (lo + hi) >>> 1; + if (timestamps[mid] > cutoff) { + hi = mid; + } else { + lo = mid + 1; + } + } + // lo is the index of the first fresh timestamp + if (lo === timestamps.length) { + // All stale this.requestTimestamps.delete(ip); - } else if (fresh.length !== timestamps.length) { - this.requestTimestamps.set(ip, fresh); + } else if (lo > 0) { + // Some stale, some fresh — keep only the fresh ones + this.requestTimestamps.set(ip, timestamps.slice(lo)); } } if (this.requestTimestamps.size > this.maxTrackedIPs) { + // Evict IPs with the oldest last-activity (max timestamp). + // Timestamps are sorted ascending, so the last element is the max — + // O(1) per comparison instead of O(T) via Math.max(...arr). const entries = Array.from(this.requestTimestamps.entries()); entries.sort((a, b) => { - const aMax = Math.max(...a[1]); - const bMax = Math.max(...b[1]); + const aMax = a[1][a[1].length - 1]; + const bMax = b[1][b[1].length - 1]; return aMax - bMax; }); const toEvict = this.requestTimestamps.size - this.maxTrackedIPs; @@ -85,4 +105,4 @@ export class RateLimiter { return true; } -} +} \ No newline at end of file diff --git a/tests/manager_test.ts b/tests/manager_test.ts index 1054183..2239020 100644 --- a/tests/manager_test.ts +++ b/tests/manager_test.ts @@ -32,6 +32,7 @@ Deno.test("manager save flushes all queues to persist", () => { assertEquals(events.every((p: any) => p.enqueue === true), true); assertEquals(events.filter((p: any) => p.queue === "q1").length, 2); assertEquals(events.filter((p: any) => p.queue === "q2").length, 1); + persist.close(); }); Deno.test("manager save overwrites previous persist data", () => { @@ -47,6 +48,7 @@ Deno.test("manager save overwrites previous persist data", () => { const events = persist.loadState(); assertEquals(events.length, 2); + persist.close(); }); Deno.test("manager save with empty queues writes nothing", () => { @@ -57,6 +59,7 @@ Deno.test("manager save with empty queues writes nothing", () => { mgr.save(); assertEquals(persist.loadState(), []); + persist.close(); }); Deno.test("manager save preserves queue order", () => { @@ -73,6 +76,7 @@ Deno.test("manager save preserves queue order", () => { const events = persist.loadState(); const payloads = events.map((e: any) => e.payload); assertEquals(payloads, ["first", "second", "third"]); + persist.close(); }); Deno.test("manager save then load round-trips data", () => { @@ -94,6 +98,7 @@ Deno.test("manager save then load round-trips data", () => { assertEquals(mgr2.dequeue("x"), "one"); assertEquals(mgr2.dequeue("x"), "two"); assertEquals(mgr2.dequeue("y"), "three"); + persist.close(); }); async function startServer(env: Record): Promise<{ child: Deno.ChildProcess; port: number }> { @@ -194,6 +199,7 @@ Deno.test("manager persistency", () => { assertEquals("bar", mgr.dequeue("foo")); assertEquals(1, mgr.length("fee")); assertEquals("gat", mgr.dequeue("fee")); + persist.close(); }); Deno.test("json persistency", () => { @@ -209,6 +215,7 @@ Deno.test("json persistency", () => { assertEquals([], persist.loadState()); assertEquals(1, mgr.length("foo")); assertEquals(payload, mgr.dequeue("foo")); + persist.close(); }); Deno.test("persistency", () => { @@ -225,6 +232,7 @@ Deno.test("persistency", () => { assertNotEquals("", load()); persist.clear(); + persist.close(); }); Deno.test("concurrent enqueue writes to file", async () => { @@ -256,6 +264,7 @@ Deno.test("concurrent enqueue writes to file", async () => { }); persist.clear(); + persist.close(); }); Deno.test("concurrent manager operations maintain data integrity", async () => { @@ -288,6 +297,7 @@ Deno.test("concurrent manager operations maintain data integrity", async () => { assertEquals(0, mgr.length("q1")); persist.clear(); + persist.close(); }); Deno.test("persist operations handle I/O errors gracefully", () => { @@ -303,6 +313,7 @@ Deno.test("persist operations handle I/O errors gracefully", () => { } assertEquals(true, errorThrown); + persist.close(); }); @@ -378,6 +389,7 @@ Deno.test("manager load cleans up empty queues from persistence", () => { mgr.enqueue("other", "value"); assertEquals("value", mgr.dequeue("other")); + persist.close(); }); Deno.test("empty dequeue does not add entry to persistence log", () => { @@ -391,6 +403,7 @@ Deno.test("empty dequeue does not add entry to persistence log", () => { assertEquals(events.length, 0); persist.clear(); + persist.close(); }); Deno.test("non-empty dequeue still adds entry to persistence log", () => { @@ -411,6 +424,7 @@ Deno.test("non-empty dequeue still adds entry to persistence log", () => { assertEquals(deqEntry.enqueue, false); persist.clear(); + persist.close(); }); Deno.test("Manager enqueue and dequeue numbers", () => { diff --git a/tests/persist_test.ts b/tests/persist_test.ts index 572c6fc..ebd2946 100644 --- a/tests/persist_test.ts +++ b/tests/persist_test.ts @@ -53,6 +53,7 @@ Deno.test("persist FileStore.dir() with trailing slash writes files in that dire p.clear(); p.saveEvent("q", "hello", true); assertEquals(p.loadState()[0].payload, "hello"); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -63,6 +64,7 @@ Deno.test("persist FileStore.dir() without trailing slash still works", () => { p.clear(); p.saveEvent("q", "world", true); assertEquals(p.loadState()[0].payload, "world"); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -75,6 +77,7 @@ Deno.test("persist FileStore.dir() multi-segment path does not corrupt the file p.clear(); p.saveEvent("q", "nested", true); assertEquals(p.loadState()[0].payload, "nested"); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -91,6 +94,7 @@ Deno.test("persist FileStore.dir() second call replaces first", () => { let dir1HasFile = false; try { Deno.statSync(dir1 + "/persist.dat"); dir1HasFile = true; } catch { /* expected */ } assertEquals(dir1HasFile, false); + p.close(); Deno.removeSync(dir1, { recursive: true }); Deno.removeSync(dir2, { recursive: true }); }); @@ -107,6 +111,7 @@ Deno.test("persist FileStore.saveEvent() does not overwrite existing content", ( const events = p.loadState(); assertEquals(events[0].payload, "line1"); assertEquals(events[1].payload, "line2"); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -170,6 +175,7 @@ Deno.test("persist FileStore.saveEvent() waits for an existing file lock", async const store = new Persistency.FileStore(); store.dir(tmpDir); assertEquals(store.loadState()[0].payload, "from-child"); + store.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -181,6 +187,7 @@ Deno.test("persist FileStore.clear() creates file when it does not exist", () => p.dir(tmpDir + "/"); p.clear(); assertEquals(p.loadState(), []); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -192,6 +199,7 @@ Deno.test("persist FileStore.clear() truncates existing content", () => { p.saveEvent("q", "existing", true); p.clear(); assertEquals(p.loadState(), []); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -209,6 +217,7 @@ Deno.test("persist FileStore.loadState() reads large file correctly", () => { assertEquals(result.length, 100); assertEquals(result[0].payload.i, 0); assertEquals(result[99].payload.i, 99); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -219,6 +228,7 @@ Deno.test("persist FileStore.loadState() returns empty array when file does not const p = new Persistency.FileStore(); p.dir(tmpDir + "/"); assertEquals(p.loadState(), []); + p.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -242,6 +252,7 @@ Deno.test("manager save() writes enqueue:true dequeue:false for each item", () = assertEquals(events[0].dequeue, false); assertEquals(events[1].enqueue, true); assertEquals(events[1].dequeue, false); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -260,6 +271,7 @@ Deno.test("manager save() writes correct payloads in queue order", () => { const events = persist.loadState(); assertEquals(events[0].payload, "first"); assertEquals(events[1].payload, "second"); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -273,6 +285,7 @@ Deno.test("manager save() on empty manager writes nothing", () => { mgr.save(); assertEquals(persist.loadState(), []); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -290,6 +303,7 @@ Deno.test("manager enqueue log: enqueue=true dequeue=false", () => { const events = persist.loadState(); assertEquals(events[0].enqueue, true); assertEquals(events[0].dequeue, false); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -307,6 +321,7 @@ Deno.test("manager dequeue log: enqueue=false dequeue=true", () => { const events = persist.loadState(); assertEquals(events[0].enqueue, false); assertEquals(events[0].dequeue, true); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -322,6 +337,7 @@ Deno.test("manager load() enqueue entry adds item to queue", () => { const mgr = new QueueManager(persist); mgr.load(); assertEquals(mgr.dequeue("q"), "x"); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -336,5 +352,6 @@ Deno.test("manager load() dequeue entry removes item from queue", () => { const mgr = new QueueManager(persist); mgr.load(); assertEquals(mgr.length("q"), 0); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); From c19c65c50e9c5d761b4f5111a131dcbc76101c23 Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 16:07:18 +0100 Subject: [PATCH 2/7] chore: trigger CI From 26eb25e3aba7226b573f3aa770336593dab95e84 Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 16:11:22 +0100 Subject: [PATCH 3/7] chore: trigger CI synchronize event --- src/router.ts | 1 + 1 file changed, 1 insertion(+) diff --git a/src/router.ts b/src/router.ts index 0116c56..ccdb1b0 100644 --- a/src/router.ts +++ b/src/router.ts @@ -41,3 +41,4 @@ export class Router { return new Response("Not found.", { status: 404 }); }; } +// Performance fixes for #53-60 From ea2d9e841ac4f3cd54b9d61ab8e468b2d6a11564 Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 16:12:09 +0100 Subject: [PATCH 4/7] ci: add workflow_dispatch for manual triggering --- .github/workflows/ci.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3ba4156..81f21cf 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -8,6 +8,7 @@ on: tags: - "v*" pull_request: + workflow_dispatch: permissions: contents: read From 31902a26d000478d360dc3c971a71df296d9a6f2 Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 16:40:44 +0100 Subject: [PATCH 5/7] test: add mutation coverage for persist.ts and rate_limiter.ts - saveBatch([]) does not create file (kills early-return mutation) - saveBatch() writes multiple events (kills ensureOpen removal) - loadState() handles data without trailing newline (kills leftover flush mutations) - MemoryStore.saveEvent dequeue field (kills boolean negation mutation) - MemoryStore.saveBatch appends events (kills loop body mutations) - cleanup skips IPs with empty timestamp arrays (kills fast-path mutations) - eviction sort orders by max timestamp regardless of insertion order (kills sort mutations) - Remove leftover trigger comment from router.ts Co-Authored-By: Claude --- src/router.ts | 1 - tests/persist_test.ts | 75 ++++++++++++++++++++++++++++++++++++++ tests/rate_limiter_test.ts | 41 +++++++++++++++++++++ 3 files changed, 116 insertions(+), 1 deletion(-) diff --git a/src/router.ts b/src/router.ts index ccdb1b0..0116c56 100644 --- a/src/router.ts +++ b/src/router.ts @@ -41,4 +41,3 @@ export class Router { return new Response("Not found.", { status: 404 }); }; } -// Performance fixes for #53-60 diff --git a/tests/persist_test.ts b/tests/persist_test.ts index ebd2946..a3d7759 100644 --- a/tests/persist_test.ts +++ b/tests/persist_test.ts @@ -355,3 +355,78 @@ Deno.test("manager load() dequeue entry removes item from queue", () => { persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); + +// ── saveBatch edge cases ──────────────────────────────────────────────────── + +Deno.test("persist FileStore.saveBatch([]) does not create file", () => { + const tmpDir = Deno.makeTempDirSync(); + const p = new Persistency.FileStore(); + p.dir(tmpDir + "/"); + p.saveBatch([]); + let fileExists = false; + try { Deno.statSync(tmpDir + "/persist.dat"); fileExists = true; } catch { /* expected */ } + assertEquals(fileExists, false); + p.close(); + Deno.removeSync(tmpDir, { recursive: true }); +}); + +Deno.test("persist FileStore.saveBatch() writes multiple events in one batch", () => { + const tmpDir = Deno.makeTempDirSync(); + const p = new Persistency.FileStore(); + p.dir(tmpDir + "/"); + p.clear(); + p.saveBatch([ + { queue: "q", payload: "a", enqueue: true, dequeue: false }, + { queue: "q", payload: "b", enqueue: false, dequeue: true }, + ]); + const events = p.loadState(); + assertEquals(events.length, 2); + assertEquals(events[0].payload, "a"); + assertEquals(events[0].enqueue, true); + assertEquals(events[0].dequeue, false); + assertEquals(events[1].payload, "b"); + assertEquals(events[1].enqueue, false); + assertEquals(events[1].dequeue, true); + p.close(); + Deno.removeSync(tmpDir, { recursive: true }); +}); + +// ── loadState() without trailing newline ──────────────────────────────────── + +Deno.test("persist FileStore.loadState() handles data without trailing newline", () => { + const tmpDir = Deno.makeTempDirSync(); + const line1 = JSON.stringify({ queue: "q", payload: "with-newline", enqueue: true, dequeue: false }) + "\n"; + const line2 = JSON.stringify({ queue: "q", payload: "no-newline", enqueue: true, dequeue: false }); + Deno.writeFileSync(tmpDir + "/persist.dat", new TextEncoder().encode(line1 + line2)); + + const p = new Persistency.FileStore(); + p.dir(tmpDir + "/"); + const events = p.loadState(); + assertEquals(events.length, 2); + assertEquals(events[0].payload, "with-newline"); + assertEquals(events[1].payload, "no-newline"); + p.close(); + Deno.removeSync(tmpDir, { recursive: true }); +}); + +// ── MemoryStore edge cases ────────────────────────────────────────────────── + +Deno.test("persist MemoryStore.saveEvent sets dequeue to opposite of isEnqueue", () => { + const p = new Persistency.MemoryStore(); + p.saveEvent("q", "enq", true); + assertEquals(p.loadState()[0].dequeue, false); + p.saveEvent("q", "deq", false); + assertEquals(p.loadState()[1].dequeue, true); +}); + +Deno.test("persist MemoryStore.saveBatch() appends all events", () => { + const p = new Persistency.MemoryStore(); + p.saveBatch([ + { queue: "q", payload: "a", enqueue: true, dequeue: false }, + { queue: "q", payload: "b", enqueue: true, dequeue: false }, + ]); + const events = p.loadState(); + assertEquals(events.length, 2); + assertEquals(events[0].payload, "a"); + assertEquals(events[1].payload, "b"); +}); diff --git a/tests/rate_limiter_test.ts b/tests/rate_limiter_test.ts index 244b0fc..fa1287c 100644 --- a/tests/rate_limiter_test.ts +++ b/tests/rate_limiter_test.ts @@ -426,3 +426,44 @@ Deno.test("rate limiter: eviction sort uses max timestamp (Math.max not Math.min }, 20); }); }); + +// ── Cleanup fast path: empty timestamp arrays ────────────────────────────────── + +Deno.test("rate limiter: cleanup skips IPs with empty timestamp arrays", () => { + // The fast-path `if (timestamps.length === 0 || timestamps[0] > cutoff) continue` + // must skip IPs whose timestamp array is empty. Without the continue, the + // binary search would compute lo=0===length → delete the IP entry. + const limiter = new RateLimiter(100, 60000, 1, 10000); + const internal = limiter as unknown as { requestTimestamps: Map }; + + internal.requestTimestamps.set("empty.ip", []); + limiter.isAllowed(req("trigger.ip")); + + assertEquals(internal.requestTimestamps.has("empty.ip"), true); +}); + +// ── Eviction sort: insertion order ≠ sorted order ────────────────────────────── + +Deno.test("rate limiter: eviction sort orders by max timestamp regardless of insertion order", () => { + // Plant IPs in reverse order of their max timestamp so that without the + // sort, the wrong IP (newest) would be evicted first. + const limiter = new RateLimiter(100, 60000, 1, 2); + const internal = limiter as unknown as { requestTimestamps: Map }; + + const oldTs = Date.now() - 5000; + const midTs = Date.now() - 2000; + const newTs = Date.now(); + + // Insert in reverse: newest first, oldest last + internal.requestTimestamps.set("newest.ip", [newTs]); + internal.requestTimestamps.set("mid.ip", [midTs]); + internal.requestTimestamps.set("oldest.ip", [oldTs]); + + // Trigger cleanup — size=3 > maxTrackedIPs=2 → evict 1 (oldest) + limiter.isAllowed(req("trigger.ip")); + + // With sort: oldest.ip evicted (correct) + // Without sort (insertion order): newest.ip evicted (wrong) + assertEquals(internal.requestTimestamps.has("oldest.ip"), false); + assertEquals(internal.requestTimestamps.has("newest.ip"), true); +}); From deb38ab07eae38f2b99aaba3947be541d2776892 Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 17:07:39 +0100 Subject: [PATCH 6/7] test: add mutation coverage for manager.ts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add 8 targeted tests to kill survived Stryker mutations in manager.ts (78% → expected ~88%, threshold 80%): - QueueNameTooLongError message/name string literals - canEnqueue validates queue name - enqueue throws at depth limit (condition + block + message) - persistEnabled=false skips saveEvent - load() throws at depth limit (condition + operator + block + message) - load() skips events with neither enqueue nor dequeue - dequeue on registered-empty queue doesn't delete (wasNonEmpty + logical op) - length() auto-registers unknown queue (condition + block + boolean + call) Co-Authored-By: Claude --- tests/manager_test.ts | 60 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 60 insertions(+) diff --git a/tests/manager_test.ts b/tests/manager_test.ts index 2239020..2780346 100644 --- a/tests/manager_test.ts +++ b/tests/manager_test.ts @@ -495,3 +495,63 @@ Deno.test("Manager: separate queues don't interfere (catches queue isolation)", assertEquals(mgr.length("queue-b"), 1); assertEquals(mgr.dequeue("queue-b"), "b-item"); }); + +// ── Mutation coverage for manager.ts ────────────────────────────────────────── + +Deno.test("QueueNameTooLongError has correct message and name", () => { + const error = new QueueNameTooLongError(); + assertEquals(error.message, "Queue name too long"); + assertEquals(error.name, "QueueNameTooLongError"); +}); + +Deno.test("manager canEnqueue validates queue name", () => { + const mgr = new QueueManager(new Persistency.MemoryStore()); + assertThrows(() => mgr.canEnqueue("x".repeat(129)), QueueNameTooLongError); +}); + +Deno.test("manager enqueue throws when queue depth limit reached", () => { + const mgr = new QueueManager(new Persistency.MemoryStore(), 2, 1000); + mgr.enqueue("q", "a"); + mgr.enqueue("q", "b"); + assertThrows(() => mgr.enqueue("q", "c"), Error, "Queue depth limit reached"); +}); + +Deno.test("manager persistEnabled=false skips saveEvent on enqueue", () => { + const persist = new Persistency.MemoryStore(); + const mgr = new QueueManager(persist, 10000, 1000, false); + mgr.enqueue("q", "item"); + assertEquals(persist.loadState().length, 0); +}); + +Deno.test("manager load throws when queue depth limit exceeded", () => { + const persist = new Persistency.MemoryStore(); + persist.saveEvent("q", "a", true); + persist.saveEvent("q", "b", true); + const mgr = new QueueManager(persist, 1, 1000); + assertThrows(() => mgr.load(), Error, "Queue depth limit reached"); +}); + +Deno.test("manager load skips events with neither enqueue nor dequeue", () => { + const persist = new Persistency.MemoryStore(); + persist.saveBatch([ + { queue: "q", payload: "a", enqueue: true, dequeue: false }, + { queue: "q", payload: "b", enqueue: false, dequeue: false }, + ]); + const mgr = new QueueManager(persist); + mgr.load(); + assertEquals(mgr.length("q"), 1); + assertEquals(mgr.dequeue("q"), "a"); +}); + +Deno.test("dequeue on registered-but-empty queue does not delete it", () => { + const mgr = new QueueManager(new Persistency.MemoryStore(), 10000, 1); + mgr.length("auto-created"); // auto-registers empty queue + mgr.dequeue("auto-created"); // wasRegistered=true, wasNonEmpty=false + assertEquals(false, mgr.canCreateQueue()); // queue should still be registered +}); + +Deno.test("manager length auto-registers unknown queue", () => { + const mgr = new QueueManager(new Persistency.MemoryStore(), 10000, 1); + mgr.length("new-q"); + assertEquals(false, mgr.canCreateQueue()); // queue was auto-registered, at limit +}); From 85d05c04efa57d15862cd72252eb05e19a2df95c Mon Sep 17 00:00:00 2001 From: Jonathan Baldie Date: Thu, 27 Aug 2026 18:44:54 +0100 Subject: [PATCH 7/7] fix: extract parseLine to module-level function for complexity and LCOM4 Move the line-parsing try/catch logic from FileStore.loadState into a module-level parseLine() function. This reduces CyclomaticComplexity of loadState (13 -> under 10), eliminates empty catch blocks, and fixes the LCOM4 cohesion violation (2 -> 1) by keeping parseLine out of the FileStore class. Co-Authored-By: Claude --- src/persist.ts | 34 +++++++++++++++++++--------------- 1 file changed, 19 insertions(+), 15 deletions(-) diff --git a/src/persist.ts b/src/persist.ts index 2f1f713..000a07b 100644 --- a/src/persist.ts +++ b/src/persist.ts @@ -17,6 +17,18 @@ export function isQueueEvent(value: unknown): value is QueueEvent { event.enqueue !== event.dequeue; } +function parseLine(line: string): QueueEvent | undefined { + try { + const event = JSON.parse(line); + if (isQueueEvent(event)) { + return event; + } + return undefined; + } catch { + return undefined; + } +} + export interface QueueStore { saveEvent(queueName: string, payload: T, isEnqueue: boolean): void; saveBatch(events: Array>): void; @@ -108,7 +120,7 @@ export class FileStore implements QueueStore { let leftover = ""; while (true) { const read = file.readSync(chunk); - if (read === null || read <= 0) { + if (!read) { break; } leftover += decoder.decode(chunk.subarray(0, read), { stream: true }); @@ -117,13 +129,9 @@ export class FileStore implements QueueStore { const line = leftover.slice(0, idx); leftover = leftover.slice(idx + 1); if (line.length > 0) { - try { - const event = JSON.parse(line); - if (isQueueEvent(event)) { - events.push(event); - } - } catch { - // skip malformed lines + const event = parseLine(line); + if (event) { + events.push(event); } } idx = leftover.indexOf("\n"); @@ -132,13 +140,9 @@ export class FileStore implements QueueStore { // Flush decoder and process any remaining line (no trailing newline) leftover += decoder.decode(); if (leftover.length > 0) { - try { - const event = JSON.parse(leftover); - if (isQueueEvent(event)) { - events.push(event); - } - } catch { - // skip malformed lines + const event = parseLine(leftover); + if (event) { + events.push(event); } } return events;