Skip to content
Closed
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
3 changes: 2 additions & 1 deletion apps/server/src/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9251,7 +9251,8 @@ it.layer(NodeServices.layer)("server router seam", (it) => {
);
if (oversized) {
assertTrue(threadResult._tag === "Failure");
assert.equal(threadResult.failure._tag, "OrchestrationGetSnapshotError");
assertTrue(threadResult.failure._tag === "OrchestrationThreadNotFoundError");
assert.equal(threadResult.failure.threadId, defaultThreadId);
assert.equal(
threadResult.failure.message,
`Thread ${defaultThreadId} was not found`,
Expand Down
6 changes: 2 additions & 4 deletions apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ import {
type OrchestrationShellStreamItem,
OrchestrationGetFullThreadDiffError,
OrchestrationGetSnapshotError,
OrchestrationThreadNotFoundError,
OrchestrationSearchThreadsError,
OrchestrationGetTurnDiffError,
ORCHESTRATION_WS_METHODS,
Expand Down Expand Up @@ -1736,10 +1737,7 @@ const makeWsRpcLayer = (
if (replayOnMissingSnapshot !== undefined) {
return replayOnMissingSnapshot;
}
return yield* new OrchestrationGetSnapshotError({
message: `Thread ${input.threadId} was not found`,
cause: input.threadId,
});
return yield* new OrchestrationThreadNotFoundError({ threadId: input.threadId });
}

const afterSnapshot =
Expand Down
60 changes: 57 additions & 3 deletions packages/client-runtime/src/errors/orchestration.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,55 @@
import { OrchestrationDispatchCommandError } from "@t3tools/contracts";
import { describe, expect, it } from "vite-plus/test";
import {
OrchestrationDispatchCommandError,
OrchestrationGetSnapshotError,
OrchestrationThreadNotFoundError,
ThreadId,
} from "@t3tools/contracts";
import { describe, expect, it } from "@effect/vitest";

import { wasBootstrapThreadDeleted } from "./orchestration.ts";
import { wasBootstrapThreadDeleted, isOrchestrationThreadNotFoundError } from "./orchestration.ts";

describe("isOrchestrationThreadNotFoundError", () => {
it("matches the typed not-found error", () => {
expect(
isOrchestrationThreadNotFoundError(
new OrchestrationThreadNotFoundError({ threadId: ThreadId.make("thread-1") }),
),
).toBe(true);
});

it("rejects a generic snapshot error with a matching message", () => {
expect(
isOrchestrationThreadNotFoundError(
new OrchestrationGetSnapshotError({
message: "Thread thread-1 was not found",
cause: "thread-1",
}),
),
).toBe(false);
});

it("rejects plain errors with a matching message", () => {
expect(isOrchestrationThreadNotFoundError(new Error("Thread thread-1 was not found"))).toBe(
false,
);
});

it("rejects other snapshot errors", () => {
expect(
isOrchestrationThreadNotFoundError(
new OrchestrationGetSnapshotError({
message: "Failed to load thread thread-1",
cause: "thread-1",
}),
),
).toBe(false);
});

it("rejects unrelated errors", () => {
expect(isOrchestrationThreadNotFoundError(new Error("boom"))).toBe(false);
expect(isOrchestrationThreadNotFoundError(null)).toBe(false);
});
});

describe("wasBootstrapThreadDeleted", () => {
it("accepts only a confirmed deleted bootstrap thread", () => {
Expand All @@ -13,11 +61,17 @@ describe("wasBootstrapThreadDeleted", () => {
}),
),
).toBe(true);
});

it("rejects a missing disposition", () => {
expect(
wasBootstrapThreadDeleted(
new OrchestrationDispatchCommandError({ message: "Failed to create worktree." }),
),
).toBe(false);
});

it("rejects unrelated errors", () => {
expect(wasBootstrapThreadDeleted(new Error("connection lost"))).toBe(false);
});
});
7 changes: 6 additions & 1 deletion packages/client-runtime/src/errors/orchestration.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
import { OrchestrationDispatchCommandError } from "@t3tools/contracts";
import {
OrchestrationDispatchCommandError,
OrchestrationThreadNotFoundError,
} from "@t3tools/contracts";
import * as Schema from "effect/Schema";

const isOrchestrationDispatchCommandError = Schema.is(OrchestrationDispatchCommandError);
Expand All @@ -8,3 +11,5 @@ export function wasBootstrapThreadDeleted(error: unknown): boolean {
isOrchestrationDispatchCommandError(error) && error.bootstrapThreadDisposition === "deleted"
);
}

export const isOrchestrationThreadNotFoundError = Schema.is(OrchestrationThreadNotFoundError);
99 changes: 99 additions & 0 deletions packages/client-runtime/src/rpc/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -465,6 +465,105 @@ describe("environment RPC", () => {
}),
);

it.effect("ends the subscription permanently when the failure is terminal", () =>
Effect.gen(function* () {
const notFound = new Error("thread was not found");
const subscriptionCount = yield* Ref.make(0);
const handled = yield* Ref.make<Cause.Cause<unknown> | null>(null);
const client = {
[WS_METHODS.subscribeTerminalEvents]: () =>
Stream.unwrap(
Ref.updateAndGet(subscriptionCount, (count) => count + 1).pipe(
Effect.map(() => Stream.fail(notFound)),
),
),
} as unknown as WsRpcProtocolClient;
const { activeSession, supervisor } = yield* makeHarness();

yield* SubscriptionRef.set(activeSession, Option.some(session(client)));
const subscriptionFiber = yield* subscribe(
WS_METHODS.subscribeTerminalEvents,
{},
{
terminalFailure: {
matches: (error) => error === notFound,
handle: (cause) => Ref.set(handled, cause),
},
},
).pipe(
Stream.runDrain,
Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
Effect.forkChild,
);
for (let attempt = 0; attempt < 100 && (yield* Ref.get(handled)) === null; attempt += 1) {
yield* Effect.yieldNow;
}

expect(yield* Ref.get(subscriptionCount)).toBe(1);
expect(yield* Ref.get(handled)).not.toBeNull();

// Far past any retry delay: a terminal failure must not re-attempt.
yield* TestClock.adjust("30 seconds");
yield* Effect.yieldNow;
yield* Fiber.interrupt(subscriptionFiber);

expect(yield* Ref.get(subscriptionCount)).toBe(1);
}),
);

it.effect("does not treat failures rejected by the terminal classifier as terminal", () =>
Effect.gen(function* () {
const transient = new Error("transient snapshot failure");
const subscriptionCount = yield* Ref.make(0);
const retryReady = yield* Deferred.make<void>();
const retryStarted = yield* Deferred.make<void>();
const client = {
[WS_METHODS.subscribeTerminalEvents]: () =>
Stream.unwrap(
Ref.getAndUpdate(subscriptionCount, (count) => count + 1).pipe(
Effect.map((count) =>
count === 0
? Stream.fail(transient)
: Stream.fromEffect(Deferred.succeed(retryStarted, undefined)).pipe(
Stream.drain,
Stream.concat(Stream.never),
),
),
),
),
} as unknown as WsRpcProtocolClient;
const { activeSession, supervisor } = yield* makeHarness();

yield* SubscriptionRef.set(activeSession, Option.some(session(client)));
const subscriptionFiber = yield* subscribe(
WS_METHODS.subscribeTerminalEvents,
{},
{
onExpectedFailure: () => Deferred.succeed(retryReady, undefined).pipe(Effect.asVoid),
retryExpectedFailureAfter: "100 millis",
terminalFailure: {
matches: () => false,
handle: () => Effect.void,
},
},
).pipe(
Stream.runDrain,
Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
Effect.forkChild,
);

// The handler runs immediately before the retry delay; let that fiber schedule its sleep.
yield* Deferred.await(retryReady);
yield* Effect.yieldNow;
yield* TestClock.adjust("100 millis");
yield* Deferred.await(retryStarted);
yield* Fiber.interrupt(subscriptionFiber);

// Classified as a regular expected failure: retried once, handler untouched.
expect(yield* Ref.get(subscriptionCount)).toBe(2);
}),
);

it.effect.each(["input", "stream"] as const)(
"does not classify %s subscription defects as expected failures",
(where) =>
Expand Down
32 changes: 32 additions & 0 deletions packages/client-runtime/src/rpc/client.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { ORCHESTRATION_WS_METHODS, WS_METHODS } from "@t3tools/contracts";
import * as Cause from "effect/Cause";
import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import type * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
Expand Down Expand Up @@ -180,6 +181,15 @@ interface SubscriptionOptions<TTag extends EnvironmentSubscriptionRpcTag> {
) => Effect.Effect<void, never, never>;
readonly retryExpectedFailureAfter?: Duration.Input;
readonly resubscribe?: Stream.Stream<unknown, never, never>;
/**
* Classifies an all-Fail cause as terminal: the attempt ends for this
* session with no retry, after `handle` runs. Checked after transport
* failures and before expected-failure retry.
*/
readonly terminalFailure?: {
readonly matches: (error: EnvironmentRpcStreamFailure<TTag>) => boolean;
readonly handle: (cause: Cause.Cause<EnvironmentRpcStreamFailure<TTag>>) => Effect.Effect<void>;
};
}

function subscribeDynamicMapped<TTag extends EnvironmentSubscriptionRpcTag, A>(
Expand All @@ -195,6 +205,12 @@ function subscribeDynamicMapped<TTag extends EnvironmentSubscriptionRpcTag, A>(
Effect.gen(function* () {
const supervisor = yield* EnvironmentSupervisor;
const observer = yield* EnvironmentRpcSubscriptionObserver;
// Signaled before a terminalFailure handler runs. Interrupts the
// outer session stream below so a supervisor.session replacement
// landing during the handler's cache I/O cannot re-issue the
// subscription after a terminal tombstone; the handler itself still
// runs to completion as the draining inner.
const terminalHalt = yield* Deferred.make<void>();
const sessionChanges = SubscriptionRef.changes(supervisor.session);
const sessions =
options?.resubscribe === undefined
Expand All @@ -206,6 +222,7 @@ function subscribeDynamicMapped<TTag extends EnvironmentSubscriptionRpcTag, A>(
),
);
return sessions.pipe(
Stream.interruptWhen(Deferred.await(terminalHalt)),
Comment thread
cursor[bot] marked this conversation as resolved.
Stream.switchMap(
Option.match({
onNone: () => Stream.empty,
Expand Down Expand Up @@ -268,6 +285,21 @@ function subscribeDynamicMapped<TTag extends EnvironmentSubscriptionRpcTag, A>(
),
).pipe(Stream.drain);
}
const terminal = options?.terminalFailure;
if (
hasOnlyExpectedFailures &&
terminal !== undefined &&
cause.reasons.every(
(reason) => reason._tag === "Fail" && terminal.matches(reason.error),
)
) {
return Stream.fromEffect(
Deferred.succeed(terminalHalt, undefined).pipe(
Effect.asVoid,
Effect.andThen(terminal.handle(cause)),
),
).pipe(Stream.drain);
}
if (hasOnlyExpectedFailures && options?.onExpectedFailure !== undefined) {
const handled = Stream.fromEffect(options.onExpectedFailure(cause)).pipe(
Stream.drain,
Expand Down
Loading
Loading