Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 63 additions & 0 deletions packages/client-runtime/src/state/threads-sync.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import {
EnvironmentId,
EventId,
MessageId,
ORCHESTRATION_V2_WS_METHODS,
ThreadId,
TurnItemId,
Expand Down Expand Up @@ -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 });
Expand Down
41 changes: 32 additions & 9 deletions packages/client-runtime/src/state/threads.ts
Original file line number Diff line number Diff line change
Expand Up @@ -435,17 +435,37 @@ export const makeEnvironmentThreadState = Effect.fn("EnvironmentThreadState.make
});

type EventItem = Extract<OrchestrationV2ThreadStreamItem, { kind: "event" }>;
type SequencedItem = Extract<
OrchestrationV2ThreadStreamItem,
{ kind: "event" | "unknown-event" }
>;
const applyEventsLocked = Effect.fn("EnvironmentThreadState.applyEventsLocked")(function* (
items: ReadonlyArray<EventItem>,
items: ReadonlyArray<SequencedItem>,
) {
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,
// Bounded: the type comes from a newer server and is not validated here.
eventType: item.eventType.slice(0, 64),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Truncating eventType bounds its size but still puts unchecked wire text (potentially payload or secrets) into logs. Could you record only a safe diagnostic instead?

Suggested change
eventType: item.eventType.slice(0, 64),
eventTypeLength: item.eventType.length,

Posted via Macroscope — Effect Service Conventions

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
Expand Down Expand Up @@ -593,9 +613,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;
}
Expand Down
47 changes: 47 additions & 0 deletions packages/contracts/src/orchestrationV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import {
OrchestrationV2ProviderCapabilities,
OrchestrationV2ProviderThread,
OrchestrationV2ProviderThreadJson,
OrchestrationV2RpcSchemas,
OrchestrationV2ShellSnapshot,
OrchestrationV2SubscribeThreadInput,
OrchestrationV2Subagent,
Expand Down Expand Up @@ -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({
Expand Down
45 changes: 45 additions & 0 deletions packages/contracts/src/orchestrationV2.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -2868,6 +2869,48 @@ export const OrchestrationV2ThreadHistoryPage = Schema.Struct({
});
export type OrchestrationV2ThreadHistoryPage = typeof OrchestrationV2ThreadHistoryPage.Type;

const knownDomainEventTypes: ReadonlySet<string> = 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"),
Expand All @@ -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;

Expand Down
Loading