/** * Per-run change inspection and rollback. A rollback never deletes: claims → rejected (reason rollback), provenance rows * → is_current = false (winner restored to the latest remaining current row per field), events → review_status rejected. * Entity columns are left as they are; the worker's reconciliation re-derives them from the remaining claims. */ import type { ClaimDTO, EventDTO, ProvenanceDTO } from "@dci/core"; import { pg, claimCols, eventCols, eventJoins } from "../../lib/sql.js"; import { bool, int, iso, reqStr, str, type Row } from "../../lib/rows.js"; import { claimDto, eventDto, provenanceDto } from "../../lib/dto.js"; import { entityKey, resolveEntityRefs } from "../../lib/resolve.js"; const CAP = 2000; export interface RunChanges { runId: string; provenance: Array; claims: ClaimDTO[]; events: EventDTO[]; documentVersions: Array<{ id: string; documentId: string; url: string | null; fetchedAt: string | null; significance: number; changes: number }>; counts: { provenance: number; claims: number; events: number; documentVersions: number }; } export async function runExists(id: string): Promise { const sql = pg(); return (await sql`select 1 from connector_runs where id = ${id}`).length > 0; } export async function runChanges(id: string): Promise { const sql = pg(); if (!(await runExists(id))) return null; const [prov, claims, events, versions, counts] = await Promise.all([ sql`select p.*, s.name as source_name, s.kind as source_kind from provenance p left join sources s on s.id = p.source_id where p.run_id = ${id} order by p.entity_type, p.entity_id, p.field limit ${CAP}`, sql`select ${claimCols(sql)} from claims k left join sources s on s.id = k.source_id where k.run_id = ${id} order by k.subject_type, k.subject_id, k.predicate limit ${CAP}`, sql`select ${eventCols(sql)} from events e ${eventJoins(sql)} where e.run_id = ${id} order by e.detected_at desc limit ${CAP}`, sql`select v.id, v.document_id, d.url, v.fetched_at, v.significance, jsonb_array_length(coalesce(v.detected_changes, '[]'::jsonb)) as changes from document_versions v left join documents d on d.id = v.document_id where v.run_id = ${id} order by v.fetched_at desc limit ${CAP}`, sql`select (select count(*) from provenance where run_id = ${id})::int as p, (select count(*) from claims where run_id = ${id})::int as c, (select count(*) from events where run_id = ${id})::int as e, (select count(*) from document_versions where run_id = ${id})::int as v`, ]); const refs = await resolveEntityRefs(events.map((r) => ({ type: reqStr(r.entity_type), id: str(r.entity_id) }))); const c = counts[0] ?? {}; return { runId: id, provenance: prov.map((p) => ({ ...provenanceDto(p), entityType: reqStr(p.entity_type), entityId: reqStr(p.entity_id), isCurrent: bool(p.is_current) })), claims: claims.map((k) => claimDto(k, false)), events: events.map((e) => eventDto(e, refs.get(entityKey(e.entity_type, e.entity_id)) ?? null)), documentVersions: versions.map((v) => ({ id: reqStr(v.id), documentId: reqStr(v.document_id), url: str(v.url), fetchedAt: iso(v.fetched_at), significance: int(v.significance), changes: int(v.changes) })), counts: { provenance: int(c.p), claims: int(c.c), events: int(c.e), documentVersions: int(c.v) }, }; } export interface RollbackResult { runId: string; claimsRejected: number; provenanceRetired: number; winnersRestored: number; eventsRejected: number } export async function rollbackRun(id: string): Promise { const sql = pg(); if (!(await runExists(id))) return null; return sql.begin(async (tx) => { const claims = await tx`update claims set status = 'rejected', rejection_reason = 'rollback' where run_id = ${id} and status <> 'rejected'`; const touched = await tx`select distinct entity_type, entity_id, field from provenance where run_id = ${id} and is_current`; const prov = await tx`update provenance set is_current = false, is_winner = false where run_id = ${id} and is_current`; let restored = 0; for (const t of touched) { const r = await tx`update provenance set is_winner = true where id = ( select id from provenance where entity_type = ${reqStr(t.entity_type)} and entity_id = ${reqStr(t.entity_id)} and field = ${reqStr(t.field)} and is_current order by last_observed desc, retrieved_at desc limit 1) and not exists (select 1 from provenance w where w.entity_type = ${reqStr(t.entity_type)} and w.entity_id = ${reqStr(t.entity_id)} and w.field = ${reqStr(t.field)} and w.is_current and w.is_winner)`; restored += r.count; } const events = await tx`update events set review_status = 'rejected' where run_id = ${id} and review_status <> 'rejected'`; await tx`update connector_runs set log = coalesce(log, '[]'::jsonb) || ${JSON.stringify([{ t: new Date().toISOString(), level: "warn", msg: `rolled back by admin: ${claims.count} claims rejected, ${prov.count} provenance rows retired, ${events.count} events rejected` }])}::jsonb where id = ${id}`; return { runId: id, claimsRejected: claims.count, provenanceRetired: prov.count, winnersRestored: restored, eventsRejected: events.count }; }); }