From f3dbae2ef89137cccb60381d029a2288fd61ea26 Mon Sep 17 00:00:00 2001 From: Kevin Rajan Date: Sat, 26 Sep 2026 04:38:06 +0000 Subject: [PATCH] fix(core): recover stale event sequence (#51411) Use the greater of the stored aggregate cursor and the latest persisted event sequence before assigning a local sequence so a stale cursor cannot retry an already-written sequence forever. --- packages/core/src/event.ts | 11 ++++++++-- packages/core/test/event.test.ts | 36 ++++++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+), 2 deletions(-) diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index c92ac0ac2ce3..4e29ae10e62d 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -3,7 +3,7 @@ export * as EventV2 from "./event" import { Cause, Context, Effect, Layer, Option, PubSub, Queue, 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, inArray } from "drizzle-orm" +import { and, asc, desc, eq, gt, inArray } from "drizzle-orm" import { Database } from "./database/database" import { EventSequenceTable, EventTable } from "./event/sql" import { Location } from "./location" @@ -246,7 +246,14 @@ export const layerWith = (options?: LayerOptions) => .where(eq(EventSequenceTable.aggregate_id, aggregateID)) .get() .pipe(Effect.orDie) - const latest = row?.seq ?? -1 + const latestEvent = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .orderBy(desc(EventTable.seq)) + .get() + .pipe(Effect.orDie) + const latest = Math.max(row?.seq ?? -1, latestEvent?.seq ?? -1) const encoded = Schema.encodeUnknownSync(definition.data)(event.data) as Record< string, unknown diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index e4329a2dde98..5201908936cd 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -419,6 +419,42 @@ describe("EventV2", () => { }), ) + it.effect("repairs a stale durable sequence before publishing", () => + 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: "first" }) + yield* events.publish(SyncMessage, { id: aggregateID, text: "second" }) + yield* events.claim(aggregateID, "owner-a") + yield* db + .update(EventSequenceTable) + .set({ seq: 0 }) + .where(eq(EventSequenceTable.aggregate_id, aggregateID)) + .run() + .pipe(Effect.orDie) + + const event = yield* events.publish(SyncMessage, { id: aggregateID, text: "recovered" }) + const rows = yield* db + .select({ seq: EventTable.seq }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, aggregateID)) + .all() + .pipe(Effect.orDie) + const sequence = yield* db + .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id }) + .from(EventSequenceTable) + .where(eq(EventSequenceTable.aggregate_id, aggregateID)) + .get() + .pipe(Effect.orDie) + + expect(event.durable?.seq).toBe(2) + expect(rows.map((row) => row.seq)).toEqual([0, 1, 2]) + expect(sequence).toEqual({ seq: 2, ownerID: "owner-a" }) + }), + ) + it.effect("replays durable aggregate events after a sequence and tails new events", () => Effect.gen(function* () { const events = yield* EventV2.Service