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); +}