diff --git a/.specgit.yaml b/.specgit.yaml index 6b3d72ca40..63d268f986 100644 --- a/.specgit.yaml +++ b/.specgit.yaml @@ -1,8 +1,8 @@ version: 1 -delivery: summary-diff-continue +delivery: event-idempotency-gate context: kind: branch - branch: fix/525-summary-diff-continue + branch: fix/523-event-idempotency-gate issues: - - 525 -pr: 526 + - 523 +pr: 527 diff --git a/packages/core/schema.json b/packages/core/schema.json index 2afa95c92a..96ddf9fd6a 100644 --- a/packages/core/schema.json +++ b/packages/core/schema.json @@ -1,9 +1,9 @@ { "version": "7", "dialect": "sqlite", - "id": "abadf28b-1770-46c6-bbf6-b7800a9ca874", + "id": "7a2e2a70-584a-4604-bf73-4c6e116c20a3", "prevIds": [ - "874d8e74-d354-4dcb-b98c-c893660c9371" + "abadf28b-1770-46c6-bbf6-b7800a9ca874" ], "ddl": [ { @@ -1052,6 +1052,16 @@ "entityType": "columns", "table": "event" }, + { + "type": "text", + "notNull": false, + "autoincrement": false, + "default": null, + "generated": null, + "name": "data_hash", + "entityType": "columns", + "table": "event" + }, { "type": "text", "notNull": false, diff --git a/packages/core/src/database/migration.gen.ts b/packages/core/src/database/migration.gen.ts index c141dc82ad..de37ec19ab 100644 --- a/packages/core/src/database/migration.gen.ts +++ b/packages/core/src/database/migration.gen.ts @@ -57,5 +57,6 @@ export const migrations = ( import("./migration/20260815044858_dag_graph_rev_view"), import("./migration/20260815083000_workflow_directory_convergence"), import("./migration/20260903044702_drop_session_summary_diffs"), + import("./migration/20260903062324_add_event_data_hash"), ]) ).map((module) => module.default) satisfies DatabaseMigration.Migration[] diff --git a/packages/core/src/database/migration/20260903062324_add_event_data_hash.ts b/packages/core/src/database/migration/20260903062324_add_event_data_hash.ts new file mode 100644 index 0000000000..a2ca39d724 --- /dev/null +++ b/packages/core/src/database/migration/20260903062324_add_event_data_hash.ts @@ -0,0 +1,11 @@ +import { Effect } from "effect" +import type { DatabaseMigration } from "../migration" + +export default { + id: "20260903062324_add_event_data_hash", + up(tx) { + return Effect.gen(function* () { + yield* tx.run(`ALTER TABLE \`event\` ADD \`data_hash\` text;`) + }) + }, +} satisfies DatabaseMigration.Migration diff --git a/packages/core/src/database/schema.gen.ts b/packages/core/src/database/schema.gen.ts index 851962d2eb..293e89b4b0 100644 --- a/packages/core/src/database/schema.gen.ts +++ b/packages/core/src/database/schema.gen.ts @@ -149,6 +149,7 @@ export default { \`seq\` integer NOT NULL, \`type\` text NOT NULL, \`data\` text NOT NULL, + \`data_hash\` text, CONSTRAINT \`fk_event_aggregate_id_event_sequence_aggregate_id_fk\` FOREIGN KEY (\`aggregate_id\`) REFERENCES \`event_sequence\`(\`aggregate_id\`) ON DELETE CASCADE ); `) diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index ec80977a89..d5ee2e7757 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -3,13 +3,14 @@ export * as EventV2 from "./event" import { Cause, Context, Effect, FiberSet, Layer, Option, PubSub, Schema, Stream } from "effect" import { Event } from "@opencode-ai/schema/event" import type { Data, Definition, Payload } from "@opencode-ai/schema/event" -import { and, asc, eq, gt } from "drizzle-orm" +import { and, asc, desc, eq, gt } from "drizzle-orm" import { Database } from "./database/database" import { EventSequenceTable, EventTable } from "./event/sql" import { Location } from "./location" import { LayerNode } from "./effect/layer-node" import { isDeepStrictEqual } from "node:util" import { Durable } from "@opencode-ai/schema/durable-event-manifest" +import { Hash } from "./util/hash" export const ID = Event.ID export type ID = import("@opencode-ai/schema/event").ID @@ -227,6 +228,30 @@ export const layerWith = (options?: LayerOptions) => if (input && row?.ownerID && row.ownerID !== input.ownerID) { return undefined } + const dataHash = Hash.sha256(JSON.stringify(encoded)) + if (!input) { + // Idempotency gate (#523): a fresh append that byte-for-byte repeats + // this aggregate's latest same-type event carries zero information + // delta. Skip it entirely — no seq consumed, no projectors, no commit + // hook, no durable wake — so the persisted sequence stays dense and + // both the replayAll contiguity check and gt(seq, after) readers are + // unaffected. Replay appends (input) keep their exact-seq contract, + // and legacy rows carry a NULL hash so they never match. + const previous = yield* db + .select({ dataHash: EventTable.data_hash }) + .from(EventTable) + .where( + and( + eq(EventTable.aggregate_id, aggregateID), + eq(EventTable.type, versionedType(definition.type, durable.version)), + ), + ) + .orderBy(desc(EventTable.seq)) + .limit(1) + .get() + .pipe(Effect.orDie) + if (previous && previous.dataHash === dataHash) return undefined + } const seq = input?.seq ?? latest + 1 if (input && seq !== latest + 1) { yield* Effect.die( @@ -278,6 +303,7 @@ export const layerWith = (options?: LayerOptions) => seq, type: versionedType(definition.type, durable.version), data: encoded, + data_hash: dataHash, }, ]) .run() @@ -479,7 +505,9 @@ export const layerWith = (options?: LayerOptions) => .transaction( () => Effect.gen(function* () { - const results = new Array<{ aggregateID: string; seq: number }>() + // Aligned with entries by index: a deduped entry yields + // undefined so the payload pairing below stays positional. + const results = new Array<{ aggregateID: string; seq: number } | undefined>() for (const entry of entries) { // No replay input: seq is allocated contiguously from the latest sequence inside the transaction. const result = yield* commitDurableEventInner( @@ -488,7 +516,7 @@ export const layerWith = (options?: LayerOptions) => undefined, entry.commit, ) - if (result) results.push(result) + results.push(result) } return results }), @@ -499,19 +527,19 @@ export const layerWith = (options?: LayerOptions) => return results }), ) - const payloads = entries.flatMap((entry, index) => { + const payloads = entries.map((entry, index) => { const result = committed[index] - if (!result) return [] - return [ - { - ...entry.event, - durable: { - aggregateID: result.aggregateID, - seq: result.seq, - version: entry.durable.version, - }, - } as Payload, - ] + // A deduped entry is still notified (mirrors the single-publish + // path) but stays unstamped: it occupies no sequence position. + if (!result) return entry.event + return { + ...entry.event, + durable: { + aggregateID: result.aggregateID, + seq: result.seq, + version: entry.durable.version, + }, + } as Payload }) for (const payload of payloads) { yield* notify(payload) diff --git a/packages/core/src/event/sql.ts b/packages/core/src/event/sql.ts index 38fe34f1e3..74f868e971 100644 --- a/packages/core/src/event/sql.ts +++ b/packages/core/src/event/sql.ts @@ -17,6 +17,11 @@ export const EventTable = sqliteTable( seq: integer().notNull(), type: text().notNull(), data: text({ mode: "json" }).$type>().notNull(), + // sha256 of the serialized payload, written once at append time. The + // idempotency gate compares against the latest same-type row via + // event_aggregate_type_seq_idx instead of re-hashing MiB-scale payloads. + // Nullable: legacy rows predate the column and never match the gate. + data_hash: text(), }, (table) => [ uniqueIndex("event_aggregate_seq_idx").on(table.aggregate_id, table.seq), diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index 7034c1cccc..c6722ce2f6 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -10,7 +10,7 @@ import { EventSequenceTable, EventTable } from "@opencode-ai/core/event/sql" import { Location } from "@opencode-ai/core/location" import { AbsolutePath } from "@opencode-ai/core/schema" import { WorkspaceV2 } from "@opencode-ai/core/workspace" -import { eq } from "drizzle-orm" +import { asc, eq } from "drizzle-orm" import { location } from "./fixture/location" import { testEffect } from "./lib/effect" @@ -382,6 +382,238 @@ describe("EventV2", () => { }), ) + it.effect("skips a byte-identical duplicate of the latest same-type durable event", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = EventV2.ID.create() + + const first = yield* events.publish(SyncMessage, { id: aggregateID, text: "hello" }) + const duplicate = yield* events.publish(SyncMessage, { id: aggregateID, text: "hello" }) + const rows = yield* db + .select({ seq: EventTable.seq, dataHash: EventTable.data_hash }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + + expect(first.durable?.seq).toBe(0) + expect(duplicate.durable).toBeUndefined() + expect(rows).toHaveLength(1) + expect(rows[0]?.dataHash).toHaveLength(64) + }), + ) + + it.effect("keeps the persisted sequence dense when a duplicate is skipped", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = EventV2.ID.create() + + yield* events.publish(SyncMessage, { id: aggregateID, text: "hello" }) + yield* events.publish(SyncMessage, { id: aggregateID, text: "hello" }) + const third = yield* events.publish(SyncMessage, { id: aggregateID, text: "world" }) + const rows = yield* db + .select({ seq: EventTable.seq, data: EventTable.data }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .orderBy(asc(EventTable.seq)) + .all() + .pipe(Effect.orDie) + + expect(third.durable?.seq).toBe(1) + expect(rows.map((row) => [row.seq, row.data["text"]])).toEqual([ + [0, "hello"], + [1, "world"], + ]) + }), + ) + + it.effect("appends when the payload differs from the latest same-type event", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = EventV2.ID.create() + + yield* events.publish(SyncMessage, { id: aggregateID, text: "a" }) + yield* events.publish(SyncMessage, { id: aggregateID, text: "b" }) + yield* events.publish(SyncMessage, { id: aggregateID, text: "a" }) + const rows = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + + expect(rows).toHaveLength(3) + }), + ) + + it.effect("dedupes only against the same aggregate and event type", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = EventV2.ID.create() + const otherAggregateID = EventV2.ID.create() + + yield* events.publish(SyncMessage, { id: aggregateID, text: "same" }) + yield* events.publish(SyncSent, { messageID: aggregateID, text: "same" }) + yield* events.publish(SyncMessage, { id: otherAggregateID, text: "same" }) + const own = yield* db + .select({ type: EventTable.type }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + const other = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, otherAggregateID)) + .all() + .pipe(Effect.orDie) + + expect(new Set(own.map((row) => row.type))).toHaveLength(2) + expect(other).toHaveLength(1) + }), + ) + + it.effect("skips duplicates inside a publishMany batch and keeps payloads aligned", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = EventV2.ID.create() + + const payloads = yield* events.publishMany([ + { definition: SyncMessage, data: { id: aggregateID, text: "a" } }, + { definition: SyncMessage, data: { id: aggregateID, text: "a" } }, + { definition: SyncMessage, data: { id: aggregateID, text: "b" } }, + ]) + const rows = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + + expect(rows).toHaveLength(2) + expect(payloads.map((event) => event.data)).toEqual([ + { id: aggregateID, text: "a" }, + { id: aggregateID, text: "a" }, + { id: aggregateID, text: "b" }, + ]) + expect(payloads.map((event) => event.durable?.seq)).toEqual([0, undefined, 1]) + }), + ) + + it.effect("runs projectors and commit hooks only for persisted appends", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const projected = new Array() + yield* events.project(SyncMessage, (event) => + Effect.sync(() => { + projected.push(event) + }), + ) + const commits = new Array() + const aggregateID = EventV2.ID.create() + const publishWithCommit = () => + events.publish( + SyncMessage, + { id: aggregateID, text: "hello" }, + { commit: (seq) => Effect.sync(() => commits.push(seq)) }, + ) + + yield* publishWithCommit() + yield* publishWithCommit() + + expect(projected.map((event) => event.durable?.seq)).toEqual([0]) + expect(commits).toEqual([0]) + }), + ) + + it.effect("never dedupes against legacy rows with a NULL hash", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = EventV2.ID.create() + + yield* db + .insert(EventSequenceTable) + .values([{ aggregate_id: aggregateID, seq: 0 }]) + .run() + .pipe(Effect.orDie) + yield* db + .insert(EventTable) + .values([ + { + id: EventV2.ID.create(), + aggregate_id: aggregateID, + seq: 0, + type: EventV2.versionedType(SyncMessage.type, 1), + data: { id: aggregateID, text: "legacy" }, + }, + ]) + .run() + .pipe(Effect.orDie) + + const published = yield* events.publish(SyncMessage, { id: aggregateID, text: "legacy" }) + const rows = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + + expect(published.durable?.seq).toBe(1) + expect(rows).toHaveLength(2) + }), + ) + + it.effect("replay with an explicit seq is never deduped", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const { db } = yield* Database.Service + const aggregateID = Session.ID.create() + + yield* events.publish(DurableMessage, durableData(aggregateID, "same")) + yield* events.replay({ + id: EventV2.ID.create(), + type: EventV2.versionedType(DurableMessage.type, 1), + seq: 1, + aggregateID, + data: durableData(aggregateID, "same"), + }) + const rows = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + + expect(rows).toHaveLength(2) + }), + ) + + it.effect("durable readers observe only persisted events after a skipped duplicate", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const aggregateID = Session.ID.create() + yield* events.publish(DurableMessage, durableData(aggregateID, "zero")) + const fiber = yield* events + .durable({ aggregateID }) + .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped) + yield* Effect.yieldNow + + yield* events.publish(DurableMessage, durableData(aggregateID, "zero")) + yield* events.publish(DurableMessage, durableData(aggregateID, "one")) + + expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([ + [0, durableData(aggregateID, "zero")], + [1, durableData(aggregateID, "one")], + ]) + }), + ) + it.effect("replays durable aggregate events after a sequence and tails new events", () => Effect.gen(function* () { const events = yield* EventV2.Service diff --git a/packages/opencode/src/server/routes/instance/httpapi/handlers/sync.ts b/packages/opencode/src/server/routes/instance/httpapi/handlers/sync.ts index 28fd245a63..4bae9e965f 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/handlers/sync.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/handlers/sync.ts @@ -72,7 +72,13 @@ export const syncHandlers = HttpApiBuilder.group(InstanceHttpApi, "sync", (handl const history = Effect.fn("SyncHttpApi.history")(function* (ctx: { payload: typeof HistoryPayload.Type }) { const exclude = Object.entries(ctx.payload) return yield* db - .select() + .select({ + id: EventTable.id, + aggregate_id: EventTable.aggregate_id, + seq: EventTable.seq, + type: EventTable.type, + data: EventTable.data, + }) .from(EventTable) .where( exclude.length > 0