diff --git a/lib/automations/behavior-hash.ts b/lib/automations/behavior-hash.ts index 2e8b14b..b6da937 100644 --- a/lib/automations/behavior-hash.ts +++ b/lib/automations/behavior-hash.ts @@ -107,8 +107,12 @@ export type StepsBehaviorInput = { function normalizeStep(step: Step): Record { const out: Record = {}; for (const [key, value] of Object.entries(step)) { - // Identity and UI-only fields carry no behavior. - if (key === "id" || key === "comment" || key === "note" || key === "preserveOnDuplication") { + // Identity and UI-only fields carry no behavior. `note` is only + // UI-metadata on VARIABLE steps — on kipclip-annotation it is record + // content and must stay in the hash (a production pair differs by + // exactly that field). + if (key === "id" || key === "comment") continue; + if (step.$type === "variable" && (key === "note" || key === "preserveOnDuplication")) { continue; } // Instance-local webhook values, regenerated on duplication. diff --git a/lib/automations/legacy-convert.test.ts b/lib/automations/legacy-convert.test.ts new file mode 100644 index 0000000..da73e4b --- /dev/null +++ b/lib/automations/legacy-convert.test.ts @@ -0,0 +1,222 @@ +import { describe, it, expect } from "vitest"; +import { + convertLegacyAutomation, + legacyIndexToStepId, + LegacyConvertError, +} from "./legacy-convert.ts"; +import { + triggerGatePrefixLength, + isConditionStep, + type ConditionStep, + type LoopStep, +} from "./steps.ts"; +import type { Action, FetchStep, Variable } from "../db/schema.ts"; + +const webhook = (over: Partial = {}): Action => + ({ + $type: "webhook", + callbackUrl: "https://example.com/hook", + secret: "s3cret", + ...over, + }) as Action; + +describe("convertLegacyAutomation (KTD-12)", () => { + it("emits gates first with variable inlining, keeping conditions on the trigger prefix", () => { + const steps = convertLegacyAutomation({ + conditions: [ + { field: "status", operator: "eq", value: "{{wanted}}" }, + { field: "repo", operator: "eq", value: "{{self}}" }, + ], + variables: [{ name: "wanted", value: "live" }] as Variable[], + fetches: [], + actions: [webhook()], + }); + // Gate BEFORE the variable step (gate-first ordering). + expect(steps[0]!.$type).toBe("condition"); + expect(steps[1]!.$type).toBe("variable"); + const gate = steps[0] as ConditionStep; + expect(gate.assertions[0]).toMatchObject({ + field: "event.commit.record.status", + operator: "eq", + value: "live", // {{wanted}} inlined to the literal + }); + expect(gate.assertions[1]).toMatchObject({ field: "event.did", value: "{{self}}" }); + // The whole gate stays inside the trigger-gate prefix (hot path). + expect(triggerGatePrefixLength(steps)).toBe(1); + }); + + it("splits more than 10 conditions across consecutive gates", () => { + const steps = convertLegacyAutomation({ + conditions: Array.from({ length: 13 }, (_, i) => ({ + field: `f${i}`, + operator: "exists", + value: "", + })), + variables: [], + fetches: [], + actions: [webhook()], + }); + const gates = steps.filter(isConditionStep); + expect(gates).toHaveLength(2); + expect(gates[0]!.assertions).toHaveLength(10); + expect(gates[1]!.assertions).toHaveLength(3); + expect(triggerGatePrefixLength(steps)).toBe(2); + }); + + it("maps fetch conditions to a gate after the fetch, with the bare-found special case as boolean eq", () => { + const steps = convertLegacyAutomation({ + conditions: [], + variables: [], + fetches: [ + { + kind: "record", + name: "mirror", + uri: "{{event.commit.record.subject}}", + conditions: [{ field: "found", operator: "not-exists", value: "" }], + }, + ] as FetchStep[], + actions: [webhook()], + }); + expect(steps[0]).toMatchObject({ $type: "fetch", name: "mirror" }); + const gate = steps[1] as ConditionStep; + expect(gate.assertions[0]).toEqual({ + field: "mirror.found", + operator: "eq", + value: "false", + }); + }); + + it("normalizes a fetch with no kind to record (the malformed production row)", () => { + const steps = convertLegacyAutomation({ + conditions: [], + variables: [], + fetches: [{ name: "repo", uri: "{{event.commit.record.subject}}" }] as FetchStep[], + actions: [webhook()], + }); + expect(steps[0]!.$type).toBe("fetch"); + }); + + it("wraps a forEach action in a loop named action{N} with per-item gates inside", () => { + const steps = convertLegacyAutomation({ + conditions: [], + variables: [], + fetches: [], + actions: [ + { + $type: "record", + targetCollection: "com.example.x", + recordTemplate: '{"u":"{{item.uri}}"}', + forEach: { + path: "event.commit.record.facets[].features[]", + conditions: [{ field: "$type", operator: "eq", value: "app.bsky.richtext.facet#link" }], + }, + } as Action, + { + $type: "kipclip-annotation", + subject: "{{item.uri}}", + forEach: { path: "action1.results[]" }, + } as Action, + ], + }); + const loop1 = steps[0] as LoopStep; + expect(loop1).toMatchObject({ $type: "loop", name: "action1" }); + expect(loop1.steps[0]!.$type).toBe("condition"); // per-item gate + expect((loop1.steps[0] as ConditionStep).assertions[0]).toMatchObject({ + field: "item.$type", + operator: "eq", + }); + expect(loop1.steps[1]!.$type).toBe("record"); + expect("name" in loop1.steps[1]! ? loop1.steps[1]!.name : undefined).toBeUndefined(); + // Downstream loop keeps referencing action1.results[] verbatim. + const loop2 = steps[1] as LoopStep; + expect(loop2).toMatchObject({ $type: "loop", name: "action2", path: "action1.results[]" }); + }); + + it("renames collect-links steps to their outputName", () => { + const steps = convertLegacyAutomation({ + conditions: [], + variables: [], + fetches: [], + actions: [ + { + $type: "collect-links", + outputName: "links", + source: "{{event.commit.record}}", + } as Action, + webhook(), + ], + }); + expect(steps[0]).toMatchObject({ $type: "collect-links", name: "links" }); + expect((steps[0] as Record).outputName).toBeUndefined(); + expect(steps[1]).toMatchObject({ $type: "webhook", name: "action2" }); + }); + + it("keeps schedule modifiers on the action step", () => { + const steps = convertLegacyAutomation({ + conditions: [], + variables: [], + fetches: [], + actions: [ + webhook({ + schedule: { + identifier: "remind-{{event.commit.rkey}}", + delay: { unit: "h", value: { kind: "literal", n: 2 } }, + }, + } as never), + ], + }); + expect((steps[0] as Record).schedule).toBeTruthy(); + }); + + it("fails loudly on unknown operators instead of guessing (hostile-row robustness)", () => { + expect(() => + convertLegacyAutomation({ + conditions: [{ field: "x", operator: "matches", value: "y" }], + variables: [], + fetches: [], + actions: [webhook()], + }), + ).toThrow(LegacyConvertError); + }); +}); + +describe("legacyIndexToStepId", () => { + it("maps actionIndex to the step named action{i+1}, including loops", () => { + const steps = convertLegacyAutomation({ + conditions: [{ field: "status", operator: "eq", value: "live" }], + variables: [], + fetches: [], + actions: [ + webhook(), + { + $type: "record", + targetCollection: "c.x", + recordTemplate: "{}", + forEach: { path: "event.commit.record.tags[]" }, + } as Action, + ], + }); + const id0 = legacyIndexToStepId(steps, 0); + const id1 = legacyIndexToStepId(steps, 1); + expect(id0).toBe(steps.find((s) => "name" in s && s.name === "action1")!.id); + expect(id1).toBe(steps.find((s) => s.$type === "loop")!.id); + }); + + it("falls back positionally for renamed collect-links steps", () => { + const steps = convertLegacyAutomation({ + conditions: [], + variables: [], + fetches: [], + actions: [ + { + $type: "collect-links", + outputName: "links", + source: "{{event.commit.record}}", + } as Action, + webhook(), + ], + }); + expect(legacyIndexToStepId(steps, 0)).toBe(steps[0]!.id); + expect(legacyIndexToStepId(steps, 1)).toBe(steps[1]!.id); + }); +}); diff --git a/lib/automations/legacy-convert.ts b/lib/automations/legacy-convert.ts new file mode 100644 index 0000000..2c3ab84 --- /dev/null +++ b/lib/automations/legacy-convert.ts @@ -0,0 +1,270 @@ +/** Legacy-shape → steps-model conversion (KTD-12 of the steps-rework plan). + * + * Behavior-preserving by construction: + * - Top-level conditions become LEADING gate steps (before variable steps), + * with `{{variableName}}` operand references inlined to the variable's + * literal value — legacy variables are save-time literals, so inlining is + * behavior-identical and keeps every migrated automation's old conditions + * inside the trigger-gate prefix (hot path, dry-run silence, webhook + * projection). Lists longer than the per-step assertion cap split across + * consecutive gates (identical under `all`). + * - Variables become variable steps (literal values). + * - Fetches become fetch/search steps; their "continue only if" conditions + * become a gate step immediately after, with fields prefixed by the fetch + * name. The legacy bare-`found` exists/not-exists special case (which + * tested the boolean flag) maps to a strict boolean eq assertion. + * - Actions become action steps named `action{i+1}` so stored `{{actionN.*}}` + * template references keep resolving without rewriting. A `forEach` action + * becomes a loop step carrying the legacy name and wrapping the (unnamed) + * action, with per-item conditions as a gate inside the loop — the loop + * export reproduces the legacy `actionN` entry shape (last scalars + + * results[]). + * - collect-links `outputName` becomes the step `name`. + * - A fetch with no `kind` is normalized to `record` (one malformed row in + * production predates the discriminator). + */ +import type { Action, Condition, FetchStep as LegacyFetch, Variable } from "../db/schema.js"; +import { generateTid } from "./pds.js"; +import { STEP_LIMITS } from "./limits.js"; +import type { Assertion, AssertionOperator, ConditionStep, LoopStep, Step } from "./steps.js"; + +const LEGACY_OPERATOR_MAP: Record = { + eq: "eq", + startsWith: "startsWith", + endsWith: "endsWith", + contains: "contains", + exists: "exists", + "not-exists": "notExists", +}; + +export type LegacyAutomationShape = { + conditions: Condition[]; + fetches: LegacyFetch[]; + variables: Variable[]; + actions: Action[]; +}; + +export class LegacyConvertError extends Error {} + +function mapOperator(op: string, where: string): AssertionOperator { + const mapped = LEGACY_OPERATOR_MAP[op]; + if (!mapped) throw new LegacyConvertError(`${where}: unknown legacy operator "${op}"`); + return mapped; +} + +/** Root a legacy top-level condition field: bare fields address the commit + * record, "repo" is the legacy alias for the event DID. */ +function rootConditionField(field: string): string { + if (field === "repo") return "event.did"; + if (field === "event.commit.operation") { + // normalizeConditions drops these on save, but ancient rows could carry + // one; operation filtering lives on the trigger, so dropping is exact. + throw new LegacyConvertError("event.commit.operation conditions belong on the trigger"); + } + if (field.startsWith("event.")) return field; + return `event.commit.record.${field}`; +} + +/** Inline {{variableName}} references with their literal values so trigger + * gates stay event-only. {{self}} is hot-path-resolvable and stays. */ +function inlineVariables(value: string, variables: Variable[]): string { + if (!value.includes("{{")) return value; + return value.replace(/\{\{([^}]+)\}\}/g, (match, raw: string) => { + const name = raw.trim(); + if (name === "self") return match; + const variable = variables.find((v) => v.name === name); + return variable ? variable.value : match; + }); +} + +function conditionToAssertion(c: Condition, variables: Variable[], where: string): Assertion { + const operator = mapOperator(c.operator ?? "eq", where); + const valueless = operator === "exists" || operator === "notExists"; + return { + field: rootConditionField(c.field), + operator, + ...(valueless ? {} : { value: inlineVariables(c.value, variables) }), + ...(c.comment ? { comment: c.comment } : {}), + }; +} + +function chunkGates(assertions: Assertion[]): ConditionStep[] { + const gates: ConditionStep[] = []; + for (let i = 0; i < assertions.length; i += STEP_LIMITS.assertionsPerCondition) { + gates.push({ + $type: "condition", + id: generateTid(), + assertions: assertions.slice(i, i + STEP_LIMITS.assertionsPerCondition), + }); + } + return gates; +} + +/** Fetch-attached condition → assertion against the fetch's entry. The bare + * `found` + presence special case tested the boolean flag directly; boolean + * eq is the strict equivalent in the typed model. */ +function fetchConditionToAssertion(c: Condition, fetchName: string, where: string): Assertion { + const operator = mapOperator(c.operator ?? "eq", where); + if (c.field === "found" && (operator === "exists" || operator === "notExists")) { + return { + field: `${fetchName}.found`, + operator: "eq", + value: operator === "exists" ? "true" : "false", + }; + } + const valueless = operator === "exists" || operator === "notExists"; + return { + field: `${fetchName}.${c.field}`, + operator, + ...(valueless ? {} : { value: c.value }), + ...(c.comment ? { comment: c.comment } : {}), + }; +} + +/** Per-item forEach condition → assertion against the loop item. Legacy + * fields accepted an optional `item.` prefix. */ +function itemConditionToAssertion(c: Condition, where: string): Assertion { + const operator = mapOperator(c.operator ?? "eq", where); + const field = + c.field === "item" ? "item" : c.field.startsWith("item.") ? c.field : `item.${c.field}`; + const valueless = operator === "exists" || operator === "notExists"; + return { + field, + operator, + ...(valueless ? {} : { value: c.value }), + ...(c.comment ? { comment: c.comment } : {}), + }; +} + +function actionToStep(action: Action, legacyName: string): Step { + const { forEach, comment, ...config } = action as Action & { + forEach?: { path: string; conditions?: Condition[] }; + comment?: string; + }; + + // collect-links: outputName becomes the step name. + let name = legacyName; + const cfg = { ...config } as Record; + if (action.$type === "collect-links") { + name = (cfg.outputName as string) ?? legacyName; + delete cfg.outputName; + } + + const inner = { + ...cfg, + id: generateTid(), + ...(comment ? { comment } : {}), + } as unknown as Step; + + if (!forEach) { + return { ...inner, name } as Step; + } + + // forEach → loop named after the legacy action, wrapping the unnamed + // action; per-item conditions become a gate inside the loop (iteration + // skip semantics, matching the legacy filter-continue behavior). + const bodySteps: Step[] = []; + if (forEach.conditions && forEach.conditions.length > 0) { + bodySteps.push({ + $type: "condition", + id: generateTid(), + assertions: forEach.conditions.map((c, i) => + itemConditionToAssertion(c, `${legacyName} item condition ${i + 1}`), + ), + }); + } + bodySteps.push(inner); + + const loop: LoopStep = { + $type: "loop", + id: generateTid(), + name, + path: forEach.path, + steps: bodySteps, + }; + return loop; +} + +/** Convert one legacy-shaped automation into a step tree (KTD-12 order: + * trigger gates, variables, fetches with their gates, actions). Throws + * LegacyConvertError on anything unconvertible — the caller accumulates + * per-row failures and rolls the whole transaction back. */ +export function convertLegacyAutomation(row: LegacyAutomationShape): Step[] { + const steps: Step[] = []; + + if (row.conditions.length > 0) { + const assertions = row.conditions.map((c, i) => + conditionToAssertion(c, row.variables, `condition ${i + 1}`), + ); + steps.push(...chunkGates(assertions)); + } + + for (const variable of row.variables) { + steps.push({ + $type: "variable", + id: generateTid(), + name: variable.name, + value: variable.value, + ...(variable.note ? { note: variable.note } : {}), + ...(variable.preserveOnDuplication ? { preserveOnDuplication: true } : {}), + }); + } + + for (const fetch of row.fetches) { + if (fetch.kind === "search") { + steps.push({ + $type: "search", + id: generateTid(), + name: fetch.name, + repo: fetch.repo, + collection: fetch.collection, + where: fetch.where.map((w) => ({ field: w.field, operator: "eq", value: w.value })), + limit: 1, + ...(fetch.comment ? { comment: fetch.comment } : {}), + }); + } else { + // Missing `kind` normalizes to "record" (the runtime default). + steps.push({ + $type: "fetch", + id: generateTid(), + name: fetch.name, + uri: fetch.uri, + ...(fetch.collection ? { collection: fetch.collection } : {}), + ...(fetch.comment ? { comment: fetch.comment } : {}), + }); + } + if (fetch.conditions && fetch.conditions.length > 0) { + steps.push({ + $type: "condition", + id: generateTid(), + assertions: fetch.conditions.map((c, i) => + fetchConditionToAssertion(c, fetch.name, `fetch "${fetch.name}" condition ${i + 1}`), + ), + }); + } + } + + row.actions.forEach((action, i) => { + steps.push(actionToStep(action, `action${i + 1}`)); + }); + + return steps; +} + +/** The step id that legacy delivery_logs rows at `actionIndex` map to: the + * step named `action{i+1}` (a loop when the action carried forEach). */ +export function legacyIndexToStepId(steps: Step[], actionIndex: number): string | undefined { + const name = `action${actionIndex + 1}`; + for (const step of steps) { + if ("name" in step && step.name === name) return step.id; + // collect-links kept its outputName; fall through to positional matching + // below for automations whose action names diverge. + } + // Positional fallback: the Nth top-level action-or-loop step. Covers + // collect-links (renamed to its outputName) exactly, because conversion + // preserves action order at the top level. + const actionish = steps.filter( + (s) => s.$type === "loop" || !["condition", "fetch", "search", "variable"].includes(s.$type), + ); + return actionish[actionIndex]?.id; +} diff --git a/lib/scripts/migrate-automations-to-steps.ts b/lib/scripts/migrate-automations-to-steps.ts new file mode 100644 index 0000000..2ef83cf --- /dev/null +++ b/lib/scripts/migrate-automations-to-steps.ts @@ -0,0 +1,247 @@ +/** Steps-model row rewrite (U13 of the steps-rework plan). + * + * Rewrites every legacy-shaped automation row into the steps model in ONE + * transaction: rows + behavior hashes + delivery-log step_id backfill + + * pending scheduled-row keys. Per-row failures are accumulated into a full + * report and roll the whole transaction back — the production run is + * all-or-nothing (KTD-11). Old columns are never written and their + * checksums are asserted identical (the binary-rollback escape hatch). + * + * `steps IS NULL` is the idempotency predicate: a second run is a no-op. + * + * Usage: + * DATABASE_PATH=data/airglow.db bun run lib/scripts/migrate-automations-to-steps.ts + * DATABASE_PATH=copy.db bun run lib/scripts/migrate-automations-to-steps.ts --dry-run + */ +import { createHash } from "node:crypto"; +import { sql, eq, isNull, and, inArray } from "drizzle-orm"; +import { db } from "../db/index.js"; +import { automations, deliveryLogs, scheduledActions } from "../db/schema.js"; +import { convertLegacyAutomation, legacyIndexToStepId } from "../automations/legacy-convert.js"; +import { computeStepsBehaviorHash } from "../automations/behavior-hash.js"; +import { countSteps, flattenSteps, type Step } from "../automations/steps.js"; +import { STEP_LIMITS } from "../automations/limits.js"; + +const dryRun = process.argv.includes("--dry-run"); + +/** Pseudo-rows written at the automation level (not attributable to a step): + * rate-limit disables, session-loss disable/re-enable notices, and + * automation-level dry-run skips. Their step_id stays NULL (KTD-9). */ +function isPseudoRow(row: { + actionIndex: number; + statusCode: number | null; + error: string | null; + message: string | null; + dryRun: boolean; +}): boolean { + if (row.error?.startsWith("Rate limit exceeded")) return true; + if (row.error?.startsWith("Disabled after")) return true; + if (row.message === "Re-enabled after re-authorization.") return true; + if (row.dryRun && row.actionIndex === 0 && row.message?.startsWith("Skipped:")) return true; + return false; +} + +function oldColumnsChecksum(row: { + actions: unknown; + fetches: unknown; + variables: unknown; + conditions: unknown; +}): string { + return createHash("sha256") + .update(JSON.stringify([row.actions, row.fetches, row.variables, row.conditions])) + .digest("hex"); +} + +/** Structural sanity for a converted tree: bounded, unique ids and names. + * (Full save-time validation includes live webhook verification, which a + * migration must not perform; the rehearsal script covers semantics.) */ +function verifyTree(steps: Step[], uri: string): string | null { + const flat = flattenSteps(steps); // throws StepDepthError past the cap + const total = countSteps(steps); + if (total > STEP_LIMITS.totalSteps) { + return `${uri}: converted to ${total} steps (cap ${STEP_LIMITS.totalSteps})`; + } + const ids = new Set(); + const names = new Set(); + for (const step of flat) { + if (ids.has(step.id)) return `${uri}: duplicate step id ${step.id}`; + ids.add(step.id); + const name = "name" in step ? step.name : undefined; + if (name) { + if (names.has(name)) return `${uri}: duplicate step name ${name}`; + names.add(name); + } + } + return null; +} + +async function main() { + const rows = await db.query.automations.findMany(); + const unmigrated = rows.filter((r) => r.steps === null); + console.log(`${rows.length} automations, ${unmigrated.length} to migrate`); + + const preChecksums = new Map(rows.map((r) => [r.uri, oldColumnsChecksum(r)])); + + const failures: string[] = []; + const converted = new Map(); + for (const row of unmigrated) { + try { + const steps = convertLegacyAutomation({ + conditions: row.conditions, + fetches: row.fetches, + variables: row.variables, + actions: row.actions, + }); + const structural = verifyTree(steps, row.uri); + if (structural) { + failures.push(structural); + continue; + } + converted.set(row.uri, steps); + } catch (err) { + failures.push(`${row.uri}: ${err instanceof Error ? err.message : String(err)}`); + } + } + + if (failures.length > 0) { + console.error(`\n${failures.length} row(s) failed conversion — nothing was written:`); + for (const f of failures) console.error(` - ${f}`); + process.exit(1); + } + + db.run(sql`BEGIN IMMEDIATE`); + let committed = false; + const report = { rows: 0, logsBackfilled: 0, logsPseudo: 0, logsOutOfRange: 0, scheduled: 0 }; + try { + const now = new Date(); + for (const row of unmigrated) { + const steps = converted.get(row.uri)!; + const behaviorHash = computeStepsBehaviorHash({ + lexicon: row.lexicon, + operations: row.operations, + wantedDids: row.wantedDids, + steps, + }); + await db + .update(automations) + .set({ steps, behaviorHash, pdsStaleSince: now }) + .where(eq(automations.uri, row.uri)); + report.rows++; + + // delivery_logs backfill: actionIndex i -> the step named action{i+1} + // (the loop when the action carried forEach). IDs are read back from + // what was just stored, never from a separate minting pass. + const stored = await db.query.automations.findFirst({ + where: eq(automations.uri, row.uri), + }); + const storedSteps = stored!.steps!; + const logRows = await db + .select({ + id: deliveryLogs.id, + actionIndex: deliveryLogs.actionIndex, + statusCode: deliveryLogs.statusCode, + error: deliveryLogs.error, + message: deliveryLogs.message, + dryRun: deliveryLogs.dryRun, + }) + .from(deliveryLogs) + .where(and(eq(deliveryLogs.automationUri, row.uri), isNull(deliveryLogs.stepId))); + + const byIndex = new Map(); + const backfillable: Array<{ id: number; stepId: string }> = []; + for (const log of logRows) { + if (isPseudoRow(log)) { + report.logsPseudo++; + continue; + } + if (log.actionIndex >= row.actions.length) { + report.logsOutOfRange++; + continue; + } + if (!byIndex.has(log.actionIndex)) { + byIndex.set(log.actionIndex, legacyIndexToStepId(storedSteps, log.actionIndex)); + } + const stepId = byIndex.get(log.actionIndex); + if (stepId) backfillable.push({ id: log.id, stepId }); + else report.logsOutOfRange++; + } + for (const { id, stepId } of backfillable) { + await db.update(deliveryLogs).set({ stepId }).where(eq(deliveryLogs.id, id)); + report.logsBackfilled++; + } + + // Pending scheduled rows: key by step id; the legacy index stays put so + // a binary rollback replays them exactly as before. + const pending = await db + .select({ id: scheduledActions.id, sourceActionIndex: scheduledActions.sourceActionIndex }) + .from(scheduledActions) + .where( + and( + eq(scheduledActions.automationUri, row.uri), + inArray(scheduledActions.status, ["pending", "running"]), + isNull(scheduledActions.sourceStepId), + ), + ); + for (const p of pending) { + const stepId = legacyIndexToStepId(storedSteps, p.sourceActionIndex); + if (stepId) { + await db + .update(scheduledActions) + .set({ sourceStepId: stepId }) + .where(eq(scheduledActions.id, p.id)); + report.scheduled++; + } + } + } + + // Invariants (KTD-11) — asserted before COMMIT. + const post = await db.query.automations.findMany(); + const stillNull = post.filter((r) => r.steps === null); + if (stillNull.length > 0) { + throw new Error(`invariant: ${stillNull.length} rows still have NULL steps`); + } + for (const r of post) { + if (!r.behaviorHash) throw new Error(`invariant: ${r.uri} has no behaviorHash`); + if (preChecksums.get(r.uri) !== oldColumnsChecksum(r)) { + throw new Error(`invariant: legacy columns changed for ${r.uri}`); + } + const flat = flattenSteps(r.steps!); + const idSet = new Set(flat.map((s) => s.id)); + const orphanLogs = await db + .select({ id: deliveryLogs.id, stepId: deliveryLogs.stepId }) + .from(deliveryLogs) + .where(eq(deliveryLogs.automationUri, r.uri)); + for (const log of orphanLogs) { + if (log.stepId && !idSet.has(log.stepId)) { + throw new Error(`invariant: log ${log.id} references unknown step ${log.stepId}`); + } + } + } + + if (dryRun) { + db.run(sql`ROLLBACK`); + console.log("\n--dry-run: rolled back. Report:"); + } else { + db.run(sql`COMMIT`); + committed = true; + console.log("\nCommitted. Report:"); + } + } catch (err) { + if (!committed) db.run(sql`ROLLBACK`); + console.error("Migration failed and was rolled back:", err); + process.exit(1); + } + + console.log( + ` automations migrated: ${report.rows}\n` + + ` delivery logs backfilled: ${report.logsBackfilled}\n` + + ` pseudo rows left NULL: ${report.logsPseudo}\n` + + ` out-of-range rows left NULL: ${report.logsOutOfRange}\n` + + ` scheduled rows keyed: ${report.scheduled}`, + ); + + const integrity = db.all<{ integrity_check: string }>(sql`PRAGMA integrity_check`); + console.log(` integrity_check: ${integrity[0]?.integrity_check ?? "??"}`); +} + +await main(); diff --git a/lib/scripts/rehearse-steps-migration.ts b/lib/scripts/rehearse-steps-migration.ts new file mode 100644 index 0000000..94678fe --- /dev/null +++ b/lib/scripts/rehearse-steps-migration.ts @@ -0,0 +1,289 @@ +/** Migration rehearsal (U13): proves the KTD-12 conversion is + * behavior-preserving for every row in the target DB before the production + * run. Read-only — point it at a COPY of the deploy-day backup artifact. + * + * DATABASE_PATH=data/backups/prod-copy.db bun run lib/scripts/rehearse-steps-migration.ts + * + * Four equivalence classes (plan U13): + * (a) every converted tree is structurally valid and every stored template + * root resolves in scope at its step (plus @atproto/lexicon validation + * of the serialized record shape); + * (b) condition verdict parity: the legacy matcher and the new trigger-gate + * evaluator agree on a synthetic-event battery derived from each row's + * conditions (match / mismatch / missing-field cases); + * (c) behavior-hash partitions: rows that were duplicates of each other + * under the legacy canonicalization stay duplicates under the steps + * canonicalization (blueprint divergence detection survives); + * (d) OAuth scope parity: the flattened step tree needs exactly the scope + * the legacy action list needed. + */ +import { readFileSync } from "node:fs"; +import { fileURLToPath } from "node:url"; +import { Lexicons } from "@atproto/lexicon"; +import { db } from "../db/index.js"; +import { convertLegacyAutomation } from "../automations/legacy-convert.js"; +import { computeBehaviorHash, computeStepsBehaviorHash } from "../automations/behavior-hash.js"; +import { toPdsStep } from "../automations/pds-serialize.js"; +import { matchConditions, evaluateConditionStep } from "../jetstream/matcher.js"; +import type { JetstreamEvent } from "../jetstream/matcher.js"; +import { actionsNeedFullScope } from "../auth/client.js"; +import { PLACEHOLDER_RE } from "../actions/placeholders.js"; +import { + countSteps, + flattenSteps, + isActionStep, + scopeAt, + triggerGateSteps, + stepName, + type Step, +} from "../automations/steps.js"; + +const lexiconDoc = JSON.parse( + readFileSync( + fileURLToPath(new URL("../../lexicons/run/airglow/automation.json", import.meta.url)), + "utf8", + ), +); +const lexicons = new Lexicons([lexiconDoc]); + +const BUILTIN = new Set(["event", "now", "self", "automation"]); + +function checkScopes(steps: Step[]): string | null { + for (const step of flattenSteps(steps)) { + const scope = scopeAt(steps, step.id); + if (!scope) return `step ${step.id} not found by scopeAt`; + const visible = new Set([...scope.names, ...scope.maybeUnset, ...scope.itemAliases]); + const strings: string[] = []; + for (const [key, value] of Object.entries(step)) { + if (key === "then" || key === "else" || key === "steps") continue; + if (typeof value === "string") strings.push(value); + else if (Array.isArray(value)) { + for (const v of value) { + if (typeof v === "string") strings.push(v); + else if (v && typeof v === "object") { + for (const inner of Object.values(v as Record)) { + if (typeof inner === "string") strings.push(inner); + } + } + } + } else if (value && typeof value === "object") { + for (const inner of Object.values(value as Record)) { + if (typeof inner === "string") strings.push(inner); + } + } + } + for (const s of strings) { + for (const [, raw] of s.matchAll(PLACEHOLDER_RE)) { + let path = raw!.trim(); + const fn = /^\w+\((.*)\)$/.exec(path); + if (fn) path = fn[1]!.trim(); + if (path.startsWith("secret:")) continue; + const root = path.split(".")[0]!; + if (!BUILTIN.has(root) && !visible.has(root)) { + return `step "${stepName(step) ?? step.id}" references {{${path}}} out of scope`; + } + } + } + } + return null; +} + +/** Build a synthetic event with `field` (rooted form) set to `value`. */ +function syntheticEvent(assignments: Array<{ field: string; value: unknown }>): JetstreamEvent { + const event: JetstreamEvent = { + did: "did:plc:synthetic", + time_us: 1, + kind: "commit", + commit: { + rev: "r", + operation: "create", + collection: "test.collection", + rkey: "rk", + record: {}, + }, + }; + for (const { field, value } of assignments) { + const rooted = + field === "repo" + ? "event.did" + : field.startsWith("event.") + ? field + : `event.commit.record.${field}`; + const segments = rooted.slice("event.".length).split("."); + let target: Record = event as unknown as Record; + for (let i = 0; i < segments.length - 1; i++) { + const key = segments[i]!; + if (typeof target[key] !== "object" || target[key] === null) target[key] = {}; + target = target[key] as Record; + } + target[segments[segments.length - 1]!] = value; + } + return event; +} + +async function main() { + const rows = await db.query.automations.findMany(); + console.log(`Rehearsing conversion for ${rows.length} rows\n`); + + const failures: string[] = []; + const legacyHashes = new Map(); + const stepsHashes = new Map(); + let maxAssertionsPerGate = 0; + let maxStepsPerRow = 0; + + for (const row of rows) { + let steps: Step[]; + try { + steps = row.steps ?? convertLegacyAutomation(row); + } catch (err) { + failures.push(`(a) ${row.uri}: ${err instanceof Error ? err.message : String(err)}`); + continue; + } + + maxStepsPerRow = Math.max(maxStepsPerRow, countSteps(steps)); + for (const gate of triggerGateSteps(steps)) { + maxAssertionsPerGate = Math.max(maxAssertionsPerGate, gate.assertions.length); + } + + // (a) structure + scope + lexicon validation of the serialized record. + const scopeError = checkScopes(steps); + if (scopeError) failures.push(`(a) ${row.uri}: ${scopeError}`); + const record = { + $type: "run.airglow.automation", + name: row.name || "unnamed", + trigger: { + $type: "run.airglow.automation#pdsEventTrigger", + lexicon: row.lexicon, + operations: row.operations, + ...(row.wantedDids.length > 0 ? { wantedDids: row.wantedDids } : {}), + }, + steps: steps.map(toPdsStep), + active: row.active, + createdAt: new Date().toISOString(), + }; + const lexResult = lexicons.validate("run.airglow.automation", record); + if (!lexResult.success) { + failures.push(`(a) ${row.uri}: lexicon validation failed: ${String(lexResult.error)}`); + } + + // (b) verdict parity on a synthetic-event battery. + if (row.conditions.length > 0) { + const gates = triggerGateSteps(steps); + const battery: JetstreamEvent[] = []; + // All-match event: every condition's field set to its (inlined) value. + const inline = (v: string) => + v.replace(PLACEHOLDER_RE, (m, raw: string) => { + const name = raw.trim(); + if (name === "self") return row.did; + const variable = row.variables.find((x) => x.name === name); + return variable ? variable.value : m; + }); + battery.push( + syntheticEvent( + row.conditions.map((c) => ({ + field: c.field, + value: c.operator === "not-exists" ? undefined : inline(c.value ?? "x"), + })), + ), + ); + // Per-condition mismatch events. + for (const c of row.conditions) { + battery.push( + syntheticEvent( + row.conditions.map((other) => ({ + field: other.field, + value: + other === c + ? "___divergence___" + : other.operator === "not-exists" + ? undefined + : inline(other.value ?? "x"), + })), + ), + ); + } + // Missing-everything event. + battery.push(syntheticEvent([])); + + for (const event of battery) { + const legacy = matchConditions(event, row.conditions, row.did, row.variables); + const ctx = { + event, + automation: { did: row.did, rkey: row.rkey, name: row.name }, + values: {}, + itemStack: [], + runId: "rehearse", + }; + const modern = gates.every((g) => evaluateConditionStep(g, ctx)); + if (legacy !== modern) { + failures.push( + `(b) ${row.uri}: verdict divergence (legacy=${legacy}, steps=${modern}) on ${JSON.stringify(event.commit?.record)}`, + ); + break; + } + } + } + + // (c) hash partitions (legacy computed in memory: 22 prod rows have NULL + // stored hashes and the canonicalization changed by design). + legacyHashes.set( + row.uri, + computeBehaviorHash({ + lexicon: row.lexicon, + operations: row.operations, + conditions: row.conditions, + actions: row.actions, + fetches: row.fetches, + variables: row.variables, + wantedDids: row.wantedDids, + }), + ); + stepsHashes.set( + row.uri, + computeStepsBehaviorHash({ + lexicon: row.lexicon, + operations: row.operations, + wantedDids: row.wantedDids, + steps, + }), + ); + + // (d) scope parity. + const legacyScope = actionsNeedFullScope(row.actions); + const stepsScope = actionsNeedFullScope(flattenSteps(steps).filter(isActionStep) as never[]); + if (legacyScope !== stepsScope) { + failures.push( + `(d) ${row.uri}: scope divergence (legacy=${legacyScope}, steps=${stepsScope})`, + ); + } + } + + // (c) compare the two partitions. + const partition = (hashes: Map) => { + const groups = new Map(); + for (const [uri, hash] of hashes) { + const list = groups.get(hash) ?? []; + list.push(uri); + groups.set(hash, list); + } + return [...groups.values()] + .map((uris) => uris.sort().join("|")) + .sort() + .join("\n"); + }; + if (partition(legacyHashes) !== partition(stepsHashes)) { + failures.push("(c) behavior-hash partitions diverge between canonicalizations"); + } + + console.log(`corpus max assertions per trigger gate: ${maxAssertionsPerGate} (cap 10)`); + console.log(`corpus max steps per row: ${maxStepsPerRow} (cap 25)\n`); + + if (failures.length > 0) { + console.error(`${failures.length} failure(s):`); + for (const f of failures) console.error(` - ${f}`); + process.exit(1); + } + console.log(`All ${rows.length} rows pass classes (a)-(d).`); +} + +await main();