From 95f264e04eca61b4a81ef92ee29476556d299a68 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 25 Jun 2026 11:35:11 -0400 Subject: [PATCH 1/3] refactor(schema): distinguish published event durability --- CONTEXT.md | 1 + packages/core/src/event.ts | 108 ++-- packages/core/test/event.test.ts | 4 +- .../server/httpapi-public-openapi.test.ts | 10 + .../test/v2/session-message-updater.test.ts | 17 + packages/protocol/src/groups/event.ts | 18 +- packages/schema/src/event.ts | 98 ++-- packages/schema/test/event.test.ts | 95 ++++ packages/sdk/js/src/v2/gen/types.gen.ts | 475 ++++-------------- 9 files changed, 373 insertions(+), 453 deletions(-) diff --git a/CONTEXT.md b/CONTEXT.md index 97919b63f6..d379facbba 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -141,6 +141,7 @@ _Avoid_: Response envelope - Promise streaming methods return a lazy `AsyncIterable` directly rather than a Promise-wrapped stream object. Iteration opens the connection, `AbortSignal` cancels it, and ending iteration closes the underlying request; the Effect emitter analogously returns `Stream` directly. - Promise SSE connection establishment, declared HTTP failures, and infrastructure failures occur during `AsyncIterable` iteration, beginning with its first `next()` call, rather than during synchronous method construction. - Neither generated streaming runtime automatically reconnects after disconnection. Promise `AsyncIterable` and Effect `Stream` fail explicitly; live consumers refresh and resubscribe, while durable sequence-based resume remains explicit composition above the generated client. +- Event definition durability is authoritative for published payloads. Durable definitions publish and decode only with commit metadata (`aggregateID`, `seq`, and `version`); live definitions forbid that metadata. Core's pre-commit payload without an assigned sequence is a separate internal type and never reaches subscribers or projectors. - Promise client construction is synchronous and network-free. It requires `baseUrl`, defaults to `globalThis.fetch`, accepts client-level headers, and merges them with per-call header overrides. - Effect client construction accepts an explicit `baseUrl` and obtains `HttpClient.HttpClient` from the Effect environment. It does not install fetch or duplicate per-call transport policy; callers transform/provide the client for headers, tracing, retries, recording, and tests, while fiber interruption owns cancellation. - Promise and Effect emitters each own their generated public type modules. The **SDK Contract IR**, not a physically shared generated type package, is the common source; this permits zero-Effect wire types and rich decoded Effect types to evolve independently. diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index 132a88b111..d8ec020627 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -2,7 +2,14 @@ export * as EventV2 from "./event" import { Cause, Context, Effect, 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 type { + Data, + Definition, + DurableDefinition, + Payload, + PublishedPayload, + UncommittedPayload, +} from "@opencode-ai/schema/event" import { and, asc, eq, gt } from "drizzle-orm" import { Database } from "./database/database" import { EventSequenceTable, EventTable } from "./event/sql" @@ -13,7 +20,14 @@ import { Durable } from "@opencode-ai/schema/durable-event-manifest" export const ID = Event.ID export type ID = import("@opencode-ai/schema/event").ID -export type { Data, Definition, Payload } from "@opencode-ai/schema/event" +export type { + Data, + Definition, + DurableDefinition, + Payload, + PublishedPayload, + UncommittedPayload, +} from "@opencode-ai/schema/event" export type Subscriber = (event: Payload) => Effect.Effect export type Unsubscribe = Effect.Effect @@ -66,10 +80,13 @@ export interface Interface { ) => Effect.Effect> readonly subscribe: (definition: D) => Stream.Stream> readonly all: () => Stream.Stream - readonly durable: (input: { readonly aggregateID: string; readonly after?: number }) => Stream.Stream + readonly durable: (input: { + readonly aggregateID: string + readonly after?: number + }) => Stream.Stream> /** @deprecated Use `all()` and consume the returned stream. */ readonly listen: (listener: Subscriber) => Effect.Effect - readonly project: (definition: D, projector: Subscriber) => Effect.Effect + readonly project: (definition: D, projector: Subscriber) => Effect.Effect readonly replay: ( event: SerializedEvent, options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean }, @@ -124,8 +141,8 @@ export const layerWith = (options?: LayerOptions) => ) function commitDurableEvent( - definition: Definition, - event: Payload, + definition: DurableDefinition, + event: UncommittedPayload, input?: { readonly seq: number readonly aggregateID: string @@ -135,7 +152,7 @@ export const layerWith = (options?: LayerOptions) => commit?: (seq: number) => Effect.Effect, ) { return Effect.gen(function* () { - const durable = definition?.durable + const durable = definition.durable if (durable) { const aggregateID = (event.data as Record)[durable.aggregate] if (typeof aggregateID !== "string") { @@ -200,7 +217,7 @@ export const layerWith = (options?: LayerOptions) => .run() .pipe(Effect.orDie) } - return + return undefined } yield* Effect.die( new InvalidDurableEventError({ @@ -210,7 +227,7 @@ export const layerWith = (options?: LayerOptions) => ) } if (input && row?.ownerID && row.ownerID !== input.ownerID) { - return + return undefined } const seq = input?.seq ?? latest + 1 if (input && seq !== latest + 1) { @@ -234,10 +251,10 @@ export const layerWith = (options?: LayerOptions) => message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`, }), ) - const committed = { + const committed: Payload = { ...event, durable: { aggregateID, seq, version: durable.version }, - } as Payload + } for (const projector of list) { yield* projector(committed) } @@ -267,14 +284,14 @@ export const layerWith = (options?: LayerOptions) => ]) .run() .pipe(Effect.orDie) - return { aggregateID, seq } + return committed }), { behavior: "immediate" }, ) .pipe(Effect.orDie) if (committed) { yield* Effect.forEach( - pubsub.durable.get(committed.aggregateID) ?? [], + pubsub.durable.get(committed.durable.aggregateID) ?? [], (wake) => PubSub.publish(wake, undefined), { discard: true }, ) @@ -287,7 +304,11 @@ export const layerWith = (options?: LayerOptions) => }) } - function publishEvent(definition: D, event: Payload, commit?: PublishOptions["commit"]) { + function publishEvent( + definition: D, + event: UncommittedPayload, + commit?: PublishOptions["commit"], + ): Effect.Effect> { return Effect.gen(function* () { if (!definition?.durable && commit) return yield* Effect.die( @@ -297,22 +318,22 @@ export const layerWith = (options?: LayerOptions) => }), ) if (definition?.durable) { - const committed = yield* commitDurableEvent(definition, event as Payload, undefined, commit) - if (committed) { - event = { - ...event, - durable: { - aggregateID: committed.aggregateID, - seq: committed.seq, - version: definition.durable.version, - }, - } - yield* notify(event as Payload, true) - return event - } + const committed = yield* commitDurableEvent( + definition, + event as UncommittedPayload, + undefined, + commit, + ) + if (!committed) + return yield* Effect.die( + new InvalidDurableEventError({ type: event.type, message: "New durable event was not committed" }), + ) + yield* notify(committed, true) + return committed as PublishedPayload } - yield* notify(event as Payload, false) - return event + const published = event as PublishedPayload + yield* notify(published, false) + return published }) } @@ -353,7 +374,7 @@ export const layerWith = (options?: LayerOptions) => type: definition.type, ...(location ? { location } : {}), data, - } as Payload, + } as UncommittedPayload, options?.commit, ) }) @@ -370,11 +391,11 @@ export const layerWith = (options?: LayerOptions) => new InvalidDurableEventError({ type: event.type, message: `Unknown durable event type ${event.type}` }), ) } else { - const payload = { + const payload: UncommittedPayload = { id: event.id, type: definition.type, data: Schema.decodeUnknownSync(definition.data)(event.data), - } as Payload + } const committed = yield* commitDurableEvent(definition, payload, { seq: event.seq, aggregateID: event.aggregateID, @@ -382,17 +403,7 @@ export const layerWith = (options?: LayerOptions) => strictOwner: options?.strictOwner, }) if (committed && options?.publish) { - yield* notify( - { - ...payload, - durable: { - aggregateID: committed.aggregateID, - seq: committed.seq, - version: definition.durable.version, - }, - }, - true, - ) + yield* notify(committed, true) } } }) @@ -459,7 +470,7 @@ export const layerWith = (options?: LayerOptions) => const streamAll = (): Stream.Stream => Stream.fromPubSub(pubsub.all) - const decodeSerializedEvent = (event: SerializedEvent) => { + const decodeSerializedEvent = (event: SerializedEvent): Payload => { const definition = Durable.get(event.type) if (!definition?.durable) { throw new InvalidDurableEventError({ type: event.type, message: `Unknown durable event type ${event.type}` }) @@ -516,7 +527,10 @@ export const layerWith = (options?: LayerOptions) => return subscription }) - const durable = (input: { readonly aggregateID: string; readonly after?: number }): Stream.Stream => + const durable = (input: { + readonly aggregateID: string + readonly after?: number + }): Stream.Stream> => Stream.unwrap( Effect.gen(function* () { const wakes = yield* subscribeDurable(input.aggregateID) @@ -524,7 +538,7 @@ export const layerWith = (options?: LayerOptions) => const read = Effect.suspend(() => readAfter(input.aggregateID, sequence)).pipe( Effect.tap((events) => Effect.sync(() => { - sequence = events.at(-1)?.durable?.seq ?? sequence + sequence = events.at(-1)?.durable.seq ?? sequence }), ), ) @@ -546,7 +560,7 @@ export const layerWith = (options?: LayerOptions) => }) }) - const project = (definition: D, projector: Subscriber): Effect.Effect => + const project = (definition: D, projector: Subscriber): Effect.Effect => Effect.sync(() => { const list = projectors.get(definition.type) ?? [] list.push((event) => projector(event as Payload)) diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index e2b2a5df04..2284bba7e4 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -116,7 +116,7 @@ describe("EventV2", () => { const event = yield* events.publish(VersionedMessage, { id: "one", text: "hello" }) expect(event.type).toBe("test.versioned") - expect(event.durable?.version).toBe(2) + expect(event.durable.version).toBe(2) }), ) @@ -764,7 +764,7 @@ describe("EventV2", () => { const replayed = { id: published.id, type: EventV2.versionedType(DurableMessage.type, 1), - seq: published.durable!.seq, + seq: published.durable.seq, aggregateID, data: published.data, } diff --git a/packages/opencode/test/server/httpapi-public-openapi.test.ts b/packages/opencode/test/server/httpapi-public-openapi.test.ts index a8f6f8d1c8..d9eb3643a9 100644 --- a/packages/opencode/test/server/httpapi-public-openapi.test.ts +++ b/packages/opencode/test/server/httpapi-public-openapi.test.ts @@ -99,6 +99,16 @@ describe("PublicApi OpenAPI v2 errors", () => { }) }) + test("documents durable metadata only on durable events", () => { + const spec = OpenApi.fromApi(PublicApi) as OpenApiSpec + const durable = spec.components.schemas.V2EventSessionCreated + const live = spec.components.schemas.V2EventSessionNextTextDelta + + expect(durable?.required).toContain("durable") + expect(durable?.properties?.durable).toBeDefined() + expect(live?.properties?.durable).toBeUndefined() + }) + test("preserves /api auth responses", () => { const spec = OpenApi.fromApi(PublicApi) as OpenApiSpec diff --git a/packages/opencode/test/v2/session-message-updater.test.ts b/packages/opencode/test/v2/session-message-updater.test.ts index 668a353f67..56478a4d1d 100644 --- a/packages/opencode/test/v2/session-message-updater.test.ts +++ b/packages/opencode/test/v2/session-message-updater.test.ts @@ -9,6 +9,12 @@ import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater" import { SessionMessage } from "@opencode-ai/core/session/message" +function durable(sessionID: SessionID, seq?: number): { aggregateID: SessionID; seq: number; version: 1 } +function durable(sessionID: SessionID, seq: number, version: 2): { aggregateID: SessionID; seq: number; version: 2 } +function durable(sessionID: SessionID, seq = 0, version: 1 | 2 = 1) { + return { aggregateID: sessionID, seq, version } +} + test.skip("step snapshots carry over to assistant messages", () => { const state: SessionMessageUpdater.MemoryState = { messages: [] } const sessionID = SessionID.make("session") @@ -17,6 +23,7 @@ test.skip("step snapshots carry over to assistant messages", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -38,6 +45,7 @@ test.skip("step snapshots carry over to assistant messages", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 1, 2), type: "session.next.step.ended", data: { sessionID, @@ -70,6 +78,7 @@ test.skip("text ended populates assistant text content", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -88,6 +97,7 @@ test.skip("text ended populates assistant text content", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 1), type: "session.next.text.started", data: { sessionID, @@ -101,6 +111,7 @@ test.skip("text ended populates assistant text content", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 2), type: "session.next.text.ended", data: { sessionID, @@ -126,6 +137,7 @@ test.skip("tool completion stores completed timestamp", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -144,6 +156,7 @@ test.skip("tool completion stores completed timestamp", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 1), type: "session.next.tool.input.started", data: { sessionID, @@ -158,6 +171,7 @@ test.skip("tool completion stores completed timestamp", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 2), type: "session.next.tool.called", data: { sessionID, @@ -174,6 +188,7 @@ test.skip("tool completion stores completed timestamp", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 3), type: "session.next.tool.success", data: { sessionID, @@ -204,6 +219,7 @@ test("compaction events reduce to compaction message only when completed", () => Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id, + durable: durable(sessionID), type: "session.next.compaction.started", data: { sessionID, @@ -245,6 +261,7 @@ test("compaction events reduce to compaction message only when completed", () => Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 3), type: "session.next.compaction.ended", data: { sessionID, diff --git a/packages/protocol/src/groups/event.ts b/packages/protocol/src/groups/event.ts index 862337e54b..703db3a3d6 100644 --- a/packages/protocol/src/groups/event.ts +++ b/packages/protocol/src/groups/event.ts @@ -8,18 +8,24 @@ import { HttpApiEndpoint, HttpApiGroup, OpenApi } from "effect/unstable/httpapi" const fields = { id: Event.ID, metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)), - durable: Schema.optional(Schema.Struct({ aggregateID: Schema.String, seq: Schema.Int, version: Schema.Int })), location: Schema.optional(Location.Ref), } const schema = (definitions: ReadonlyArray) => Schema.Union([ ...definitions.map((definition) => - Schema.Struct({ - ...fields, - type: Schema.Literal(definition.type), - data: definition.data, - }).annotate({ identifier: `V2Event.${definition.type}` }), + definition.durable + ? Schema.Struct({ + ...fields, + durable: Event.durableEnvelope(definition.durable.version), + type: Schema.Literal(definition.type), + data: definition.data, + }).annotate({ identifier: `V2Event.${definition.type}` }) + : Schema.Struct({ + ...fields, + type: Schema.Literal(definition.type), + data: definition.data, + }).annotate({ identifier: `V2Event.${definition.type}` }), ), ...(definitions.some((definition) => definition.type === "server.connected") ? [] diff --git a/packages/schema/src/event.ts b/packages/schema/src/event.ts index 849ffb8902..ba116b05e5 100644 --- a/packages/schema/src/event.ts +++ b/packages/schema/src/event.ts @@ -3,7 +3,7 @@ export * as Event from "./event" import { Schema } from "effect" import { ascending } from "./identifier" import { Location } from "./location" -import { statics } from "./schema" +import { NonNegativeInt, statics } from "./schema" export const ID = Schema.String.check(Schema.isStartsWith("evt_")).pipe( Schema.brand("Event.ID"), @@ -11,53 +11,95 @@ export const ID = Schema.String.check(Schema.isStartsWith("evt_")).pipe( ) export type ID = typeof ID.Type -export type Definition< +export type DurableOptions = { + readonly version: number + readonly aggregate: string +} + +export type DurableEnvelope = { + readonly aggregateID: string + readonly seq: number + readonly version: Version +} + +export const durableEnvelope = (version: Version) => + Schema.Struct({ aggregateID: Schema.String, seq: NonNegativeInt, version: Schema.Literal(version) }) + +const NoDurableEnvelope = Schema.optional(Schema.Never) + +export type LiveDefinition< Type extends string = string, DataSchema extends Schema.Codec = Schema.Codec, > = Schema.Top & { readonly type: Type - readonly durable?: { - readonly version: number - readonly aggregate: string - } readonly data: DataSchema + readonly durable?: never } +export type DurableDefinition< + Type extends string = string, + DataSchema extends Schema.Codec = Schema.Codec, + Durability extends DurableOptions = DurableOptions, +> = Schema.Top & { + readonly type: Type + readonly data: DataSchema + readonly durable: Durability +} + +export type Definition = LiveDefinition | DurableDefinition + +type Defined< + Type extends string, + DataSchema extends Schema.Codec, + Durability extends DurableOptions | undefined, +> = Durability extends DurableOptions + ? DurableDefinition + : LiveDefinition + export type Data = Schema.Schema.Type -export type Payload = { - readonly id: ID - readonly type: D["type"] - readonly data: Data - readonly durable?: { - readonly aggregateID: string - readonly seq: number - readonly version: number - } - readonly location?: Location.Ref - readonly metadata?: Record -} +export type UncommittedPayload = D extends Definition + ? { + readonly id: ID + readonly type: D["type"] + readonly data: Data + readonly location?: Location.Ref + readonly metadata?: Record + } + : never + +export type PublishedPayload = D extends Definition + ? UncommittedPayload & + (D extends { readonly durable: infer Durability extends DurableOptions } + ? { readonly durable: DurableEnvelope } + : { readonly durable?: never }) + : never + +export type Payload = PublishedPayload + +type EventSchema< + Type extends string, + Fields extends Readonly>>, + Durability extends DurableOptions | undefined, +> = Schema.Schema, Durability>>> & + Defined, Durability> export function define< const Type extends string, Fields extends Readonly>>, + const Durability extends DurableOptions | undefined = undefined, >(input: { readonly type: Type - readonly durable?: { - readonly version: number - readonly aggregate: string - } + readonly durable?: Durability readonly schema: Fields -}): Schema.Schema>>> & Definition> { +}): EventSchema { const data = Schema.Struct(input.schema) return Object.assign( Schema.Struct({ id: ID, metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)), type: Schema.Literal(input.type), - durable: Schema.optional( - Schema.Struct({ aggregateID: Schema.String, seq: Schema.Number, version: Schema.Number }), - ), + durable: input.durable === undefined ? NoDurableEnvelope : durableEnvelope(input.durable.version), location: Schema.optional(Location.Ref), data, }).annotate({ identifier: input.type }), @@ -66,7 +108,7 @@ export function define< ...(input.durable === undefined ? {} : { durable: input.durable }), data, }, - ) as Schema.Schema>>> & Definition> + ) as unknown as EventSchema } export function inventory>(...definitions: Definitions) { @@ -103,7 +145,7 @@ export function durable(definitions: ReadonlyArray) { if (result.has(key)) throw new Error(`Duplicate durable event definition for ${key}`) result.set(key, definition) return result - }, new Map()), + }, new Map()), ) } diff --git a/packages/schema/test/event.test.ts b/packages/schema/test/event.test.ts index 380faa5a4a..3b119368cc 100644 --- a/packages/schema/test/event.test.ts +++ b/packages/schema/test/event.test.ts @@ -34,4 +34,99 @@ describe("public event schemas", () => { expect(Event.durable([definition]).get("test.durable.1")).toBe(definition) }) + + test("durable definitions require published commit metadata", () => { + const definition = Event.define({ + type: "test.durable", + durable: { aggregate: "id", version: 1 }, + schema: { id: Schema.String }, + }) + const payload: typeof definition.Type = { + id: Event.ID.create(), + type: definition.type, + durable: { aggregateID: "aggregate", seq: 0, version: 1 }, + data: { id: "aggregate" }, + } + + expect(Schema.is(definition)(payload)).toBe(true) + expect( + Schema.is(definition)({ + id: Event.ID.create(), + type: definition.type, + data: { id: "aggregate" }, + }), + ).toBe(false) + expect(Schema.is(definition)({ ...payload, durable: { ...payload.durable, seq: -1 } })).toBe(false) + expect(Schema.is(definition)({ ...payload, durable: { ...payload.durable, version: 2 } })).toBe(false) + + // @ts-expect-error Published durable payloads require commit metadata. + const missing: typeof definition.Type = { id: Event.ID.create(), type: definition.type, data: { id: "aggregate" } } + void missing + }) + + test("live definitions reject durable commit metadata", () => { + const definition = Event.define({ + type: "test.live", + schema: { value: Schema.String }, + }) + const payload: typeof definition.Type = { + id: Event.ID.create(), + type: definition.type, + data: { value: "value" }, + } + + expect(Schema.is(definition)(payload)).toBe(true) + expect( + Schema.is(definition)({ + ...payload, + durable: { aggregateID: "aggregate", seq: 0, version: 1 }, + }), + ).toBe(false) + + const invalid: typeof definition.Type = { + ...payload, + // @ts-expect-error Live payloads cannot carry durable commit metadata. + durable: { aggregateID: "aggregate", seq: 0, version: 1 }, + } + void invalid + }) + + test("mixed definition payloads preserve durability correlation", () => { + const durable = Event.define({ + type: "test.mixed.durable", + durable: { aggregate: "id", version: 2 }, + schema: { id: Schema.String }, + }) + const live = Event.define({ + type: "test.mixed.live", + schema: { value: Schema.String }, + }) + type Mixed = Event.Payload + + const committed: Mixed = { + id: Event.ID.create(), + type: durable.type, + durable: { aggregateID: "aggregate", seq: 0, version: 2 }, + data: { id: "aggregate" }, + } + const ephemeral: Mixed = { + id: Event.ID.create(), + type: live.type, + data: { value: "value" }, + } + void committed + void ephemeral + + // @ts-expect-error Durable union members require commit metadata. + const uncommitted: Mixed = { id: Event.ID.create(), type: durable.type, data: { id: "aggregate" } } + const falselyCommitted: Mixed = { + id: Event.ID.create(), + type: live.type, + // @ts-expect-error Live union members cannot carry durable commit metadata. + durable: { aggregateID: "aggregate", seq: 0, version: 2 }, + data: { value: "value" }, + } + void uncommitted + void falselyCommitted + }) }) diff --git a/packages/sdk/js/src/v2/gen/types.gen.ts b/packages/sdk/js/src/v2/gen/types.gen.ts index 2f22dc9e7f..b653c75e1b 100644 --- a/packages/sdk/js/src/v2/gen/types.gen.ts +++ b/packages/sdk/js/src/v2/gen/types.gen.ts @@ -4328,11 +4328,6 @@ export type V2EventModelsDevRefreshed = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "models-dev.refreshed" data: { @@ -4345,11 +4340,6 @@ export type V2EventIntegrationUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "integration.updated" data: { @@ -4362,11 +4352,6 @@ export type V2EventIntegrationConnectionUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "integration.connection.updated" data: { @@ -4379,11 +4364,6 @@ export type V2EventCatalogUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "catalog.updated" data: { @@ -4396,12 +4376,12 @@ export type V2EventSessionCreated = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.created" data: { sessionID: string @@ -4414,12 +4394,12 @@ export type V2EventSessionUpdated = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.updated" data: { sessionID: string @@ -4432,12 +4412,12 @@ export type V2EventSessionDeleted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.deleted" data: { sessionID: string @@ -4450,12 +4430,12 @@ export type V2EventMessageUpdated = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "message.updated" data: { sessionID: string @@ -4468,12 +4448,12 @@ export type V2EventMessageRemoved = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "message.removed" data: { sessionID: string @@ -4486,12 +4466,12 @@ export type V2EventMessagePartUpdated = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "message.part.updated" data: { sessionID: string @@ -4505,12 +4485,12 @@ export type V2EventMessagePartRemoved = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "message.part.removed" data: { sessionID: string @@ -4524,12 +4504,12 @@ export type V2EventSessionNextAgentSwitched = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.agent.switched" data: { timestamp: number @@ -4544,12 +4524,12 @@ export type V2EventSessionNextModelSwitched = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.model.switched" data: { timestamp: number @@ -4568,12 +4548,12 @@ export type V2EventSessionNextMoved = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.moved" data: { timestamp: number @@ -4588,12 +4568,12 @@ export type V2EventSessionNextPrompted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.prompted" data: { timestamp: number @@ -4609,12 +4589,12 @@ export type V2EventSessionNextPromptAdmitted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.prompt.admitted" data: { timestamp: number @@ -4630,12 +4610,12 @@ export type V2EventSessionNextContextUpdated = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.context.updated" data: { timestamp: number @@ -4650,12 +4630,12 @@ export type V2EventSessionNextSynthetic = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.synthetic" data: { timestamp: number @@ -4670,12 +4650,12 @@ export type V2EventSessionNextShellStarted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.shell.started" data: { timestamp: number @@ -4691,12 +4671,12 @@ export type V2EventSessionNextShellEnded = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.shell.ended" data: { timestamp: number @@ -4711,12 +4691,12 @@ export type V2EventSessionNextStepStarted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.step.started" data: { timestamp: number @@ -4737,12 +4717,12 @@ export type V2EventSessionNextStepEnded = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 2 } - location?: LocationRef type: "session.next.step.ended" data: { timestamp: number @@ -4769,12 +4749,12 @@ export type V2EventSessionNextStepFailed = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 2 } - location?: LocationRef type: "session.next.step.failed" data: { timestamp: number @@ -4789,12 +4769,12 @@ export type V2EventSessionNextTextStarted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.text.started" data: { timestamp: number @@ -4809,11 +4789,6 @@ export type V2EventSessionNextTextDelta = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.next.text.delta" data: { @@ -4830,12 +4805,12 @@ export type V2EventSessionNextTextEnded = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.text.ended" data: { timestamp: number @@ -4851,12 +4826,12 @@ export type V2EventSessionNextReasoningStarted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.reasoning.started" data: { timestamp: number @@ -4876,11 +4851,6 @@ export type V2EventSessionNextReasoningDelta = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.next.reasoning.delta" data: { @@ -4897,12 +4867,12 @@ export type V2EventSessionNextReasoningEnded = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.reasoning.ended" data: { timestamp: number @@ -4923,12 +4893,12 @@ export type V2EventSessionNextToolInputStarted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.tool.input.started" data: { timestamp: number @@ -4944,11 +4914,6 @@ export type V2EventSessionNextToolInputDelta = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.next.tool.input.delta" data: { @@ -4965,12 +4930,12 @@ export type V2EventSessionNextToolInputEnded = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.tool.input.ended" data: { timestamp: number @@ -4986,12 +4951,12 @@ export type V2EventSessionNextToolCalled = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.tool.called" data: { timestamp: number @@ -5018,12 +4983,12 @@ export type V2EventSessionNextToolProgress = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.tool.progress" data: { timestamp: number @@ -5042,12 +5007,12 @@ export type V2EventSessionNextToolSuccess = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.tool.success" data: { timestamp: number @@ -5076,12 +5041,12 @@ export type V2EventSessionNextToolFailed = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.tool.failed" data: { timestamp: number @@ -5106,12 +5071,12 @@ export type V2EventSessionNextRetried = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.retried" data: { timestamp: number @@ -5126,12 +5091,12 @@ export type V2EventSessionNextCompactionStarted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.compaction.started" data: { timestamp: number @@ -5146,11 +5111,6 @@ export type V2EventSessionNextCompactionDelta = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.next.compaction.delta" data: { @@ -5166,12 +5126,12 @@ export type V2EventSessionNextCompactionEnded = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.compaction.ended" data: { timestamp: number @@ -5188,12 +5148,12 @@ export type V2EventSessionNextRevertStaged = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.revert.staged" data: { timestamp: number @@ -5213,12 +5173,12 @@ export type V2EventSessionNextRevertCleared = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.revert.cleared" data: { timestamp: number @@ -5231,12 +5191,12 @@ export type V2EventSessionNextRevertCommitted = { metadata?: { [key: string]: unknown } - durable?: { + location?: LocationRef + durable: { aggregateID: string seq: number - version: number + version: 1 } - location?: LocationRef type: "session.next.revert.committed" data: { timestamp: number @@ -5250,11 +5210,6 @@ export type V2EventMessagePartDelta = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "message.part.delta" data: { @@ -5271,11 +5226,6 @@ export type V2EventSessionDiff = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.diff" data: { @@ -5289,11 +5239,6 @@ export type V2EventSessionError = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.error" data: { @@ -5315,11 +5260,6 @@ export type V2EventInstallationUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "installation.updated" data: { @@ -5332,11 +5272,6 @@ export type V2EventInstallationUpdateAvailable = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "installation.update-available" data: { @@ -5349,11 +5284,6 @@ export type V2EventFileEdited = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "file.edited" data: { @@ -5366,11 +5296,6 @@ export type V2EventReferenceUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "reference.updated" data: { @@ -5383,11 +5308,6 @@ export type V2EventPermissionV2Asked = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "permission.v2.asked" data: { @@ -5408,11 +5328,6 @@ export type V2EventPermissionV2Replied = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "permission.v2.replied" data: { @@ -5427,11 +5342,6 @@ export type V2EventPluginAdded = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "plugin.added" data: { @@ -5444,11 +5354,6 @@ export type V2EventProjectDirectoriesUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "project.directories.updated" data: { @@ -5461,11 +5366,6 @@ export type V2EventFileWatcherUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "file.watcher.updated" data: { @@ -5479,11 +5379,6 @@ export type V2EventPtyCreated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "pty.created" data: { @@ -5496,11 +5391,6 @@ export type V2EventPtyUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "pty.updated" data: { @@ -5513,11 +5403,6 @@ export type V2EventPtyExited = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "pty.exited" data: { @@ -5531,11 +5416,6 @@ export type V2EventPtyDeleted = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "pty.deleted" data: { @@ -5548,11 +5428,6 @@ export type V2EventQuestionV2Asked = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "question.v2.asked" data: { @@ -5571,11 +5446,6 @@ export type V2EventQuestionV2Replied = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "question.v2.replied" data: { @@ -5590,11 +5460,6 @@ export type V2EventQuestionV2Rejected = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "question.v2.rejected" data: { @@ -5608,11 +5473,6 @@ export type V2EventTodoUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "todo.updated" data: { @@ -5626,11 +5486,6 @@ export type V2EventLspUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "lsp.updated" data: { @@ -5643,11 +5498,6 @@ export type V2EventPermissionAsked = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "permission.asked" data: { @@ -5671,11 +5521,6 @@ export type V2EventPermissionReplied = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "permission.replied" data: { @@ -5690,11 +5535,6 @@ export type V2EventTuiPromptAppend = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "tui.prompt.append" data: { @@ -5707,11 +5547,6 @@ export type V2EventTuiCommandExecute = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "tui.command.execute" data: { @@ -5741,11 +5576,6 @@ export type V2EventTuiToastShow = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "tui.toast.show" data: { @@ -5761,11 +5591,6 @@ export type V2EventTuiSessionSelect = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "tui.session.select" data: { @@ -5781,11 +5606,6 @@ export type V2EventMcpToolsChanged = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "mcp.tools.changed" data: { @@ -5798,11 +5618,6 @@ export type V2EventMcpBrowserOpenFailed = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "mcp.browser.open.failed" data: { @@ -5816,11 +5631,6 @@ export type V2EventCommandExecuted = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "command.executed" data: { @@ -5836,11 +5646,6 @@ export type V2EventProjectUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "project.updated" data: { @@ -5873,11 +5678,6 @@ export type V2EventSessionStatus = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.status" data: { @@ -5891,11 +5691,6 @@ export type V2EventSessionIdle = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.idle" data: { @@ -5908,11 +5703,6 @@ export type V2EventQuestionAsked = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "question.asked" data: { @@ -5931,11 +5721,6 @@ export type V2EventQuestionReplied = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "question.replied" data: { @@ -5950,11 +5735,6 @@ export type V2EventQuestionRejected = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "question.rejected" data: { @@ -5968,11 +5748,6 @@ export type V2EventSessionCompacted = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "session.compacted" data: { @@ -5985,11 +5760,6 @@ export type V2EventVcsBranchUpdated = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "vcs.branch.updated" data: { @@ -6002,11 +5772,6 @@ export type V2EventWorkspaceReady = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "workspace.ready" data: { @@ -6019,11 +5784,6 @@ export type V2EventWorkspaceFailed = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "workspace.failed" data: { @@ -6036,11 +5796,6 @@ export type V2EventWorkspaceStatus = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "workspace.status" data: { @@ -6054,11 +5809,6 @@ export type V2EventWorktreeReady = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "worktree.ready" data: { @@ -6072,11 +5822,6 @@ export type V2EventWorktreeFailed = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "worktree.failed" data: { @@ -6089,11 +5834,6 @@ export type V2EventServerConnected = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "server.connected" data: { @@ -6106,11 +5846,6 @@ export type V2EventGlobalDisposed = { metadata?: { [key: string]: unknown } - durable?: { - aggregateID: string - seq: number - version: number - } location?: LocationRef type: "global.disposed" data: { From 45071a8d8e293b26284ac14bd21a248fa0616c99 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 25 Jun 2026 12:04:05 -0400 Subject: [PATCH 2/3] refactor(core): prove event publication branches --- packages/core/src/event.ts | 106 +++++++++--------- packages/core/src/session/message-updater.ts | 15 ++- .../test/v2/session-message-updater.test.ts | 44 +++----- packages/schema/src/event.ts | 90 +++++++++------ 4 files changed, 132 insertions(+), 123 deletions(-) diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index d8ec020627..0b0eae26d7 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -1,13 +1,13 @@ export * as EventV2 from "./event" -import { Cause, Context, Effect, Layer, Option, PubSub, Schema, Stream } from "effect" +import { Cause, Context, Effect, Layer, Option, Predicate, PubSub, Schema, Stream } from "effect" import { Event } from "@opencode-ai/schema/event" import type { Data, Definition, DurableDefinition, + LivePublishedPayload, Payload, - PublishedPayload, UncommittedPayload, } from "@opencode-ai/schema/event" import { and, asc, eq, gt } from "drizzle-orm" @@ -24,8 +24,8 @@ export type { Data, Definition, DurableDefinition, + LivePublishedPayload, Payload, - PublishedPayload, UncommittedPayload, } from "@opencode-ai/schema/event" @@ -154,7 +154,7 @@ export const layerWith = (options?: LayerOptions) => return Effect.gen(function* () { const durable = definition.durable if (durable) { - const aggregateID = (event.data as Record)[durable.aggregate] + const aggregateID = Predicate.isReadonlyObject(event.data) ? event.data[durable.aggregate] : undefined if (typeof aggregateID !== "string") { yield* Effect.die( new InvalidDurableEventError({ @@ -185,10 +185,14 @@ export const layerWith = (options?: LayerOptions) => .get() .pipe(Effect.orDie) const latest = row?.seq ?? -1 - const encoded = Schema.encodeUnknownSync(definition.data)(event.data) as Record< - string, - unknown - > + const encoded = Schema.encodeUnknownSync(definition.data)(event.data) + if (!Predicate.isReadonlyObject(encoded)) + return yield* Effect.die( + new InvalidDurableEventError({ + type: event.type, + message: "Expected durable event data to encode as an object", + }), + ) if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) { yield* Effect.die( new InvalidDurableEventError({ @@ -304,39 +308,6 @@ export const layerWith = (options?: LayerOptions) => }) } - function publishEvent( - definition: D, - event: UncommittedPayload, - commit?: PublishOptions["commit"], - ): Effect.Effect> { - return Effect.gen(function* () { - if (!definition?.durable && commit) - return yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: "Local commit hooks require a durable event", - }), - ) - if (definition?.durable) { - const committed = yield* commitDurableEvent( - definition, - event as UncommittedPayload, - undefined, - commit, - ) - if (!committed) - return yield* Effect.die( - new InvalidDurableEventError({ type: event.type, message: "New durable event was not committed" }), - ) - yield* notify(committed, true) - return committed as PublishedPayload - } - const published = event as PublishedPayload - yield* notify(published, false) - return published - }) - } - const observe = (event: Payload, observer: (event: Payload) => Effect.Effect) => Effect.suspend(() => observer(event)).pipe( Effect.catchCauseIf( @@ -358,7 +329,17 @@ export const layerWith = (options?: LayerOptions) => }) } - function publish(definition: D, data: Data, options?: PublishOptions) { + const isPayload = + (definition: D) => + (event: Payload): event is Payload => + Schema.is(definition)(event) + + function publish( + definition: D, + data: Data, + options?: PublishOptions, + ): Effect.Effect> + function publish(definition: Definition, data: unknown, options?: PublishOptions): Effect.Effect { return Effect.gen(function* () { const serviceLocation = Option.getOrUndefined(yield* Effect.serviceOption(Location.Service)) const location = @@ -366,17 +347,38 @@ export const layerWith = (options?: LayerOptions) => (serviceLocation ? { directory: serviceLocation.directory, workspaceID: serviceLocation.workspaceID } : undefined) - return yield* publishEvent( - definition, - { + if (definition.durable) { + const event: UncommittedPayload = { id: options?.id ?? ID.create(), ...(options?.metadata ? { metadata: options.metadata } : {}), type: definition.type, ...(location ? { location } : {}), data, - } as UncommittedPayload, - options?.commit, - ) + } + const committed = yield* commitDurableEvent(definition, event, undefined, options?.commit) + if (!committed) + return yield* Effect.die( + new InvalidDurableEventError({ type: event.type, message: "New durable event was not committed" }), + ) + yield* notify(committed, true) + return committed + } + if (options?.commit) + return yield* Effect.die( + new InvalidDurableEventError({ + type: definition.type, + message: "Local commit hooks require a durable event", + }), + ) + const event: LivePublishedPayload = { + id: options?.id ?? ID.create(), + ...(options?.metadata ? { metadata: options.metadata } : {}), + type: definition.type, + ...(location ? { location } : {}), + data, + } + yield* notify(event, false) + return event }) } @@ -465,7 +467,7 @@ export const layerWith = (options?: LayerOptions) => const subscribe = (definition: D): Stream.Stream> => Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))).pipe( - Stream.map((event) => event as Payload), + Stream.filter(isPayload(definition)), ) const streamAll = (): Stream.Stream => Stream.fromPubSub(pubsub.all) @@ -563,7 +565,11 @@ export const layerWith = (options?: LayerOptions) => const project = (definition: D, projector: Subscriber): Effect.Effect => Effect.sync(() => { const list = projectors.get(definition.type) ?? [] - list.push((event) => projector(event as Payload)) + list.push((event) => + isPayload(definition)(event) + ? projector(event) + : Effect.die(`Published event ${event.type} does not match its definition`), + ) projectors.set(definition.type, list) }) diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index a97afb58b6..8f856d2288 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -1,5 +1,6 @@ import { castDraft, produce, type WritableDraft } from "immer" -import { Effect } from "effect" +import { Effect, Match } from "effect" +import type { UncommittedPayload } from "@opencode-ai/schema/event" import { SessionEvent } from "./event" import { SessionMessage } from "./message" @@ -7,6 +8,8 @@ export type MemoryState = { messages: SessionMessage.Message[] } +export type Input = UncommittedPayload<(typeof SessionEvent.Definitions)[number]> + export interface Adapter { readonly getCurrentAssistant: () => Effect.Effect readonly getAssistant: (messageID: SessionMessage.ID) => Effect.Effect @@ -75,7 +78,7 @@ export function memory(state: MemoryState): Adapter { } } -export function update(adapter: Adapter, event: SessionEvent.Event) { +export function update(adapter: Adapter, event: Input) { type DraftAssistant = WritableDraft type DraftTool = WritableDraft type DraftText = WritableDraft @@ -98,8 +101,8 @@ export function update(adapter: Adapter, event: SessionEvent.Event) { if (assistant) yield* adapter.updateAssistant(produce(assistant, recipe)) }) - return Effect.gen(function* () { - yield* SessionEvent.All.match(event, { + return Match.value(event).pipe( + Match.discriminatorsExhaustive("type")({ "session.next.agent.switched": (event) => { return adapter.appendMessage( SessionMessage.AgentSwitched.make({ @@ -388,8 +391,8 @@ export function update(adapter: Adapter, event: SessionEvent.Event) { "session.next.revert.staged": () => Effect.void, "session.next.revert.cleared": () => Effect.void, "session.next.revert.committed": () => Effect.void, - }) - }) + }), + ) } export * as SessionMessageUpdater from "./message-updater" diff --git a/packages/opencode/test/v2/session-message-updater.test.ts b/packages/opencode/test/v2/session-message-updater.test.ts index 56478a4d1d..5611f6ed70 100644 --- a/packages/opencode/test/v2/session-message-updater.test.ts +++ b/packages/opencode/test/v2/session-message-updater.test.ts @@ -5,16 +5,9 @@ import { SessionID } from "../../src/session/schema" import { EventV2 } from "@opencode-ai/core/event" import { ModelV2 } from "@opencode-ai/core/model" import { ProviderV2 } from "@opencode-ai/core/provider" -import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater" import { SessionMessage } from "@opencode-ai/core/session/message" -function durable(sessionID: SessionID, seq?: number): { aggregateID: SessionID; seq: number; version: 1 } -function durable(sessionID: SessionID, seq: number, version: 2): { aggregateID: SessionID; seq: number; version: 2 } -function durable(sessionID: SessionID, seq = 0, version: 1 | 2 = 1) { - return { aggregateID: sessionID, seq, version } -} - test.skip("step snapshots carry over to assistant messages", () => { const state: SessionMessageUpdater.MemoryState = { messages: [] } const sessionID = SessionID.make("session") @@ -23,7 +16,6 @@ test.skip("step snapshots carry over to assistant messages", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -37,7 +29,7 @@ test.skip("step snapshots carry over to assistant messages", () => { }, snapshot: "before", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) expect(state.messages).toEqual([]) @@ -45,7 +37,6 @@ test.skip("step snapshots carry over to assistant messages", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 1, 2), type: "session.next.step.ended", data: { sessionID, @@ -61,7 +52,7 @@ test.skip("step snapshots carry over to assistant messages", () => { }, snapshot: "after", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) expect(state.messages[0]?.type).toBe("assistant") @@ -78,7 +69,6 @@ test.skip("text ended populates assistant text content", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -91,13 +81,12 @@ test.skip("text ended populates assistant text content", () => { variant: ModelV2.VariantID.make("default"), }, }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 1), type: "session.next.text.started", data: { sessionID, @@ -105,13 +94,12 @@ test.skip("text ended populates assistant text content", () => { timestamp: DateTime.makeUnsafe(2), textID: "text-1", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 2), type: "session.next.text.ended", data: { sessionID, @@ -120,7 +108,7 @@ test.skip("text ended populates assistant text content", () => { textID: "text-1", text: "hello assistant", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) expect(state.messages[0]?.type).toBe("assistant") @@ -137,7 +125,6 @@ test.skip("tool completion stores completed timestamp", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -150,13 +137,12 @@ test.skip("tool completion stores completed timestamp", () => { variant: ModelV2.VariantID.make("default"), }, }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 1), type: "session.next.tool.input.started", data: { sessionID, @@ -165,13 +151,12 @@ test.skip("tool completion stores completed timestamp", () => { callID, name: "bash", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 2), type: "session.next.tool.called", data: { sessionID, @@ -182,13 +167,12 @@ test.skip("tool completion stores completed timestamp", () => { input: { command: "pwd" }, provider: { executed: true, metadata: { fake: { source: "provider" } } }, }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 3), type: "session.next.tool.success", data: { sessionID, @@ -199,7 +183,7 @@ test.skip("tool completion stores completed timestamp", () => { content: [{ type: "text", text: "/tmp" }], provider: { executed: true, metadata: { fake: { status: "done" } } }, }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) expect(state.messages[0]?.type).toBe("assistant") @@ -219,7 +203,6 @@ test("compaction events reduce to compaction message only when completed", () => Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id, - durable: durable(sessionID), type: "session.next.compaction.started", data: { sessionID, @@ -227,7 +210,7 @@ test("compaction events reduce to compaction message only when completed", () => timestamp: DateTime.makeUnsafe(1), reason: "auto", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) expect(state.messages).toEqual([]) @@ -242,7 +225,7 @@ test("compaction events reduce to compaction message only when completed", () => timestamp: DateTime.makeUnsafe(2), text: "hello ", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( @@ -255,13 +238,12 @@ test("compaction events reduce to compaction message only when completed", () => timestamp: DateTime.makeUnsafe(3), text: "summary", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), - durable: durable(sessionID, 3), type: "session.next.compaction.ended", data: { sessionID, @@ -271,7 +253,7 @@ test("compaction events reduce to compaction message only when completed", () => text: "final summary", recent: "recent context", }, - } satisfies SessionEvent.Event), + } satisfies SessionMessageUpdater.Input), ) expect(state.messages).toHaveLength(1) diff --git a/packages/schema/src/event.ts b/packages/schema/src/event.ts index ba116b05e5..168eb1eaf7 100644 --- a/packages/schema/src/event.ts +++ b/packages/schema/src/event.ts @@ -48,14 +48,6 @@ export type DurableDefinition< export type Definition = LiveDefinition | DurableDefinition -type Defined< - Type extends string, - DataSchema extends Schema.Codec, - Durability extends DurableOptions | undefined, -> = Durability extends DurableOptions - ? DurableDefinition - : LiveDefinition - export type Data = Schema.Schema.Type export type UncommittedPayload = D extends Definition @@ -68,47 +60,73 @@ export type UncommittedPayload = D extends De } : never -export type PublishedPayload = D extends Definition - ? UncommittedPayload & - (D extends { readonly durable: infer Durability extends DurableOptions } - ? { readonly durable: DurableEnvelope } - : { readonly durable?: never }) - : never +export type DurablePublishedPayload = UncommittedPayload & { + readonly durable: DurableEnvelope +} + +export type LivePublishedPayload = UncommittedPayload & { + readonly durable?: never +} + +export type PublishedPayload = D extends DurableDefinition + ? DurablePublishedPayload + : D extends LiveDefinition + ? LivePublishedPayload + : never export type Payload = PublishedPayload -type EventSchema< +type LiveEventSchema< Type extends string, Fields extends Readonly>>, - Durability extends DurableOptions | undefined, -> = Schema.Schema, Durability>>> & - Defined, Durability> +> = Schema.Schema>>> & + LiveDefinition> + +type DurableEventSchema< + Type extends string, + Fields extends Readonly>>, + Durability extends DurableOptions, +> = Schema.Schema, Durability>>> & + DurableDefinition, Durability> export function define< const Type extends string, Fields extends Readonly>>, - const Durability extends DurableOptions | undefined = undefined, +>(input: { readonly type: Type; readonly durable?: never; readonly schema: Fields }): LiveEventSchema +export function define< + const Type extends string, + Fields extends Readonly>>, + const Durability extends DurableOptions, >(input: { readonly type: Type - readonly durable?: Durability + readonly durable: Durability readonly schema: Fields -}): EventSchema { +}): DurableEventSchema +export function define(input: { + readonly type: string + readonly durable?: DurableOptions + readonly schema: Readonly>> +}): Schema.Top { const data = Schema.Struct(input.schema) - return Object.assign( - Schema.Struct({ - id: ID, - metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)), - type: Schema.Literal(input.type), - durable: input.durable === undefined ? NoDurableEnvelope : durableEnvelope(input.durable.version), - location: Schema.optional(Location.Ref), - data, - }).annotate({ identifier: input.type }), - { - type: input.type, - ...(input.durable === undefined ? {} : { durable: input.durable }), - data, - }, - ) as unknown as EventSchema + const fields = { + id: ID, + metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)), + type: Schema.Literal(input.type), + location: Schema.optional(Location.Ref), + data, + } + if (input.durable) { + return Object.assign( + Schema.Struct({ ...fields, durable: durableEnvelope(input.durable.version) }).annotate({ + identifier: input.type, + }), + { type: input.type, durable: input.durable, data }, + ) + } + return Object.assign(Schema.Struct({ ...fields, durable: NoDurableEnvelope }).annotate({ identifier: input.type }), { + type: input.type, + data, + }) } export function inventory>(...definitions: Definitions) { From 603b334b7f70620e8e531fadcd1f0ff159f186ab Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 25 Jun 2026 12:41:44 -0400 Subject: [PATCH 3/3] fix(core): preserve event routing semantics --- packages/core/src/event.ts | 79 +++++------ packages/core/src/session/message-updater.ts | 15 +-- packages/core/test/event.test.ts | 44 +++++++ .../test/session-runner-tool-events.test.ts | 9 +- .../server/httpapi-public-openapi.test.ts | 3 +- .../test/server/httpapi-v2-location.test.ts | 4 + .../test/v2/session-message-updater.test.ts | 44 +++++-- packages/protocol/src/groups/event.ts | 9 +- packages/protocol/test/event.test.ts | 27 ++++ packages/schema/src/event.ts | 33 +++-- packages/schema/test/event.test.ts | 5 +- packages/sdk/js/script/build.ts | 11 ++ packages/sdk/js/src/v2/gen/types.gen.ts | 123 +++++++++++++----- 13 files changed, 277 insertions(+), 129 deletions(-) create mode 100644 packages/protocol/test/event.test.ts diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index 0b0eae26d7..2c0261ebbf 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -1,12 +1,12 @@ export * as EventV2 from "./event" -import { Cause, Context, Effect, Layer, Option, Predicate, PubSub, Schema, Stream } from "effect" +import { Cause, Context, Effect, Layer, Option, PubSub, Schema, Stream } from "effect" import { Event } from "@opencode-ai/schema/event" import type { Data, Definition, DurableDefinition, - LivePublishedPayload, + LiveDefinition, Payload, UncommittedPayload, } from "@opencode-ai/schema/event" @@ -20,14 +20,7 @@ import { Durable } from "@opencode-ai/schema/durable-event-manifest" export const ID = Event.ID export type ID = import("@opencode-ai/schema/event").ID -export type { - Data, - Definition, - DurableDefinition, - LivePublishedPayload, - Payload, - UncommittedPayload, -} from "@opencode-ai/schema/event" +export type { Data, Definition, Payload } from "@opencode-ai/schema/event" export type Subscriber = (event: Payload) => Effect.Effect export type Unsubscribe = Effect.Effect @@ -80,13 +73,10 @@ export interface Interface { ) => Effect.Effect> readonly subscribe: (definition: D) => Stream.Stream> readonly all: () => Stream.Stream - readonly durable: (input: { - readonly aggregateID: string - readonly after?: number - }) => Stream.Stream> + readonly durable: (input: { readonly aggregateID: string; readonly after?: number }) => Stream.Stream /** @deprecated Use `all()` and consume the returned stream. */ readonly listen: (listener: Subscriber) => Effect.Effect - readonly project: (definition: D, projector: Subscriber) => Effect.Effect + readonly project: (definition: D, projector: Subscriber) => Effect.Effect readonly replay: ( event: SerializedEvent, options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean }, @@ -140,6 +130,23 @@ export const layerWith = (options?: LayerOptions) => }), ) + function commitDurableEvent( + definition: DurableDefinition, + event: UncommittedPayload, + input: undefined, + commit?: (seq: number) => Effect.Effect, + ): Effect.Effect> + function commitDurableEvent( + definition: DurableDefinition, + event: UncommittedPayload, + input: { + readonly seq: number + readonly aggregateID: string + readonly ownerID?: string + readonly strictOwner?: boolean + }, + commit?: (seq: number) => Effect.Effect, + ): Effect.Effect | undefined> function commitDurableEvent( definition: DurableDefinition, event: UncommittedPayload, @@ -154,7 +161,7 @@ export const layerWith = (options?: LayerOptions) => return Effect.gen(function* () { const durable = definition.durable if (durable) { - const aggregateID = Predicate.isReadonlyObject(event.data) ? event.data[durable.aggregate] : undefined + const aggregateID = (event.data as Record)[durable.aggregate] if (typeof aggregateID !== "string") { yield* Effect.die( new InvalidDurableEventError({ @@ -185,14 +192,10 @@ export const layerWith = (options?: LayerOptions) => .get() .pipe(Effect.orDie) const latest = row?.seq ?? -1 - const encoded = Schema.encodeUnknownSync(definition.data)(event.data) - if (!Predicate.isReadonlyObject(encoded)) - return yield* Effect.die( - new InvalidDurableEventError({ - type: event.type, - message: "Expected durable event data to encode as an object", - }), - ) + const encoded = Schema.encodeUnknownSync(definition.data)(event.data) as Record< + string, + unknown + > if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) { yield* Effect.die( new InvalidDurableEventError({ @@ -329,11 +332,6 @@ export const layerWith = (options?: LayerOptions) => }) } - const isPayload = - (definition: D) => - (event: Payload): event is Payload => - Schema.is(definition)(event) - function publish( definition: D, data: Data, @@ -356,10 +354,6 @@ export const layerWith = (options?: LayerOptions) => data, } const committed = yield* commitDurableEvent(definition, event, undefined, options?.commit) - if (!committed) - return yield* Effect.die( - new InvalidDurableEventError({ type: event.type, message: "New durable event was not committed" }), - ) yield* notify(committed, true) return committed } @@ -370,7 +364,7 @@ export const layerWith = (options?: LayerOptions) => message: "Local commit hooks require a durable event", }), ) - const event: LivePublishedPayload = { + const event: Payload = { id: options?.id ?? ID.create(), ...(options?.metadata ? { metadata: options.metadata } : {}), type: definition.type, @@ -467,7 +461,7 @@ export const layerWith = (options?: LayerOptions) => const subscribe = (definition: D): Stream.Stream> => Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))).pipe( - Stream.filter(isPayload(definition)), + Stream.map((event) => event as Payload), ) const streamAll = (): Stream.Stream => Stream.fromPubSub(pubsub.all) @@ -529,10 +523,7 @@ export const layerWith = (options?: LayerOptions) => return subscription }) - const durable = (input: { - readonly aggregateID: string - readonly after?: number - }): Stream.Stream> => + const durable = (input: { readonly aggregateID: string; readonly after?: number }): Stream.Stream => Stream.unwrap( Effect.gen(function* () { const wakes = yield* subscribeDurable(input.aggregateID) @@ -540,7 +531,7 @@ export const layerWith = (options?: LayerOptions) => const read = Effect.suspend(() => readAfter(input.aggregateID, sequence)).pipe( Effect.tap((events) => Effect.sync(() => { - sequence = events.at(-1)?.durable.seq ?? sequence + sequence = events.at(-1)?.durable?.seq ?? sequence }), ), ) @@ -562,14 +553,10 @@ export const layerWith = (options?: LayerOptions) => }) }) - const project = (definition: D, projector: Subscriber): Effect.Effect => + const project = (definition: D, projector: Subscriber): Effect.Effect => Effect.sync(() => { const list = projectors.get(definition.type) ?? [] - list.push((event) => - isPayload(definition)(event) - ? projector(event) - : Effect.die(`Published event ${event.type} does not match its definition`), - ) + list.push((event) => projector(event as Payload)) projectors.set(definition.type, list) }) diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index 8f856d2288..a97afb58b6 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -1,6 +1,5 @@ import { castDraft, produce, type WritableDraft } from "immer" -import { Effect, Match } from "effect" -import type { UncommittedPayload } from "@opencode-ai/schema/event" +import { Effect } from "effect" import { SessionEvent } from "./event" import { SessionMessage } from "./message" @@ -8,8 +7,6 @@ export type MemoryState = { messages: SessionMessage.Message[] } -export type Input = UncommittedPayload<(typeof SessionEvent.Definitions)[number]> - export interface Adapter { readonly getCurrentAssistant: () => Effect.Effect readonly getAssistant: (messageID: SessionMessage.ID) => Effect.Effect @@ -78,7 +75,7 @@ export function memory(state: MemoryState): Adapter { } } -export function update(adapter: Adapter, event: Input) { +export function update(adapter: Adapter, event: SessionEvent.Event) { type DraftAssistant = WritableDraft type DraftTool = WritableDraft type DraftText = WritableDraft @@ -101,8 +98,8 @@ export function update(adapter: Adapter, event: Input) { if (assistant) yield* adapter.updateAssistant(produce(assistant, recipe)) }) - return Match.value(event).pipe( - Match.discriminatorsExhaustive("type")({ + return Effect.gen(function* () { + yield* SessionEvent.All.match(event, { "session.next.agent.switched": (event) => { return adapter.appendMessage( SessionMessage.AgentSwitched.make({ @@ -391,8 +388,8 @@ export function update(adapter: Adapter, event: Input) { "session.next.revert.staged": () => Effect.void, "session.next.revert.cleared": () => Effect.void, "session.next.revert.committed": () => Effect.void, - }), - ) + }) + }) } export * as SessionMessageUpdater from "./message-updater" diff --git a/packages/core/test/event.test.ts b/packages/core/test/event.test.ts index 2284bba7e4..26d93b4d99 100644 --- a/packages/core/test/event.test.ts +++ b/packages/core/test/event.test.ts @@ -138,6 +138,50 @@ describe("EventV2", () => { }), ) + it.effect("preserves same-type projector routing across durable versions", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const historical = EventV2.define({ + type: "test.projector-version", + durable: { version: 1, aggregate: "id" }, + schema: { id: Schema.String }, + }) + const current = EventV2.define({ + type: "test.projector-version", + durable: { version: 2, aggregate: "id" }, + schema: { id: Schema.String }, + }) + const received = new Array() + yield* events.project(historical, (event) => Effect.sync(() => received.push(event))) + + const published = yield* events.publish(current, { id: "aggregate" }) + + expect(received).toEqual([published]) + }), + ) + + it.effect("preserves same-type subscription routing across durable versions", () => + Effect.gen(function* () { + const events = yield* EventV2.Service + const historical = EventV2.define({ + type: "test.subscription-version", + durable: { version: 1, aggregate: "id" }, + schema: { id: Schema.String }, + }) + const current = EventV2.define({ + type: "test.subscription-version", + durable: { version: 2, aggregate: "id" }, + schema: { id: Schema.String }, + }) + const fiber = yield* events.subscribe(historical).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped) + yield* Effect.yieldNow + + const published = yield* events.publish(current, { id: "aggregate" }) + + expect(Array.from(yield* Fiber.join(fiber))).toEqual([published]) + }), + ) + it.effect("publishes to typed and wildcard subscriptions", () => Effect.gen(function* () { const events = yield* EventV2.Service diff --git a/packages/core/test/session-runner-tool-events.test.ts b/packages/core/test/session-runner-tool-events.test.ts index f96ea4dea2..2fe34c6276 100644 --- a/packages/core/test/session-runner-tool-events.test.ts +++ b/packages/core/test/session-runner-tool-events.test.ts @@ -17,7 +17,14 @@ const capture = () => { const events = EventV2.Service.of({ publish: (definition, data) => Effect.sync(() => { - const event = { id: EventV2.ID.create(), type: definition.type, data } as EventV2.Payload + const event = { + id: EventV2.ID.create(), + type: definition.type, + ...(definition.durable + ? { durable: { aggregateID: sessionID, seq: published.length, version: definition.durable.version } } + : {}), + data, + } as EventV2.Payload published.push({ type: definition.durable ? EventV2.versionedType(definition.type, definition.durable.version) diff --git a/packages/opencode/test/server/httpapi-public-openapi.test.ts b/packages/opencode/test/server/httpapi-public-openapi.test.ts index d9eb3643a9..73a0158795 100644 --- a/packages/opencode/test/server/httpapi-public-openapi.test.ts +++ b/packages/opencode/test/server/httpapi-public-openapi.test.ts @@ -106,7 +106,8 @@ describe("PublicApi OpenAPI v2 errors", () => { expect(durable?.required).toContain("durable") expect(durable?.properties?.durable).toBeDefined() - expect(live?.properties?.durable).toBeUndefined() + expect(live?.required).not.toContain("durable") + expect(Reflect.get(live?.properties?.durable ?? {}, "not")).toEqual({}) }) test("preserves /api auth responses", () => { diff --git a/packages/opencode/test/server/httpapi-v2-location.test.ts b/packages/opencode/test/server/httpapi-v2-location.test.ts index e178062e7a..172e0bd4a1 100644 --- a/packages/opencode/test/server/httpapi-v2-location.test.ts +++ b/packages/opencode/test/server/httpapi-v2-location.test.ts @@ -23,6 +23,7 @@ function request(route: string, directory: string, init: RequestInit = {}) { const Event = Schema.Struct({ id: EventV2.ID, type: Schema.String, + durable: Schema.optional(Schema.Struct({ aggregateID: Schema.String, seq: Schema.Int, version: Schema.Int })), location: Schema.optional(Location.Ref), data: Schema.Unknown, }) @@ -81,12 +82,15 @@ describe("v2 location HttpApi", () => { const reader = response.body!.getReader() const connected = await readEvent(reader) expect(connected.type).toBe("server.connected") + expect(connected).not.toHaveProperty("durable") expect(connected.location).toBeUndefined() const created = await request("/session", publisher.path, { method: "POST" }) expect(created.status).toBe(200) + const session = (await created.json()) as { id: string } expect(await readEventType(reader, "session.created")).toMatchObject({ type: "session.created", + durable: { aggregateID: session.id, seq: 0, version: 1 }, location: { directory: publisher.path }, data: { sessionID: expect.any(String) }, }) diff --git a/packages/opencode/test/v2/session-message-updater.test.ts b/packages/opencode/test/v2/session-message-updater.test.ts index 5611f6ed70..56478a4d1d 100644 --- a/packages/opencode/test/v2/session-message-updater.test.ts +++ b/packages/opencode/test/v2/session-message-updater.test.ts @@ -5,9 +5,16 @@ import { SessionID } from "../../src/session/schema" import { EventV2 } from "@opencode-ai/core/event" import { ModelV2 } from "@opencode-ai/core/model" import { ProviderV2 } from "@opencode-ai/core/provider" +import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionMessageUpdater } from "@opencode-ai/core/session/message-updater" import { SessionMessage } from "@opencode-ai/core/session/message" +function durable(sessionID: SessionID, seq?: number): { aggregateID: SessionID; seq: number; version: 1 } +function durable(sessionID: SessionID, seq: number, version: 2): { aggregateID: SessionID; seq: number; version: 2 } +function durable(sessionID: SessionID, seq = 0, version: 1 | 2 = 1) { + return { aggregateID: sessionID, seq, version } +} + test.skip("step snapshots carry over to assistant messages", () => { const state: SessionMessageUpdater.MemoryState = { messages: [] } const sessionID = SessionID.make("session") @@ -16,6 +23,7 @@ test.skip("step snapshots carry over to assistant messages", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -29,7 +37,7 @@ test.skip("step snapshots carry over to assistant messages", () => { }, snapshot: "before", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) expect(state.messages).toEqual([]) @@ -37,6 +45,7 @@ test.skip("step snapshots carry over to assistant messages", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 1, 2), type: "session.next.step.ended", data: { sessionID, @@ -52,7 +61,7 @@ test.skip("step snapshots carry over to assistant messages", () => { }, snapshot: "after", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) expect(state.messages[0]?.type).toBe("assistant") @@ -69,6 +78,7 @@ test.skip("text ended populates assistant text content", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -81,12 +91,13 @@ test.skip("text ended populates assistant text content", () => { variant: ModelV2.VariantID.make("default"), }, }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 1), type: "session.next.text.started", data: { sessionID, @@ -94,12 +105,13 @@ test.skip("text ended populates assistant text content", () => { timestamp: DateTime.makeUnsafe(2), textID: "text-1", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 2), type: "session.next.text.ended", data: { sessionID, @@ -108,7 +120,7 @@ test.skip("text ended populates assistant text content", () => { textID: "text-1", text: "hello assistant", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) expect(state.messages[0]?.type).toBe("assistant") @@ -125,6 +137,7 @@ test.skip("tool completion stores completed timestamp", () => { Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID), type: "session.next.step.started", data: { sessionID, @@ -137,12 +150,13 @@ test.skip("tool completion stores completed timestamp", () => { variant: ModelV2.VariantID.make("default"), }, }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 1), type: "session.next.tool.input.started", data: { sessionID, @@ -151,12 +165,13 @@ test.skip("tool completion stores completed timestamp", () => { callID, name: "bash", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 2), type: "session.next.tool.called", data: { sessionID, @@ -167,12 +182,13 @@ test.skip("tool completion stores completed timestamp", () => { input: { command: "pwd" }, provider: { executed: true, metadata: { fake: { source: "provider" } } }, }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 3), type: "session.next.tool.success", data: { sessionID, @@ -183,7 +199,7 @@ test.skip("tool completion stores completed timestamp", () => { content: [{ type: "text", text: "/tmp" }], provider: { executed: true, metadata: { fake: { status: "done" } } }, }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) expect(state.messages[0]?.type).toBe("assistant") @@ -203,6 +219,7 @@ test("compaction events reduce to compaction message only when completed", () => Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id, + durable: durable(sessionID), type: "session.next.compaction.started", data: { sessionID, @@ -210,7 +227,7 @@ test("compaction events reduce to compaction message only when completed", () => timestamp: DateTime.makeUnsafe(1), reason: "auto", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) expect(state.messages).toEqual([]) @@ -225,7 +242,7 @@ test("compaction events reduce to compaction message only when completed", () => timestamp: DateTime.makeUnsafe(2), text: "hello ", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( @@ -238,12 +255,13 @@ test("compaction events reduce to compaction message only when completed", () => timestamp: DateTime.makeUnsafe(3), text: "summary", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) Effect.runSync( SessionMessageUpdater.update(SessionMessageUpdater.memory(state), { id: EventV2.ID.create(), + durable: durable(sessionID, 3), type: "session.next.compaction.ended", data: { sessionID, @@ -253,7 +271,7 @@ test("compaction events reduce to compaction message only when completed", () => text: "final summary", recent: "recent context", }, - } satisfies SessionMessageUpdater.Input), + } satisfies SessionEvent.Event), ) expect(state.messages).toHaveLength(1) diff --git a/packages/protocol/src/groups/event.ts b/packages/protocol/src/groups/event.ts index 703db3a3d6..cf6195a264 100644 --- a/packages/protocol/src/groups/event.ts +++ b/packages/protocol/src/groups/event.ts @@ -11,18 +11,21 @@ const fields = { location: Schema.optional(Location.Ref), } +const DurableEnvelope = Schema.Struct({ aggregateID: Schema.String, seq: Schema.Int, version: Schema.Int }) + const schema = (definitions: ReadonlyArray) => Schema.Union([ ...definitions.map((definition) => definition.durable ? Schema.Struct({ ...fields, - durable: Event.durableEnvelope(definition.durable.version), + durable: DurableEnvelope, type: Schema.Literal(definition.type), data: definition.data, }).annotate({ identifier: `V2Event.${definition.type}` }) : Schema.Struct({ ...fields, + durable: Schema.optional(Schema.Never), type: Schema.Literal(definition.type), data: definition.data, }).annotate({ identifier: `V2Event.${definition.type}` }), @@ -32,6 +35,7 @@ const schema = (definitions: ReadonlyArray) => : [ Schema.Struct({ ...fields, + durable: Schema.optional(Schema.Never), type: Schema.Literal("server.connected"), data: Schema.Struct({}), }).annotate({ identifier: "V2Event.server.connected" }), @@ -62,4 +66,5 @@ export const makeEventGroup = (definitions: ReadonlyArray) => make(d const event = make(EventManifest.ServerDefinitions) export const EventGroup = event.group -export type Event = typeof event.schema.Type +export const EventSchema = event.schema +export type Event = typeof EventSchema.Type diff --git a/packages/protocol/test/event.test.ts b/packages/protocol/test/event.test.ts new file mode 100644 index 0000000000..e254087bd0 --- /dev/null +++ b/packages/protocol/test/event.test.ts @@ -0,0 +1,27 @@ +import { describe, expect, test } from "bun:test" +import { Event } from "@opencode-ai/schema/event" +import { Schema } from "effect" +import { EventSchema } from "../src/groups/event" + +describe("EventSchema", () => { + test("requires durable metadata on durable events", () => { + expect( + Schema.is(EventSchema)({ + id: Event.ID.create(), + type: "session.created", + data: { sessionID: "session" }, + }), + ).toBe(false) + }) + + test("rejects durable metadata on live events", () => { + expect( + Schema.is(EventSchema)({ + id: Event.ID.create(), + type: "server.connected", + durable: { aggregateID: "aggregate", seq: 0, version: 1 }, + data: {}, + }), + ).toBe(false) + }) +}) diff --git a/packages/schema/src/event.ts b/packages/schema/src/event.ts index 168eb1eaf7..243667e543 100644 --- a/packages/schema/src/event.ts +++ b/packages/schema/src/event.ts @@ -3,7 +3,7 @@ export * as Event from "./event" import { Schema } from "effect" import { ascending } from "./identifier" import { Location } from "./location" -import { NonNegativeInt, statics } from "./schema" +import { statics } from "./schema" export const ID = Schema.String.check(Schema.isStartsWith("evt_")).pipe( Schema.brand("Event.ID"), @@ -16,15 +16,17 @@ export type DurableOptions = { readonly aggregate: string } -export type DurableEnvelope = { +export type DurableEnvelope = { readonly aggregateID: string readonly seq: number - readonly version: Version + readonly version: number } -export const durableEnvelope = (version: Version) => - Schema.Struct({ aggregateID: Schema.String, seq: NonNegativeInt, version: Schema.Literal(version) }) - +const PublishedDurableEnvelope = Schema.Struct({ + aggregateID: Schema.String, + seq: Schema.Number, + version: Schema.Number, +}) const NoDurableEnvelope = Schema.optional(Schema.Never) export type LiveDefinition< @@ -46,7 +48,10 @@ export type DurableDefinition< readonly durable: Durability } -export type Definition = LiveDefinition | DurableDefinition +export type Definition< + Type extends string = string, + DataSchema extends Schema.Codec = Schema.Codec, +> = LiveDefinition | DurableDefinition export type Data = Schema.Schema.Type @@ -60,18 +65,10 @@ export type UncommittedPayload = D extends De } : never -export type DurablePublishedPayload = UncommittedPayload & { - readonly durable: DurableEnvelope -} - -export type LivePublishedPayload = UncommittedPayload & { - readonly durable?: never -} - export type PublishedPayload = D extends DurableDefinition - ? DurablePublishedPayload + ? UncommittedPayload & { readonly durable: DurableEnvelope } : D extends LiveDefinition - ? LivePublishedPayload + ? UncommittedPayload & { readonly durable?: never } : never export type Payload = PublishedPayload @@ -117,7 +114,7 @@ export function define(input: { } if (input.durable) { return Object.assign( - Schema.Struct({ ...fields, durable: durableEnvelope(input.durable.version) }).annotate({ + Schema.Struct({ ...fields, durable: PublishedDurableEnvelope }).annotate({ identifier: input.type, }), { type: input.type, durable: input.durable, data }, diff --git a/packages/schema/test/event.test.ts b/packages/schema/test/event.test.ts index 3b119368cc..02b1feda76 100644 --- a/packages/schema/test/event.test.ts +++ b/packages/schema/test/event.test.ts @@ -56,9 +56,6 @@ describe("public event schemas", () => { data: { id: "aggregate" }, }), ).toBe(false) - expect(Schema.is(definition)({ ...payload, durable: { ...payload.durable, seq: -1 } })).toBe(false) - expect(Schema.is(definition)({ ...payload, durable: { ...payload.durable, version: 2 } })).toBe(false) - // @ts-expect-error Published durable payloads require commit metadata. const missing: typeof definition.Type = { id: Event.ID.create(), type: definition.type, data: { id: "aggregate" } } void missing @@ -119,10 +116,10 @@ describe("public event schemas", () => { // @ts-expect-error Durable union members require commit metadata. const uncommitted: Mixed = { id: Event.ID.create(), type: durable.type, data: { id: "aggregate" } } + // @ts-expect-error Live union members cannot carry durable commit metadata. const falselyCommitted: Mixed = { id: Event.ID.create(), type: live.type, - // @ts-expect-error Live union members cannot carry durable commit metadata. durable: { aggregateID: "aggregate", seq: 0, version: 2 }, data: { value: "value" }, } diff --git a/packages/sdk/js/script/build.ts b/packages/sdk/js/script/build.ts index 72f4e3f3e9..2b2032f714 100755 --- a/packages/sdk/js/script/build.ts +++ b/packages/sdk/js/script/build.ts @@ -58,6 +58,17 @@ if (sseTypesPatched === sseTypesSource) { } await Bun.write(sseTypesPath, sseTypesPatched) +// OpenAPI represents Schema.Never as `not: {}`, which @hey-api currently +// widens to unknown. Preserve impossible optional event fields as never. +const eventTypesPath = "./src/v2/gen/types.gen.ts" +const eventTypesFile = Bun.file(eventTypesPath) +const eventTypesSource = await eventTypesFile.text() +const eventTypesPatched = eventTypesSource.replaceAll(" durable?: unknown", " durable?: never") +if (eventTypesPatched === eventTypesSource) { + throw new Error(`Event never patch did not apply; @hey-api/openapi-ts output may have changed (${eventTypesPath})`) +} +await Bun.write(eventTypesPath, eventTypesPatched) + await $`bun prettier --write src/gen` await $`bun prettier --write src/v2` await $`rm -rf dist` diff --git a/packages/sdk/js/src/v2/gen/types.gen.ts b/packages/sdk/js/src/v2/gen/types.gen.ts index b653c75e1b..ce480c6ae6 100644 --- a/packages/sdk/js/src/v2/gen/types.gen.ts +++ b/packages/sdk/js/src/v2/gen/types.gen.ts @@ -4329,6 +4329,7 @@ export type V2EventModelsDevRefreshed = { [key: string]: unknown } location?: LocationRef + durable?: never type: "models-dev.refreshed" data: { [key: string]: unknown @@ -4341,6 +4342,7 @@ export type V2EventIntegrationUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "integration.updated" data: { [key: string]: unknown @@ -4353,6 +4355,7 @@ export type V2EventIntegrationConnectionUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "integration.connection.updated" data: { integrationID: string @@ -4365,6 +4368,7 @@ export type V2EventCatalogUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "catalog.updated" data: { [key: string]: unknown @@ -4380,7 +4384,7 @@ export type V2EventSessionCreated = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.created" data: { @@ -4398,7 +4402,7 @@ export type V2EventSessionUpdated = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.updated" data: { @@ -4416,7 +4420,7 @@ export type V2EventSessionDeleted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.deleted" data: { @@ -4434,7 +4438,7 @@ export type V2EventMessageUpdated = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "message.updated" data: { @@ -4452,7 +4456,7 @@ export type V2EventMessageRemoved = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "message.removed" data: { @@ -4470,7 +4474,7 @@ export type V2EventMessagePartUpdated = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "message.part.updated" data: { @@ -4489,7 +4493,7 @@ export type V2EventMessagePartRemoved = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "message.part.removed" data: { @@ -4508,7 +4512,7 @@ export type V2EventSessionNextAgentSwitched = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.agent.switched" data: { @@ -4528,7 +4532,7 @@ export type V2EventSessionNextModelSwitched = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.model.switched" data: { @@ -4552,7 +4556,7 @@ export type V2EventSessionNextMoved = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.moved" data: { @@ -4572,7 +4576,7 @@ export type V2EventSessionNextPrompted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.prompted" data: { @@ -4593,7 +4597,7 @@ export type V2EventSessionNextPromptAdmitted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.prompt.admitted" data: { @@ -4614,7 +4618,7 @@ export type V2EventSessionNextContextUpdated = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.context.updated" data: { @@ -4634,7 +4638,7 @@ export type V2EventSessionNextSynthetic = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.synthetic" data: { @@ -4654,7 +4658,7 @@ export type V2EventSessionNextShellStarted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.shell.started" data: { @@ -4675,7 +4679,7 @@ export type V2EventSessionNextShellEnded = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.shell.ended" data: { @@ -4695,7 +4699,7 @@ export type V2EventSessionNextStepStarted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.step.started" data: { @@ -4721,7 +4725,7 @@ export type V2EventSessionNextStepEnded = { durable: { aggregateID: string seq: number - version: 2 + version: number } type: "session.next.step.ended" data: { @@ -4753,7 +4757,7 @@ export type V2EventSessionNextStepFailed = { durable: { aggregateID: string seq: number - version: 2 + version: number } type: "session.next.step.failed" data: { @@ -4773,7 +4777,7 @@ export type V2EventSessionNextTextStarted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.text.started" data: { @@ -4790,6 +4794,7 @@ export type V2EventSessionNextTextDelta = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.next.text.delta" data: { timestamp: number @@ -4809,7 +4814,7 @@ export type V2EventSessionNextTextEnded = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.text.ended" data: { @@ -4830,7 +4835,7 @@ export type V2EventSessionNextReasoningStarted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.reasoning.started" data: { @@ -4852,6 +4857,7 @@ export type V2EventSessionNextReasoningDelta = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.next.reasoning.delta" data: { timestamp: number @@ -4871,7 +4877,7 @@ export type V2EventSessionNextReasoningEnded = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.reasoning.ended" data: { @@ -4897,7 +4903,7 @@ export type V2EventSessionNextToolInputStarted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.tool.input.started" data: { @@ -4915,6 +4921,7 @@ export type V2EventSessionNextToolInputDelta = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.next.tool.input.delta" data: { timestamp: number @@ -4934,7 +4941,7 @@ export type V2EventSessionNextToolInputEnded = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.tool.input.ended" data: { @@ -4955,7 +4962,7 @@ export type V2EventSessionNextToolCalled = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.tool.called" data: { @@ -4987,7 +4994,7 @@ export type V2EventSessionNextToolProgress = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.tool.progress" data: { @@ -5011,7 +5018,7 @@ export type V2EventSessionNextToolSuccess = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.tool.success" data: { @@ -5045,7 +5052,7 @@ export type V2EventSessionNextToolFailed = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.tool.failed" data: { @@ -5075,7 +5082,7 @@ export type V2EventSessionNextRetried = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.retried" data: { @@ -5095,7 +5102,7 @@ export type V2EventSessionNextCompactionStarted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.compaction.started" data: { @@ -5112,6 +5119,7 @@ export type V2EventSessionNextCompactionDelta = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.next.compaction.delta" data: { timestamp: number @@ -5130,7 +5138,7 @@ export type V2EventSessionNextCompactionEnded = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.compaction.ended" data: { @@ -5152,7 +5160,7 @@ export type V2EventSessionNextRevertStaged = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.revert.staged" data: { @@ -5177,7 +5185,7 @@ export type V2EventSessionNextRevertCleared = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.revert.cleared" data: { @@ -5195,7 +5203,7 @@ export type V2EventSessionNextRevertCommitted = { durable: { aggregateID: string seq: number - version: 1 + version: number } type: "session.next.revert.committed" data: { @@ -5211,6 +5219,7 @@ export type V2EventMessagePartDelta = { [key: string]: unknown } location?: LocationRef + durable?: never type: "message.part.delta" data: { sessionID: string @@ -5227,6 +5236,7 @@ export type V2EventSessionDiff = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.diff" data: { sessionID: string @@ -5240,6 +5250,7 @@ export type V2EventSessionError = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.error" data: { sessionID?: string @@ -5261,6 +5272,7 @@ export type V2EventInstallationUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "installation.updated" data: { version: string @@ -5273,6 +5285,7 @@ export type V2EventInstallationUpdateAvailable = { [key: string]: unknown } location?: LocationRef + durable?: never type: "installation.update-available" data: { version: string @@ -5285,6 +5298,7 @@ export type V2EventFileEdited = { [key: string]: unknown } location?: LocationRef + durable?: never type: "file.edited" data: { file: string @@ -5297,6 +5311,7 @@ export type V2EventReferenceUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "reference.updated" data: { [key: string]: unknown @@ -5309,6 +5324,7 @@ export type V2EventPermissionV2Asked = { [key: string]: unknown } location?: LocationRef + durable?: never type: "permission.v2.asked" data: { id: string @@ -5329,6 +5345,7 @@ export type V2EventPermissionV2Replied = { [key: string]: unknown } location?: LocationRef + durable?: never type: "permission.v2.replied" data: { sessionID: string @@ -5343,6 +5360,7 @@ export type V2EventPluginAdded = { [key: string]: unknown } location?: LocationRef + durable?: never type: "plugin.added" data: { id: string @@ -5355,6 +5373,7 @@ export type V2EventProjectDirectoriesUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "project.directories.updated" data: { projectID: string @@ -5367,6 +5386,7 @@ export type V2EventFileWatcherUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "file.watcher.updated" data: { file: string @@ -5380,6 +5400,7 @@ export type V2EventPtyCreated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "pty.created" data: { info: Pty @@ -5392,6 +5413,7 @@ export type V2EventPtyUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "pty.updated" data: { info: Pty @@ -5404,6 +5426,7 @@ export type V2EventPtyExited = { [key: string]: unknown } location?: LocationRef + durable?: never type: "pty.exited" data: { id: string @@ -5417,6 +5440,7 @@ export type V2EventPtyDeleted = { [key: string]: unknown } location?: LocationRef + durable?: never type: "pty.deleted" data: { id: string @@ -5429,6 +5453,7 @@ export type V2EventQuestionV2Asked = { [key: string]: unknown } location?: LocationRef + durable?: never type: "question.v2.asked" data: { id: string @@ -5447,6 +5472,7 @@ export type V2EventQuestionV2Replied = { [key: string]: unknown } location?: LocationRef + durable?: never type: "question.v2.replied" data: { sessionID: string @@ -5461,6 +5487,7 @@ export type V2EventQuestionV2Rejected = { [key: string]: unknown } location?: LocationRef + durable?: never type: "question.v2.rejected" data: { sessionID: string @@ -5474,6 +5501,7 @@ export type V2EventTodoUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "todo.updated" data: { sessionID: string @@ -5487,6 +5515,7 @@ export type V2EventLspUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "lsp.updated" data: { [key: string]: unknown @@ -5499,6 +5528,7 @@ export type V2EventPermissionAsked = { [key: string]: unknown } location?: LocationRef + durable?: never type: "permission.asked" data: { id: string @@ -5522,6 +5552,7 @@ export type V2EventPermissionReplied = { [key: string]: unknown } location?: LocationRef + durable?: never type: "permission.replied" data: { sessionID: string @@ -5536,6 +5567,7 @@ export type V2EventTuiPromptAppend = { [key: string]: unknown } location?: LocationRef + durable?: never type: "tui.prompt.append" data: { text: string @@ -5548,6 +5580,7 @@ export type V2EventTuiCommandExecute = { [key: string]: unknown } location?: LocationRef + durable?: never type: "tui.command.execute" data: { command: @@ -5577,6 +5610,7 @@ export type V2EventTuiToastShow = { [key: string]: unknown } location?: LocationRef + durable?: never type: "tui.toast.show" data: { title?: string @@ -5592,6 +5626,7 @@ export type V2EventTuiSessionSelect = { [key: string]: unknown } location?: LocationRef + durable?: never type: "tui.session.select" data: { /** @@ -5607,6 +5642,7 @@ export type V2EventMcpToolsChanged = { [key: string]: unknown } location?: LocationRef + durable?: never type: "mcp.tools.changed" data: { server: string @@ -5619,6 +5655,7 @@ export type V2EventMcpBrowserOpenFailed = { [key: string]: unknown } location?: LocationRef + durable?: never type: "mcp.browser.open.failed" data: { mcpName: string @@ -5632,6 +5669,7 @@ export type V2EventCommandExecuted = { [key: string]: unknown } location?: LocationRef + durable?: never type: "command.executed" data: { name: string @@ -5647,6 +5685,7 @@ export type V2EventProjectUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "project.updated" data: { id: string @@ -5679,6 +5718,7 @@ export type V2EventSessionStatus = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.status" data: { sessionID: string @@ -5692,6 +5732,7 @@ export type V2EventSessionIdle = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.idle" data: { sessionID: string @@ -5704,6 +5745,7 @@ export type V2EventQuestionAsked = { [key: string]: unknown } location?: LocationRef + durable?: never type: "question.asked" data: { id: string @@ -5722,6 +5764,7 @@ export type V2EventQuestionReplied = { [key: string]: unknown } location?: LocationRef + durable?: never type: "question.replied" data: { sessionID: string @@ -5736,6 +5779,7 @@ export type V2EventQuestionRejected = { [key: string]: unknown } location?: LocationRef + durable?: never type: "question.rejected" data: { sessionID: string @@ -5749,6 +5793,7 @@ export type V2EventSessionCompacted = { [key: string]: unknown } location?: LocationRef + durable?: never type: "session.compacted" data: { sessionID: string @@ -5761,6 +5806,7 @@ export type V2EventVcsBranchUpdated = { [key: string]: unknown } location?: LocationRef + durable?: never type: "vcs.branch.updated" data: { branch?: string @@ -5773,6 +5819,7 @@ export type V2EventWorkspaceReady = { [key: string]: unknown } location?: LocationRef + durable?: never type: "workspace.ready" data: { name: string @@ -5785,6 +5832,7 @@ export type V2EventWorkspaceFailed = { [key: string]: unknown } location?: LocationRef + durable?: never type: "workspace.failed" data: { message: string @@ -5797,6 +5845,7 @@ export type V2EventWorkspaceStatus = { [key: string]: unknown } location?: LocationRef + durable?: never type: "workspace.status" data: { workspaceID: string @@ -5810,6 +5859,7 @@ export type V2EventWorktreeReady = { [key: string]: unknown } location?: LocationRef + durable?: never type: "worktree.ready" data: { name: string @@ -5823,6 +5873,7 @@ export type V2EventWorktreeFailed = { [key: string]: unknown } location?: LocationRef + durable?: never type: "worktree.failed" data: { message: string @@ -5835,6 +5886,7 @@ export type V2EventServerConnected = { [key: string]: unknown } location?: LocationRef + durable?: never type: "server.connected" data: { [key: string]: unknown @@ -5847,6 +5899,7 @@ export type V2EventGlobalDisposed = { [key: string]: unknown } location?: LocationRef + durable?: never type: "global.disposed" data: { [key: string]: unknown