diff --git a/packages/core/schema.json b/packages/core/schema.json index ef6fef485..718b65eef 100644 --- a/packages/core/schema.json +++ b/packages/core/schema.json @@ -1,9 +1,9 @@ { "version": "7", "dialect": "sqlite", - "id": "54667ea5-ab5c-486b-b9d0-a398dcc360d6", + "id": "f984d12a-8384-4efa-a337-4a1518eba1f1", "prevIds": [ - "847009df-a964-45cd-b79f-6d07bf06e3e2" + "d9a74873-59fb-4487-9bf1-76477c0f831c" ], "ddl": [ { @@ -1406,6 +1406,26 @@ "entityType": "columns", "table": "session_receipt_operation" }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": "0", + "generated": null, + "name": "reserved_receipts", + "entityType": "columns", + "table": "session_receipt_operation" + }, + { + "type": "integer", + "notNull": true, + "autoincrement": false, + "default": "0", + "generated": null, + "name": "reserved_metadata_bytes", + "entityType": "columns", + "table": "session_receipt_operation" + }, { "type": "text", "notNull": true, diff --git a/packages/core/src/database/migration.gen.ts b/packages/core/src/database/migration.gen.ts index c4097baa8..19ca0c19f 100644 --- a/packages/core/src/database/migration.gen.ts +++ b/packages/core/src/database/migration.gen.ts @@ -45,5 +45,6 @@ export const migrations = ( import("./migration/20260828201050_normal_stryfe"), import("./migration/20260913205004_session-lineage"), import("./migration/20260913212936_session-receipt-storage"), + import("./migration/20260913221452_session-receipt-budget-reservation"), ]) ).map((module) => module.default) satisfies DatabaseMigration.Migration[] diff --git a/packages/core/src/database/migration/20260913221452_session-receipt-budget-reservation.ts b/packages/core/src/database/migration/20260913221452_session-receipt-budget-reservation.ts new file mode 100644 index 000000000..00ec72c10 --- /dev/null +++ b/packages/core/src/database/migration/20260913221452_session-receipt-budget-reservation.ts @@ -0,0 +1,32 @@ +import { Effect } from "effect" +import type { DatabaseMigration } from "../migration" + +export default { + id: "20260913221452_session-receipt-budget-reservation", + up(tx) { + return Effect.gen(function* () { + yield* tx.run(`PRAGMA foreign_keys=OFF;`) + yield* tx.run(` + CREATE TABLE \`__new_session_receipt_operation\` ( + \`id\` text PRIMARY KEY, + \`root_id\` text NOT NULL, + \`session_id\` text NOT NULL, + \`origin\` text NOT NULL, + \`reserved_receipts\` integer DEFAULT 0 NOT NULL, + \`reserved_metadata_bytes\` integer DEFAULT 0 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( + `INSERT INTO \`__new_session_receipt_operation\`(\`id\`, \`root_id\`, \`session_id\`, \`origin\`, \`reserved_receipts\`, \`reserved_metadata_bytes\`, \`state\`) SELECT \`id\`, \`root_id\`, \`session_id\`, \`origin\`, \`reserved_receipts\`, \`reserved_metadata_bytes\`, \`state\` FROM \`session_receipt_operation\`;`, + ) + yield* tx.run(`DROP TABLE \`session_receipt_operation\`;`) + yield* tx.run(`ALTER TABLE \`__new_session_receipt_operation\` RENAME TO \`session_receipt_operation\`;`) + yield* tx.run(`PRAGMA foreign_keys=ON;`) + yield* tx.run( + `CREATE INDEX \`session_receipt_operation_root_idx\` ON \`session_receipt_operation\` (\`root_id\`);`, + ) + }) + }, +} satisfies DatabaseMigration.Migration diff --git a/packages/core/src/database/schema.gen.ts b/packages/core/src/database/schema.gen.ts index 744437099..7ce18169e 100644 --- a/packages/core/src/database/schema.gen.ts +++ b/packages/core/src/database/schema.gen.ts @@ -233,6 +233,8 @@ export default { \`root_id\` text NOT NULL, \`session_id\` text NOT NULL, \`origin\` text NOT NULL, + \`reserved_receipts\` integer DEFAULT 0 NOT NULL, + \`reserved_metadata_bytes\` integer DEFAULT 0 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 ); diff --git a/packages/core/src/session/sql.ts b/packages/core/src/session/sql.ts index 383a9fb33..fb15326ef 100644 --- a/packages/core/src/session/sql.ts +++ b/packages/core/src/session/sql.ts @@ -120,6 +120,8 @@ export const SessionReceiptOperationTable = sqliteTable( .references(() => SessionTable.id, { onDelete: "cascade" }), session_id: text().$type().notNull(), origin: text().notNull(), + reserved_receipts: integer().notNull().default(0), + reserved_metadata_bytes: integer().notNull().default(0), state: text().$type<"prepared" | "evidence_ready" | "committed">().notNull(), }, (table) => [index("session_receipt_operation_root_idx").on(table.root_id)], diff --git a/packages/core/test/database-migration.test.ts b/packages/core/test/database-migration.test.ts index d4fcc7140..7fbb24e21 100644 --- a/packages/core/test/database-migration.test.ts +++ b/packages/core/test/database-migration.test.ts @@ -117,7 +117,7 @@ describe("DatabaseMigration", () => { 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')`, + sql`INSERT INTO session_receipt_operation (id, root_id, session_id, origin, reserved_receipts, reserved_metadata_bytes, state) VALUES ('operation', 'root', 'root', 'agent', 1, 64, '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)`, @@ -134,12 +134,19 @@ describe("DatabaseMigration", () => { yield* DatabaseMigration.apply(db) expect( yield* db.get(sql` - SELECT operation.state, receipt.creation_seq AS sequence, assessment.revision, assessment.expires_at AS expiresAt + SELECT operation.state, operation.reserved_receipts AS reservedReceipts, operation.reserved_metadata_bytes AS reservedMetadataBytes, 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 }) + ).toEqual({ + state: "committed", + reservedReceipts: 1, + reservedMetadataBytes: 64, + sequence: 1, + revision: 1, + expiresAt: 3, + }) }), ) }) diff --git a/packages/opencode/src/session/receipt.ts b/packages/opencode/src/session/receipt.ts index 3028f9b08..c6ee5c3c7 100644 --- a/packages/opencode/src/session/receipt.ts +++ b/packages/opencode/src/session/receipt.ts @@ -10,6 +10,11 @@ import { Effect } from "effect" import { SessionID } from "./schema" export namespace SessionReceipt { + export type Budget = { + maxReceipts: number + maxMetadataBytes: number + } + export type Fact = { id: string resource: string @@ -29,10 +34,15 @@ export namespace SessionReceipt { timeCreated: number } - export function publish( - database: Database.Interface, - input: { id: string; sessionID: SessionID; origin: string; receipts: ReadonlyArray }, - ) { + type ReservationInput = { + id: string + sessionID: SessionID + origin: string + receipts: ReadonlyArray + budget: Budget + } + + export function reserve(database: Database.Interface, input: ReservationInput) { return database.db.transaction( (tx) => Effect.gen(function* () { @@ -43,12 +53,30 @@ export namespace SessionReceipt { .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() + const receiptCount = input.receipts.length + const metadataBytes = metadataSize(input.receipts) + if (receiptCount > input.budget.maxReceipts || metadataBytes > input.budget.maxMetadataBytes) + return yield* Effect.fail(new Error(`Receipt budget exceeded for ${lineage.rootID}`)) + + const reservations = yield* tx + .select({ + receipts: SessionReceiptOperationTable.reserved_receipts, + metadataBytes: SessionReceiptOperationTable.reserved_metadata_bytes, + }) + .from(SessionReceiptOperationTable) + .where(eq(SessionReceiptOperationTable.root_id, lineage.rootID)) + .all() + const reservedReceipts = reservations.reduce( + (total, reservation) => total + reservation.receipts, + receiptCount, + ) + const reservedMetadataBytes = reservations.reduce( + (total, reservation) => total + reservation.metadataBytes, + metadataBytes, + ) + if (reservedReceipts > input.budget.maxReceipts || reservedMetadataBytes > input.budget.maxMetadataBytes) + return yield* Effect.fail(new Error(`Receipt budget exhausted for ${lineage.rootID}`)) + yield* tx .insert(SessionReceiptOperationTable) .values({ @@ -56,68 +84,119 @@ export namespace SessionReceipt { root_id: lineage.rootID, session_id: input.sessionID, origin: input.origin, + reserved_receipts: receiptCount, + reserved_metadata_bytes: metadataBytes, state: "prepared", }) .run() + }), + { behavior: "immediate" }, + ) + } + + export function abort(database: Database.Interface, id: string) { + return database.db.transaction( + (tx) => + Effect.gen(function* () { 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)) + .delete(SessionReceiptOperationTable) + .where(and(eq(SessionReceiptOperationTable.id, id), eq(SessionReceiptOperationTable.state, "prepared"))) .run() }), { behavior: "immediate" }, ) } - export function appendAssessment(database: Database.Interface, input: Assessment) { + export function commit(database: Database.Interface, input: { id: string; receipts: ReadonlyArray }) { return database.db .transaction( (tx) => Effect.gen(function* () { - const receipt = yield* tx - .select({ rootID: SessionReceiptTable.root_id }) + const reservation = yield* tx + .select() + .from(SessionReceiptOperationTable) + .where(eq(SessionReceiptOperationTable.id, input.id)) + .get() + if (!reservation || reservation.state !== "prepared") + return yield* Effect.fail(new Error(`No prepared receipt reservation ${input.id}`)) + if ( + reservation.reserved_receipts !== input.receipts.length || + reservation.reserved_metadata_bytes !== metadataSize(input.receipts) + ) + return yield* Effect.fail(new Error(`Receipt reservation mismatch for ${input.id}`)) + + const latest = yield* tx + .select({ sequence: SessionReceiptTable.creation_seq }) .from(SessionReceiptTable) - .where(eq(SessionReceiptTable.id, input.receiptID)) + .where(eq(SessionReceiptTable.root_id, reservation.root_id)) + .orderBy(desc(SessionReceiptTable.creation_seq)) .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, - }) + .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: reservation.root_id, + 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" }, ) + .pipe(Effect.tapError(() => abort(database, input.id))) + } + + export function publish(database: Database.Interface, input: ReservationInput) { + return Effect.gen(function* () { + yield* reserve(database, input) + yield* commit(database, input) + }) + } + + 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) { @@ -154,13 +233,23 @@ export namespace SessionReceipt { const operations = yield* database.db .select() .from(SessionReceiptOperationTable) - .where(and(eq(SessionReceiptOperationTable.root_id, lineage.rootID), eq(SessionReceiptOperationTable.state, "committed"))) + .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))) + .where( + inArray( + SessionReceiptTable.operation_id, + operations.map((operation) => operation.id), + ), + ) .orderBy(asc(SessionReceiptTable.creation_seq)) .all() : [] @@ -184,3 +273,7 @@ export namespace SessionReceipt { }) } } + +function metadataSize(receipts: ReadonlyArray) { + return new TextEncoder().encode(JSON.stringify(receipts)).byteLength +} diff --git a/packages/opencode/test/session/session.test.ts b/packages/opencode/test/session/session.test.ts index 8fc1a3de0..17d581288 100644 --- a/packages/opencode/test/session/session.test.ts +++ b/packages/opencode/test/session/session.test.ts @@ -213,16 +213,137 @@ describe("step-finish token propagation via event", () => { }) describe("Session", () => { + it.instance("reserves root-wide receipt and metadata capacity before publishing", () => + Effect.gen(function* () { + const session = yield* SessionNs.Service + const database = yield* Database.Service + const root = yield* session.create({ title: "budget root" }) + const budget = { maxReceipts: 2, maxMetadataBytes: 1_000 } + const first = [ + { id: "budget_first", resource: "file:///first", operation: "write", outcome: "applied", timeCreated: 1 }, + ] + const second = [ + { id: "budget_second", resource: "file:///second", operation: "write", outcome: "applied", timeCreated: 2 }, + ] + + yield* SessionReceipt.publish(database, { + id: "budget_op_first", + sessionID: root.id, + origin: "agent", + receipts: first, + budget, + }) + yield* SessionReceipt.publish(database, { + id: "budget_op_second", + sessionID: root.id, + origin: "agent", + receipts: second, + budget, + }) + + const receiptOverflow = yield* Effect.exit( + SessionReceipt.publish(database, { + id: "budget_op_receipt_overflow", + sessionID: root.id, + origin: "agent", + receipts: [ + { id: "budget_third", resource: "file:///third", operation: "write", outcome: "applied", timeCreated: 3 }, + ], + budget, + }), + ) + expect(receiptOverflow._tag).toBe("Failure") + expect(yield* SessionReceipt.committed(database, root.id)).toHaveLength(2) + + const metadataOverflow = yield* Effect.exit( + SessionReceipt.reserve(database, { + id: "budget_op_metadata_overflow", + sessionID: root.id, + origin: "agent", + receipts: [ + { + id: "budget_metadata", + resource: "file:///metadata", + operation: "write", + outcome: "applied", + timeCreated: 3, + }, + ], + budget: { maxReceipts: 3, maxMetadataBytes: 1 }, + }), + ) + expect(metadataOverflow._tag).toBe("Failure") + }), + ) + + it.instance("keeps concurrent reservations within a root budget", () => + Effect.gen(function* () { + const session = yield* SessionNs.Service + const database = yield* Database.Service + const root = yield* session.create({ title: "reservation root" }) + const budget = { maxReceipts: 1, maxMetadataBytes: 1_000 } + const reserve = (id: string) => + SessionReceipt.reserve(database, { + id, + sessionID: root.id, + origin: "agent", + receipts: [ + { id: `${id}_receipt`, resource: `file:///${id}`, operation: "write", outcome: "applied", timeCreated: 1 }, + ], + budget, + }) + + const reservations = yield* Effect.all( + [Effect.exit(reserve("reservation_one")), Effect.exit(reserve("reservation_two"))], + { + concurrency: "unbounded", + }, + ) + expect(reservations.filter((result) => result._tag === "Success")).toHaveLength(1) + }), + ) + + it.instance("releases aborted reservations without consuming root capacity", () => + Effect.gen(function* () { + const session = yield* SessionNs.Service + const database = yield* Database.Service + const root = yield* session.create({ title: "aborted reservation root" }) + const budget = { maxReceipts: 1, maxMetadataBytes: 1_000 } + const receipt = (id: string) => [ + { id, resource: `file:///${id}`, operation: "write", outcome: "applied", timeCreated: 1 }, + ] + + yield* SessionReceipt.reserve(database, { + id: "reservation_aborted", + sessionID: root.id, + origin: "agent", + receipts: receipt("reservation_aborted_receipt"), + budget, + }) + yield* SessionReceipt.abort(database, "reservation_aborted") + yield* SessionReceipt.reserve(database, { + id: "reservation_after_abort", + sessionID: root.id, + origin: "agent", + receipts: receipt("reservation_after_abort_receipt"), + budget, + }) + expect(yield* SessionReceipt.committed(database, root.id)).toEqual([]) + }), + ) + 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" }) + const budget = { maxReceipts: 4, maxMetadataBytes: 100_000 } yield* SessionReceipt.publish(database, { id: "op_first", sessionID: root.id, origin: "agent", + budget, 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 }, @@ -237,8 +358,22 @@ describe("Session", () => { 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 }, + { + 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, + }, ], }, ]) @@ -247,10 +382,20 @@ describe("Session", () => { id: "op_followup", sessionID: root.id, origin: "agent", - receipts: [{ id: "receipt_third", resource: "file:///third", operation: "delete", outcome: "applied", timeCreated: 3 }], + budget, + 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 }, + { + id: "receipt_third", + sequence: 3, + resource: "file:///third", + operation: "delete", + outcome: "applied", + timeCreated: 3, + }, ]) const duplicate = yield* Effect.exit( @@ -258,12 +403,32 @@ describe("Session", () => { id: "op_rolled_back", sessionID: root.id, origin: "agent", - receipts: [{ id: "receipt_first", resource: "file:///third", operation: "write", outcome: "applied", timeCreated: 3 }], + budget, + 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) + yield* SessionReceipt.publish(database, { + id: "op_after_failure", + sessionID: root.id, + origin: "agent", + budget, + receipts: [ + { + id: "receipt_after_failure", + resource: "file:///fourth", + operation: "write", + outcome: "applied", + timeCreated: 4, + }, + ], + }) + expect(yield* SessionReceipt.committed(database, root.id)).toHaveLength(3) + const rewrite = yield* Effect.exit( database.db .update(SessionReceiptTable) @@ -281,11 +446,21 @@ describe("Session", () => { const session = yield* SessionNs.Service const database = yield* Database.Service const root = yield* session.create({ title: "assessment root" }) + const budget = { maxReceipts: 100, maxMetadataBytes: 100_000 } 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 }], + budget, + receipts: [ + { + id: "receipt_assessed", + resource: "file:///assessed", + operation: "write", + outcome: "applied", + timeCreated: 1, + }, + ], }) yield* SessionReceipt.appendAssessment(database, {