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
108 changes: 107 additions & 1 deletion packages/client-runtime/src/rpc/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -319,7 +319,11 @@ describe("environment RPC", () => {
} as unknown as WsRpcProtocolClient;
const { activeSession, retryCount, supervisor } = yield* makeHarness();

const subscriptionFiber = yield* subscribe(WS_METHODS.subscribeTerminalEvents, {}).pipe(
const subscriptionFiber = yield* subscribe(
WS_METHODS.subscribeTerminalEvents,
{},
{ resubscribeOnEndAfter: "250 millis" },
).pipe(
Stream.runDrain,
Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
Effect.forkChild,
Expand All @@ -328,6 +332,8 @@ describe("environment RPC", () => {
for (let attempt = 0; attempt < 100 && subscriptions.length < 1; attempt += 1) {
yield* Effect.yieldNow;
}
yield* TestClock.adjust("1 second");
expect(subscriptions).toEqual(["first"]);
yield* SubscriptionRef.set(activeSession, Option.none());
yield* SubscriptionRef.set(activeSession, Option.some(session(secondClient)));

Expand All @@ -341,6 +347,105 @@ describe("environment RPC", () => {
}),
);

it.effect.each(["completed", "interrupted"] as const)(
"reopens a %s automation stream after a delay without replacing its session",
(ending) =>
Effect.gen(function* () {
const requests = yield* Queue.unbounded<never, Cause.Done>();
const subscribed = yield* Deferred.make<void>();
const replacementReceived = yield* Deferred.make<void>();
let subscriptions = 0;
const client = {
[WS_METHODS.previewAutomationConnect]: () => {
subscriptions += 1;
return subscriptions === 1
? Stream.fromEffect(Deferred.succeed(subscribed, undefined)).pipe(
Stream.drain,
Stream.concat(Stream.fromQueue(requests)),
)
: Stream.succeed({ type: "connected", connectionId: "replacement" }).pipe(
Stream.concat(Stream.never),
);
},
} as unknown as WsRpcProtocolClient;
const { activeSession, retryCount, supervisor } = yield* makeHarness();
yield* SubscriptionRef.set(activeSession, Option.some(session(client)));
const fiber = yield* subscribe(
WS_METHODS.previewAutomationConnect,
{ clientId: "desktop", environmentId: TARGET.environmentId },
{ resubscribeOnEndAfter: "250 millis" },
).pipe(
Stream.runForEach((event) => {
expect(event).toEqual({ type: "connected", connectionId: "replacement" });
return Deferred.succeed(replacementReceived, undefined);
}),
Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
Effect.forkChild,
);
yield* Deferred.await(subscribed);
// RPC delivers a server-side interruption by failing its response queue.
yield* ending === "completed"
? Queue.end(requests)
: Queue.failCause(requests, Cause.interrupt());
yield* TestClock.adjust("249 millis");
expect(subscriptions).toBe(1);
yield* TestClock.adjust("1 millis");
expect(subscriptions).toBe(2);
yield* Deferred.await(replacementReceived);
yield* Fiber.interrupt(fiber);
expect(yield* Ref.get(retryCount)).toBe(0);
}),
);

it.effect.each(["unmount", "session replacement"] as const)(
"cancels a pending automation resubscription on %s",
(cancellation) =>
Effect.gen(function* () {
const firstSubscribed = yield* Deferred.make<void>();
const secondSubscribed = yield* Deferred.make<void>();
const subscriptions: string[] = [];
const firstClient = {
[WS_METHODS.previewAutomationConnect]: () => {
subscriptions.push("first");
return Stream.fromEffect(Deferred.succeed(firstSubscribed, undefined)).pipe(
Stream.drain,
);
},
} as unknown as WsRpcProtocolClient;
const secondClient = {
[WS_METHODS.previewAutomationConnect]: () => {
subscriptions.push("second");
return Stream.fromEffect(Deferred.succeed(secondSubscribed, undefined)).pipe(
Stream.drain,
Stream.concat(Stream.never),
);
},
} as unknown as WsRpcProtocolClient;
const { activeSession, supervisor } = yield* makeHarness();
yield* SubscriptionRef.set(activeSession, Option.some(session(firstClient)));
const fiber = yield* subscribe(
WS_METHODS.previewAutomationConnect,
{ clientId: "desktop", environmentId: TARGET.environmentId },
{ resubscribeOnEndAfter: "250 millis" },
).pipe(
Stream.runDrain,
Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
Effect.forkChild,
);
yield* Deferred.await(firstSubscribed);
yield* TestClock.adjust("100 millis");
if (cancellation === "unmount") {
yield* Fiber.interrupt(fiber);
} else {
yield* SubscriptionRef.set(activeSession, Option.some(session(secondClient)));
yield* Deferred.await(secondSubscribed);
}
yield* TestClock.adjust("250 millis");
expect(subscriptions).toEqual(cancellation === "unmount" ? ["first"] : ["first", "second"]);
yield* Fiber.interrupt(fiber);
}),
);

it.effect("surfaces domain subscription failures without reconnecting", () =>
Effect.gen(function* () {
const domainError = new Error("terminal subscription rejected");
Expand Down Expand Up @@ -499,6 +604,7 @@ describe("environment RPC", () => {
expectedFailureCount += 1;
}),
retryExpectedFailureAfter: "250 millis",
resubscribeOnEndAfter: "250 millis",
},
).pipe(
Stream.runDrain,
Expand Down
16 changes: 16 additions & 0 deletions packages/client-runtime/src/rpc/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import type * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
import * as Schema from "effect/Schema";
import * as Schedule from "effect/Schedule";
import * as Stream from "effect/Stream";
import * as SubscriptionRef from "effect/SubscriptionRef";
import { RpcClientError } from "effect/unstable/rpc";
Expand Down Expand Up @@ -183,6 +184,8 @@ interface SubscriptionOptions<TTag extends EnvironmentSubscriptionRpcTag> {
) => Effect.Effect<void, never, never>;
readonly retryExpectedFailureAfter?: Duration.Input;
readonly resubscribe?: Stream.Stream<unknown, never, never>;
/** Reopen a remotely ended/interrupted stream while its session remains active. */
readonly resubscribeOnEndAfter?: Duration.Input;
}

function subscribeDynamicMapped<TTag extends EnvironmentSubscriptionRpcTag, A>(
Expand Down Expand Up @@ -238,6 +241,19 @@ function subscribeDynamicMapped<TTag extends EnvironmentSubscriptionRpcTag, A>(
);
}),
).pipe(
(stream) =>
options?.resubscribeOnEndAfter === undefined
? stream
: stream.pipe(
// Broker eviction arrives as a remote interruption. Local
// cancellation (unmount/session switch) still stops this stream.
Stream.catchCause((cause) =>
Cause.hasInterruptsOnly(cause)
? Stream.empty
: Stream.failCause(cause),
),
Stream.repeat(Schedule.spaced(options.resubscribeOnEndAfter)),
),
Stream.tapCause((cause) =>
options?.onDefect !== undefined &&
cause.reasons.some(
Expand Down
11 changes: 8 additions & 3 deletions packages/client-runtime/src/state/preview.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,14 @@
import { WS_METHODS } from "@t3tools/contracts";
import { type PreviewAutomationHost, WS_METHODS } from "@t3tools/contracts";
import { Atom } from "effect/unstable/reactivity";

import type { EnvironmentRegistry } from "../connection/registry.ts";
import { subscribe } from "../rpc/client.ts";
import {
createAtomCommandScheduler,
createEnvironmentRpcCommand,
createEnvironmentRpcQueryAtomFamily,
createEnvironmentRpcSubscriptionAtomFamily,
createEnvironmentSubscriptionAtomFamily,
} from "./runtime.ts";

export const previewAutomationHostFocusConcurrencyKey = (value: {
Expand Down Expand Up @@ -45,9 +47,12 @@ export function createPreviewEnvironmentAtoms<R, E>(
// unmounted projects stop contributing probe candidates on the server.
idleTtlMs: 0,
}),
automationRequests: createEnvironmentRpcSubscriptionAtomFamily(runtime, {
automationRequests: createEnvironmentSubscriptionAtomFamily(runtime, {
label: "environment-data:preview:automation-requests",
tag: WS_METHODS.previewAutomationConnect,
subscribe: (input: PreviewAutomationHost) =>
subscribe(WS_METHODS.previewAutomationConnect, input, {
resubscribeOnEndAfter: "250 millis",
}),
// Automation requests are commands, not cached query data. Dispose the
// stream immediately with its owner so stale requests cannot replay when
// a thread remounts and the server can clear disconnected hosts promptly.
Expand Down
Loading