From ec6e47b40c125ce830ea97b88dd72bd84193d159 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 23:52:40 -0700 Subject: [PATCH 1/3] fix(contracts): clients tolerate turn item types they don't know yet Turn item types are a closed union, so a server that adds one (such as the upcoming secret_request) made every older client fail to decode the thread snapshot and show "Could not synchronize the thread". Thread projections (turnItems, visibleTurnItems) and history pages now drop rows whose turn item type this build does not know, under both the runtime and JSON codecs. A turn-item.updated event carrying an unknown type decodes to the existing unknown-event case, so the client skips it and still advances its resume cursor. A known type with a broken payload still fails. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../ThreadLaunchService.test.ts | 5 +- .../contracts/src/orchestrationV2.test.ts | 129 ++++++++++++++++++ packages/contracts/src/orchestrationV2.ts | 94 +++++++++++-- 3 files changed, 212 insertions(+), 16 deletions(-) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 5e53250915f8..284d6d4b8e33 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -329,7 +329,10 @@ it.effect.each( assert.equal(wire.messages[0]?.text, task.prompt); assert.equal(wire.messages[0]?.scheduledTaskId, task.id); assert.equal(wire.messages[0]?.createdBy, createdBy); - const turnItem = wire.turnItems.find((item) => item.type === "user_message"); + const turnItem = wire.turnItems.find( + (item): item is Extract => + item.type === "user_message", + ); assert.equal(turnItem?.text, task.prompt); assert.equal(turnItem?.scheduledTaskId, task.id); }).pipe(Effect.provide(Layer.mergeAll(harness.layer, scheduledTasks))); diff --git a/packages/contracts/src/orchestrationV2.test.ts b/packages/contracts/src/orchestrationV2.test.ts index aae9cc40bd8c..543386ca20e3 100644 --- a/packages/contracts/src/orchestrationV2.test.ts +++ b/packages/contracts/src/orchestrationV2.test.ts @@ -34,6 +34,7 @@ import { OrchestrationV2ShellSnapshot, OrchestrationV2SubscribeThreadInput, OrchestrationV2Subagent, + OrchestrationV2ThreadHistoryPage, OrchestrationV2ThreadProjection, OrchestrationV2ThreadStreamItem, OrchestrationV2ThreadShell, @@ -175,6 +176,134 @@ describe("orchestration V2 contracts", () => { ).toThrow(); }); + it("skips turn item types from a newer server in snapshots and turn-item events", () => { + const decodeWireItems = Schema.decodeUnknownSync( + Schema.toCodecJson(Schema.Array(OrchestrationV2RpcSchemas.subscribeThread.output)), + ); + const decodeHistoryPage = Schema.decodeUnknownSync( + Schema.toCodecJson(OrchestrationV2ThreadHistoryPage), + ); + const item = (id: string, type: string, extra: Record) => ({ + id, + type, + threadId: "thread-1", + runId: null, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "completed", + title: null, + startedAt: null, + completedAt: null, + updatedAt: DateTime.formatIso(now), + ...extra, + }); + const known = item("item-known", "system_notice", { message: "Hello" }); + // A type no build of this client knows, standing in for a newer server's item. + const future = item("item-future", "secret_request", { secretRef: "ref-1" }); + const projected = (position: number, turnItem: { readonly id: string }) => ({ + position, + visibility: "local", + sourceThreadId: "thread-1", + sourceItemId: turnItem.id, + item: turnItem, + }); + const projection = { + thread: { + createdBy: "user", + creationSource: "web", + id: "thread-1", + projectId: "project-1", + title: "Thread", + providerInstanceId: "codex", + modelSelection: { instanceId: "codex", model: "gpt-5-codex" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: null, + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: "thread-1" }, + forkedFrom: null, + createdAt: DateTime.formatIso(now), + updatedAt: DateTime.formatIso(now), + archivedAt: null, + deletedAt: null, + }, + runs: [], + attempts: [], + nodes: [], + subagents: [], + providerSessions: [], + providerThreads: [], + providerTurns: [], + runtimeRequests: [], + messages: [], + plans: [], + turnItems: [future, known], + checkpointScopes: [], + checkpoints: [], + contextHandoffs: [], + contextTransfers: [], + visibleTurnItems: [projected(0, future), projected(1, known)], + updatedAt: DateTime.formatIso(now), + }; + const turnItemEvent = (sequence: number, payload: unknown) => ({ + kind: "event", + sequence, + event: { + id: `event-${sequence}`, + type: "turn-item.updated", + threadId: "thread-1", + occurredAt: DateTime.formatIso(now), + payload, + }, + }); + + const [snapshot, futureEvent, knownEvent] = decodeWireItems([ + { kind: "snapshot", snapshotSequence: 1, projection }, + turnItemEvent(2, future), + turnItemEvent(3, known), + ]); + + expect(snapshot).toMatchObject({ kind: "snapshot" }); + if (snapshot?.kind !== "snapshot") throw new Error("expected a snapshot"); + expect(snapshot.projection.turnItems.map((turnItem) => turnItem.id)).toEqual(["item-known"]); + expect(snapshot.projection.visibleTurnItems.map((row) => row.item.id)).toEqual(["item-known"]); + expect(futureEvent).toEqual({ + kind: "unknown-event", + sequence: 2, + eventType: "turn-item.updated", + }); + expect(knownEvent).toMatchObject({ + kind: "event", + event: { type: "turn-item.updated", payload: { id: "item-known", type: "system_notice" } }, + }); + expect( + decodeHistoryPage({ + snapshotSequence: 1, + items: [projected(0, future), projected(1, known)], + nextCursor: null, + hasMoreHistory: false, + }).items.map((row) => row.item.id), + ).toEqual(["item-known"]); + + // A known turn item type with a broken payload is a real defect, not a newer item. + const broken = item("item-broken", "system_notice", {}); + expect(() => decodeWireItems([turnItemEvent(4, broken)])).toThrow(); + expect(() => + decodeWireItems([ + { + kind: "snapshot", + snapshotSequence: 1, + projection: { ...projection, turnItems: [broken] }, + }, + ]), + ).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 675faa4b7f4f..a650970ba840 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -3,6 +3,7 @@ import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; import * as SchemaAST from "effect/SchemaAST"; import * as SchemaGetter from "effect/SchemaGetter"; +import * as SchemaTransformation from "effect/SchemaTransformation"; import { CheckpointId, @@ -1494,6 +1495,65 @@ export const OrchestrationV2TurnItem = Schema.Union([ ]); export type OrchestrationV2TurnItem = typeof OrchestrationV2TurnItem.Type; +const knownTurnItemTypes: ReadonlySet = new Set( + OrchestrationV2TurnItem.members.map((member) => member.fields.type.literal), +); + +const isUnknownTurnItem = (value: unknown): boolean => + typeof value === "object" && + value !== null && + "type" in value && + !knownTurnItemTypes.has(value.type); + +/** Matches only a turn item whose type this build does not know. */ +const UnknownTurnItem = Schema.Struct({ + type: Schema.String.check( + Schema.makeFilter( + (type: string) => + !knownTurnItemTypes.has(type) || "A known turn item type must decode in full.", + ), + ), +}); + +/** + * A turn item array that skips items whose `type` this build does not know. + * Newer servers add turn item types; older clients drop those rows instead of + * failing the whole snapshot. A known type that does not decode still fails. + * Unlike ForwardCompatibleArray, elements keep their own codec, so this holds + * under the JSON wire codec too. Encoding is the plain array encoding. + */ +const TurnItemArray = ( + element: Element, + unknownElement: Unknown, + isUnknown: (value: unknown) => boolean, +) => + Schema.Array(Schema.Union([element, unknownElement])).pipe( + Schema.decodeTo( + Schema.Array(Schema.toType(element)), + SchemaTransformation.transform< + ReadonlyArray, + ReadonlyArray + >({ + decode: (values) => + values.filter((value) => !isUnknown(value)) as ReadonlyArray, + encode: (values) => values, + }), + ), + ); + +const turnItemArray = (element: Element) => + TurnItemArray(element, UnknownTurnItem, isUnknownTurnItem); + +const UnknownProjectedTurnItem = Schema.Struct({ item: UnknownTurnItem }); + +const projectedTurnItemArray = (element: Element) => + TurnItemArray( + element, + UnknownProjectedTurnItem, + (row) => + typeof row === "object" && row !== null && "item" in row && isUnknownTurnItem(row.item), + ); + export const OrchestrationV2ProjectedTurnItem = Schema.Struct({ position: NonNegativeInt, visibility: Schema.Literals(["local", "inherited", "synthetic"]), @@ -1681,12 +1741,12 @@ export const OrchestrationV2ThreadProjection = Schema.Struct({ runtimeRequests: Schema.Array(OrchestrationV2RuntimeRequest), messages: Schema.Array(OrchestrationV2ConversationMessage), plans: Schema.Array(OrchestrationV2PlanArtifact), - turnItems: Schema.Array(OrchestrationV2TurnItem), + turnItems: turnItemArray(OrchestrationV2TurnItem), checkpointScopes: Schema.Array(OrchestrationV2CheckpointScope), checkpoints: Schema.Array(OrchestrationV2Checkpoint), contextHandoffs: Schema.Array(OrchestrationV2ContextHandoff), contextTransfers: Schema.Array(OrchestrationV2ContextTransfer), - visibleTurnItems: Schema.Array(OrchestrationV2ProjectedTurnItem), + visibleTurnItems: projectedTurnItemArray(OrchestrationV2ProjectedTurnItem), updatedAt: Schema.DateTimeUtc, }); export type OrchestrationV2ThreadProjection = typeof OrchestrationV2ThreadProjection.Type; @@ -2249,12 +2309,12 @@ export const OrchestrationV2ThreadProjectionJson = OrchestrationV2ThreadProjecti runtimeRequests: Schema.Array(OrchestrationV2RuntimeRequestJson), messages: Schema.Array(OrchestrationV2ConversationMessageJson), plans: Schema.Array(OrchestrationV2PlanArtifact), - turnItems: Schema.Array(OrchestrationV2TurnItemJson), + turnItems: turnItemArray(OrchestrationV2TurnItemJson), checkpointScopes: Schema.Array(OrchestrationV2CheckpointScopeJson), checkpoints: Schema.Array(OrchestrationV2CheckpointJson), contextHandoffs: Schema.Array(OrchestrationV2ContextHandoffJson), contextTransfers: Schema.Array(OrchestrationV2ContextTransferJson), - visibleTurnItems: Schema.Array(OrchestrationV2ProjectedTurnItemJson), + visibleTurnItems: projectedTurnItemArray(OrchestrationV2ProjectedTurnItemJson), updatedAt: Schema.DateTimeUtcFromString, }), ); @@ -3129,7 +3189,7 @@ export type OrchestrationV2ThreadBoundedSnapshot = typeof OrchestrationV2ThreadB /** Older timeline page for progressive history. Rows are chronological. */ export const OrchestrationV2ThreadHistoryPage = Schema.Struct({ snapshotSequence: NonNegativeInt, - items: Schema.Array(OrchestrationV2ProjectedTurnItem), + items: projectedTurnItemArray(OrchestrationV2ProjectedTurnItem), nextCursor: Schema.NullOr(TrimmedNonEmptyString), hasMoreHistory: Schema.Boolean, }); @@ -3143,22 +3203,26 @@ const knownDomainEventTypes: ReadonlySet = new Set( ); /** - * 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. + * A thread event whose type this build does not know, or a turn-item.updated + * carrying a turn item type it does not know. Newer servers add both; 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.", - ), + type: Schema.String, + payload: Schema.optional(Schema.Unknown), + }).check( + Schema.makeFilter( + (event) => + !knownDomainEventTypes.has(event.type) || + (event.type === "turn-item.updated" && isUnknownTurnItem(event.payload)) || + "A known event type must decode in full.", ), - }), + ), }).pipe( Schema.decodeTo( Schema.Struct({ From 5349dc9f71b04a3a8bf557a4da5e8665ef70bb6f Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 08:48:29 -0700 Subject: [PATCH 2/3] feat(contracts): shared forward-compatible unions, and arrays that keep transformations ForwardCompatibleArray checked each element with its own schema but passed the raw value through, so an element with a transformation (a DateTimeUtc, a decoding default) was dropped or threw. It now decodes each element through its schema and drops only the ones that fail. ForwardCompatibleUnion adds a catch-all for a tagged union's members a newer server may add: an unknown tag decodes to UnknownUnionMember for the caller to handle, while a known member that does not decode still fails. ForwardCompatibleUnionArray drops unknown members from an array. Turn items now use these shared helpers instead of a local copy. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/contracts/src/baseSchemas.test.ts | 95 +++++++++++ packages/contracts/src/baseSchemas.ts | 178 +++++++++++++++++++-- packages/contracts/src/orchestrationV2.ts | 101 ++++++------ 3 files changed, 305 insertions(+), 69 deletions(-) create mode 100644 packages/contracts/src/baseSchemas.test.ts diff --git a/packages/contracts/src/baseSchemas.test.ts b/packages/contracts/src/baseSchemas.test.ts new file mode 100644 index 000000000000..cd4842e0df83 --- /dev/null +++ b/packages/contracts/src/baseSchemas.test.ts @@ -0,0 +1,95 @@ +import { describe, expect, it } from "@effect/vitest"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Schema from "effect/Schema"; + +import { + ForwardCompatibleArray, + ForwardCompatibleUnion, + ForwardCompatibleUnionArray, + isUnknownUnionMember, +} from "./baseSchemas.ts"; + +const at = "2026-10-05T00:00:00.000Z"; +const Shape = Schema.Union([ + Schema.Struct({ kind: Schema.Literal("circle"), radius: Schema.Number, at: Schema.DateTimeUtc }), + Schema.Struct({ kind: Schema.Literal("square"), side: Schema.Number }), +]); +const members = Shape.members; +/** How clients decode: the JSON wire codec over the runtime schema. */ +const fromWire = (schema: S) => + Schema.decodeUnknownSync(Schema.toCodecJson(schema) as never) as (value: unknown) => S["Type"]; + +describe("ForwardCompatibleArray", () => { + it("decodes elements through their own codec and drops ones it cannot read", () => { + const decoded = fromWire(ForwardCompatibleArray(Shape))([ + { kind: "circle", radius: 1, at }, + { kind: "triangle" }, + { kind: "square", side: "wide" }, + ]); + expect(decoded).toHaveLength(1); + expect(DateTime.isDateTime((decoded[0] as { at: unknown }).at)).toBe(true); + }); + + it("applies decoding defaults instead of failing", () => { + const WithDefault = Schema.Struct({ + name: Schema.String, + count: Schema.optionalKey(Schema.Number).pipe( + Schema.withDecodingDefaultKey(Effect.succeed(7)), + ), + }); + expect(Schema.decodeUnknownSync(ForwardCompatibleArray(WithDefault))([{ name: "a" }])).toEqual([ + { name: "a", count: 7 }, + ]); + }); +}); + +describe("ForwardCompatibleUnion", () => { + const decode = fromWire(ForwardCompatibleUnion(members, "kind")); + + it("decodes a member from a newer server as unknown", () => { + const decoded = decode({ kind: "triangle", corners: 3 }); + expect(isUnknownUnionMember(decoded)).toBe(true); + expect(decoded).toEqual({ _unknown: true, tag: "kind", value: "triangle" }); + }); + + it("decodes known members in full, transformations included", () => { + const decoded = decode({ kind: "circle", radius: 1, at }); + expect(isUnknownUnionMember(decoded)).toBe(false); + expect(DateTime.isDateTime((decoded as { at: unknown }).at)).toBe(true); + }); + + it("still fails a known member whose payload is broken", () => { + expect(() => decode({ kind: "square", side: "wide" })).toThrow(); + }); + + it("refuses to encode an unknown member", () => { + const encode = Schema.encodeSync(ForwardCompatibleUnion(members, "kind") as never); + expect(() => encode({ _unknown: true, tag: "kind", value: "triangle" } as never)).toThrow(); + }); +}); + +describe("ForwardCompatibleUnionArray", () => { + const decode = fromWire(ForwardCompatibleUnionArray(members, "kind")); + + it("drops members a newer server added and keeps the rest", () => { + const decoded = decode([ + { kind: "triangle" }, + { kind: "circle", radius: 1, at }, + { kind: "square", side: 2 }, + ]); + expect(decoded.map((shape) => shape.kind)).toEqual(["circle", "square"]); + }); + + it("fails when a known member is broken, so real bugs stay visible", () => { + expect(() => decode([{ kind: "square", side: "wide" }])).toThrow(); + }); + + it("encodes as a plain array", () => { + const encode = Schema.encodeSync( + Schema.toCodecJson(ForwardCompatibleUnionArray(members, "kind")), + ); + const square = { kind: "square" as const, side: 2 }; + expect(encode([square])).toEqual([square]); + }); +}); diff --git a/packages/contracts/src/baseSchemas.ts b/packages/contracts/src/baseSchemas.ts index c81b81426b6a..70a01e5e5b8f 100644 --- a/packages/contracts/src/baseSchemas.ts +++ b/packages/contracts/src/baseSchemas.ts @@ -1,6 +1,7 @@ import * as Effect from "effect/Effect"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; +import * as SchemaGetter from "effect/SchemaGetter"; import * as SchemaTransformation from "effect/SchemaTransformation"; export const TrimmedString = Schema.String.pipe( @@ -35,14 +36,6 @@ export type DpopFailureReason = typeof DpopFailureReason.Type; export const IsoDateTime = Schema.String; export type IsoDateTime = typeof IsoDateTime.Type; -/** - * Wire codec for server→client arrays whose element unions grow over time - * (new literal members, new struct variants). Decoding drops elements the - * current build cannot decode instead of failing the whole payload — a client - * has to keep decoding configs sent by servers newer than itself, and - * rejecting the payload would take down the connection over data the client - * couldn't act on anyway. Encoding is the plain array encoding. - */ /** * Same idea for one optional value whose literal set grows over time: a * member this build does not know decodes as absent rather than failing the @@ -107,20 +100,175 @@ export const OmittedWhenNull = (value: Value) => { ); }; -export const ForwardCompatibleArray = (element: Element) => { - const decodeElement = Schema.decodeUnknownOption(element as never); - return Schema.Array(Schema.Unknown).pipe( +/** + * Wire codec for a server→client value whose shape grows over time: a union + * a newer server may add members to, or an array of such values. Clients keep + * decoding payloads from servers newer than themselves instead of failing the + * connection over data they could not act on anyway. + * + * Decoding runs each value through its own schema, so transformations (dates, + * trimming, decoding defaults) apply as usual. Only values that schema + * rejects are dropped; encoding is the plain encoding. + * + * For a tagged union, prefer {@link ForwardCompatibleUnion}: it drops only + * values whose tag this build does not know, so a known member with a broken + * payload still fails loudly. + */ +export const ForwardCompatibleArray = (element: Element) => + Schema.Array( + Schema.UndefinedOr(element).pipe( + // An element this build cannot read becomes a hole, filtered out below. + Schema.catchDecoding(() => Effect.succeedSome(undefined)), + ), + ).pipe( + Schema.decodeTo( + Schema.Array(Schema.toType(element)), + SchemaTransformation.transform< + ReadonlyArray, + ReadonlyArray + >({ + decode: (values) => values.filter((value) => value !== undefined), + encode: (values) => values, + }), + ), + ) as unknown as ForwardCompatibleArray; +export type ForwardCompatibleArray = Schema.Codec< + ReadonlyArray, + ReadonlyArray, + Element["DecodingServices"], + Element["EncodingServices"] +>; + +/** A member of a {@link ForwardCompatibleUnion} whose tag this build does not know. */ +export interface UnknownUnionMember { + readonly _unknown: true; + readonly tag: Tag; + readonly value: string; +} + +/** + * Wire codec for a server→client tagged union a newer server may add members + * to. A value whose `tag` field holds a value this build does not know decodes + * to {@link UnknownUnionMember} instead of failing; a known member decodes + * through its own schema, so a broken known payload still fails. Callers + * decide what an unknown member means: skip it, or show a fallback. + * + * Unknown members are decode-only; encoding one fails. + */ +export const ForwardCompatibleUnion = < + const Members extends ReadonlyArray, + const Tag extends string, +>( + members: Members, + tag: Tag, +) => { + const known = knownTags(members, tag); + const unknownMember = Schema.Struct({ + [tag]: Schema.String.check( + Schema.makeFilter( + (value: string) => !known.has(value) || `A known ${tag} must decode in full.`, + ), + ), + } as Record).pipe( Schema.decodeTo( - Schema.Array(element), - SchemaTransformation.transform, ReadonlyArray>({ + Schema.Struct({ + _unknown: Schema.Literal(true), + tag: Schema.Literal(tag), + value: Schema.String, + }), + { + decode: SchemaGetter.transform((raw: Record) => ({ + _unknown: true as const, + tag, + value: raw[tag]!, + })), + encode: SchemaGetter.forbidden(() => `Unknown ${tag} values are never sent.`), + }, + ), + ); + return Schema.Union([...members, unknownMember]) as unknown as ForwardCompatibleUnion< + Members, + Tag + >; +}; +export type ForwardCompatibleUnion< + Members extends ReadonlyArray, + Tag extends string, +> = Schema.Codec< + Members[number]["Type"] | UnknownUnionMember, + Members[number]["Encoded"], + Members[number]["DecodingServices"], + Members[number]["EncodingServices"] +>; + +/** + * Whether a raw, undecoded value carries a `tag` this build does not know + * among `members`. For checks outside a {@link ForwardCompatibleUnion}, such + * as an envelope that must skip a payload of an unknown kind. + */ +export const hasUnknownUnionTag = ( + members: ReadonlyArray, + tag: string, +): ((value: unknown) => boolean) => { + const known = knownTags(members, tag); + return (value) => + typeof value === "object" && + value !== null && + tag in value && + !known.has((value as Record)[tag] as string); +}; + +/** Whether a decoded {@link ForwardCompatibleUnion} value is a member this build does not know. */ +export const isUnknownUnionMember = ( + value: A, +): value is Extract> => + typeof value === "object" && value !== null && "_unknown" in value && value._unknown === true; + +/** + * An array of a growing tagged union that drops members this build does not + * know. Known members decode through their own schema and still fail when + * broken. Encoding is the plain array encoding. + */ +export const ForwardCompatibleUnionArray = < + const Members extends ReadonlyArray, + const Tag extends string, +>( + members: Members, + tag: Tag, +) => + Schema.Array(ForwardCompatibleUnion(members, tag)).pipe( + Schema.decodeTo( + Schema.Array(Schema.toType(Schema.Union(members))), + SchemaTransformation.transform< + ReadonlyArray, + ReadonlyArray> + >({ decode: (values) => - values.filter((value) => Option.isSome(decodeElement(value))) as ReadonlyArray< - Element["Encoded"] + values.filter((value) => !isUnknownUnionMember(value)) as ReadonlyArray< + Members[number]["Type"] >, encode: (values) => values, }), ), ); + +const knownTags = ( + members: ReadonlyArray, + tag: string, +): ReadonlySet => { + const tags = new Set(); + for (const member of members) { + const field = (member.fields as Record)[tag] as + | { readonly literal?: unknown; readonly literals?: ReadonlyArray } + | undefined; + for (const literal of field?.literals ?? [field?.literal]) { + if (typeof literal !== "string") { + throw new Error(`Every ForwardCompatibleUnion member needs a string literal ${tag}.`); + } + tags.add(literal); + } + } + return tags; }; /** diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index a650970ba840..9498fdb00002 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -16,6 +16,10 @@ import { IsoDateTime, MessageId, NodeId, + ForwardCompatibleUnion, + ForwardCompatibleUnionArray, + hasUnknownUnionTag, + isUnknownUnionMember, NonNegativeInt, PlanId, PositiveInt, @@ -1495,65 +1499,36 @@ export const OrchestrationV2TurnItem = Schema.Union([ ]); export type OrchestrationV2TurnItem = typeof OrchestrationV2TurnItem.Type; -const knownTurnItemTypes: ReadonlySet = new Set( - OrchestrationV2TurnItem.members.map((member) => member.fields.type.literal), -); - -const isUnknownTurnItem = (value: unknown): boolean => - typeof value === "object" && - value !== null && - "type" in value && - !knownTurnItemTypes.has(value.type); - -/** Matches only a turn item whose type this build does not know. */ -const UnknownTurnItem = Schema.Struct({ - type: Schema.String.check( - Schema.makeFilter( - (type: string) => - !knownTurnItemTypes.has(type) || "A known turn item type must decode in full.", - ), - ), -}); - /** - * A turn item array that skips items whose `type` this build does not know. - * Newer servers add turn item types; older clients drop those rows instead of - * failing the whole snapshot. A known type that does not decode still fails. - * Unlike ForwardCompatibleArray, elements keep their own codec, so this holds - * under the JSON wire codec too. Encoding is the plain array encoding. + * Turn item types grow over time, so clients decode them forward-compatibly: + * an item whose type this build does not know is dropped from snapshots and + * history instead of failing the thread. A known type that does not decode + * still fails. Arrays of projected rows filter on the row's nested item. */ -const TurnItemArray = ( - element: Element, - unknownElement: Unknown, - isUnknown: (value: unknown) => boolean, +const isUnknownTurnItem = hasUnknownUnionTag(OrchestrationV2TurnItem.members, "type"); + +const turnItemArray = >( + union: Schema.Union, +) => ForwardCompatibleUnionArray(union.members, "type"); + +/** Projected rows whose nested item may be of a type this build does not know. */ +const projectedTurnItemArray = ( + row: Row, + rowWithUnknownItem: Item, ) => - Schema.Array(Schema.Union([element, unknownElement])).pipe( + Schema.Array(rowWithUnknownItem).pipe( Schema.decodeTo( - Schema.Array(Schema.toType(element)), - SchemaTransformation.transform< - ReadonlyArray, - ReadonlyArray - >({ - decode: (values) => - values.filter((value) => !isUnknown(value)) as ReadonlyArray, - encode: (values) => values, + Schema.Array(Schema.toType(row)), + SchemaTransformation.transform, ReadonlyArray>({ + decode: (rows) => + rows.filter( + (projected) => !isUnknownUnionMember((projected as { readonly item: unknown }).item), + ) as ReadonlyArray, + encode: (rows) => rows, }), ), ); -const turnItemArray = (element: Element) => - TurnItemArray(element, UnknownTurnItem, isUnknownTurnItem); - -const UnknownProjectedTurnItem = Schema.Struct({ item: UnknownTurnItem }); - -const projectedTurnItemArray = (element: Element) => - TurnItemArray( - element, - UnknownProjectedTurnItem, - (row) => - typeof row === "object" && row !== null && "item" in row && isUnknownTurnItem(row.item), - ); - export const OrchestrationV2ProjectedTurnItem = Schema.Struct({ position: NonNegativeInt, visibility: Schema.Literals(["local", "inherited", "synthetic"]), @@ -1746,7 +1721,13 @@ export const OrchestrationV2ThreadProjection = Schema.Struct({ checkpoints: Schema.Array(OrchestrationV2Checkpoint), contextHandoffs: Schema.Array(OrchestrationV2ContextHandoff), contextTransfers: Schema.Array(OrchestrationV2ContextTransfer), - visibleTurnItems: projectedTurnItemArray(OrchestrationV2ProjectedTurnItem), + visibleTurnItems: projectedTurnItemArray( + OrchestrationV2ProjectedTurnItem, + OrchestrationV2ProjectedTurnItem.mapFields((fields) => ({ + ...fields, + item: ForwardCompatibleUnion(OrchestrationV2TurnItem.members, "type"), + })), + ), updatedAt: Schema.DateTimeUtc, }); export type OrchestrationV2ThreadProjection = typeof OrchestrationV2ThreadProjection.Type; @@ -2314,7 +2295,13 @@ export const OrchestrationV2ThreadProjectionJson = OrchestrationV2ThreadProjecti checkpoints: Schema.Array(OrchestrationV2CheckpointJson), contextHandoffs: Schema.Array(OrchestrationV2ContextHandoffJson), contextTransfers: Schema.Array(OrchestrationV2ContextTransferJson), - visibleTurnItems: projectedTurnItemArray(OrchestrationV2ProjectedTurnItemJson), + visibleTurnItems: projectedTurnItemArray( + OrchestrationV2ProjectedTurnItemJson, + OrchestrationV2ProjectedTurnItemJson.mapFields((fields) => ({ + ...fields, + item: ForwardCompatibleUnion(OrchestrationV2TurnItemJson.members, "type"), + })), + ), updatedAt: Schema.DateTimeUtcFromString, }), ); @@ -3189,7 +3176,13 @@ export type OrchestrationV2ThreadBoundedSnapshot = typeof OrchestrationV2ThreadB /** Older timeline page for progressive history. Rows are chronological. */ export const OrchestrationV2ThreadHistoryPage = Schema.Struct({ snapshotSequence: NonNegativeInt, - items: projectedTurnItemArray(OrchestrationV2ProjectedTurnItem), + items: projectedTurnItemArray( + OrchestrationV2ProjectedTurnItem, + OrchestrationV2ProjectedTurnItem.mapFields((fields) => ({ + ...fields, + item: ForwardCompatibleUnion(OrchestrationV2TurnItem.members, "type"), + })), + ), nextCursor: Schema.NullOr(TrimmedNonEmptyString), hasMoreHistory: Schema.Boolean, }); From 302872a89d8d83b3990f4a6764092bffe04204a9 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 08:52:05 -0700 Subject: [PATCH 3/3] fix(contracts): an unknown union member can't be mistaken for a known one UnknownUnionMember was a plain object marked with _unknown: true, so a known member whose schema had the same field was dropped as unknown. It is now a class only the codec constructs, checked with instanceof. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/contracts/src/baseSchemas.test.ts | 16 ++++++++-- packages/contracts/src/baseSchemas.ts | 37 ++++++++++------------ 2 files changed, 31 insertions(+), 22 deletions(-) diff --git a/packages/contracts/src/baseSchemas.test.ts b/packages/contracts/src/baseSchemas.test.ts index cd4842e0df83..10ece5be6b4a 100644 --- a/packages/contracts/src/baseSchemas.test.ts +++ b/packages/contracts/src/baseSchemas.test.ts @@ -8,6 +8,7 @@ import { ForwardCompatibleUnion, ForwardCompatibleUnionArray, isUnknownUnionMember, + UnknownUnionMember, } from "./baseSchemas.ts"; const at = "2026-10-05T00:00:00.000Z"; @@ -50,7 +51,7 @@ describe("ForwardCompatibleUnion", () => { it("decodes a member from a newer server as unknown", () => { const decoded = decode({ kind: "triangle", corners: 3 }); expect(isUnknownUnionMember(decoded)).toBe(true); - expect(decoded).toEqual({ _unknown: true, tag: "kind", value: "triangle" }); + expect(decoded).toEqual(new UnknownUnionMember("kind", "triangle")); }); it("decodes known members in full, transformations included", () => { @@ -65,11 +66,22 @@ describe("ForwardCompatibleUnion", () => { it("refuses to encode an unknown member", () => { const encode = Schema.encodeSync(ForwardCompatibleUnion(members, "kind") as never); - expect(() => encode({ _unknown: true, tag: "kind", value: "triangle" } as never)).toThrow(); + expect(() => encode(new UnknownUnionMember("kind", "triangle") as never)).toThrow(); }); }); describe("ForwardCompatibleUnionArray", () => { + it("keeps a known member however its fields are named", () => { + const Flagged = Schema.Struct({ + kind: Schema.Literal("flag"), + _unknown: Schema.Literal(true), + tag: Schema.String, + value: Schema.String, + }); + const flag = { kind: "flag", _unknown: true, tag: "kind", value: "x" }; + expect(fromWire(ForwardCompatibleUnionArray([Flagged], "kind"))([flag])).toEqual([flag]); + }); + const decode = fromWire(ForwardCompatibleUnionArray(members, "kind")); it("drops members a newer server added and keeps the rest", () => { diff --git a/packages/contracts/src/baseSchemas.ts b/packages/contracts/src/baseSchemas.ts index 70a01e5e5b8f..19d7bf67cdfe 100644 --- a/packages/contracts/src/baseSchemas.ts +++ b/packages/contracts/src/baseSchemas.ts @@ -139,11 +139,18 @@ export type ForwardCompatibleArray = Schema.Codec< Element["EncodingServices"] >; -/** A member of a {@link ForwardCompatibleUnion} whose tag this build does not know. */ -export interface UnknownUnionMember { - readonly _unknown: true; +/** + * A member of a {@link ForwardCompatibleUnion} whose tag this build does not + * know. A class, so only the codec can make one: no decoded payload, however + * it is shaped, is ever mistaken for it. + */ +export class UnknownUnionMember { readonly tag: Tag; readonly value: string; + constructor(tag: Tag, value: string) { + this.tag = tag; + this.value = value; + } } /** @@ -170,21 +177,12 @@ export const ForwardCompatibleUnion = < ), ), } as Record).pipe( - Schema.decodeTo( - Schema.Struct({ - _unknown: Schema.Literal(true), - tag: Schema.Literal(tag), - value: Schema.String, - }), - { - decode: SchemaGetter.transform((raw: Record) => ({ - _unknown: true as const, - tag, - value: raw[tag]!, - })), - encode: SchemaGetter.forbidden(() => `Unknown ${tag} values are never sent.`), - }, - ), + Schema.decodeTo(Schema.instanceOf(UnknownUnionMember), { + decode: SchemaGetter.transform( + (raw: Record) => new UnknownUnionMember(tag, raw[tag]!), + ), + encode: SchemaGetter.forbidden(() => `Unknown ${tag} values are never sent.`), + }), ); return Schema.Union([...members, unknownMember]) as unknown as ForwardCompatibleUnion< Members, @@ -221,8 +219,7 @@ export const hasUnknownUnionTag = ( /** Whether a decoded {@link ForwardCompatibleUnion} value is a member this build does not know. */ export const isUnknownUnionMember = ( value: A, -): value is Extract> => - typeof value === "object" && value !== null && "_unknown" in value && value._unknown === true; +): value is Extract> => value instanceof UnknownUnionMember; /** * An array of a growing tagged union that drops members this build does not