From b7f36dc9f3ffebec93e16516852265378d177bc7 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Sun, 27 Sep 2026 02:22:25 -0700 Subject: [PATCH 1/2] fix(clients): older clients skip event types they don't know instead of breaking A thread stream event with a type the client build does not know used to fail the chunk decode. The RPC client turns that into a defect, which ends the thread subscription, and resubscribing replays the same event. The thread stream item union now has a decode-only fallback arm. It takes an `event` item whose `type` is not a known domain event type and decodes it as `{ kind: "unknown-event", sequence, eventType }`. A known type with a payload that does not decode still fails. The client-runtime reducer skips these items with a debug log and still advances its resume cursor past them. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/state/threads-sync.test.ts | 63 +++++++++++++++++++ packages/client-runtime/src/state/threads.ts | 40 +++++++++--- .../contracts/src/orchestrationV2.test.ts | 47 ++++++++++++++ packages/contracts/src/orchestrationV2.ts | 45 +++++++++++++ 4 files changed, 186 insertions(+), 9 deletions(-) diff --git a/packages/client-runtime/src/state/threads-sync.test.ts b/packages/client-runtime/src/state/threads-sync.test.ts index bf925382e941..32036b150a35 100644 --- a/packages/client-runtime/src/state/threads-sync.test.ts +++ b/packages/client-runtime/src/state/threads-sync.test.ts @@ -1,6 +1,7 @@ import { EnvironmentId, EventId, + MessageId, ORCHESTRATION_V2_WS_METHODS, ThreadId, TurnItemId, @@ -1771,6 +1772,68 @@ describe("EnvironmentThreads", () => { }), ); + it.effect("skips an unknown event type and resumes after it", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ cached: BASE_PROJECTION, completionMarker: true }); + const unknown = (sequence: number): OrchestrationV2ThreadStreamItem => ({ + kind: "unknown-event", + sequence, + eventType: "run.background-work-cancelled", + }); + const occurredAt = DateTime.makeUnsafe("2026-06-20T01:00:00.000Z"); + yield* Queue.offerAll(harness.inputs, [ + { + kind: "event", + sequence: CACHED_SNAPSHOT_SEQUENCE + 1, + event: { + id: EventId.make("event-message"), + type: "message.updated", + threadId: THREAD_ID, + occurredAt, + payload: { + id: MessageId.make("message-before"), + threadId: THREAD_ID, + runId: null, + nodeId: null, + role: "assistant", + text: "Before", + streaming: false, + attachments: [], + createdBy: "agent", + creationSource: "provider", + createdAt: occurredAt, + updatedAt: occurredAt, + }, + }, + }, + unknown(CACHED_SNAPSHOT_SEQUENCE + 2), + titleUpdated("After", CACHED_SNAPSHOT_SEQUENCE + 3), + // A trailing unknown event must still advance the resume cursor. + unknown(CACHED_SNAPSHOT_SEQUENCE + 4), + synchronized(), + ]); + const live = yield* awaitThreadState( + harness.observed, + (value) => + value.status === "live" && + Option.isSome(value.data) && + value.data.value.thread.title === "After", + ); + expect(Option.isNone(live.error)).toBe(true); + expect(Option.getOrThrow(live.data).messages.map((message) => message.text)).toEqual([ + "Before", + ]); + + yield* harness.replaceSession; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(harness.subscriptionCount)) >= 2) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(harness.subscriptionCount)).toBe(2); + expect(yield* Ref.get(harness.lastSubscribeAfterSequence)).toBe(CACHED_SNAPSHOT_SEQUENCE + 4); + }), + ); + it.effect("resumes replacement sessions from the latest applied sequence", () => Effect.gen(function* () { const harness = yield* makeHarness({ cached: BASE_PROJECTION, completionMarker: true }); diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index 5997dfa69c33..52d127ed4593 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -435,17 +435,36 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make }); type EventItem = Extract; + type SequencedItem = Extract< + OrchestrationV2ThreadStreamItem, + { kind: "event" | "unknown-event" } + >; const applyEventsLocked = Effect.fn("EnvironmentThreadState.applyEventsLocked")(function* ( - items: ReadonlyArray, + items: ReadonlyArray, ) { - let sequence = yield* SubscriptionRef.get(lastSequence); - const fresh = items.filter((item) => { - if (item.sequence <= sequence) return false; + const appliedSequence = yield* SubscriptionRef.get(lastSequence); + let sequence = appliedSequence; + const fresh: EventItem[] = []; + for (const item of items) { + if (item.sequence <= sequence) continue; + // An event type from a newer server still moves the resume cursor past it. sequence = item.sequence; - return true; - }); - if (fresh.length === 0) return; + if (item.kind === "event") { + fresh.push(item); + continue; + } + yield* Effect.logDebug("Skipped a thread event type this client does not know.").pipe( + Effect.annotateLogs({ + environmentId, + threadId, + eventType: item.eventType, + sequence: item.sequence, + }), + ); + } + if (sequence === appliedSequence) return; yield* SubscriptionRef.set(lastSequence, sequence); + if (fresh.length === 0) return; const waiting = yield* Ref.get(awaitingCompletion); // Apply against the latest projection/history in one update so a concurrent @@ -593,9 +612,12 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make ) { yield* applyLock.withPermits(1)( Effect.gen(function* () { - let events: EventItem[] = []; + let events: SequencedItem[] = []; for (const item of items) { - if (item.kind === "event" && item.event.type !== "thread.deleted") { + if ( + item.kind === "unknown-event" || + (item.kind === "event" && item.event.type !== "thread.deleted") + ) { events.push(item); continue; } diff --git a/packages/contracts/src/orchestrationV2.test.ts b/packages/contracts/src/orchestrationV2.test.ts index 3677dbd10175..b6acbace84bf 100644 --- a/packages/contracts/src/orchestrationV2.test.ts +++ b/packages/contracts/src/orchestrationV2.test.ts @@ -29,6 +29,7 @@ import { OrchestrationV2ProviderCapabilities, OrchestrationV2ProviderThread, OrchestrationV2ProviderThreadJson, + OrchestrationV2RpcSchemas, OrchestrationV2ShellSnapshot, OrchestrationV2SubscribeThreadInput, OrchestrationV2Subagent, @@ -123,6 +124,52 @@ describe("orchestration V2 contracts", () => { } }); + it("decodes thread event types from a newer server as skippable items", () => { + const decodeWireItems = Schema.decodeUnknownSync( + Schema.toCodecJson(Schema.Array(OrchestrationV2RpcSchemas.subscribeThread.output)), + ); + const detached = (id: string, sequence: number) => ({ + kind: "event", + sequence, + event: { + id, + type: "provider-session.detached", + threadId: "thread-1", + occurredAt: DateTime.formatIso(now), + payload: { providerSessionId: "provider-session-1", detachedAt: DateTime.formatIso(now) }, + }, + }); + + const items = decodeWireItems([ + detached("event-1", 1), + { + kind: "event", + sequence: 2, + event: { + id: "event-2", + type: "run.background-work-cancelled", + threadId: "thread-1", + occurredAt: DateTime.formatIso(now), + payload: { runId: "run-1", restartCancelledBackgroundWork: [] }, + }, + }, + detached("event-3", 3), + ]); + + expect(items.map((item) => item.kind)).toEqual(["event", "unknown-event", "event"]); + expect(items[1]).toEqual({ + kind: "unknown-event", + sequence: 2, + eventType: "run.background-work-cancelled", + }); + // A known type with a broken payload is a real defect, not a newer event. + expect(() => + decodeWireItems([ + { ...detached("event-4", 4), event: { ...detached("event-4", 4).event, payload: {} } }, + ]), + ).toThrow(); + }); + it("negotiates bounded socket snapshots as an optional capability", () => { expect( decodeOrchestrationV2SubscribeThreadInput({ diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 5b298b5e6db2..bed162e3880e 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -1,6 +1,7 @@ import { OrchestrationMessageContext } from "./composerContext.ts"; import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; +import * as SchemaGetter from "effect/SchemaGetter"; import { CheckpointId, @@ -2868,6 +2869,48 @@ export const OrchestrationV2ThreadHistoryPage = Schema.Struct({ }); export type OrchestrationV2ThreadHistoryPage = typeof OrchestrationV2ThreadHistoryPage.Type; +const knownDomainEventTypes: ReadonlySet = new Set( + OrchestrationV2DomainEvent.members.flatMap((member) => { + const type = member.fields.type; + return "literals" in type ? type.literals : [type.literal]; + }), +); + +/** + * A thread event whose type this build does not know. Newer servers add event + * types; older clients decode them to this case and skip them, still advancing + * their resume cursor, instead of failing the whole subscription. A known type + * whose payload does not decode still fails. Decode-only: servers never send it. + */ +const OrchestrationV2UnknownThreadStreamEvent = Schema.Struct({ + kind: Schema.Literal("event"), + sequence: NonNegativeInt, + event: Schema.Struct({ + type: Schema.String.check( + Schema.makeFilter( + (type: string) => + !knownDomainEventTypes.has(type) || "A known event type must decode in full.", + ), + ), + }), +}).pipe( + Schema.decodeTo( + Schema.Struct({ + kind: Schema.Literal("unknown-event"), + sequence: NonNegativeInt, + eventType: Schema.String, + }), + { + decode: SchemaGetter.transform((item) => ({ + kind: "unknown-event" as const, + sequence: item.sequence, + eventType: item.event.type, + })), + encode: SchemaGetter.forbidden(() => "Servers never send unknown thread events."), + }, + ), +); + export const OrchestrationV2ThreadStreamItem = Schema.Union([ Schema.Struct({ kind: Schema.Literal("synchronized"), @@ -2890,6 +2933,8 @@ export const OrchestrationV2ThreadStreamItem = Schema.Union([ sequence: NonNegativeInt, event: OrchestrationV2DomainEvent, }), + // After the known arm: union members are tried in order. + OrchestrationV2UnknownThreadStreamEvent, ]); export type OrchestrationV2ThreadStreamItem = typeof OrchestrationV2ThreadStreamItem.Type; From c926cbb633e097c8c49afb6d45461102389d2fe7 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Sun, 27 Sep 2026 14:11:24 -0700 Subject: [PATCH 2/2] fix(clients): bound the unknown event type in the skip log Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/client-runtime/src/state/threads.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/packages/client-runtime/src/state/threads.ts b/packages/client-runtime/src/state/threads.ts index 52d127ed4593..5231958b346e 100644 --- a/packages/client-runtime/src/state/threads.ts +++ b/packages/client-runtime/src/state/threads.ts @@ -457,7 +457,8 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make Effect.annotateLogs({ environmentId, threadId, - eventType: item.eventType, + // Bounded: the type comes from a newer server and is not validated here. + eventType: item.eventType.slice(0, 64), sequence: item.sequence, }), );