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 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 142f0ef..5380df9 100644 --- a/src/manager.ts +++ b/src/manager.ts @@ -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,7 +119,9 @@ export default class Manager { } queue.push(payload); - this.store.saveEvent(name, payload, true); + if (this.persistEnabled) { + this.store.saveEvent(name, payload, true); + } return this; } @@ -88,7 +140,7 @@ export default class Manager { this.queues.delete(name); } - if (payload !== undefined) { + if (payload !== undefined && this.persistEnabled) { this.store.saveEvent(name, payload, false); } @@ -103,7 +155,7 @@ export default class Manager { return undefined; } - return queue[0]; + return queue.peek(); } public length(name: string): number { @@ -118,11 +170,13 @@ 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 { @@ -147,7 +201,7 @@ export default class Manager { if (!existing && !this.canCreateQueue()) { return; } - const queue = existing || []; + const queue = existing || new FIFOQueue(); if (queue.length >= this.queueDepthLimit) { return; } @@ -168,4 +222,4 @@ export default class Manager { this.queues.delete(event.queue); } } -} +} \ No newline at end of file diff --git a/src/persist.ts b/src/persist.ts index fb5373d..000a07b 100644 --- a/src/persist.ts +++ b/src/persist.ts @@ -17,15 +17,31 @@ 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; 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"; @@ -38,33 +54,57 @@ export class FileStore implements QueueStore { Deno.mkdirSync(this.directory, { recursive: true }); } + // 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.ensureDirectory(); + this.writeHandle = Deno.openSync(this.path, { write: true, create: true, append: true }); + } + } + public saveEvent(queueName: string, payload: T, isEnqueue: boolean): void { - this.ensureDirectory(); + 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 { - file.writeSync(new TextEncoder().encode(line + "\n")); + 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 { + 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 { - this.ensureDirectory(); - 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(); } } @@ -73,34 +113,39 @@ 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) { + if (!read) { 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) { + const event = parseLine(line); + if (event) { + events.push(event); + } + } + 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) { + const event = parseLine(leftover); + if (event) { + events.push(event); + } } - const content = new TextDecoder().decode(buf); - return content.split("\n") - .filter((line: string) => line.length > 0) - .flatMap((line: string) => { - try { - const event = JSON.parse(line); - return isQueueEvent(event) ? [event] : []; - } catch { - return []; - } - }); + return events; } finally { file.unlockSync(); file.close(); @@ -114,8 +159,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 { @@ -130,6 +186,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 = []; } @@ -139,4 +201,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 f339406..4b0e1fd 100644 --- a/src/rate_limiter.ts +++ b/src/rate_limiter.ts @@ -27,19 +27,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; @@ -83,4 +103,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 4957d84..2e5beb0 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(); }); @@ -411,6 +422,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", () => { @@ -424,6 +436,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", () => { @@ -444,6 +457,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", () => { @@ -514,3 +528,52 @@ 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 skips events exceeding queue depth limit", () => { + const persist = new Persistency.MemoryStore(); + persist.saveEvent("q", "a", true); + persist.saveEvent("q", "b", true); + const mgr = new QueueManager(persist, 1, 1000); + mgr.load(); + assertEquals(mgr.length("q"), 1); + assertEquals(mgr.dequeue("q"), "a"); +}); + +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"); +}); diff --git a/tests/persist_test.ts b/tests/persist_test.ts index df6173f..10278b2 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,9 +352,87 @@ 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 }); }); +// ── 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"); +}); + +// ── FileStore directory and malformed line handling ────────────────────────── + Deno.test("persist FileStore.saveEvent creates a missing directory", () => { const tmpDir = Deno.makeTempDirSync(); const missing = tmpDir + "/does-not-exist"; @@ -346,6 +440,7 @@ Deno.test("persist FileStore.saveEvent creates a missing directory", () => { persist.dir(missing); persist.saveEvent("q", "x", true); assertEquals(persist.loadState()[0].payload, "x"); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -363,6 +458,7 @@ Deno.test("persist FileStore.loadState skips malformed lines", () => { const events = persist.loadState(); assertEquals(events.length, 1); assertEquals(events[0].payload, "kept"); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); @@ -419,5 +515,6 @@ Deno.test("FileStore.loadState skips non-event JSON records", () => { const events = persist.loadState(); assertEquals(events.length, 1); assertEquals(events[0].payload, "kept"); + persist.close(); Deno.removeSync(tmpDir, { recursive: true }); }); diff --git a/tests/rate_limiter_test.ts b/tests/rate_limiter_test.ts index 5c74c02..49deb62 100644 --- a/tests/rate_limiter_test.ts +++ b/tests/rate_limiter_test.ts @@ -439,3 +439,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); +});