diff --git a/src/lib/outbox.test.ts b/src/lib/outbox.test.ts index c7f88b1..996e72d 100644 --- a/src/lib/outbox.test.ts +++ b/src/lib/outbox.test.ts @@ -228,6 +228,134 @@ describe('MetadataOutbox', () => { expect(JSON.stringify([...storage.values])).not.toContain('secret'); }); + it('serializes only validated import fields and never invokes caller serialization hooks', async () => { + const toJSON = vi.fn(() => ({ poisoned: true })); + const item = { + ...importRequest.items[0], + ignored: 'attacker-controlled extra field', + toJSON, + }; + const outbox = new MetadataOutbox('alice', { + storage, locks: locks as unknown as LockManager, generation: generations(), + }); + + await outbox.enqueueImportLinks({ source: importRequest.source, items: [item] }); + + expect(toJSON).not.toHaveBeenCalled(); + expect(outbox.records()).toEqual([{ + version: 1, + generation: 'generation-1', + kind: 'import', + state: 'queued', + source: importRequest.source, + item: importRequest.items[0], + }]); + expect(JSON.stringify([...storage.values])).not.toContain('ignored'); + expect(JSON.stringify([...storage.values])).not.toContain('poisoned'); + }); + + it.each([ + ['non-object request', null], + ['blank source', { ...importRequest, source: ' ' }], + ['non-array items', { ...importRequest, items: {} }], + ['too many items', { ...importRequest, items: Array(101).fill(importRequest.items[0]) }], + ['non-object item', { ...importRequest, items: [null] }], + ['blank field', { + ...importRequest, + items: [{ ...importRequest.items[0], link: ' ' }], + }], + ['line break', { + ...importRequest, + items: [{ ...importRequest.items[0], title: 'forged\nmetadata' }], + }], + ['carriage return', { + ...importRequest, + items: [{ ...importRequest.items[0], link: 'forged\rmetadata' }], + }], + ['oversized external key', { + ...importRequest, + items: [{ ...importRequest.items[0], external_key: 'é'.repeat(1025) }], + }], + ['oversized URL', { + ...importRequest, + items: [{ ...importRequest.items[0], url: 'é'.repeat(1025) }], + }], + ['oversized title', { + ...importRequest, + items: [{ ...importRequest.items[0], title: 'é'.repeat(251) }], + }], + ['non-string field', { + ...importRequest, + items: [{ ...importRequest.items[0], url: 7 }], + }], + ])('rejects invalid import metadata atomically: %s', async (_label, request) => { + const outbox = new MetadataOutbox('alice', { + storage, locks: locks as unknown as LockManager, generation: generations(), + }); + + await expect(outbox.enqueueImportLinks( + request as unknown as Parameters[0], + )).rejects.toThrow('invalid outbox'); + expect(storage.values.size).toBe(0); + }); + + it('sanitizes server-controlled error metadata before persisting it', async () => { + const malicious = `invalid\r\n${'x'.repeat(300)}`; + const statuses = vi.fn(); + const outbox = new MetadataOutbox('alice', { + storage, + locks: locks as unknown as LockManager, + generation: generations(), + sendImport: () => Promise.reject(Object.assign(new Error(malicious), { status: 400 })), + onStatus: statuses, + }); + + await outbox.enqueueImportLinks(importRequest); + await outbox.drain(identity); + + expect(outbox.records()[0]).toMatchObject({ + state: 'blocked', + error: expect.not.stringContaining('\n'), + }); + expect(outbox.records()[0]?.error).toHaveLength(200); + expect(statuses).toHaveBeenCalledWith({ kind: 'blocked', message: malicious }); + }); + + it('omits empty or non-string error metadata from durable records', async () => { + const emptyError = new MetadataOutbox('alice', { + storage, + locks: locks as unknown as LockManager, + generation: generations(), + sendImport: () => Promise.reject(Object.assign(new Error(''), { status: 400 })), + }); + await emptyError.enqueueImportLinks(importRequest); + await emptyError.drain(identity); + expect(emptyError.records()[0]).not.toHaveProperty('error'); + + const key = [...storage.values.keys()][0]!; + storage.setItem(key, JSON.stringify({ + ...JSON.parse(storage.getItem(key)!), + state: 'sending', + error: 7, + })); + await emptyError.reconcile(cancelled, new Map()); + expect(emptyError.records()[0]).toMatchObject({ state: 'retry' }); + expect(emptyError.records()[0]).not.toHaveProperty('error'); + }); + + it.each([ + ['blank task id', ' ', 'valid reason'], + ['line break in reason', 'client-1', 'forged\nmetadata'], + ['oversized reason', 'client-1', 'é'.repeat(1001)], + ])('rejects invalid tombstone metadata before storage: %s', async (_label, taskId, reason) => { + const outbox = new MetadataOutbox('alice', { + storage, locks: locks as unknown as LockManager, generation: generations(), + }); + + await expect(outbox.enqueueTombstone(taskId, reason)).rejects.toThrow('invalid outbox'); + expect(storage.values.size).toBe(0); + }); + it.each([ ['2xx', undefined], ['network', new TypeError('offline')], @@ -366,6 +494,88 @@ describe('MetadataOutbox', () => { ); }); + it('canonicalizes legacy records instead of copying attacker-controlled fields', () => { + const logicalKey = `import:${importRequest.items[0]!.external_key}`; + const legacyKey = `kb.outbox.v1.alice.${encodeURIComponent(logicalKey)}`; + storage.values.set(legacyKey, JSON.stringify({ + version: 1, + generation: 'legacy-generation', + kind: 'import', + state: 'queued', + source: importRequest.source, + item: { + ...importRequest.items[0], + ignored: 'attacker-controlled extra field', + toJSON: { poisoned: true }, + }, + ignored: 'attacker-controlled record field', + })); + + const outbox = new MetadataOutbox('alice', { + storage, locks: locks as unknown as LockManager, generation: generations(), + }); + + expect(outbox.records()).toEqual([{ + version: 1, + generation: 'legacy-generation', + kind: 'import', + state: 'queued', + source: importRequest.source, + item: importRequest.items[0], + }]); + const migrated = [...storage.values.entries()].find(([key]) => key !== legacyKey); + expect(migrated).toBeDefined(); + expect(migrated![1]).not.toContain('ignored'); + expect(migrated![1]).not.toContain('toJSON'); + expect(migrated![1]).not.toContain('poisoned'); + }); + + it.each([ + ['line break', { ...importRequest.items[0], title: 'forged\nmetadata' }], + ['oversized URL', { ...importRequest.items[0], url: 'é'.repeat(1025) }], + ])('retains but does not migrate invalid legacy metadata: %s', (_label, item) => { + const logicalKey = `import:${item.external_key}`; + const legacyKey = `kb.outbox.v1.alice.${encodeURIComponent(logicalKey)}`; + const raw = JSON.stringify({ + version: 1, + generation: 'legacy-generation', + kind: 'import', + state: 'queued', + source: importRequest.source, + item, + }); + storage.values.set(legacyKey, raw); + + const outbox = new MetadataOutbox('alice', { + storage, locks: locks as unknown as LockManager, generation: generations(), + }); + + expect(outbox.records()).toEqual([]); + expect(storage.values).toEqual(new Map([[legacyKey, raw]])); + }); + + it('retains a valid legacy record when canonical migration storage fails', () => { + const logicalKey = `import:${importRequest.items[0]!.external_key}`; + const legacyKey = `kb.outbox.v1.alice.${encodeURIComponent(logicalKey)}`; + const raw = JSON.stringify({ + version: 1, + generation: 'legacy-generation', + kind: 'import', + state: 'queued', + source: importRequest.source, + item: importRequest.items[0], + }); + storage.values.set(legacyKey, raw); + storage.failSet = true; + + const outbox = new MetadataOutbox('alice', { + storage, locks: locks as unknown as LockManager, generation: generations(), + }); + + expect(outbox.records()).toEqual([]); + expect(storage.values).toEqual(new Map([[legacyKey, raw]])); + }); + it('does not lose interleaved two-instance additions and removals', async () => { const first = new MetadataOutbox('alice', { storage, locks: locks as unknown as LockManager, generation: generations(), diff --git a/src/lib/outbox.ts b/src/lib/outbox.ts index b5c081a..643a74d 100644 --- a/src/lib/outbox.ts +++ b/src/lib/outbox.ts @@ -11,6 +11,11 @@ import { const PREFIX = 'kb.outbox.v1'; const LOCK_PREFIX = 'kb:outbox:'; +const MAX_IMPORT_ITEMS = 100; +const MAX_EXTERNAL_KEY_BYTES = 2048; +const MAX_URL_BYTES = 2048; +const MAX_TITLE_BYTES = 500; +const MAX_REASON_BYTES = 2000; type OutboxState = 'awaiting_canonical' | 'queued' | 'sending' | 'retry' | 'blocked'; @@ -152,6 +157,91 @@ function isRecord(value: unknown): value is Record { return typeof value === 'object' && value !== null && !Array.isArray(value); } +function validatedText( + value: unknown, + field: string, + maxBytes?: number, +): string { + if ( + typeof value !== 'string' || + value.trim() === '' || + value.includes('\r') || + value.includes('\n') || + (maxBytes !== undefined && new TextEncoder().encode(value).byteLength > maxBytes) + ) { + throw new TypeError(`invalid outbox ${field}`); + } + return value; +} + +function validatedImportRequest(req: RecordImportLinksRequest): RecordImportLinksRequest { + if (!isRecord(req)) throw new TypeError('invalid outbox import request'); + const source = validatedText(req.source, 'import source'); + if (!Array.isArray(req.items) || req.items.length > MAX_IMPORT_ITEMS) { + throw new TypeError('invalid outbox import items'); + } + const items = req.items.map((item) => { + if (!isRecord(item)) throw new TypeError('invalid outbox import item'); + return { + external_key: validatedText( + item.external_key, + 'import external key', + MAX_EXTERNAL_KEY_BYTES, + ), + link: validatedText(item.link, 'import link'), + url: validatedText(item.url, 'import URL', MAX_URL_BYTES), + title: validatedText(item.title, 'import title', MAX_TITLE_BYTES), + }; + }); + return { source, items }; +} + +function storedError(value: unknown): string | undefined { + if (typeof value !== 'string') return undefined; + const safe = value.replace(/[\r\n]+/g, ' ').slice(0, 200); + return safe === '' ? undefined : safe; +} + +function validatedStorageRecord(record: OutboxRecord): OutboxRecord { + const generation = validatedText(record.generation, 'generation'); + const error = storedError(record.error); + if (record.kind === 'import') { + const validated = validatedImportRequest({ source: record.source, items: [record.item] }); + return { + version: 1, + generation, + kind: 'import', + state: record.state, + source: validated.source, + item: validated.items[0]!, + ...(error === undefined ? {} : { error }), + }; + } + const clientTaskId = validatedText(record.clientTaskId, 'tombstone task ID'); + const reason = validatedText(record.reason, 'tombstone reason', MAX_REASON_BYTES); + if (record.state === 'awaiting_canonical') { + return { + version: 1, + generation, + kind: 'tombstone', + state: 'awaiting_canonical', + clientTaskId, + reason, + ...(error === undefined ? {} : { error }), + }; + } + return { + version: 1, + generation, + kind: 'tombstone', + state: record.state, + clientTaskId, + canonicalTaskId: validatedText(record.canonicalTaskId, 'canonical task ID'), + reason, + ...(error === undefined ? {} : { error }), + }; +} + function parseRecord(raw: string | null): OutboxRecord | null { if (raw === null) return null; try { @@ -244,8 +334,12 @@ export class MetadataOutbox { if (key !== expected) continue; const target = recordKey(this.ns, logicalKey(record)); if (this.storage.getItem(target) === null) { - const raw = this.storage.getItem(key); - if (raw !== null) this.storage.setItem(target, raw); + try { + this.write(target, record); + } catch (error) { + if (error instanceof TypeError) continue; + throw error; + } } } } catch { @@ -296,7 +390,7 @@ export class MetadataOutbox { } private write(key: string, record: OutboxRecord): void { - this.storage.setItem(key, JSON.stringify(record)); + this.storage.setItem(key, JSON.stringify(validatedStorageRecord(record))); } private fresh(record: NewOutboxRecord): OutboxRecord { @@ -309,12 +403,14 @@ export class MetadataOutbox { /** Persist the user's reason before the board PUT can acknowledge an ID. */ async enqueueTombstone(clientTaskId: string, reason: string): Promise { - const key = recordKey(this.ns, tombstoneLogicalKey(clientTaskId)); + const safeClientTaskId = validatedText(clientTaskId, 'tombstone task ID'); + const safeReason = validatedText(reason, 'tombstone reason', MAX_REASON_BYTES); + const key = recordKey(this.ns, tombstoneLogicalKey(safeClientTaskId)); const record = this.fresh({ kind: 'tombstone', state: 'awaiting_canonical', - clientTaskId, - reason, + clientTaskId: safeClientTaskId, + reason: safeReason, }); const written = await this.locked(() => this.write(key, record)); if (written === undefined && !this.locks?.request) { @@ -326,10 +422,13 @@ export class MetadataOutbox { } async enqueueImportLinks(req: RecordImportLinksRequest): Promise { + const validated = validatedImportRequest(req); const writeAll = () => { - for (const item of req.items) { + for (const item of validated.items) { const key = recordKey(this.ns, importLogicalKey(item.external_key)); - this.write(key, this.fresh({ kind: 'import', state: 'queued', source: req.source, item })); + this.write(key, this.fresh({ + kind: 'import', state: 'queued', source: validated.source, item, + })); } }; const written = await this.locked(writeAll); diff --git a/src/lib/remote.test.ts b/src/lib/remote.test.ts index f382396..6a99b9a 100644 --- a/src/lib/remote.test.ts +++ b/src/lib/remote.test.ts @@ -2463,6 +2463,31 @@ describe('RemoteStore concurrency', () => { })).toBe(false); }); + it.each([ + ['mixed-case', ['server-A', 'server-a', 'Server-b']], + ['numeric-like', ['server-2', 'server-10', 'server-01']], + ['non-ASCII', ['server-ä', 'server-Ω', 'server-😀']], + ])('compares %s deleted canonical IDs independently of insertion order', (_kind, ids) => { + const value = board('unchanged'); + const canonicalTaskIDs = new Map([[value.tasks[0]!.id, 'server-live']]); + const base = { + board: value, + canonicalTaskIDs, + deletedCanonicalIDs: new Set(ids), + migratedRaw: false, + pendingBoardWrite: null, + }; + + expect(sameBoardSemantics(base, { + ...base, + deletedCanonicalIDs: new Set([...ids].reverse()), + })).toBe(true); + expect(sameBoardSemantics(base, { + ...base, + deletedCanonicalIDs: new Set([...ids.slice(0, -1), `${ids.at(-1)}-different`]), + })).toBe(false); + }); + it('exposes a current-epoch guard across an awaited success callback', async () => { const store = new RemoteStore(); const release = deferred(); diff --git a/src/lib/remote.ts b/src/lib/remote.ts index 62fa1ff..d3dbda1 100644 --- a/src/lib/remote.ts +++ b/src/lib/remote.ts @@ -195,10 +195,8 @@ export function sameBoardSemantics( canonicalSequence(current.board, current.canonicalTaskIDs), canonicalSequence(target.board, target.canonicalTaskIDs), ) && - same( - [...current.deletedCanonicalIDs].sort(), - [...target.deletedCanonicalIDs].sort(), - ) && + current.deletedCanonicalIDs.size === target.deletedCanonicalIDs.size && + [...current.deletedCanonicalIDs].every((id) => target.deletedCanonicalIDs.has(id)) && current.migratedRaw === target.migratedRaw && same(current.pendingBoardWrite, target.pendingBoardWrite) );