diff --git a/packages/core/schema.json b/packages/core/schema.json index ff6946241e..ef6fef4855 100644 --- a/packages/core/schema.json +++ b/packages/core/schema.json @@ -1,9 +1,9 @@ { "version": "7", "dialect": "sqlite", - "id": "847009df-a964-45cd-b79f-6d07bf06e3e2", + "id": "54667ea5-ab5c-486b-b9d0-a398dcc360d6", "prevIds": [ - "68fc67ad-c3bd-4c0f-923d-db9817bf7475" + "847009df-a964-45cd-b79f-6d07bf06e3e2" ], "ddl": [ { @@ -82,6 +82,18 @@ "name": "session_message", "entityType": "tables" }, + { + "name": "session_receipt_assessment", + "entityType": "tables" + }, + { + "name": "session_receipt_operation", + "entityType": "tables" + }, + { + "name": "session_receipt", + "entityType": "tables" + }, { "name": "session", "entityType": "tables" @@ -1264,6 +1276,226 @@ "entityType": "columns", "table": "session_message" }, + { + "type": "text", + "notNull": false, + "autoincrement": false, + "default": null, + "generated": null, + "name": "id", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "receipt_id", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "root_id", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "confidence", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "net_state", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "evidence_state", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "revision", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "integer", + "notNull": false, + "autoincrement": false, + "default": null, + "generated": null, + "name": "expires_at", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "time_created", + "entityType": "columns", + "table": "session_receipt_assessment" + }, + { + "type": "text", + "notNull": false, + "autoincrement": false, + "default": null, + "generated": null, + "name": "id", + "entityType": "columns", + "table": "session_receipt_operation" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "root_id", + "entityType": "columns", + "table": "session_receipt_operation" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "session_id", + "entityType": "columns", + "table": "session_receipt_operation" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "origin", + "entityType": "columns", + "table": "session_receipt_operation" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "state", + "entityType": "columns", + "table": "session_receipt_operation" + }, + { + "type": "text", + "notNull": false, + "autoincrement": false, + "default": null, + "generated": null, + "name": "id", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "operation_id", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "root_id", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "creation_seq", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "resource", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "operation", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "text", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "outcome", + "entityType": "columns", + "table": "session_receipt" + }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": null, + "generated": null, + "name": "time_created", + "entityType": "columns", + "table": "session_receipt" + }, { "type": "text", "notNull": false, @@ -1874,6 +2106,81 @@ "entityType": "fks", "table": "session_message" }, + { + "columns": [ + "receipt_id" + ], + "tableTo": "session_receipt", + "columnsTo": [ + "id" + ], + "onUpdate": "NO ACTION", + "onDelete": "CASCADE", + "nameExplicit": false, + "name": "fk_session_receipt_assessment_receipt_id_session_receipt_id_fk", + "entityType": "fks", + "table": "session_receipt_assessment" + }, + { + "columns": [ + "root_id" + ], + "tableTo": "session", + "columnsTo": [ + "id" + ], + "onUpdate": "NO ACTION", + "onDelete": "CASCADE", + "nameExplicit": false, + "name": "fk_session_receipt_assessment_root_id_session_id_fk", + "entityType": "fks", + "table": "session_receipt_assessment" + }, + { + "columns": [ + "root_id" + ], + "tableTo": "session", + "columnsTo": [ + "id" + ], + "onUpdate": "NO ACTION", + "onDelete": "CASCADE", + "nameExplicit": false, + "name": "fk_session_receipt_operation_root_id_session_id_fk", + "entityType": "fks", + "table": "session_receipt_operation" + }, + { + "columns": [ + "operation_id" + ], + "tableTo": "session_receipt_operation", + "columnsTo": [ + "id" + ], + "onUpdate": "NO ACTION", + "onDelete": "CASCADE", + "nameExplicit": false, + "name": "fk_session_receipt_operation_id_session_receipt_operation_id_fk", + "entityType": "fks", + "table": "session_receipt" + }, + { + "columns": [ + "root_id" + ], + "tableTo": "session", + "columnsTo": [ + "id" + ], + "onUpdate": "NO ACTION", + "onDelete": "CASCADE", + "nameExplicit": false, + "name": "fk_session_receipt_root_id_session_id_fk", + "entityType": "fks", + "table": "session_receipt" + }, { "columns": [ "project_id" @@ -2103,6 +2410,33 @@ "table": "session_message", "entityType": "pks" }, + { + "columns": [ + "id" + ], + "nameExplicit": false, + "name": "session_receipt_assessment_pk", + "table": "session_receipt_assessment", + "entityType": "pks" + }, + { + "columns": [ + "id" + ], + "nameExplicit": false, + "name": "session_receipt_operation_pk", + "table": "session_receipt_operation", + "entityType": "pks" + }, + { + "columns": [ + "id" + ], + "nameExplicit": false, + "name": "session_receipt_pk", + "table": "session_receipt", + "entityType": "pks" + }, { "columns": [ "id" @@ -2443,6 +2777,84 @@ "entityType": "indexes", "table": "session_message" }, + { + "columns": [ + { + "value": "receipt_id", + "isExpression": false + }, + { + "value": "revision", + "isExpression": false + } + ], + "isUnique": true, + "where": null, + "origin": "manual", + "name": "session_receipt_assessment_receipt_revision_idx", + "entityType": "indexes", + "table": "session_receipt_assessment" + }, + { + "columns": [ + { + "value": "root_id", + "isExpression": false + } + ], + "isUnique": false, + "where": null, + "origin": "manual", + "name": "session_receipt_assessment_root_idx", + "entityType": "indexes", + "table": "session_receipt_assessment" + }, + { + "columns": [ + { + "value": "root_id", + "isExpression": false + } + ], + "isUnique": false, + "where": null, + "origin": "manual", + "name": "session_receipt_operation_root_idx", + "entityType": "indexes", + "table": "session_receipt_operation" + }, + { + "columns": [ + { + "value": "root_id", + "isExpression": false + }, + { + "value": "creation_seq", + "isExpression": false + } + ], + "isUnique": true, + "where": null, + "origin": "manual", + "name": "session_receipt_root_creation_seq_idx", + "entityType": "indexes", + "table": "session_receipt" + }, + { + "columns": [ + { + "value": "operation_id", + "isExpression": false + } + ], + "isUnique": false, + "where": null, + "origin": "manual", + "name": "session_receipt_operation_idx", + "entityType": "indexes", + "table": "session_receipt" + }, { "columns": [ { diff --git a/packages/core/script/migration.ts b/packages/core/script/migration.ts index 4f383f5f8f..cb1d2b5c92 100644 --- a/packages/core/script/migration.ts +++ b/packages/core/script/migration.ts @@ -147,12 +147,23 @@ export default { up(tx) { return Effect.gen(function* () { ${renderStatements(sql)} +${renderReceiptTriggers()} }) }, } satisfies Omit ` } +function renderReceiptTriggers() { + return [ + "CREATE TRIGGER session_receipt_immutable BEFORE UPDATE ON session_receipt BEGIN SELECT RAISE(ABORT, 'session receipt facts are immutable'); END;", + "CREATE TRIGGER session_receipt_assessment_append_only BEFORE UPDATE ON session_receipt_assessment BEGIN SELECT RAISE(ABORT, 'session receipt assessments are append-only'); END;", + "CREATE TRIGGER session_receipt_operation_state BEFORE UPDATE OF state ON session_receipt_operation WHEN NOT ((OLD.state = 'prepared' AND NEW.state = 'evidence_ready') OR (OLD.state = 'evidence_ready' AND NEW.state = 'committed')) BEGIN SELECT RAISE(ABORT, 'invalid session receipt operation state transition'); END;", + ] + .map((statement) => ` yield* tx.run(${JSON.stringify(statement)})`) + .join("\n") +} + function renderStatements(sql: string) { return sql .split("--> statement-breakpoint") diff --git a/packages/core/src/database/database.ts b/packages/core/src/database/database.ts index 4442393e54..b26af00642 100644 --- a/packages/core/src/database/database.ts +++ b/packages/core/src/database/database.ts @@ -61,6 +61,9 @@ const MERGE_TABLES = [ "session", "session_lineage", "session_lineage_origin", + "session_receipt_operation", + "session_receipt", + "session_receipt_assessment", "session_message", "session_input", "session_context_epoch", diff --git a/packages/core/src/database/migration.gen.ts b/packages/core/src/database/migration.gen.ts index 9d026c645b..c4097baa8a 100644 --- a/packages/core/src/database/migration.gen.ts +++ b/packages/core/src/database/migration.gen.ts @@ -44,5 +44,6 @@ export const migrations = ( import("./migration/20260820000001_add_session_directories"), import("./migration/20260828201050_normal_stryfe"), import("./migration/20260913205004_session-lineage"), + import("./migration/20260913212936_session-receipt-storage"), ]) ).map((module) => module.default) satisfies DatabaseMigration.Migration[] diff --git a/packages/core/src/database/migration/20260913212936_session-receipt-storage.ts b/packages/core/src/database/migration/20260913212936_session-receipt-storage.ts new file mode 100644 index 0000000000..81024a3b52 --- /dev/null +++ b/packages/core/src/database/migration/20260913212936_session-receipt-storage.ts @@ -0,0 +1,87 @@ +import { Effect } from "effect" +import type { DatabaseMigration } from "../migration" + +export default { + id: "20260913212936_session-receipt-storage", + up(tx) { + return Effect.gen(function* () { + yield* tx.run(` + CREATE TABLE \`session_receipt_assessment\` ( + \`id\` text PRIMARY KEY, + \`receipt_id\` text NOT NULL, + \`root_id\` text NOT NULL, + \`confidence\` text NOT NULL, + \`net_state\` text NOT NULL, + \`evidence_state\` text NOT NULL, + \`revision\` integer NOT NULL, + \`expires_at\` integer, + \`time_created\` integer NOT NULL, + CONSTRAINT \`fk_session_receipt_assessment_receipt_id_session_receipt_id_fk\` FOREIGN KEY (\`receipt_id\`) REFERENCES \`session_receipt\`(\`id\`) ON DELETE CASCADE, + CONSTRAINT \`fk_session_receipt_assessment_root_id_session_id_fk\` FOREIGN KEY (\`root_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE + ); + `) + yield* tx.run(` + CREATE TABLE \`session_receipt_operation\` ( + \`id\` text PRIMARY KEY, + \`root_id\` text NOT NULL, + \`session_id\` text NOT NULL, + \`origin\` text NOT NULL, + \`state\` text NOT NULL, + CONSTRAINT \`fk_session_receipt_operation_root_id_session_id_fk\` FOREIGN KEY (\`root_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE + ); + `) + yield* tx.run(` + CREATE TABLE \`session_receipt\` ( + \`id\` text PRIMARY KEY, + \`operation_id\` text NOT NULL, + \`root_id\` text NOT NULL, + \`creation_seq\` integer NOT NULL, + \`resource\` text NOT NULL, + \`operation\` text NOT NULL, + \`outcome\` text NOT NULL, + \`time_created\` integer NOT NULL, + CONSTRAINT \`fk_session_receipt_operation_id_session_receipt_operation_id_fk\` FOREIGN KEY (\`operation_id\`) REFERENCES \`session_receipt_operation\`(\`id\`) ON DELETE CASCADE, + CONSTRAINT \`fk_session_receipt_root_id_session_id_fk\` FOREIGN KEY (\`root_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE + ); + `) + yield* tx.run( + `CREATE UNIQUE INDEX \`session_receipt_assessment_receipt_revision_idx\` ON \`session_receipt_assessment\` (\`receipt_id\`,\`revision\`);`, + ) + yield* tx.run( + `CREATE INDEX \`session_receipt_assessment_root_idx\` ON \`session_receipt_assessment\` (\`root_id\`);`, + ) + yield* tx.run( + `CREATE INDEX \`session_receipt_operation_root_idx\` ON \`session_receipt_operation\` (\`root_id\`);`, + ) + yield* tx.run( + `CREATE UNIQUE INDEX \`session_receipt_root_creation_seq_idx\` ON \`session_receipt\` (\`root_id\`,\`creation_seq\`);`, + ) + yield* tx.run(`CREATE INDEX \`session_receipt_operation_idx\` ON \`session_receipt\` (\`operation_id\`);`) + yield* tx.run(` + CREATE TRIGGER \`session_receipt_immutable\` + BEFORE UPDATE ON \`session_receipt\` + BEGIN + SELECT RAISE(ABORT, 'session receipt facts are immutable'); + END; + `) + yield* tx.run(` + CREATE TRIGGER \`session_receipt_assessment_append_only\` + BEFORE UPDATE ON \`session_receipt_assessment\` + BEGIN + SELECT RAISE(ABORT, 'session receipt assessments are append-only'); + END; + `) + yield* tx.run(` + CREATE TRIGGER \`session_receipt_operation_state\` + BEFORE UPDATE OF \`state\` ON \`session_receipt_operation\` + WHEN NOT ( + (OLD.\`state\` = 'prepared' AND NEW.\`state\` = 'evidence_ready') + OR (OLD.\`state\` = 'evidence_ready' AND NEW.\`state\` = 'committed') + ) + BEGIN + SELECT RAISE(ABORT, 'invalid session receipt operation state transition'); + END; + `) + }) + }, +} satisfies DatabaseMigration.Migration diff --git a/packages/core/src/database/schema.gen.ts b/packages/core/src/database/schema.gen.ts index b687865f28..7444370993 100644 --- a/packages/core/src/database/schema.gen.ts +++ b/packages/core/src/database/schema.gen.ts @@ -212,6 +212,45 @@ export default { CONSTRAINT \`fk_session_message_session_id_session_id_fk\` FOREIGN KEY (\`session_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE ); `) + yield* tx.run(` + CREATE TABLE \`session_receipt_assessment\` ( + \`id\` text PRIMARY KEY, + \`receipt_id\` text NOT NULL, + \`root_id\` text NOT NULL, + \`confidence\` text NOT NULL, + \`net_state\` text NOT NULL, + \`evidence_state\` text NOT NULL, + \`revision\` integer NOT NULL, + \`expires_at\` integer, + \`time_created\` integer NOT NULL, + CONSTRAINT \`fk_session_receipt_assessment_receipt_id_session_receipt_id_fk\` FOREIGN KEY (\`receipt_id\`) REFERENCES \`session_receipt\`(\`id\`) ON DELETE CASCADE, + CONSTRAINT \`fk_session_receipt_assessment_root_id_session_id_fk\` FOREIGN KEY (\`root_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE + ); + `) + yield* tx.run(` + CREATE TABLE \`session_receipt_operation\` ( + \`id\` text PRIMARY KEY, + \`root_id\` text NOT NULL, + \`session_id\` text NOT NULL, + \`origin\` text NOT NULL, + \`state\` text NOT NULL, + CONSTRAINT \`fk_session_receipt_operation_root_id_session_id_fk\` FOREIGN KEY (\`root_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE + ); + `) + yield* tx.run(` + CREATE TABLE \`session_receipt\` ( + \`id\` text PRIMARY KEY, + \`operation_id\` text NOT NULL, + \`root_id\` text NOT NULL, + \`creation_seq\` integer NOT NULL, + \`resource\` text NOT NULL, + \`operation\` text NOT NULL, + \`outcome\` text NOT NULL, + \`time_created\` integer NOT NULL, + CONSTRAINT \`fk_session_receipt_operation_id_session_receipt_operation_id_fk\` FOREIGN KEY (\`operation_id\`) REFERENCES \`session_receipt_operation\`(\`id\`) ON DELETE CASCADE, + CONSTRAINT \`fk_session_receipt_root_id_session_id_fk\` FOREIGN KEY (\`root_id\`) REFERENCES \`session\`(\`id\`) ON DELETE CASCADE + ); + `) yield* tx.run(` CREATE TABLE \`session\` ( \`id\` text PRIMARY KEY, @@ -306,10 +345,32 @@ export default { `CREATE INDEX \`session_message_session_time_created_id_idx\` ON \`session_message\` (\`session_id\`,\`time_created\`,\`id\`);`, ) yield* tx.run(`CREATE INDEX \`session_message_time_created_idx\` ON \`session_message\` (\`time_created\`);`) + yield* tx.run( + `CREATE UNIQUE INDEX \`session_receipt_assessment_receipt_revision_idx\` ON \`session_receipt_assessment\` (\`receipt_id\`,\`revision\`);`, + ) + yield* tx.run( + `CREATE INDEX \`session_receipt_assessment_root_idx\` ON \`session_receipt_assessment\` (\`root_id\`);`, + ) + yield* tx.run( + `CREATE INDEX \`session_receipt_operation_root_idx\` ON \`session_receipt_operation\` (\`root_id\`);`, + ) + yield* tx.run( + `CREATE UNIQUE INDEX \`session_receipt_root_creation_seq_idx\` ON \`session_receipt\` (\`root_id\`,\`creation_seq\`);`, + ) + yield* tx.run(`CREATE INDEX \`session_receipt_operation_idx\` ON \`session_receipt\` (\`operation_id\`);`) yield* tx.run(`CREATE INDEX \`session_project_idx\` ON \`session\` (\`project_id\`);`) yield* tx.run(`CREATE INDEX \`session_workspace_idx\` ON \`session\` (\`workspace_id\`);`) yield* tx.run(`CREATE INDEX \`session_parent_idx\` ON \`session\` (\`parent_id\`);`) yield* tx.run(`CREATE INDEX \`todo_session_idx\` ON \`todo\` (\`session_id\`);`) + yield* tx.run( + "CREATE TRIGGER session_receipt_immutable BEFORE UPDATE ON session_receipt BEGIN SELECT RAISE(ABORT, 'session receipt facts are immutable'); END;", + ) + yield* tx.run( + "CREATE TRIGGER session_receipt_assessment_append_only BEFORE UPDATE ON session_receipt_assessment BEGIN SELECT RAISE(ABORT, 'session receipt assessments are append-only'); END;", + ) + yield* tx.run( + "CREATE TRIGGER session_receipt_operation_state BEFORE UPDATE OF state ON session_receipt_operation WHEN NOT ((OLD.state = 'prepared' AND NEW.state = 'evidence_ready') OR (OLD.state = 'evidence_ready' AND NEW.state = 'committed')) BEGIN SELECT RAISE(ABORT, 'invalid session receipt operation state transition'); END;", + ) }) }, } satisfies Omit diff --git a/packages/core/src/session/sql.ts b/packages/core/src/session/sql.ts index b0eea148ce..383a9fb332 100644 --- a/packages/core/src/session/sql.ts +++ b/packages/core/src/session/sql.ts @@ -109,6 +109,71 @@ export const SessionLineageOriginTable = sqliteTable( ], ) +/** Immutable root-owned facts for Files Changed operation publication. */ +export const SessionReceiptOperationTable = sqliteTable( + "session_receipt_operation", + { + id: text().primaryKey(), + root_id: text() + .$type() + .notNull() + .references(() => SessionTable.id, { onDelete: "cascade" }), + session_id: text().$type().notNull(), + origin: text().notNull(), + state: text().$type<"prepared" | "evidence_ready" | "committed">().notNull(), + }, + (table) => [index("session_receipt_operation_root_idx").on(table.root_id)], +) + +/** Immutable resource facts. Creation sequence is scoped to the lineage root. */ +export const SessionReceiptTable = sqliteTable( + "session_receipt", + { + id: text().primaryKey(), + operation_id: text() + .notNull() + .references(() => SessionReceiptOperationTable.id, { onDelete: "cascade" }), + root_id: text() + .$type() + .notNull() + .references(() => SessionTable.id, { onDelete: "cascade" }), + creation_seq: integer().notNull(), + resource: text().notNull(), + operation: text().notNull(), + outcome: text().notNull(), + time_created: integer().notNull(), + }, + (table) => [ + uniqueIndex("session_receipt_root_creation_seq_idx").on(table.root_id, table.creation_seq), + index("session_receipt_operation_idx").on(table.operation_id), + ], +) + +/** Append-only observations deliberately separate from immutable receipt facts. */ +export const SessionReceiptAssessmentTable = sqliteTable( + "session_receipt_assessment", + { + id: text().primaryKey(), + receipt_id: text() + .notNull() + .references(() => SessionReceiptTable.id, { onDelete: "cascade" }), + root_id: text() + .$type() + .notNull() + .references(() => SessionTable.id, { onDelete: "cascade" }), + confidence: text().notNull(), + net_state: text().notNull(), + evidence_state: text().notNull(), + revision: integer().notNull(), + expires_at: integer(), + time_created: integer().notNull(), + }, + (table) => [ + uniqueIndex("session_receipt_assessment_receipt_revision_idx").on(table.receipt_id, table.revision), + index("session_receipt_assessment_root_idx").on(table.root_id), + ], +) + export const MessageTable = sqliteTable( "message", { diff --git a/packages/core/test/database-migration.test.ts b/packages/core/test/database-migration.test.ts index b381cc7418..d4fcc71408 100644 --- a/packages/core/test/database-migration.test.ts +++ b/packages/core/test/database-migration.test.ts @@ -35,6 +35,9 @@ const run = (effect: Effect.Effect) => effect.pipe(Effect.provide(SqliteClient.layer({ filename: ":memory:", disableWAL: true })), Effect.scoped), ) +const runAtPath = (filename: string, effect: Effect.Effect) => + Effect.runPromise(effect.pipe(Effect.provide(SqliteClient.layer({ filename, disableWAL: true })), Effect.scoped)) + const makeDb = EffectDrizzleSqlite.makeWithDefaults() describe("DatabaseMigration", () => { @@ -99,6 +102,48 @@ describe("DatabaseMigration", () => { ) }) + test("preserves committed receipt groups, immutable receipts, and assessment history across restart", async () => { + await using tmp = await tmpdir() + const filename = path.join(tmp.path, "receipts.sqlite") + await runAtPath( + filename, + Effect.gen(function* () { + const db = yield* makeDb + yield* DatabaseMigration.apply(db) + yield* db.run( + sql`INSERT INTO project (id, worktree, time_created, time_updated, sandboxes) VALUES ('project', '/project', 1, 1, '[]')`, + ) + yield* db.run( + sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES ('root', 'project', 'root', '/project', 'Root', 'test', 1, 1)`, + ) + yield* db.run( + sql`INSERT INTO session_receipt_operation (id, root_id, session_id, origin, state) VALUES ('operation', 'root', 'root', 'agent', 'committed')`, + ) + yield* db.run( + sql`INSERT INTO session_receipt (id, operation_id, root_id, creation_seq, resource, operation, outcome, time_created) VALUES ('receipt', 'operation', 'root', 1, 'file:///root', 'write', 'applied', 2)`, + ) + yield* db.run( + sql`INSERT INTO session_receipt_assessment (id, receipt_id, root_id, confidence, net_state, evidence_state, revision, expires_at, time_created) VALUES ('assessment', 'receipt', 'root', 'verified', 'changed', 'available', 1, 3, 3)`, + ) + }), + ) + await runAtPath( + filename, + Effect.gen(function* () { + const db = yield* makeDb + yield* DatabaseMigration.apply(db) + expect( + yield* db.get(sql` + SELECT operation.state, receipt.creation_seq AS sequence, assessment.revision, assessment.expires_at AS expiresAt + FROM session_receipt_operation operation + JOIN session_receipt receipt ON receipt.operation_id = operation.id + JOIN session_receipt_assessment assessment ON assessment.receipt_id = receipt.id + `), + ).toEqual({ state: "committed", sequence: 1, revision: 1, expiresAt: 3 }) + }), + ) + }) + test("rejects a non-empty database without a session table", async () => { await expect( run( diff --git a/packages/opencode/src/session/receipt.ts b/packages/opencode/src/session/receipt.ts new file mode 100644 index 0000000000..3028f9b089 --- /dev/null +++ b/packages/opencode/src/session/receipt.ts @@ -0,0 +1,186 @@ +import { Database } from "@opencode-ai/core/database/database" +import { + SessionLineageTable, + SessionReceiptAssessmentTable, + SessionReceiptOperationTable, + SessionReceiptTable, +} from "@opencode-ai/core/session/sql" +import { and, asc, desc, eq, inArray } from "drizzle-orm" +import { Effect } from "effect" +import { SessionID } from "./schema" + +export namespace SessionReceipt { + export type Fact = { + id: string + resource: string + operation: string + outcome: string + timeCreated: number + } + + export type Assessment = { + id: string + receiptID: string + confidence: string + netState: string + evidenceState: string + revision: number + expiresAt?: number + timeCreated: number + } + + export function publish( + database: Database.Interface, + input: { id: string; sessionID: SessionID; origin: string; receipts: ReadonlyArray }, + ) { + return database.db.transaction( + (tx) => + Effect.gen(function* () { + const lineage = yield* tx + .select({ rootID: SessionLineageTable.root_id }) + .from(SessionLineageTable) + .where(eq(SessionLineageTable.session_id, input.sessionID)) + .get() + if (!lineage) return yield* Effect.fail(new Error(`Missing lineage root for ${input.sessionID}`)) + + const latest = yield* tx + .select({ sequence: SessionReceiptTable.creation_seq }) + .from(SessionReceiptTable) + .where(eq(SessionReceiptTable.root_id, lineage.rootID)) + .orderBy(desc(SessionReceiptTable.creation_seq)) + .get() + yield* tx + .insert(SessionReceiptOperationTable) + .values({ + id: input.id, + root_id: lineage.rootID, + session_id: input.sessionID, + origin: input.origin, + state: "prepared", + }) + .run() + yield* tx + .update(SessionReceiptOperationTable) + .set({ state: "evidence_ready" }) + .where(eq(SessionReceiptOperationTable.id, input.id)) + .run() + if (input.receipts.length > 0) + yield* tx + .insert(SessionReceiptTable) + .values( + input.receipts.map((receipt, index) => ({ + id: receipt.id, + operation_id: input.id, + root_id: lineage.rootID, + creation_seq: (latest?.sequence ?? 0) + index + 1, + resource: receipt.resource, + operation: receipt.operation, + outcome: receipt.outcome, + time_created: receipt.timeCreated, + })), + ) + .run() + yield* tx + .update(SessionReceiptOperationTable) + .set({ state: "committed" }) + .where(eq(SessionReceiptOperationTable.id, input.id)) + .run() + }), + { behavior: "immediate" }, + ) + } + + export function appendAssessment(database: Database.Interface, input: Assessment) { + return database.db + .transaction( + (tx) => + Effect.gen(function* () { + const receipt = yield* tx + .select({ rootID: SessionReceiptTable.root_id }) + .from(SessionReceiptTable) + .where(eq(SessionReceiptTable.id, input.receiptID)) + .get() + if (!receipt) return yield* Effect.fail(new Error(`Missing receipt ${input.receiptID}`)) + yield* tx + .insert(SessionReceiptAssessmentTable) + .values({ + id: input.id, + receipt_id: input.receiptID, + root_id: receipt.rootID, + confidence: input.confidence, + net_state: input.netState, + evidence_state: input.evidenceState, + revision: input.revision, + expires_at: input.expiresAt, + time_created: input.timeCreated, + }) + .run() + }), + { behavior: "immediate" }, + ) + } + + export function assessments(database: Database.Interface, receiptID: string) { + return database.db + .select() + .from(SessionReceiptAssessmentTable) + .where(eq(SessionReceiptAssessmentTable.receipt_id, receiptID)) + .orderBy(asc(SessionReceiptAssessmentTable.revision)) + .all() + .pipe( + Effect.map((assessments) => + assessments.map((assessment) => ({ + id: assessment.id, + receiptID: assessment.receipt_id, + confidence: assessment.confidence, + netState: assessment.net_state, + evidenceState: assessment.evidence_state, + revision: assessment.revision, + ...(assessment.expires_at === null ? {} : { expiresAt: assessment.expires_at }), + timeCreated: assessment.time_created, + })), + ), + ) + } + + export function committed(database: Database.Interface, sessionID: SessionID) { + return Effect.gen(function* () { + const lineage = yield* database.db + .select({ rootID: SessionLineageTable.root_id }) + .from(SessionLineageTable) + .where(eq(SessionLineageTable.session_id, sessionID)) + .get() + if (!lineage) return [] + const operations = yield* database.db + .select() + .from(SessionReceiptOperationTable) + .where(and(eq(SessionReceiptOperationTable.root_id, lineage.rootID), eq(SessionReceiptOperationTable.state, "committed"))) + .all() + const receipts = operations.length + ? yield* database.db + .select() + .from(SessionReceiptTable) + .where(inArray(SessionReceiptTable.operation_id, operations.map((operation) => operation.id))) + .orderBy(asc(SessionReceiptTable.creation_seq)) + .all() + : [] + return operations.map((operation) => ({ + id: operation.id, + rootID: operation.root_id, + sessionID: operation.session_id, + origin: operation.origin, + state: operation.state, + receipts: receipts + .filter((receipt) => receipt.operation_id === operation.id) + .map((receipt) => ({ + id: receipt.id, + sequence: receipt.creation_seq, + resource: receipt.resource, + operation: receipt.operation, + outcome: receipt.outcome, + timeCreated: receipt.time_created, + })), + })) + }) + } +} diff --git a/packages/opencode/test/session/session.test.ts b/packages/opencode/test/session/session.test.ts index 0870a747f6..8fc1a3de08 100644 --- a/packages/opencode/test/session/session.test.ts +++ b/packages/opencode/test/session/session.test.ts @@ -1,7 +1,7 @@ import { describe, expect } from "bun:test" import { SessionV1 } from "@opencode-ai/core/v1/session" import { Database } from "@opencode-ai/core/database/database" -import { SessionLineageTable } from "@opencode-ai/core/session/sql" +import { SessionLineageTable, SessionReceiptAssessmentTable, SessionReceiptTable } from "@opencode-ai/core/session/sql" import { EventV2 } from "@opencode-ai/core/event" import { SessionProjector } from "@opencode-ai/core/session/projector" import { Deferred, Effect, Exit, Layer } from "effect" @@ -19,6 +19,7 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { InstanceStore } from "@/project/instance-store" import { InstanceBootstrap } from "@/project/bootstrap" import { ExternalDiff } from "@/session/external-diff" +import { SessionReceipt } from "@/session/receipt" import path from "path" import { eq } from "drizzle-orm" @@ -212,6 +213,134 @@ describe("step-finish token propagation via event", () => { }) describe("Session", () => { + it.instance("atomically publishes one committed receipt group for a lineage root", () => + Effect.gen(function* () { + const session = yield* SessionNs.Service + const database = yield* Database.Service + const root = yield* session.create({ title: "receipt root" }) + + yield* SessionReceipt.publish(database, { + id: "op_first", + sessionID: root.id, + origin: "agent", + receipts: [ + { id: "receipt_first", resource: "file:///first", operation: "write", outcome: "applied", timeCreated: 1 }, + { id: "receipt_second", resource: "file:///second", operation: "write", outcome: "applied", timeCreated: 2 }, + ], + }) + + expect(yield* SessionReceipt.committed(database, root.id)).toEqual([ + { + id: "op_first", + rootID: root.id, + sessionID: root.id, + origin: "agent", + state: "committed", + receipts: [ + { id: "receipt_first", sequence: 1, resource: "file:///first", operation: "write", outcome: "applied", timeCreated: 1 }, + { id: "receipt_second", sequence: 2, resource: "file:///second", operation: "write", outcome: "applied", timeCreated: 2 }, + ], + }, + ]) + + yield* SessionReceipt.publish(database, { + id: "op_followup", + sessionID: root.id, + origin: "agent", + receipts: [{ id: "receipt_third", resource: "file:///third", operation: "delete", outcome: "applied", timeCreated: 3 }], + }) + expect((yield* SessionReceipt.committed(database, root.id))[1]?.receipts).toEqual([ + { id: "receipt_third", sequence: 3, resource: "file:///third", operation: "delete", outcome: "applied", timeCreated: 3 }, + ]) + + const duplicate = yield* Effect.exit( + SessionReceipt.publish(database, { + id: "op_rolled_back", + sessionID: root.id, + origin: "agent", + receipts: [{ id: "receipt_first", resource: "file:///third", operation: "write", outcome: "applied", timeCreated: 3 }], + }), + ) + expect(duplicate._tag).toBe("Failure") + expect(yield* SessionReceipt.committed(database, root.id)).toHaveLength(2) + + const rewrite = yield* Effect.exit( + database.db + .update(SessionReceiptTable) + .set({ outcome: "rewritten" }) + .where(eq(SessionReceiptTable.id, "receipt_first")) + .run(), + ) + expect(rewrite._tag).toBe("Failure") + expect((yield* SessionReceipt.committed(database, root.id))[0]?.receipts[0]?.outcome).toBe("applied") + }), + ) + + it.instance("appends receipt assessments without rewriting immutable facts", () => + Effect.gen(function* () { + const session = yield* SessionNs.Service + const database = yield* Database.Service + const root = yield* session.create({ title: "assessment root" }) + yield* SessionReceipt.publish(database, { + id: "op_assessed", + sessionID: root.id, + origin: "agent", + receipts: [{ id: "receipt_assessed", resource: "file:///assessed", operation: "write", outcome: "applied", timeCreated: 1 }], + }) + + yield* SessionReceipt.appendAssessment(database, { + id: "assessment_first", + receiptID: "receipt_assessed", + confidence: "observed", + netState: "changed", + evidenceState: "available", + revision: 1, + expiresAt: 10, + timeCreated: 2, + }) + yield* SessionReceipt.appendAssessment(database, { + id: "assessment_second", + receiptID: "receipt_assessed", + confidence: "verified", + netState: "restored", + evidenceState: "unavailable", + revision: 2, + timeCreated: 3, + }) + + expect(yield* SessionReceipt.assessments(database, "receipt_assessed")).toEqual([ + { + id: "assessment_first", + receiptID: "receipt_assessed", + confidence: "observed", + netState: "changed", + evidenceState: "available", + revision: 1, + expiresAt: 10, + timeCreated: 2, + }, + { + id: "assessment_second", + receiptID: "receipt_assessed", + confidence: "verified", + netState: "restored", + evidenceState: "unavailable", + revision: 2, + timeCreated: 3, + }, + ]) + const rewrite = yield* Effect.exit( + database.db + .update(SessionReceiptAssessmentTable) + .set({ net_state: "rewritten" }) + .where(eq(SessionReceiptAssessmentTable.id, "assessment_first")) + .run(), + ) + expect(rewrite._tag).toBe("Failure") + expect((yield* SessionReceipt.committed(database, root.id))[0]?.receipts[0]?.outcome).toBe("applied") + }), + ) + it.instance("creates one full lineage root for each new session", () => Effect.gen(function* () { const session = yield* SessionNs.Service