From ff5e4b831341d4d0a2bec970032bec5bed135280 Mon Sep 17 00:00:00 2001 From: "extro.rook.host" Date: Mon, 27 Jul 2026 17:35:14 -0600 Subject: [PATCH] fix(explore): tolerate malformed external beacons and prune empty aggregates Normalize indexed cap and vouch beacons plus API filters so malformed external values are isolated, warned once, and never create invalid aggregate identities. The vouch path previously bound record.beacon raw into the indexed column while aggregates used the guarded value. Object and array beacons therefore threw at the D1 bind and stalled the Jetstream window. Prune zero-count aggregates transactionally and share the real-SQLite D1 adapter for ingestion and API coverage. --- explore/src/api.js | 40 ++- explore/src/jetstream.js | 72 +++- ...xplore-canonical-beacon-api-260727.test.js | 183 ++++++++++ ...ore-canonical-beacon-ingest-260727.test.js | 339 ++++++++++++++++++ test/explore-cursor.test.js | 60 +--- test/explore-d1.js | 57 +++ 6 files changed, 673 insertions(+), 78 deletions(-) create mode 100644 test/explore-canonical-beacon-api-260727.test.js create mode 100644 test/explore-canonical-beacon-ingest-260727.test.js create mode 100644 test/explore-d1.js diff --git a/explore/src/api.js b/explore/src/api.js index 5a6b03b..f465cd5 100644 --- a/explore/src/api.js +++ b/explore/src/api.js @@ -1,6 +1,8 @@ // SPDX-License-Identifier: MIT // Copyright (c) 2026 sol pbc +import { BEACON_ACCEPTED_FORMS, tryNormalizeBeacon } from '../../src/lib/beacon.js'; + const CORS_HEADERS = { 'Access-Control-Allow-Origin': '*', 'Access-Control-Allow-Methods': 'GET, OPTIONS', @@ -31,6 +33,26 @@ function parseCursor(value) { return Number.isFinite(parsed) && parsed > 0 ? parsed : null; } +function parseBeaconFilter(value) { + if (!value) { + return []; + } + + const beacons = []; + const seen = new Set(); + for (const member of value.split(',')) { + const beacon = tryNormalizeBeacon(member.trim()); + if (!beacon) { + return null; + } + if (!seen.has(beacon)) { + seen.add(beacon); + beacons.push(beacon); + } + } + return beacons; +} + export async function handleRequest(request, env) { if (request.method === 'OPTIONS') { return new Response(null, { status: 204, headers: CORS_HEADERS }); @@ -46,15 +68,18 @@ export async function handleRequest(request, env) { if (pathname === '/api/caps') { const cursor = parseCursor(searchParams.get('cursor')); const limit = parseLimit(searchParams.get('limit')); - const beacon = searchParams.get('beacon'); + const beacons = parseBeaconFilter(searchParams.get('beacon')); const kind = searchParams.get('kind'); const sort = searchParams.get('sort'); + if (beacons === null) { + return json({ error: `invalid beacon filter: expected ${BEACON_ACCEPTED_FORMS}.` }, 400); + } + const conditions = []; const bindings = []; - if (beacon) { - const beacons = beacon.split(',').filter(Boolean); + if (beacons.length > 0) { const placeholders = beacons.map(() => '?').join(', '); conditions.push(`c.beacon IN (${placeholders})`); bindings.push(...beacons); @@ -94,7 +119,7 @@ export async function handleRequest(request, env) { if (pathname === '/api/cap') { const ref = searchParams.get('ref'); const uri = searchParams.get('uri'); - const beacon = searchParams.get('beacon'); + const beacons = parseBeaconFilter(searchParams.get('beacon')); if (!ref && !uri) { return json({ error: 'ref or uri is required' }, 400); @@ -104,6 +129,10 @@ export async function handleRequest(request, env) { return json({ error: 'provide ref or uri, not both' }, 400); } + if (beacons === null) { + return json({ error: `invalid beacon filter: expected ${BEACON_ACCEPTED_FORMS}.` }, 400); + } + const conditions = []; const bindings = []; @@ -116,8 +145,7 @@ export async function handleRequest(request, env) { conditions.push('c.ref = ?'); bindings.push(ref); - if (beacon) { - const beacons = beacon.split(',').filter(Boolean); + if (beacons.length > 0) { const placeholders = beacons.map(() => '?').join(', '); conditions.push(`c.beacon IN (${placeholders})`); bindings.push(...beacons); diff --git a/explore/src/jetstream.js b/explore/src/jetstream.js index 823106f..d8a6071 100644 --- a/explore/src/jetstream.js +++ b/explore/src/jetstream.js @@ -2,6 +2,7 @@ // Copyright (c) 2026 sol pbc import { resolveHandles } from './resolve.js'; +import { tryNormalizeBeacon } from '../../src/lib/beacon.js'; const CAP_COLLECTION = 'org.v-it.cap'; const VOUCH_COLLECTION = 'org.v-it.vouch'; @@ -14,9 +15,36 @@ const SKILL_COLLECTION = 'org.v-it.skill'; // live-tail (see index.js scheduled handler). const JETSTREAM_URL = 'wss://jetstream1.us-east.bsky.network/subscribe'; const STREAM_DURATION_MS = 55_000; +const BEACON_WARNING_PREFIX = 'explore: beacon normalization failed class=invalid-beacon record='; +const MAX_BEACON_WARNING_BYTES = 512; +const TRUNCATION_MARKER = '...'; + +function formatBeaconWarning(uri) { + let renderedUri; + try { + renderedUri = encodeURI(String(uri)); + } catch { + renderedUri = '[invalid-uri]'; + } + + const available = MAX_BEACON_WARNING_BYTES - BEACON_WARNING_PREFIX.length; + if (renderedUri.length <= available) { + return BEACON_WARNING_PREFIX + renderedUri; + } -function beaconValue(value) { - return typeof value === 'string' && value.length > 0 ? value : null; + let truncated = renderedUri.slice(0, available - TRUNCATION_MARKER.length); + while (/%[0-9A-F]?$/i.test(truncated)) { + truncated = truncated.slice(0, -1); + } + return BEACON_WARNING_PREFIX + truncated + TRUNCATION_MARKER; +} + +function normalizeRecordBeacon(record, uri) { + const beacon = tryNormalizeBeacon(record?.beacon); + if (record?.beacon != null && beacon === null) { + console.warn(formatBeaconWarning(uri)); + } + return beacon; } function enqueueRecordTask(recordTasks, key, task) { @@ -73,16 +101,22 @@ function decrementVouchBeaconStatement(env, beacon) { ).bind(beacon); } +function deleteEmptyBeaconStatement(env, beacon) { + return env.DB.prepare( + 'DELETE FROM beacons WHERE name = ? AND cap_count = 0 AND vouch_count = 0', + ).bind(beacon); +} + export async function processCapEvent(env, did, commit) { const { operation, rkey, record, cid } = commit; const uri = `at://${did}/${CAP_COLLECTION}/${rkey}`; if (operation === 'create' || operation === 'update') { - const nextBeacon = beaconValue(record?.beacon); + const nextBeacon = normalizeRecordBeacon(record, uri); const existing = await env.DB.prepare('SELECT beacon FROM caps WHERE did = ? AND rkey = ?') .bind(did, rkey) .first(); - const prevBeacon = beaconValue(existing?.beacon); + const prevBeacon = tryNormalizeBeacon(existing?.beacon); const capKind = typeof record?.kind === 'string' && record.kind.length > 0 ? record.kind : null; @@ -119,6 +153,7 @@ export async function processCapEvent(env, did, commit) { } else if (existing && prevBeacon !== nextBeacon) { if (prevBeacon) { stmts.push(decrementCapBeaconStatement(env, prevBeacon)); + stmts.push(deleteEmptyBeaconStatement(env, prevBeacon)); } if (nextBeacon) { stmts.push(...incrementCapBeaconStatements(env, nextBeacon)); @@ -134,14 +169,13 @@ export async function processCapEvent(env, did, commit) { .bind(did, rkey) .first(); - const stmts = [ - env.DB.prepare('DELETE FROM caps WHERE did = ? AND rkey = ?').bind(did, rkey), - ]; - - const prevBeacon = beaconValue(existing?.beacon); + const stmts = []; + const prevBeacon = tryNormalizeBeacon(existing?.beacon); if (prevBeacon) { - stmts.unshift(decrementCapBeaconStatement(env, prevBeacon)); + stmts.push(decrementCapBeaconStatement(env, prevBeacon)); + stmts.push(deleteEmptyBeaconStatement(env, prevBeacon)); } + stmts.push(env.DB.prepare('DELETE FROM caps WHERE did = ? AND rkey = ?').bind(did, rkey)); await env.DB.batch(stmts); } @@ -152,11 +186,11 @@ export async function processVouchEvent(env, did, commit) { const uri = `at://${did}/${VOUCH_COLLECTION}/${rkey}`; if (operation === 'create' || operation === 'update') { - const nextBeacon = beaconValue(record?.beacon); + const nextBeacon = normalizeRecordBeacon(record, uri); const existing = await env.DB.prepare('SELECT beacon FROM vouches WHERE did = ? AND rkey = ?') .bind(did, rkey) .first(); - const prevBeacon = beaconValue(existing?.beacon); + const prevBeacon = tryNormalizeBeacon(existing?.beacon); const vouchKind = typeof record?.kind === 'string' && record.kind.length > 0 ? record.kind : 'endorse'; @@ -179,7 +213,7 @@ export async function processVouchEvent(env, did, commit) { cid ?? null, record.subject?.uri, record.ref, - record.beacon ?? null, + nextBeacon, vouchKind, JSON.stringify(record), record.createdAt, @@ -191,6 +225,7 @@ export async function processVouchEvent(env, did, commit) { } else if (existing && prevBeacon !== nextBeacon) { if (prevBeacon) { stmts.push(decrementVouchBeaconStatement(env, prevBeacon)); + stmts.push(deleteEmptyBeaconStatement(env, prevBeacon)); } if (nextBeacon) { stmts.push(...incrementVouchBeaconStatements(env, nextBeacon)); @@ -206,14 +241,13 @@ export async function processVouchEvent(env, did, commit) { .bind(did, rkey) .first(); - const stmts = [ - env.DB.prepare('DELETE FROM vouches WHERE did = ? AND rkey = ?').bind(did, rkey), - ]; - - const prevBeacon = beaconValue(existing?.beacon); + const stmts = []; + const prevBeacon = tryNormalizeBeacon(existing?.beacon); if (prevBeacon) { - stmts.unshift(decrementVouchBeaconStatement(env, prevBeacon)); + stmts.push(decrementVouchBeaconStatement(env, prevBeacon)); + stmts.push(deleteEmptyBeaconStatement(env, prevBeacon)); } + stmts.push(env.DB.prepare('DELETE FROM vouches WHERE did = ? AND rkey = ?').bind(did, rkey)); await env.DB.batch(stmts); } diff --git a/test/explore-canonical-beacon-api-260727.test.js b/test/explore-canonical-beacon-api-260727.test.js new file mode 100644 index 0000000..83b8a96 --- /dev/null +++ b/test/explore-canonical-beacon-api-260727.test.js @@ -0,0 +1,183 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 sol pbc + +import { describe, expect, test } from 'bun:test'; +import { handleRequest } from '../explore/src/api.js'; +import { processCapEvent } from '../explore/src/jetstream.js'; +import { createSqliteEnv } from './explore-d1.js'; + +const CANONICAL = 'vit:github.com/solpbc/thermals'; +const HTTPS_ALIAS = 'https://github.com/solpbc/thermals'; +const OTHER = 'vit:tangled.org/solpbc.org/rookery'; +const ERROR_BODY = { + error: 'invalid beacon filter: expected vit:host/owner/repo or a git URL (a scheme URL such as https://host/owner/repo, ssh://git@host/owner/repo, or git://host/owner/repo; SCP-style git@host:owner/repo; or host/owner/repo).', +}; + +function capCommit(rkey, beacon, createdAt, ref = 'shared-api-ref') { + return { + operation: 'create', + rkey, + cid: `cid-${rkey}`, + record: { + title: `Cap ${rkey}`, + description: `Description ${rkey}`, + ref, + beacon, + kind: 'test', + createdAt, + }, + }; +} + +function apiRequest(pathname, params = {}) { + const url = new URL(`https://explore.example${pathname}`); + for (const [name, value] of Object.entries(params)) { + url.searchParams.set(name, value); + } + return new Request(url); +} + +async function responseJson(pathname, params, env) { + const response = await handleRequest(apiRequest(pathname, params), env); + return { response, body: await response.json() }; +} + +async function seedCaps(env) { + await processCapEvent( + env, + 'did:plc:api-one', + capCommit('api-one', HTTPS_ALIAS, '2026-07-27T10:00:00.000Z'), + ); + await processCapEvent( + env, + 'did:plc:api-two', + capCommit('api-two', CANONICAL, '2026-07-27T11:00:00.000Z'), + ); + await processCapEvent( + env, + 'did:plc:api-three', + capCommit('api-three', OTHER, '2026-07-27T12:00:00.000Z', 'other-api-ref'), + ); +} + +describe('Explore canonical beacon API filters', () => { + test('canonical and HTTPS filters return the same pinned cap rows', async () => { + const { db, env } = createSqliteEnv(); + try { + await seedCaps(env); + + const canonicalCaps = await responseJson('/api/caps', { beacon: CANONICAL }, env); + const aliasCaps = await responseJson('/api/caps', { beacon: HTTPS_ALIAS }, env); + expect(canonicalCaps.response.status).toBe(200); + expect(aliasCaps.response.status).toBe(200); + expect(canonicalCaps.body.caps.map((cap) => cap.id)).toEqual([2, 1]); + expect(aliasCaps.body.caps.map((cap) => cap.id)).toEqual([2, 1]); + expect(aliasCaps.body).toEqual(canonicalCaps.body); + expect(aliasCaps.body.caps.every((cap) => cap.beacon === CANONICAL)).toBe(true); + + const canonicalCap = await responseJson( + '/api/cap', + { ref: 'shared-api-ref', beacon: CANONICAL }, + env, + ); + const aliasCap = await responseJson( + '/api/cap', + { ref: 'shared-api-ref', beacon: HTTPS_ALIAS }, + env, + ); + expect(canonicalCap.response.status).toBe(200); + expect(aliasCap.response.status).toBe(200); + expect(canonicalCap.body.cap.id).toBe(2); + expect(aliasCap.body).toEqual(canonicalCap.body); + expect(aliasCap.body.cap.beacon).toBe(CANONICAL); + } finally { + db.close(); + } + }); + + test('trims and deduplicates comma-separated aliases into one predicate', async () => { + const { db, env } = createSqliteEnv(); + try { + await seedCaps(env); + + const statements = []; + const prepare = env.DB.prepare.bind(env.DB); + env.DB.prepare = (sql) => { + const statement = prepare(sql); + return { + bind(...args) { + statements.push({ sql, args }); + return statement.bind(...args); + }, + first() { + statements.push({ sql, args: [] }); + return statement.first(); + }, + all() { + statements.push({ sql, args: [] }); + return statement.all(); + }, + }; + }; + + const filter = ` ${HTTPS_ALIAS} , ${CANONICAL} , ssh://git@github.com/solpbc/thermals.git `; + const result = await responseJson('/api/caps', { beacon: filter }, env); + + expect(result.response.status).toBe(200); + expect(result.body.caps.map((cap) => cap.id)).toEqual([2, 1]); + expect(statements).toHaveLength(1); + expect(statements[0].sql).toContain('c.beacon IN (?)'); + expect(statements[0].args).toEqual([CANONICAL, 50]); + } finally { + db.close(); + } + }); + + test('rejects an invalid member in a mixed filter before querying D1', async () => { + const { db, env } = createSqliteEnv(); + try { + await seedCaps(env); + env.DB.prepare = () => { + throw new Error('D1 must not be queried for an invalid beacon filter'); + }; + + const mixed = `${CANONICAL},not a url`; + const caps = await responseJson('/api/caps', { beacon: mixed }, env); + expect(caps.response.status).toBe(400); + expect(caps.body).toEqual(ERROR_BODY); + + const cap = await responseJson( + '/api/cap', + { ref: 'shared-api-ref', beacon: mixed }, + env, + ); + expect(cap.response.status).toBe(400); + expect(cap.body).toEqual(ERROR_BODY); + } finally { + db.close(); + } + }); + + test('keeps an empty beacon parameter equivalent to no filter', async () => { + const { db, env } = createSqliteEnv(); + try { + await seedCaps(env); + + const emptyCaps = await responseJson('/api/caps', { beacon: '' }, env); + const allCaps = await responseJson('/api/caps', {}, env); + expect(emptyCaps.response.status).toBe(200); + expect(emptyCaps.body).toEqual(allCaps.body); + expect(emptyCaps.body.caps.map((cap) => cap.id)).toEqual([3, 2, 1]); + + const emptyCap = await responseJson( + '/api/cap', + { ref: 'shared-api-ref', beacon: '' }, + env, + ); + expect(emptyCap.response.status).toBe(200); + expect(emptyCap.body.cap.id).toBe(2); + } finally { + db.close(); + } + }); +}); diff --git a/test/explore-canonical-beacon-ingest-260727.test.js b/test/explore-canonical-beacon-ingest-260727.test.js new file mode 100644 index 0000000..0465cb0 --- /dev/null +++ b/test/explore-canonical-beacon-ingest-260727.test.js @@ -0,0 +1,339 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 sol pbc + +import { describe, expect, spyOn, test } from 'bun:test'; +import { processCapEvent, processVouchEvent } from '../explore/src/jetstream.js'; +import { createSqliteEnv } from './explore-d1.js'; + +const PROJECT_A = 'vit:github.com/solpbc/thermals'; +const PROJECT_A_ALIAS = 'https://github.com/solpbc/thermals'; +const PROJECT_B = 'vit:tangled.org/solpbc.org/rookery'; +const INVALID_BEACON = 'not a url'; + +function capCommit(operation, rkey, beacon, overrides = {}) { + if (operation === 'delete') { + return { operation, rkey }; + } + + return { + operation, + rkey, + cid: overrides.cid ?? `cid-${rkey}`, + record: { + title: overrides.title ?? `Cap ${rkey}`, + description: overrides.description ?? `Description ${rkey}`, + ref: overrides.ref ?? `ref-${rkey}`, + beacon, + kind: overrides.kind ?? 'test', + createdAt: overrides.createdAt ?? '2026-07-27T12:00:00.000Z', + }, + }; +} + +function vouchCommit(operation, rkey, beacon, overrides = {}) { + if (operation === 'delete') { + return { operation, rkey }; + } + + return { + operation, + rkey, + cid: overrides.cid ?? `cid-${rkey}`, + record: { + subject: { + uri: overrides.capUri ?? 'at://did:plc:author/org.v-it.cap/source', + }, + ref: overrides.ref ?? `ref-${rkey}`, + beacon, + kind: overrides.kind ?? 'endorse', + createdAt: overrides.createdAt ?? '2026-07-27T12:00:01.000Z', + }, + }; +} + +function aggregate(db, name) { + return db + .query('SELECT name, cap_count, vouch_count, last_activity FROM beacons WHERE name = ?') + .get(name); +} + +function expectConsistentAggregates(db) { + const rows = db.query('SELECT name, cap_count, vouch_count FROM beacons ORDER BY name').all(); + for (const row of rows) { + const counts = db.query( + `SELECT + (SELECT COUNT(*) FROM caps WHERE beacon = ?) AS cap_count, + (SELECT COUNT(*) FROM vouches WHERE beacon = ?) AS vouch_count`, + ).get(row.name, row.name); + expect(row.cap_count).toBe(counts.cap_count); + expect(row.vouch_count).toBe(counts.vouch_count); + expect(row.cap_count).toBeGreaterThanOrEqual(0); + expect(row.vouch_count).toBeGreaterThanOrEqual(0); + expect(row.cap_count + row.vouch_count).toBeGreaterThan(0); + } + + const missing = db.query( + `SELECT beacon FROM ( + SELECT beacon FROM caps WHERE beacon IS NOT NULL + UNION + SELECT beacon FROM vouches WHERE beacon IS NOT NULL + ) + WHERE beacon NOT IN (SELECT name FROM beacons)`, + ).all(); + expect(missing).toEqual([]); +} + +function expectIndexedBeacon(db, table, did, rkey, expected) { + const row = db.query(`SELECT beacon FROM ${table} WHERE did = ? AND rkey = ?`).get(did, rkey); + expect(row?.beacon ?? null).toBe(expected); +} + +describe('Explore canonical beacon ingestion', () => { + test('canonicalizes cap and vouch aliases into one aggregate row', async () => { + const { db, env } = createSqliteEnv(); + try { + const capDid = 'did:plc:caps'; + const vouchDid = 'did:plc:vouches'; + const aliasCap = capCommit('create', 'cap-alias', PROJECT_A_ALIAS); + const canonicalCap = capCommit('create', 'cap-canonical', PROJECT_A); + const vouch = vouchCommit('create', 'vouch-canonical', PROJECT_A); + + await processCapEvent(env, capDid, aliasCap); + await processCapEvent(env, capDid, canonicalCap); + await processVouchEvent(env, vouchDid, vouch); + + expect(db.query('SELECT id, rkey, beacon, record_json FROM caps ORDER BY id').all()).toEqual([ + { + id: 1, + rkey: 'cap-alias', + beacon: PROJECT_A, + record_json: JSON.stringify(aliasCap.record), + }, + { + id: 2, + rkey: 'cap-canonical', + beacon: PROJECT_A, + record_json: JSON.stringify(canonicalCap.record), + }, + ]); + expect(db.query('SELECT id, rkey, beacon, record_json FROM vouches ORDER BY id').all()).toEqual([ + { + id: 1, + rkey: 'vouch-canonical', + beacon: PROJECT_A, + record_json: JSON.stringify(vouch.record), + }, + ]); + + const rows = db.query( + 'SELECT name, cap_count, vouch_count, last_activity FROM beacons ORDER BY id', + ).all(); + expect(rows).toHaveLength(1); + expect(rows[0]).toEqual({ + name: PROJECT_A, + cap_count: 2, + vouch_count: 1, + last_activity: db.query("SELECT datetime('now') AS value").get().value, + }); + expect(aggregate(db, PROJECT_A_ALIAS)).toBeNull(); + expectConsistentAggregates(db); + } finally { + db.close(); + } + }); + + test('keeps cap aggregates exact across moves, malformed updates, retries, and deletes', async () => { + const { db, env } = createSqliteEnv(); + const did = 'did:plc:cap-transitions'; + const rkey = 'cap-transition'; + try { + await processCapEvent(env, did, capCommit('create', rkey, PROJECT_A_ALIAS)); + expectIndexedBeacon(db, 'caps', did, rkey, PROJECT_A); + expectConsistentAggregates(db); + + await processCapEvent(env, did, capCommit('update', rkey, PROJECT_B)); + expectIndexedBeacon(db, 'caps', did, rkey, PROJECT_B); + expect(aggregate(db, PROJECT_A)).toBeNull(); + expectConsistentAggregates(db); + + await processCapEvent(env, did, capCommit('update', rkey, INVALID_BEACON)); + expectIndexedBeacon(db, 'caps', did, rkey, null); + expect(aggregate(db, PROJECT_B)).toBeNull(); + expectConsistentAggregates(db); + + const validAgain = capCommit('update', rkey, PROJECT_A); + await processCapEvent(env, did, validAgain); + expectIndexedBeacon(db, 'caps', did, rkey, PROJECT_A); + expectConsistentAggregates(db); + + await processCapEvent(env, did, validAgain); + expect(aggregate(db, PROJECT_A)).toMatchObject({ cap_count: 1, vouch_count: 0 }); + expectConsistentAggregates(db); + + await processCapEvent(env, did, capCommit('delete', rkey)); + expectIndexedBeacon(db, 'caps', did, rkey, null); + expect(aggregate(db, PROJECT_A)).toBeNull(); + expectConsistentAggregates(db); + + await processCapEvent(env, did, capCommit('delete', rkey)); + expect(aggregate(db, PROJECT_A)).toBeNull(); + expectConsistentAggregates(db); + } finally { + db.close(); + } + }); + + test('keeps vouch aggregates exact across moves, malformed updates, retries, and deletes', async () => { + const { db, env } = createSqliteEnv(); + const did = 'did:plc:vouch-transitions'; + const rkey = 'vouch-transition'; + try { + await processVouchEvent(env, did, vouchCommit('create', rkey, PROJECT_A_ALIAS)); + expectIndexedBeacon(db, 'vouches', did, rkey, PROJECT_A); + expectConsistentAggregates(db); + + await processVouchEvent(env, did, vouchCommit('update', rkey, PROJECT_B)); + expectIndexedBeacon(db, 'vouches', did, rkey, PROJECT_B); + expect(aggregate(db, PROJECT_A)).toBeNull(); + expectConsistentAggregates(db); + + await processVouchEvent(env, did, vouchCommit('update', rkey, INVALID_BEACON)); + expectIndexedBeacon(db, 'vouches', did, rkey, null); + expect(aggregate(db, PROJECT_B)).toBeNull(); + expectConsistentAggregates(db); + + const validAgain = vouchCommit('update', rkey, PROJECT_A); + await processVouchEvent(env, did, validAgain); + expectIndexedBeacon(db, 'vouches', did, rkey, PROJECT_A); + expectConsistentAggregates(db); + + await processVouchEvent(env, did, validAgain); + expect(aggregate(db, PROJECT_A)).toMatchObject({ cap_count: 0, vouch_count: 1 }); + expectConsistentAggregates(db); + + await processVouchEvent(env, did, vouchCommit('delete', rkey)); + expectIndexedBeacon(db, 'vouches', did, rkey, null); + expect(aggregate(db, PROJECT_A)).toBeNull(); + expectConsistentAggregates(db); + + await processVouchEvent(env, did, vouchCommit('delete', rkey)); + expect(aggregate(db, PROJECT_A)).toBeNull(); + expectConsistentAggregates(db); + } finally { + db.close(); + } + }); + + test('removes an aggregate only after both independent counts reach zero', async () => { + const { db, env } = createSqliteEnv(); + try { + await processCapEvent(env, 'did:plc:cap-a', capCommit('create', 'cap-a', PROJECT_A)); + await processVouchEvent(env, 'did:plc:vouch-a', vouchCommit('create', 'vouch-a', PROJECT_A)); + await processCapEvent(env, 'did:plc:cap-a', capCommit('delete', 'cap-a')); + expect(aggregate(db, PROJECT_A)).toMatchObject({ cap_count: 0, vouch_count: 1 }); + await processVouchEvent(env, 'did:plc:vouch-a', vouchCommit('delete', 'vouch-a')); + expect(aggregate(db, PROJECT_A)).toBeNull(); + + await processCapEvent(env, 'did:plc:cap-b', capCommit('create', 'cap-b', PROJECT_B)); + await processVouchEvent(env, 'did:plc:vouch-b', vouchCommit('create', 'vouch-b', PROJECT_B)); + await processVouchEvent(env, 'did:plc:vouch-b', vouchCommit('delete', 'vouch-b')); + expect(aggregate(db, PROJECT_B)).toMatchObject({ cap_count: 1, vouch_count: 0 }); + await processCapEvent(env, 'did:plc:cap-b', capCommit('delete', 'cap-b')); + expect(aggregate(db, PROJECT_B)).toBeNull(); + expectConsistentAggregates(db); + } finally { + db.close(); + } + }); + + test('stores hostile beacon fields only in record_json and emits one bounded warning', async () => { + const invalidCases = [ + { label: 'object', value: { hostile: 'object-beacon-secret' } }, + { label: 'array', value: ['array-beacon-secret'] }, + { label: 'whitespace', value: ' ' }, + { label: 'oversized', value: 'oversized-beacon-secret-'.repeat(600) }, + ]; + + for (const [caseIndex, invalidCase] of invalidCases.entries()) { + for (const kind of ['cap', 'vouch']) { + const { db, env } = createSqliteEnv(); + const did = `did:plc:${kind}-${invalidCase.label}`; + const rkey = invalidCase.label === 'oversized' ? `record-${'r'.repeat(900)}` : `record-${caseIndex}`; + const commit = kind === 'cap' + ? capCommit('create', rkey, invalidCase.value) + : vouchCommit('create', rkey, invalidCase.value); + const warn = spyOn(console, 'warn').mockImplementation(() => {}); + + try { + if (kind === 'cap') { + await processCapEvent(env, did, commit); + } else { + await processVouchEvent(env, did, commit); + } + + const table = kind === 'cap' ? 'caps' : 'vouches'; + const row = db + .query(`SELECT beacon, record_json FROM ${table} WHERE did = ? AND rkey = ?`) + .get(did, rkey); + expect(row).toEqual({ + beacon: null, + record_json: JSON.stringify(commit.record), + }); + expect(db.query('SELECT * FROM beacons').all()).toEqual([]); + expect(warn).toHaveBeenCalledTimes(1); + + const message = warn.mock.calls[0][0]; + expect(message).toStartWith( + 'explore: beacon normalization failed class=invalid-beacon record=at://', + ); + expect(new TextEncoder().encode(message).byteLength).toBeLessThanOrEqual(512); + expect(message).not.toContain(JSON.stringify(invalidCase.value)); + if (typeof invalidCase.value === 'string' && invalidCase.value.trim()) { + expect(message).not.toContain(invalidCase.value.slice(0, 32)); + } + expectConsistentAggregates(db); + } finally { + warn.mockRestore(); + db.close(); + } + } + } + }); + + test('malformed updates preserve source records while clearing indexed beacons', async () => { + for (const kind of ['cap', 'vouch']) { + const { db, env } = createSqliteEnv(); + const did = `did:plc:${kind}-malformed-update`; + const rkey = `${kind}-malformed-update`; + const valid = kind === 'cap' + ? capCommit('create', rkey, PROJECT_A) + : vouchCommit('create', rkey, PROJECT_A); + const malformed = kind === 'cap' + ? capCommit('update', rkey, INVALID_BEACON) + : vouchCommit('update', rkey, INVALID_BEACON); + const handler = kind === 'cap' ? processCapEvent : processVouchEvent; + const warn = spyOn(console, 'warn').mockImplementation(() => {}); + + try { + await handler(env, did, valid); + warn.mockClear(); + await handler(env, did, malformed); + + const table = kind === 'cap' ? 'caps' : 'vouches'; + expect(db.query(`SELECT beacon, record_json FROM ${table}`).get()).toEqual({ + beacon: null, + record_json: JSON.stringify(malformed.record), + }); + expect(db.query('SELECT * FROM beacons').all()).toEqual([]); + expect(warn).toHaveBeenCalledTimes(1); + const message = warn.mock.calls[0][0]; + expect(message).not.toContain(INVALID_BEACON); + expect(new TextEncoder().encode(message).byteLength).toBeLessThanOrEqual(512); + expectConsistentAggregates(db); + } finally { + warn.mockRestore(); + db.close(); + } + } + }); +}); diff --git a/test/explore-cursor.test.js b/test/explore-cursor.test.js index f5976b3..704a476 100644 --- a/test/explore-cursor.test.js +++ b/test/explore-cursor.test.js @@ -2,11 +2,9 @@ // Copyright (c) 2026 sol pbc import { describe, expect, test } from 'bun:test'; -import { readFileSync } from 'node:fs'; -import { join } from 'node:path'; -import { Database } from 'bun:sqlite'; import { runScheduled } from '../explore/src/index.js'; import { processCapEvent, processSkillEvent, processVouchEvent } from '../explore/src/jetstream.js'; +import { createSqliteEnv } from './explore-d1.js'; function createCursorEnv({ stored = '', @@ -68,50 +66,6 @@ function createStreamReader({ result = { observedCursor: null }, error = null } return streamReader; } -function d1Statement(db, sql, args = []) { - return { - sql, - args, - bind(...nextArgs) { - return d1Statement(db, sql, nextArgs); - }, - first() { - return db.query(sql).get(...args) ?? null; - }, - run() { - return db.query(sql).run(...args); - }, - all() { - return { results: db.query(sql).all(...args) }; - }, - }; -} - -function createD1(db) { - const executeBatch = db.transaction((statements) => statements.map((stmt) => { - if (/^\s*select\b/i.test(stmt.sql)) { - return stmt.all(); - } - return stmt.run(); - })); - - return { - prepare(sql) { - return d1Statement(db, sql); - }, - batch(statements) { - return executeBatch(statements); - }, - }; -} - -function createSqliteEnv() { - const db = new Database(':memory:'); - const schemaPath = join(import.meta.dir, '..', 'explore', 'schema.sql'); - db.exec(readFileSync(schemaPath, 'utf8')); - return { db, env: { DB: createD1(db) } }; -} - describe('explore scheduled cursor', () => { test('passes stored valid cursor and writes newer observed cursor', async () => { const { env, store } = createCursorEnv({ stored: '12345' }); @@ -270,7 +224,7 @@ describe('explore event idempotency', () => { title: 'Idempotent Cap', description: 'A cap inserted twice for testing.', ref: 'idempotent-cap-test', - beacon: 'vit:example/repo', + beacon: 'vit:example.com/org/repo', kind: 'test', createdAt: '2026-07-07T00:00:00.000Z', }, @@ -279,11 +233,11 @@ describe('explore event idempotency', () => { await processCapEvent(env, capDid, capCommit); const capAfterOne = db .query("SELECT (SELECT COUNT(*) FROM caps) AS rows, (SELECT cap_count FROM beacons WHERE name = ?) AS count") - .get('vit:example/repo'); + .get('vit:example.com/org/repo'); await processCapEvent(env, capDid, capCommit); const capAfterTwo = db .query("SELECT (SELECT COUNT(*) FROM caps) AS rows, (SELECT cap_count FROM beacons WHERE name = ?) AS count") - .get('vit:example/repo'); + .get('vit:example.com/org/repo'); expect(capAfterOne).toEqual({ rows: 1, count: 1 }); expect(capAfterTwo).toEqual(capAfterOne); @@ -296,7 +250,7 @@ describe('explore event idempotency', () => { record: { subject: { uri: 'at://did:plc:capauthor/org.v-it.cap/3lcap' }, ref: 'idempotent-cap-test', - beacon: 'vit:example/repo', + beacon: 'vit:example.com/org/repo', kind: 'want', createdAt: '2026-07-07T00:00:01.000Z', }, @@ -305,11 +259,11 @@ describe('explore event idempotency', () => { await processVouchEvent(env, vouchDid, vouchCommit); const vouchAfterOne = db .query("SELECT (SELECT COUNT(*) FROM vouches) AS rows, (SELECT vouch_count FROM beacons WHERE name = ?) AS count") - .get('vit:example/repo'); + .get('vit:example.com/org/repo'); await processVouchEvent(env, vouchDid, vouchCommit); const vouchAfterTwo = db .query("SELECT (SELECT COUNT(*) FROM vouches) AS rows, (SELECT vouch_count FROM beacons WHERE name = ?) AS count") - .get('vit:example/repo'); + .get('vit:example.com/org/repo'); expect(vouchAfterOne).toEqual({ rows: 1, count: 1 }); expect(vouchAfterTwo).toEqual(vouchAfterOne); diff --git a/test/explore-d1.js b/test/explore-d1.js new file mode 100644 index 0000000..eeeb557 --- /dev/null +++ b/test/explore-d1.js @@ -0,0 +1,57 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 sol pbc + +import { readFileSync } from 'node:fs'; +import { join } from 'node:path'; +import { Database } from 'bun:sqlite'; + +function d1Statement(db, sql, args = []) { + return { + sql, + args, + bind(...nextArgs) { + return d1Statement(db, sql, nextArgs); + }, + first() { + return db.query(sql).get(...args) ?? null; + }, + run() { + return db.query(sql).run(...args); + }, + all() { + return { results: db.query(sql).all(...args) }; + }, + }; +} + +export function createD1(db) { + const executeBatch = db.transaction((statements) => statements.map((stmt) => { + if (/^\s*select\b/i.test(stmt.sql)) { + return stmt.all(); + } + return stmt.run(); + })); + + return { + prepare(sql) { + return d1Statement(db, sql); + }, + batch(statements) { + return executeBatch(statements); + }, + }; +} + +export function createSqliteEnv() { + const db = new Database(':memory:'); + const schemaPath = join(import.meta.dir, '..', 'explore', 'schema.sql'); + db.exec(readFileSync(schemaPath, 'utf8')); + return { db, env: { DB: createD1(db) } }; +} + +export function splitSqlStatements(text) { + return text + .split(';') + .map((statement) => statement.trim()) + .filter(Boolean); +} -- 2.51.2