Skip to content
Open
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
35 changes: 35 additions & 0 deletions packages/client-runtime/src/rpc/client.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
import {
DEFAULT_SERVER_SETTINGS,
EnvironmentId,
type RelayClientInstallProgressEvent,
type ServerConfigStreamEvent,
WS_METHODS,
} from "@t3tools/contracts";
import { describe, expect, it } from "@effect/vitest";
Expand Down Expand Up @@ -77,6 +79,39 @@ const makeHarness = Effect.fn("TestEnvironmentRpc.makeHarness")(function* () {
});

describe("environment RPC", () => {
it.effect("reuses the session config stream instead of opening a duplicate subscription", () =>
Effect.gen(function* () {
const event: ServerConfigStreamEvent = {
version: 1,
type: "settingsUpdated",
payload: { settings: DEFAULT_SERVER_SETTINGS },
};
let duplicateSubscriptions = 0;
const client = {
[WS_METHODS.subscribeServerConfig]: () => {
duplicateSubscriptions += 1;
return Stream.never;
},
} as unknown as WsRpcProtocolClient;
const { activeSession, supervisor } = yield* makeHarness();
yield* SubscriptionRef.set(
activeSession,
Option.some({
...session(client),
serverConfigEvents: Stream.succeed(event),
}),
);

const received = yield* subscribe(WS_METHODS.subscribeServerConfig, {}).pipe(
Stream.runHead,
Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor),
);

expect(received).toEqual(Option.some(event));
expect(duplicateSubscriptions).toBe(0);
}),
);

it.effect("observes unary requests until they complete", () =>
Effect.gen(function* () {
const observations: string[] = [];
Expand Down
6 changes: 5 additions & 1 deletion packages/client-runtime/src/rpc/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,7 +203,11 @@ export function subscribeDynamic<TTag extends EnvironmentSubscriptionRpcTag>(
Option.match({
onNone: () => Stream.empty,
onSome: (session) => {
const method = session.client[tag] as (
const method = (
tag === WS_METHODS.subscribeServerConfig && session.serverConfigEvents !== undefined
Comment thread
Adamulek123 marked this conversation as resolved.
? () => session.serverConfigEvents
: session.client[tag]
) as (
input: EnvironmentRpcInput<TTag>,
) => Stream.Stream<
EnvironmentRpcStreamValue<TTag>,
Expand Down
218 changes: 199 additions & 19 deletions packages/client-runtime/src/rpc/session.test.ts
Original file line number Diff line number Diff line change
@@ -1,15 +1,20 @@
import {
DEFAULT_SERVER_SETTINGS,
EnvironmentId,
ProviderDriverKind,
ProviderInstanceId,
ServerConfig,
type ServerConfig as ServerConfigType,
ServerConfigStreamEvent,
type ServerConfigStreamEvent as ServerConfigStreamEventType,
WS_METHODS,
} from "@t3tools/contracts";
import { describe, expect, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";
import * as Socket from "effect/unstable/socket/Socket";

Expand Down Expand Up @@ -139,9 +144,11 @@ const RpcRequest = Schema.TaggedStruct("Request", {
tag: Schema.String,
});
const decodeJson = Schema.decodeUnknownSync(Schema.fromJsonString(Schema.Unknown));
const decodeRpcRequest = Schema.decodeUnknownSync(RpcRequest);
const isRpcRequest = Schema.is(RpcRequest);
const isPing = Schema.is(Schema.Struct({ _tag: Schema.Literal("Ping") }));
const encodeJson = Schema.encodeUnknownSync(Schema.fromJsonString(Schema.Unknown));
const encodeServerConfig = Schema.encodeSync(ServerConfig);
const encodeServerConfigStreamEvent = Schema.encodeSync(ServerConfigStreamEvent);
const ENCODED_SERVER_CONFIG = encodeServerConfig(SERVER_CONFIG);
const LEGACY_SERVER_CONFIG = {
...ENCODED_SERVER_CONFIG,
Expand Down Expand Up @@ -183,9 +190,9 @@ const awaitRequest = Effect.fn("TestRpcSessionFactory.awaitRequest")(function* (
index = 0,
) {
for (let attempt = 0; attempt < 100; attempt += 1) {
const request = socket.sent[index];
const request = socket.sent.map((message) => decodeJson(message)).filter(isRpcRequest)[index];
if (request) {
return decodeRpcRequest(decodeJson(request));
return request;
}
yield* Effect.yieldNow;
}
Expand All @@ -199,17 +206,14 @@ const completeInitialConfig = Effect.fn("TestRpcSessionFactory.completeInitialCo
const request = yield* awaitRequest(socket);
expect(request).toMatchObject({
_tag: "Request",
tag: WS_METHODS.serverGetConfig,
tag: WS_METHODS.subscribeServerConfig,
payload: {},
});
socket.serverMessage(
encodeJson({
_tag: "Exit",
_tag: "Chunk",
requestId: request.id,
exit: {
_tag: "Success",
value: config,
},
values: [{ version: 1, type: "snapshot", config }],
}),
);
});
Expand All @@ -229,7 +233,9 @@ describe("RpcSessionFactory", () => {

const config = yield* session.initialConfig;
expect(config).toEqual(SERVER_CONFIG);
expect(socket.sent).toHaveLength(1);
expect(socket.sent.map((message) => decodeJson(message)).filter(isRpcRequest)).toHaveLength(
1,
);

const probeFiber = yield* Effect.forkChild(session.probe);
const probeRequest = yield* awaitRequest(socket, 1);
Expand All @@ -250,19 +256,26 @@ describe("RpcSessionFactory", () => {
);
yield* Fiber.join(probeFiber);

expect(socket.sent.map((request) => decodeRpcRequest(decodeJson(request)).tag)).toEqual([
WS_METHODS.serverGetConfig,
WS_METHODS.serverProbe,
]);
expect(
socket.sent
.map((message) => decodeJson(message))
.filter(isRpcRequest)
.map((request) => request.tag),
).toEqual([WS_METHODS.subscribeServerConfig, WS_METHODS.serverProbe]);

socket.close(1012, "service restart");
const error = yield* Effect.flip(session.closed);
const configStreamError = yield* session.serverConfigEvents!.pipe(
Stream.runDrain,
Effect.flip,
);

expect(error).toBeInstanceOf(ConnectionTransientError);
expect(error).toMatchObject({
reason: "transport",
message: "Test environment disconnected.",
});
expect(configStreamError).toMatchObject({ _tag: "RpcClientError" });
yield* Effect.yieldNow;
expect(sockets).toHaveLength(1);
}),
Expand All @@ -287,6 +300,171 @@ describe("RpcSessionFactory", () => {
}),
);

it.effect("replays current config and broadcasts updates to every subscriber", () =>
Effect.scoped(
Effect.gen(function* () {
const { factory, sockets } = yield* makeFactory();
const session = yield* factory.connect(PREPARED);
const readyFiber = yield* Effect.forkChild(session.ready);
const socket = yield* awaitSocket(sockets);
socket.open();
yield* completeInitialConfig(socket);
yield* Fiber.join(readyFiber);

const collectTwo = session.serverConfigEvents!.pipe(Stream.take(2), Stream.runCollect);
const firstSubscriber = yield* Effect.forkChild(collectTwo);
const secondSubscriber = yield* Effect.forkChild(collectTwo);
yield* Effect.yieldNow;

const shortcut = {
key: "k",
metaKey: false,
ctrlKey: false,
shiftKey: false,
altKey: false,
modKey: true,
};
const request = yield* awaitRequest(socket);
socket.serverMessage(
encodeJson({
_tag: "Chunk",
requestId: request.id,
values: [
{
version: 1,
type: "keybindingsUpdated",
payload: {
keybindings: [{ command: "terminal.toggle", shortcut }],
issues: [],
},
},
],
}),
);

const firstEvents = Array.from(yield* Fiber.join(firstSubscriber));
const secondEvents = Array.from(yield* Fiber.join(secondSubscriber));
expect(firstEvents.map((event) => event.type)).toEqual(["snapshot", "keybindingsUpdated"]);
expect(secondEvents).toEqual(firstEvents);

const replay = yield* session.serverConfigEvents!.pipe(Stream.runHead);
expect(replay).toMatchObject({
_tag: "Some",
value: {
type: "snapshot",
config: { keybindings: [{ command: "terminal.toggle", shortcut }] },
},
});
}),
),
);

it.effect.each<{
readonly event: ServerConfigStreamEventType;
readonly expectedConfig: Partial<ServerConfigType>;
}>([
{
event: {
version: 1,
type: "providerStatuses",
payload: {
providers: [
{
instanceId: ProviderInstanceId.make("codex"),
driver: ProviderDriverKind.make("codex"),
enabled: true,
installed: true,
version: "1.0.0",
status: "ready",
auth: { status: "authenticated" },
checkedAt: "2026-08-27T00:00:00.000Z",
models: [],
slashCommands: [],
skills: [],
},
],
},
},
expectedConfig: {
providers: [
{
instanceId: ProviderInstanceId.make("codex"),
driver: ProviderDriverKind.make("codex"),
enabled: true,
installed: true,
version: "1.0.0",
status: "ready",
auth: { status: "authenticated" },
checkedAt: "2026-08-27T00:00:00.000Z",
models: [],
slashCommands: [],
skills: [],
},
],
},
},
{
event: {
version: 1,
type: "settingsUpdated",
payload: {
settings: {
...DEFAULT_SERVER_SETTINGS,
newWorktreesStartFromOrigin: !DEFAULT_SERVER_SETTINGS.newWorktreesStartFromOrigin,
},
},
},
expectedConfig: {
settings: {
...DEFAULT_SERVER_SETTINGS,
newWorktreesStartFromOrigin: !DEFAULT_SERVER_SETTINGS.newWorktreesStartFromOrigin,
},
},
},
])(
"preserves $event.type events and includes them in replay snapshots",
({ event, expectedConfig }) =>
Effect.scoped(
Effect.gen(function* () {
const { factory, sockets } = yield* makeFactory();
const session = yield* factory.connect(PREPARED);
const readyFiber = yield* Effect.forkChild(session.ready);
const socket = yield* awaitSocket(sockets);
socket.open();
yield* completeInitialConfig(socket);
yield* Fiber.join(readyFiber);

const subscriber = yield* session.serverConfigEvents!.pipe(
Stream.take(2),
Stream.runCollect,
Effect.forkChild,
);
yield* Effect.yieldNow;

const request = yield* awaitRequest(socket);
socket.serverMessage(
encodeJson({
_tag: "Chunk",
requestId: request.id,
values: [encodeServerConfigStreamEvent(event)],
}),
);

const events = Array.from(yield* Fiber.join(subscriber));
expect(events[1]).toEqual(event);

const replay = yield* session.serverConfigEvents!.pipe(Stream.runHead);
expect(replay).toMatchObject({
_tag: "Some",
value: {
type: "snapshot",
config: expectedConfig,
},
});
}),
),
);

it.effect("tolerates two missed pong windows before closing the session", () =>
Effect.gen(function* () {
const { factory, sockets } = yield* makeFactory();
Expand All @@ -301,7 +479,7 @@ describe("RpcSessionFactory", () => {

yield* TestClock.adjust("15 seconds");
expect(closedFiber.pollUnsafe()).toBeUndefined();
expect(socket.sent.slice(1).map((request) => decodeJson(request))).toEqual([
expect(socket.sent.map((message) => decodeJson(message)).filter(isPing)).toEqual([
{ _tag: "Ping" },
{ _tag: "Ping" },
{ _tag: "Ping" },
Expand Down Expand Up @@ -379,10 +557,12 @@ describe("RpcSessionFactory", () => {
);
yield* Fiber.join(probeFiber);

expect(socket.sent.map((request) => decodeRpcRequest(decodeJson(request)).tag)).toEqual([
WS_METHODS.serverGetConfig,
WS_METHODS.serverGetConfig,
]);
expect(
socket.sent
.map((message) => decodeJson(message))
.filter(isRpcRequest)
.map((request) => request.tag),
).toEqual([WS_METHODS.subscribeServerConfig, WS_METHODS.serverGetConfig]);
}),
),
);
Expand Down
Loading
Loading