Experimental Bluesky client for agents
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803import { Database } from 'bun:sqlite';import { createHash, randomUUID } from 'node:crypto';import { dirname } from 'node:path';import { mkdirSync } from 'node:fs';import { targetKey, type ReadingView, type SessionInfo, type Snapshot, type Target,} from './model.ts';import type { FrozenAttachment, PreparedPublication, PublicationClaim, PublicationPlan, PublicationPostProgress, PublicationPostState, PublicationReceipt, PublicationRecordInspection, StrongReference,} from './social.ts';
interface SessionRow { name: string; current_view: string | null; history_json: string; revision: number;}
interface ViewRow { view_json: string;}
interface PublicationPlanRow { plan_json: string;}
interface PublicationPostRow { post_index: number; rkey: string; state: PublicationPostState; record_json: string | null; uri: string | null; cid: string | null; attempted_at: string | null; confirmed_at: string | null; observed_at: string | null; last_error: string | null; claim_token: string | null; claimed_at: string | null;}
interface PublicationAttachmentRow { attachment_id: string; metadata_json: string; bytes: Uint8Array;}
export interface SessionPosition { session: SessionInfo; revision: number; view: ReadingView | null;}
export interface ResolvedReference { view: ReadingView; ref: string; target: Target;}
export class StoreConflictError extends Error { constructor(session: string) { super(`Session "${session}" changed while this operation was in progress; retry it.`); this.name = 'StoreConflictError'; }}
export class UnknownReferenceError extends Error { constructor(ref: string) { super(`Unknown reference: ${ref}`); this.name = 'UnknownReferenceError'; }}
function freezeDeep<T>(value: T): T { if (value !== null && typeof value === 'object' && !Object.isFrozen(value)) { for (const key in value) { freezeDeep((value as Record<string, unknown>)[key]); } Object.freeze(value); } return value;}
function parseView(json: string): ReadingView { return freezeDeep(JSON.parse(json) as ReadingView);}
function validateSessionName(name: string): void { if (name.length === 0 || name !== name.trim() || name.length > 128 || /[\u0000-\u001f\u007f]/u.test(name)) { throw new Error('Session names must be 1-128 characters, without surrounding whitespace or control characters.'); }}
function buildReferences(viewId: string, snapshot: Snapshot): Record<string, Target> { const refs: Record<string, Target> = {}; const byTarget = new Map<string, string>(); const counters: Record<string, number> = {};
const add = (role: string, target: Target): void => { const key = targetKey(target); if (byTarget.has(key)) return; const number = (counters[role] ?? 0) + 1; counters[role] = number; const ref = `${viewId}-${role}${number}`; refs[ref] = target; byTarget.set(key, ref); };
if (snapshot.kind === 'thread') { for (const node of snapshot.nodes) { add('p', { kind: 'post', value: node.uri, label: node.author ? `Post by @${node.author.handle}` : 'Post', }); if (node.author) { add('a', { kind: 'profile', value: node.author.did, label: node.author.displayName || `@${node.author.handle}`, }); } for (const facet of node.facets) { for (const target of facet.targets) add('f', { ...target }); } for (let index = 0; index < node.quotes.length; index += 1) { add('q', { kind: 'post', value: node.quotes[index], label: 'Quoted post', ownerUri: node.uri, index, }); } for (let index = 0; index < node.cards.length; index += 1) { const card = node.cards[index]; add('c', { kind: 'link', value: card.uri, label: card.title || card.uri, ownerUri: node.uri, index, }); } for (let index = 0; index < node.media.length; index += 1) { const media = node.media[index]; add('m', { kind: 'media', value: media.url ?? media.thumbnail ?? `${node.uri}#${media.kind}-${index}`, label: media.kind === 'video' ? 'Video media' : 'Image media', ownerUri: node.uri, index, }); } } // Allocate edge-only posts after hydrated nodes so existing post ref numbering stays stable. for (const node of snapshot.nodes) { if (node.parentUri) add('p', { kind: 'post', value: node.parentUri, label: 'Known parent post' }); if (node.rootUri) add('p', { kind: 'post', value: node.rootUri, label: 'Known root post' }); for (const replyUri of node.replies) { add('p', { kind: 'post', value: replyUri, label: 'Known reply post' }); } } } else if (snapshot.kind === 'feed') { const nodesByUri = new Map(snapshot.nodes.map((node) => [node.uri, node])); for (const entry of snapshot.entries) { const node = nodesByUri.get(entry.postUri); add('e', { kind: 'occurrence', value: entry.id, label: node?.author ? `Feed occurrence ${entry.id}: post by @${node.author.handle}` : `Feed occurrence ${entry.id}`, }); } for (const node of snapshot.nodes) { add('p', { kind: 'post', value: node.uri, label: node.author ? `Post by @${node.author.handle}` : 'Post', }); if (node.author) { add('a', { kind: 'profile', value: node.author.did, label: node.author.displayName || `@${node.author.handle}`, }); } for (const facet of node.facets) { for (const target of facet.targets) add('f', { ...target }); } for (let index = 0; index < node.quotes.length; index += 1) { add('q', { kind: 'post', value: node.quotes[index], label: 'Quoted post', ownerUri: node.uri, index, }); } for (let index = 0; index < node.cards.length; index += 1) { const card = node.cards[index]; add('c', { kind: 'link', value: card.uri, label: card.title || card.uri, ownerUri: node.uri, index, }); } for (let index = 0; index < node.media.length; index += 1) { const media = node.media[index]; add('m', { kind: 'media', value: media.url ?? media.thumbnail ?? `${node.uri}#${media.kind}-${index}`, label: media.kind === 'video' ? 'Video media' : 'Image media', ownerUri: node.uri, index, }); } } for (const entry of snapshot.entries) { const actor = entry.reason?.by; if (actor) { add('a', { kind: 'profile', value: actor.did, label: actor.displayName || `@${actor.handle}`, }); } } if (snapshot.actor) { add('a', { kind: 'profile', value: snapshot.actor.did, label: snapshot.actor.displayName || `@${snapshot.actor.handle}`, }); } for (const node of snapshot.nodes) { if (node.parentUri) add('p', { kind: 'post', value: node.parentUri, label: 'Known parent post' }); if (node.rootUri) add('p', { kind: 'post', value: node.rootUri, label: 'Known root post' }); for (const replyUri of node.replies) { add('p', { kind: 'post', value: replyUri, label: 'Known reply post' }); } } for (const entry of snapshot.entries) { if (entry.parentUri) add('p', { kind: 'post', value: entry.parentUri, label: 'Known parent post' }); if (entry.rootUri) add('p', { kind: 'post', value: entry.rootUri, label: 'Known root post' }); } } else if (snapshot.kind === 'profile') { add('a', { kind: 'profile', value: snapshot.author.did, label: snapshot.author.displayName || `@${snapshot.author.handle}`, }); add('r', { kind: 'author-feed', value: snapshot.author.did, label: `Recent posts by @${snapshot.author.handle}`, }); } else { const role = ({ post: 'p', profile: 'a', link: 'l', tag: 't', media: 'm', occurrence: 'e', 'author-feed': 'r', } as const)[snapshot.target.kind]; add(role, { ...snapshot.target }); }
return refs;}
const PUBLICATION_CLAIM_TTL_MS = 5 * 60_000;
function parsePublicationPlan(row: PublicationPlanRow | null, id: string): PublicationPlan { if (!row) throw new Error(`Unknown publication: ${id}`); const plan = JSON.parse(row.plan_json) as PublicationPlan; if (plan.version !== 1 || plan.id !== id) throw new Error(`Stored publication ${id} has an unsupported plan format.`); return plan;}
function publicationProgress(plan: PublicationPlan, rows: PublicationPostRow[]): PublicationPostProgress[] { if (rows.length !== plan.posts.length) throw new Error(`Publication ${plan.id} has incomplete durable post state.`); return rows.map((row, offset) => { const prepared = plan.posts[offset]; if (!prepared || prepared.index !== row.post_index || prepared.rkey !== row.rkey) { throw new Error(`Publication ${plan.id} post state does not match its immutable plan.`); } return { index: row.post_index, rkey: row.rkey, plannedUri: prepared.plannedUri, state: row.state, ...(row.record_json !== null ? { record: JSON.parse(row.record_json) as Record<string, unknown> } : {}), ...(row.uri !== null ? { uri: row.uri } : {}), ...(row.cid !== null ? { cid: row.cid } : {}), ...(row.attempted_at !== null ? { attemptedAt: row.attempted_at } : {}), ...(row.confirmed_at !== null ? { confirmedAt: row.confirmed_at } : {}), ...(row.observed_at !== null ? { observedAt: row.observed_at } : {}), ...(row.last_error !== null ? { lastError: row.last_error } : {}), }; });}
function publicationReceipt(plan: PublicationPlan, posts: PublicationPostProgress[]): PublicationReceipt { const confirmed = posts.filter((post) => post.state === 'confirmed').length; const status: PublicationReceipt['status'] = confirmed === posts.length ? 'complete' : posts.some((post) => post.state === 'dispatching' || post.state === 'uncertain') ? 'uncertain' : posts.some((post) => post.state === 'assembling') ? 'publishing' : confirmed > 0 ? 'partial' : 'prepared'; return freezeDeep({ id: plan.id, status, plan, posts });}
export class ReadingStore { readonly path: string; private readonly database: Database;
constructor(path: string) { if (path !== ':memory:') mkdirSync(dirname(path), { recursive: true, mode: 0o700 }); this.path = path; this.database = new Database(path, { create: true, strict: true }); this.database.exec('PRAGMA busy_timeout = 5000'); this.database.exec('PRAGMA foreign_keys = ON'); if (path !== ':memory:') this.database.exec('PRAGMA journal_mode = WAL'); this.database.exec(` CREATE TABLE IF NOT EXISTS store_sequence ( singleton INTEGER PRIMARY KEY CHECK (singleton = 1), last_view INTEGER NOT NULL ); INSERT OR IGNORE INTO store_sequence(singleton, last_view) VALUES (1, 0); CREATE TABLE IF NOT EXISTS views ( id TEXT PRIMARY KEY, sequence INTEGER NOT NULL UNIQUE, view_json TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS sessions ( name TEXT PRIMARY KEY, current_view TEXT REFERENCES views(id), history_json TEXT NOT NULL, revision INTEGER NOT NULL ); CREATE TABLE IF NOT EXISTS publication_plans ( id TEXT PRIMARY KEY, owner_did TEXT NOT NULL, service TEXT NOT NULL, prepared_at TEXT NOT NULL, plan_json TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS publication_posts ( plan_id TEXT NOT NULL REFERENCES publication_plans(id), post_index INTEGER NOT NULL, rkey TEXT NOT NULL, state TEXT NOT NULL CHECK (state IN ('prepared', 'assembling', 'dispatching', 'uncertain', 'confirmed')), record_json TEXT, uri TEXT, cid TEXT, attempted_at TEXT, confirmed_at TEXT, observed_at TEXT, last_error TEXT, claim_token TEXT, claimed_at TEXT, PRIMARY KEY (plan_id, post_index), UNIQUE (plan_id, rkey) ); CREATE TABLE IF NOT EXISTS publication_attachments ( plan_id TEXT NOT NULL REFERENCES publication_plans(id), attachment_id TEXT NOT NULL, metadata_json TEXT NOT NULL, bytes BLOB NOT NULL, PRIMARY KEY (plan_id, attachment_id) ); `); }
close(): void { this.database.close(); }
private immediate<T>(operation: () => T): T { this.database.exec('BEGIN IMMEDIATE'); try { const result = operation(); this.database.exec('COMMIT'); return result; } catch (error) { this.database.exec('ROLLBACK'); throw error; } }
private readSession(name: string): SessionRow | null { return this.database .query<SessionRow, [string]>('SELECT name, current_view, history_json, revision FROM sessions WHERE name = ?') .get(name) ?? null; }
private requireSession(name: string): SessionRow { const row = this.readSession(name); if (!row) throw new Error(`Unknown session: ${name}`); return row; }
private readPublicationPlan(id: string): PublicationPlan { const row = this.database .query<PublicationPlanRow, [string]>('SELECT plan_json FROM publication_plans WHERE id = ?') .get(id) ?? null; return parsePublicationPlan(row, id); }
private readPublicationPosts(id: string): PublicationPostRow[] { return this.database.query<PublicationPostRow, [string]>(` SELECT post_index, rkey, state, record_json, uri, cid, attempted_at, confirmed_at, observed_at, last_error, claim_token, claimed_at FROM publication_posts WHERE plan_id = ? ORDER BY post_index `).all(id); }
getView(id: string): ReadingView | null { const row = this.database.query<ViewRow, [string]>('SELECT view_json FROM views WHERE id = ?').get(id); return row ? parseView(row.view_json) : null; }
position(name: string): SessionPosition { validateSessionName(name); this.database.query('INSERT OR IGNORE INTO sessions(name, current_view, history_json, revision) VALUES (?, NULL, ?, 0)') .run(name, '[]'); const row = this.requireSession(name); const history = JSON.parse(row.history_json) as string[]; return { session: freezeDeep({ name: row.name, viewId: row.current_view, history }), revision: row.revision, view: row.current_view ? this.getView(row.current_view) : null, }; }
current(name: string): ReadingView | null { return this.position(name).view; }
sessions(): SessionInfo[] { const rows = this.database .query<SessionRow, []>('SELECT name, current_view, history_json, revision FROM sessions ORDER BY name') .all(); return rows.map((row) => freezeDeep({ name: row.name, viewId: row.current_view, history: JSON.parse(row.history_json) as string[], })); }
resolveRef(ref: string): ResolvedReference { const separator = ref.indexOf('-'); if (separator < 2) throw new UnknownReferenceError(ref); const viewId = ref.slice(0, separator); if (!/^v[1-9]\d*$/u.test(viewId)) throw new UnknownReferenceError(ref); const view = this.getView(viewId); const target = view?.refs[ref]; if (!view || !target) throw new UnknownReferenceError(ref); return freezeDeep({ view, ref, target }); }
commit(name: string, snapshot: Snapshot, expectedRevision: number): ReadingView { validateSessionName(name); const frozenSnapshot = freezeDeep(JSON.parse(JSON.stringify(snapshot)) as Snapshot); return this.immediate(() => { const session = this.requireSession(name); if (session.revision !== expectedRevision) throw new StoreConflictError(name);
const sequenceRow = this.database .query<{ last_view: number }, []>('SELECT last_view FROM store_sequence WHERE singleton = 1') .get(); if (!sequenceRow) throw new Error('Reading store sequence is missing.'); const sequence = sequenceRow.last_view + 1; const id = `v${sequence}`; const view = freezeDeep({ id, snapshot: frozenSnapshot, refs: buildReferences(id, frozenSnapshot) }); const history = JSON.parse(session.history_json) as string[]; if (session.current_view) history.push(session.current_view);
this.database.query('UPDATE store_sequence SET last_view = ? WHERE singleton = 1').run(sequence); this.database.query('INSERT INTO views(id, sequence, view_json) VALUES (?, ?, ?)') .run(id, sequence, JSON.stringify(view)); const result = this.database.query(` UPDATE sessions SET current_view = ?, history_json = ?, revision = revision + 1 WHERE name = ? AND revision = ? `).run(id, JSON.stringify(history), name, expectedRevision); if (result.changes !== 1) throw new StoreConflictError(name); return view; }); }
back(name: string, expectedRevision: number): ReadingView { validateSessionName(name); return this.immediate(() => { const session = this.requireSession(name); if (session.revision !== expectedRevision) throw new StoreConflictError(name); const history = JSON.parse(session.history_json) as string[]; const previous = history.pop(); if (!previous) throw new Error(`Session "${name}" has no previous view.`); const view = this.getView(previous); if (!view) throw new Error(`Stored view ${previous} is missing.`); const result = this.database.query(` UPDATE sessions SET current_view = ?, history_json = ?, revision = revision + 1 WHERE name = ? AND revision = ? `).run(previous, JSON.stringify(history), name, expectedRevision); if (result.changes !== 1) throw new StoreConflictError(name); return view; }); }
createPublication(prepared: PreparedPublication): PublicationReceipt { const plan = JSON.parse(JSON.stringify(prepared.plan)) as PublicationPlan; if (plan.version !== 1 || !plan.id || !plan.ownerDid || !plan.service || plan.posts.length === 0) { throw new Error('Publication plan is incomplete.'); } const payloadById = new Map(prepared.attachments.map((attachment) => [attachment.metadata.id, attachment])); if (payloadById.size !== plan.attachments.length || prepared.attachments.length !== plan.attachments.length) { throw new Error('Publication attachment payloads do not match the immutable plan.'); } for (const metadata of plan.attachments) { const payload = payloadById.get(metadata.id); if (!payload || JSON.stringify(payload.metadata) !== JSON.stringify(metadata)) { throw new Error(`Publication attachment ${metadata.id} metadata does not match its plan.`); } if ( payload.bytes.length !== metadata.byteLength || createHash('sha256').update(payload.bytes).digest('hex') !== metadata.sha256 ) throw new Error(`Publication attachment ${metadata.id} bytes do not match its plan.`); }
this.immediate(() => { this.database.query(` INSERT INTO publication_plans(id, owner_did, service, prepared_at, plan_json) VALUES (?, ?, ?, ?, ?) `).run(plan.id, plan.ownerDid, plan.service, plan.preparedAt, JSON.stringify(plan)); const insertPost = this.database.query(` INSERT INTO publication_posts(plan_id, post_index, rkey, state) VALUES (?, ?, ?, 'prepared') `); plan.posts.forEach((post, offset) => { if (post.index !== offset + 1 || !post.rkey || !post.plannedUri) { throw new Error(`Publication ${plan.id} post ordering or identity is invalid.`); } insertPost.run(plan.id, post.index, post.rkey); }); const insertAttachment = this.database.query(` INSERT INTO publication_attachments(plan_id, attachment_id, metadata_json, bytes) VALUES (?, ?, ?, ?) `); for (const attachment of prepared.attachments) { insertAttachment.run( plan.id, attachment.metadata.id, JSON.stringify(attachment.metadata), attachment.bytes, ); } }); return this.publication(plan.id); }
publication(id: string): PublicationReceipt { const plan = this.readPublicationPlan(id); return publicationReceipt(plan, publicationProgress(plan, this.readPublicationPosts(id))); }
assertPublicationBinding( id: string, source: { service: string; viewerDid: string }, ): PublicationReceipt { const receipt = this.publication(id); if (receipt.plan.ownerDid !== source.viewerDid || receipt.plan.service !== source.service) { throw new Error( `Publication ${id} belongs to DID ${receipt.plan.ownerDid} at ${receipt.plan.service}, not the current authenticated account and service.`, ); } return receipt; }
claimPublicationPost( id: string, source: { service: string; viewerDid: string }, ): PublicationClaim | null { return this.immediate(() => { const plan = this.readPublicationPlan(id); if (plan.ownerDid !== source.viewerDid || plan.service !== source.service) { throw new Error(`Publication ${id} belongs to a different authenticated DID or service.`); } let rows = this.readPublicationPosts(id); if (rows.some((row) => row.state === 'dispatching' || row.state === 'uncertain')) { throw new Error(`Publication ${id} has an uncertain post; run reconcile before continuing.`); } const assembling = rows.find((row) => row.state === 'assembling'); if (assembling) { const claimedAt = assembling.claimed_at === null ? Number.NaN : Date.parse(assembling.claimed_at); if (Number.isFinite(claimedAt) && Date.now() - claimedAt < PUBLICATION_CLAIM_TTL_MS) { throw new Error(`Publication ${id} post ${assembling.post_index} is already claimed by another invocation.`); } this.database.query(` UPDATE publication_posts SET state = 'prepared', claim_token = NULL, claimed_at = NULL, last_error = 'Recovered an expired pre-dispatch assembly claim; no record dispatch had begun.' WHERE plan_id = ? AND post_index = ? AND state = 'assembling' `).run(id, assembling.post_index); rows = this.readPublicationPosts(id); } const next = rows.find((row) => row.state !== 'confirmed'); if (!next) return null; const predecessors = rows.filter((row) => row.post_index < next.post_index); if (predecessors.some((row) => row.state !== 'confirmed' || !row.uri || !row.cid)) { throw new Error(`Publication ${id} post ${next.post_index} does not have a fully confirmed prefix.`); } if (next.state !== 'prepared') throw new Error(`Publication ${id} post ${next.post_index} cannot be claimed from ${next.state}.`); const token = randomUUID(); const claimedAt = new Date().toISOString(); const claimed = this.database.query(` UPDATE publication_posts SET state = 'assembling', claim_token = ?, claimed_at = ?, last_error = NULL WHERE plan_id = ? AND post_index = ? AND state = 'prepared' `).run(token, claimedAt, id, next.post_index); if (claimed.changes !== 1) throw new Error(`Publication ${id} post ${next.post_index} was claimed concurrently.`);
const attachmentRows = this.database.query<PublicationAttachmentRow, [string]>(` SELECT attachment_id, metadata_json, bytes FROM publication_attachments WHERE plan_id = ? ORDER BY attachment_id `).all(id); const attachments = attachmentRows.flatMap((row) => { const metadata = JSON.parse(row.metadata_json) as FrozenAttachment['metadata']; if (metadata.postIndex !== next.post_index) return []; return [{ metadata, bytes: new Uint8Array(row.bytes) }]; }); const post = plan.posts[next.post_index - 1]; if (!post || post.rkey !== next.rkey) throw new Error(`Publication ${id} post plan is inconsistent.`); const confirmed: StrongReference[] = predecessors.map((row) => ({ uri: row.uri!, cid: row.cid! })); const record = next.record_json === null ? undefined : JSON.parse(next.record_json) as Record<string, unknown>; return { token, plan, post, attachments, confirmed, ...(record !== undefined ? { record } : {}) }; }); }
beginPublicationDispatch( id: string, postIndex: number, token: string, record: Record<string, unknown>, ): void { const attemptedAt = new Date().toISOString(); const recordJson = JSON.stringify(record); const result = this.database.query(` UPDATE publication_posts SET state = 'dispatching', record_json = ?, attempted_at = ?, observed_at = NULL, last_error = NULL WHERE plan_id = ? AND post_index = ? AND state = 'assembling' AND claim_token = ? AND (record_json IS NULL OR record_json = ?) `).run(recordJson, attemptedAt, id, postIndex, token, recordJson); if (result.changes !== 1) { throw new Error(`Publication ${id} post ${postIndex} lost its claim or changed its durable dispatched record.`); } }
confirmPublicationPost( id: string, postIndex: number, token: string, result: StrongReference, ): PublicationReceipt { const confirmedAt = new Date().toISOString(); const updated = this.database.query(` UPDATE publication_posts SET state = 'confirmed', uri = ?, cid = ?, confirmed_at = ?, observed_at = ?, last_error = NULL, claim_token = NULL, claimed_at = NULL WHERE plan_id = ? AND post_index = ? AND state = 'dispatching' AND claim_token = ? `).run(result.uri, result.cid, confirmedAt, confirmedAt, id, postIndex, token); if (updated.changes !== 1) throw new Error(`Publication ${id} post ${postIndex} could not record its confirmation.`); return this.publication(id); }
markPublicationUncertain( id: string, postIndex: number, token: string, error: string, ): PublicationReceipt { const updated = this.database.query(` UPDATE publication_posts SET state = 'uncertain', last_error = ?, claim_token = NULL, claimed_at = NULL WHERE plan_id = ? AND post_index = ? AND state = 'dispatching' AND claim_token = ? `).run(error.slice(0, 1000), id, postIndex, token); if (updated.changes !== 1) throw new Error(`Publication ${id} post ${postIndex} could not record uncertainty.`); return this.publication(id); }
releasePublicationAssembly( id: string, postIndex: number, token: string, error: string, ): PublicationReceipt { const updated = this.database.query(` UPDATE publication_posts SET state = 'prepared', last_error = ?, claim_token = NULL, claimed_at = NULL WHERE plan_id = ? AND post_index = ? AND state = 'assembling' AND claim_token = ? `).run(error.slice(0, 1000), id, postIndex, token); if (updated.changes !== 1) throw new Error(`Publication ${id} post ${postIndex} could not release its assembly claim.`); return this.publication(id); }
reconcilePublicationPost( id: string, postIndex: number, inspection: PublicationRecordInspection, retryAbsent: boolean, ): PublicationReceipt { return this.immediate(() => { const rows = this.readPublicationPosts(id); const row = rows.find((candidate) => candidate.post_index === postIndex); if (!row || (row.state !== 'dispatching' && row.state !== 'uncertain') || row.record_json === null) { throw new Error(`Publication ${id} post ${postIndex} is not an uncertain dispatched record.`); } if (inspection.state === 'matching') { if (!inspection.uri || !inspection.cid) throw new Error('Matching reconciliation omitted URI or CID.'); this.database.query(` UPDATE publication_posts SET state = 'confirmed', uri = ?, cid = ?, confirmed_at = ?, observed_at = ?, last_error = NULL, claim_token = NULL, claimed_at = NULL WHERE plan_id = ? AND post_index = ? AND state IN ('dispatching', 'uncertain') `).run(inspection.uri, inspection.cid, inspection.observedAt, inspection.observedAt, id, postIndex); } else if (inspection.state === 'absent' && retryAbsent) { this.database.query(` UPDATE publication_posts SET state = 'prepared', uri = NULL, cid = NULL, observed_at = ?, last_error = ?, claim_token = NULL, claimed_at = NULL WHERE plan_id = ? AND post_index = ? AND state IN ('dispatching', 'uncertain') `).run( inspection.observedAt, 'Exact PDS record was absent; explicit --retry-absent reset the same create-only rkey for a separate publish invocation.', id, postIndex, ); } else { const detail = inspection.detail ?? ( inspection.state === 'absent' ? 'Exact PDS record is currently absent; uncertainty remains until explicitly reset with --retry-absent.' : 'Exact PDS record did not match the durable prepared record.' ); this.database.query(` UPDATE publication_posts SET state = 'uncertain', observed_at = ?, last_error = ?, claim_token = NULL, claimed_at = NULL WHERE plan_id = ? AND post_index = ? AND state IN ('dispatching', 'uncertain') `).run(inspection.observedAt, detail.slice(0, 1000), id, postIndex); } return this.publication(id); }); }
fork(sourceName: string, destinationName: string): SessionPosition { validateSessionName(sourceName); validateSessionName(destinationName); if (sourceName === destinationName) throw new Error('Fork destination must have a different session name.'); return this.immediate(() => { const source = this.requireSession(sourceName); if (!source.current_view) throw new Error(`Session "${sourceName}" has no current view to fork.`); if (this.readSession(destinationName)) throw new Error(`Session "${destinationName}" already exists.`); this.database.query('INSERT INTO sessions(name, current_view, history_json, revision) VALUES (?, ?, ?, 0)') .run(destinationName, source.current_view, source.history_json); const history = JSON.parse(source.history_json) as string[]; const view = this.getView(source.current_view); if (!view) throw new Error(`Stored view ${source.current_view} is missing.`); return { session: freezeDeep({ name: destinationName, viewId: source.current_view, history }), revision: 0, view, }; }); }}