spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Per-run change inspection and rollback. A rollback never deletes: claims → rejected (reason rollback), provenance rows3 * → is_current = false (winner restored to the latest remaining current row per field), events → review_status rejected.4 * Entity columns are left as they are; the worker's reconciliation re-derives them from the remaining claims.5 */6import type { ClaimDTO, EventDTO, ProvenanceDTO } from "@dci/core";7import { pg, claimCols, eventCols, eventJoins } from "../../lib/sql.js";8import { bool, int, iso, reqStr, str, type Row } from "../../lib/rows.js";9import { claimDto, eventDto, provenanceDto } from "../../lib/dto.js";10import { entityKey, resolveEntityRefs } from "../../lib/resolve.js";1112const CAP = 2000;1314export interface RunChanges {15 runId: string;16 provenance: Array<ProvenanceDTO & { entityType: string; entityId: string; isCurrent: boolean }>;17 claims: ClaimDTO[];18 events: EventDTO[];19 documentVersions: Array<{ id: string; documentId: string; url: string | null; fetchedAt: string | null; significance: number; changes: number }>;20 counts: { provenance: number; claims: number; events: number; documentVersions: number };21}2223export async function runExists(id: string): Promise<boolean> {24 const sql = pg();25 return (await sql`select 1 from connector_runs where id = ${id}`).length > 0;26}2728export async function runChanges(id: string): Promise<RunChanges | null> {29 const sql = pg();30 if (!(await runExists(id))) return null;31 const [prov, claims, events, versions, counts] = await Promise.all([32 sql<Row[]>`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}`,33 sql<Row[]>`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}`,34 sql<Row[]>`select ${eventCols(sql)} from events e ${eventJoins(sql)} where e.run_id = ${id} order by e.detected_at desc limit ${CAP}`,35 sql<Row[]>`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}`,36 sql<Row[]>`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`,37 ]);38 const refs = await resolveEntityRefs(events.map((r) => ({ type: reqStr(r.entity_type), id: str(r.entity_id) })));39 const c = counts[0] ?? {};40 return {41 runId: id,42 provenance: prov.map((p) => ({ ...provenanceDto(p), entityType: reqStr(p.entity_type), entityId: reqStr(p.entity_id), isCurrent: bool(p.is_current) })),43 claims: claims.map((k) => claimDto(k, false)),44 events: events.map((e) => eventDto(e, refs.get(entityKey(e.entity_type, e.entity_id)) ?? null)),45 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) })),46 counts: { provenance: int(c.p), claims: int(c.c), events: int(c.e), documentVersions: int(c.v) },47 };48}4950export interface RollbackResult { runId: string; claimsRejected: number; provenanceRetired: number; winnersRestored: number; eventsRejected: number }5152export async function rollbackRun(id: string): Promise<RollbackResult | null> {53 const sql = pg();54 if (!(await runExists(id))) return null;55 return sql.begin(async (tx) => {56 const claims = await tx`update claims set status = 'rejected', rejection_reason = 'rollback' where run_id = ${id} and status <> 'rejected'`;57 const touched = await tx<Row[]>`select distinct entity_type, entity_id, field from provenance where run_id = ${id} and is_current`;58 const prov = await tx`update provenance set is_current = false, is_winner = false where run_id = ${id} and is_current`;59 let restored = 0;60 for (const t of touched) {61 const r = await tx`update provenance set is_winner = true where id = (62 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)63 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)`;64 restored += r.count;65 }66 const events = await tx`update events set review_status = 'rejected' where run_id = ${id} and review_status <> 'rejected'`;67 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}`;68 return { runId: id, claimsRejected: claims.count, provenanceRetired: prov.count, winnersRestored: restored, eventsRejected: events.count };69 });70}71