Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 22 additions & 2 deletions packages/core/schema.json
Original file line number Diff line number Diff line change
@@ -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": [
{
Expand Down Expand Up @@ -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,
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/database/migration.gen.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
@@ -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
2 changes: 2 additions & 0 deletions packages/core/src/database/schema.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
);
Expand Down
2 changes: 2 additions & 0 deletions packages/core/src/session/sql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,8 @@ export const SessionReceiptOperationTable = sqliteTable(
.references(() => SessionTable.id, { onDelete: "cascade" }),
session_id: text().$type<SessionSchema.ID>().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)],
Expand Down
13 changes: 10 additions & 3 deletions packages/core/test/database-migration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)`,
Expand All @@ -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,
})
}),
)
})
Expand Down
199 changes: 146 additions & 53 deletions packages/opencode/src/session/receipt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -29,10 +34,15 @@ export namespace SessionReceipt {
timeCreated: number
}

export function publish(
database: Database.Interface,
input: { id: string; sessionID: SessionID; origin: string; receipts: ReadonlyArray<Fact> },
) {
type ReservationInput = {
id: string
sessionID: SessionID
origin: string
receipts: ReadonlyArray<Fact>
budget: Budget
}

export function reserve(database: Database.Interface, input: ReservationInput) {
return database.db.transaction(
(tx) =>
Effect.gen(function* () {
Expand All @@ -43,81 +53,150 @@ 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({
id: input.id,
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<Fact> }) {
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) {
Expand Down Expand Up @@ -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()
: []
Expand All @@ -184,3 +273,7 @@ export namespace SessionReceipt {
})
}
}

function metadataSize(receipts: ReadonlyArray<SessionReceipt.Fact>) {
return new TextEncoder().encode(JSON.stringify(receipts)).byteLength
}
Loading
Loading