Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
210 changes: 210 additions & 0 deletions src/lib/outbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<MetadataOutbox['enqueueImportLinks']>[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')],
Expand Down Expand Up @@ -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(),
Expand Down
115 changes: 107 additions & 8 deletions src/lib/outbox.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,11 @@

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';

Expand Down Expand Up @@ -152,6 +157,91 @@
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 {
Expand Down Expand Up @@ -229,7 +319,7 @@
* Copy only records whose payload proves the exact old key. Prefix matching
* alone cannot distinguish namespace `alice` from `alice.work`.
*/
private migrateLegacyRecords(): void {

Check failure on line 322 in src/lib/outbox.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 18 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=RandomCodeSpace_kb&issues=AZ-9uKBU1vMCp7SkRG7F&open=AZ-9uKBU1vMCp7SkRG7F&pullRequest=18
try {
const candidates: string[] = [];
for (let i = 0; i < this.storage.length; i += 1) {
Expand All @@ -244,8 +334,12 @@
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 {
Expand Down Expand Up @@ -296,7 +390,7 @@
}

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 {
Expand All @@ -309,12 +403,14 @@

/** Persist the user's reason before the board PUT can acknowledge an ID. */
async enqueueTombstone(clientTaskId: string, reason: string): Promise<boolean> {
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) {
Expand All @@ -326,10 +422,13 @@
}

async enqueueImportLinks(req: RecordImportLinksRequest): Promise<boolean> {
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);
Expand Down
Loading