From 74597155b51a32d6bf7b30acf786f5a29d842b81 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Thu, 27 Aug 2026 00:20:52 +0200 Subject: [PATCH 1/5] perf(client-runtime): reuse server config subscription --- .../client-runtime/src/rpc/client.test.ts | 35 +++++++++++++ packages/client-runtime/src/rpc/client.ts | 6 ++- .../client-runtime/src/rpc/session.test.ts | 42 ++++++++------- packages/client-runtime/src/rpc/session.ts | 51 ++++++++++++++++--- 4 files changed, 108 insertions(+), 26 deletions(-) diff --git a/packages/client-runtime/src/rpc/client.test.ts b/packages/client-runtime/src/rpc/client.test.ts index 507d137caccb..a33dcf9b22a4 100644 --- a/packages/client-runtime/src/rpc/client.test.ts +++ b/packages/client-runtime/src/rpc/client.test.ts @@ -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"; @@ -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[] = []; diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index bfe57a6c0dd5..e95e2298622e 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -203,7 +203,11 @@ export function subscribeDynamic( Option.match({ onNone: () => Stream.empty, onSome: (session) => { - const method = session.client[tag] as ( + const method = ( + tag === WS_METHODS.subscribeServerConfig && session.serverConfigEvents !== undefined + ? () => session.serverConfigEvents + : session.client[tag] + ) as ( input: EnvironmentRpcInput, ) => Stream.Stream< EnvironmentRpcStreamValue, diff --git a/packages/client-runtime/src/rpc/session.test.ts b/packages/client-runtime/src/rpc/session.test.ts index 0af5850bf6c7..4e934ea8b028 100644 --- a/packages/client-runtime/src/rpc/session.test.ts +++ b/packages/client-runtime/src/rpc/session.test.ts @@ -139,7 +139,8 @@ 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 ENCODED_SERVER_CONFIG = encodeServerConfig(SERVER_CONFIG); @@ -183,9 +184,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; } @@ -199,17 +200,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 }], }), ); }); @@ -229,7 +227,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); @@ -250,10 +250,12 @@ 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); @@ -301,7 +303,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" }, @@ -379,10 +381,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]); }), ), ); diff --git a/packages/client-runtime/src/rpc/session.ts b/packages/client-runtime/src/rpc/session.ts index 9625effa406f..f56c81526870 100644 --- a/packages/client-runtime/src/rpc/session.ts +++ b/packages/client-runtime/src/rpc/session.ts @@ -1,11 +1,20 @@ -import { type ServerConfig, WS_METHODS } from "@t3tools/contracts"; +import { + type EnvironmentAuthorizationError, + type KeybindingsConfigError, + type ServerConfig, + type ServerConfigStreamEvent, + type ServerSettingsError, + WS_METHODS, +} from "@t3tools/contracts"; import * as Context from "effect/Context"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Schedule from "effect/Schedule"; import type * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; import * as RpcClient from "effect/unstable/rpc/RpcClient"; +import type * as RpcClientError from "effect/unstable/rpc/RpcClientError"; import * as RpcSerialization from "effect/unstable/rpc/RpcSerialization"; import * as Socket from "effect/unstable/socket/Socket"; @@ -25,6 +34,10 @@ const SOCKET_OPEN_TIMEOUT = "15 seconds"; export interface RpcSession { readonly client: WsRpcProtocolClient; readonly initialConfig: Effect.Effect; + readonly serverConfigEvents?: Stream.Stream< + ServerConfigStreamEvent, + ServerConfigSubscriptionError + >; readonly ready: Effect.Effect; readonly probe: Effect.Effect; readonly closed: Effect.Effect; @@ -43,8 +56,15 @@ type InitialConfigError = Effect.Error< ReturnType >; type ProbeError = Effect.Error>; +type ServerConfigSubscriptionError = + | EnvironmentAuthorizationError + | KeybindingsConfigError + | ServerSettingsError + | RpcClientError.RpcClientError; -function mapSessionRpcError(error: InitialConfigError | ProbeError): ConnectionAttemptError { +function mapSessionRpcError( + error: InitialConfigError | ProbeError | ServerConfigSubscriptionError, +): ConnectionAttemptError { switch (error._tag) { case "EnvironmentAuthorizationError": return new ConnectionBlockedError({ @@ -114,11 +134,29 @@ export const make = Effect.gen(function* () { Effect.withSpan("environment.websocket.connect"), ); const client = yield* makeWsRpcProtocolClient.pipe(Effect.provide(protocolContext)); - const initialConfig = yield* Effect.cached( - client[WS_METHODS.serverGetConfig]({}).pipe( - Effect.mapError(mapSessionRpcError), - Effect.withSpan("environment.initialSync"), + const initialConfigDeferred = yield* Deferred.make(); + const serverConfigEvents = client[WS_METHODS.subscribeServerConfig]({}).pipe( + Stream.tap((event) => + event.type === "snapshot" + ? Deferred.succeed(initialConfigDeferred, event.config).pipe(Effect.asVoid) + : Effect.void, + ), + Stream.tapError((error) => + Deferred.fail(initialConfigDeferred, mapSessionRpcError(error)).pipe(Effect.asVoid), ), + Stream.ensuring( + Deferred.fail( + initialConfigDeferred, + new ConnectionTransientErrorClass({ + reason: "remote-unavailable", + detail: `${connection.label} config subscription ended before its initial snapshot.`, + }), + ).pipe(Effect.asVoid), + ), + ); + const serverConfigQueue = yield* Stream.toQueue(serverConfigEvents, { capacity: 64 }); + const initialConfig = Deferred.await(initialConfigDeferred).pipe( + Effect.withSpan("environment.initialSync"), ); const probe = initialConfig.pipe( Effect.flatMap((config) => @@ -134,6 +172,7 @@ export const make = Effect.gen(function* () { return { client, initialConfig, + serverConfigEvents: Stream.fromQueue(serverConfigQueue), ready: Deferred.await(connected).pipe( Effect.andThen(initialConfig), Effect.asVoid, From b753046eba672039a2f8d8db93af1a11c02e2159 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Thu, 27 Aug 2026 10:43:27 +0200 Subject: [PATCH 2/5] fix(client-runtime): replay config to every subscriber --- .../client-runtime/src/rpc/session.test.ts | 60 ++++++++++++ packages/client-runtime/src/rpc/session.ts | 94 +++++++++++++++++-- 2 files changed, 145 insertions(+), 9 deletions(-) diff --git a/packages/client-runtime/src/rpc/session.test.ts b/packages/client-runtime/src/rpc/session.test.ts index 4e934ea8b028..5112ecdf0fe7 100644 --- a/packages/client-runtime/src/rpc/session.test.ts +++ b/packages/client-runtime/src/rpc/session.test.ts @@ -10,6 +10,7 @@ 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"; @@ -289,6 +290,65 @@ 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("tolerates two missed pong windows before closing the session", () => Effect.gen(function* () { const { factory, sockets } = yield* makeFactory(); diff --git a/packages/client-runtime/src/rpc/session.ts b/packages/client-runtime/src/rpc/session.ts index f56c81526870..23c070966be7 100644 --- a/packages/client-runtime/src/rpc/session.ts +++ b/packages/client-runtime/src/rpc/session.ts @@ -10,6 +10,8 @@ import * as Context from "effect/Context"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as PubSub from "effect/PubSub"; +import * as Ref from "effect/Ref"; import * as Schedule from "effect/Schedule"; import type * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; @@ -62,6 +64,36 @@ type ServerConfigSubscriptionError = | ServerSettingsError | RpcClientError.RpcClientError; +interface ServerConfigReplayState { + readonly config: ServerConfig; + readonly revision: number; +} + +interface BufferedServerConfigEvent { + readonly event: ServerConfigStreamEvent; + readonly revision: number; +} + +function applyServerConfigEvent( + config: ServerConfig, + event: ServerConfigStreamEvent, +): ServerConfig { + switch (event.type) { + case "snapshot": + return event.config; + case "keybindingsUpdated": + return { + ...config, + keybindings: event.payload.keybindings, + issues: event.payload.issues, + }; + case "providerStatuses": + return { ...config, providers: event.payload.providers }; + case "settingsUpdated": + return { ...config, settings: event.payload.settings }; + } +} + function mapSessionRpcError( error: InitialConfigError | ProbeError | ServerConfigSubscriptionError, ): ConnectionAttemptError { @@ -135,16 +167,39 @@ export const make = Effect.gen(function* () { ); const client = yield* makeWsRpcProtocolClient.pipe(Effect.provide(protocolContext)); const initialConfigDeferred = yield* Deferred.make(); - const serverConfigEvents = client[WS_METHODS.subscribeServerConfig]({}).pipe( - Stream.tap((event) => - event.type === "snapshot" - ? Deferred.succeed(initialConfigDeferred, event.config).pipe(Effect.asVoid) - : Effect.void, + const serverConfigState = yield* Ref.make(undefined); + const serverConfigUpdates = yield* PubSub.bounded(64); + const serverConfigSource = client[WS_METHODS.subscribeServerConfig]({}).pipe( + Stream.runForEach((event) => + Effect.gen(function* () { + if (event.type === "snapshot") { + yield* Deferred.succeed(initialConfigDeferred, event.config); + } + const buffered = yield* Ref.modify(serverConfigState, (current) => { + let config: ServerConfig; + if (current === undefined) { + if (event.type !== "snapshot") { + return [undefined, current] as const; + } + config = event.config; + } else { + config = applyServerConfigEvent(current.config, event); + } + const next = { + config, + revision: (current?.revision ?? 0) + 1, + }; + return [{ event, revision: next.revision }, next] as const; + }); + if (buffered !== undefined) { + yield* PubSub.publish(serverConfigUpdates, buffered); + } + }), ), - Stream.tapError((error) => + Effect.tapError((error) => Deferred.fail(initialConfigDeferred, mapSessionRpcError(error)).pipe(Effect.asVoid), ), - Stream.ensuring( + Effect.ensuring( Deferred.fail( initialConfigDeferred, new ConnectionTransientErrorClass({ @@ -154,10 +209,31 @@ export const make = Effect.gen(function* () { ).pipe(Effect.asVoid), ), ); - const serverConfigQueue = yield* Stream.toQueue(serverConfigEvents, { capacity: 64 }); + yield* serverConfigSource.pipe(Effect.forkScoped); const initialConfig = Deferred.await(initialConfigDeferred).pipe( Effect.withSpan("environment.initialSync"), ); + const serverConfigEvents = Stream.unwrap( + Effect.gen(function* () { + const subscription = yield* PubSub.subscribe(serverConfigUpdates); + yield* initialConfig.pipe(Effect.option); + const snapshot = yield* Ref.get(serverConfigState); + if (snapshot === undefined) { + return Stream.empty; + } + return Stream.concat( + Stream.succeed({ + version: 1 as const, + type: "snapshot" as const, + config: snapshot.config, + }), + Stream.fromSubscription(subscription).pipe( + Stream.filter((buffered) => buffered.revision > snapshot.revision), + Stream.map((buffered) => buffered.event), + ), + ); + }), + ); const probe = initialConfig.pipe( Effect.flatMap((config) => (config.environment.capabilities.connectionProbe === true @@ -172,7 +248,7 @@ export const make = Effect.gen(function* () { return { client, initialConfig, - serverConfigEvents: Stream.fromQueue(serverConfigQueue), + serverConfigEvents, ready: Deferred.await(connected).pipe( Effect.andThen(initialConfig), Effect.asVoid, From 12efd16677d7d1c9418963d535623e751130321d Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Thu, 27 Aug 2026 10:51:54 +0200 Subject: [PATCH 3/5] fix(client-runtime): keep config fanout nonblocking --- .../client-runtime/src/rpc/session.test.ts | 7 ++- packages/client-runtime/src/rpc/session.ts | 45 ++++++++++++------- 2 files changed, 36 insertions(+), 16 deletions(-) diff --git a/packages/client-runtime/src/rpc/session.test.ts b/packages/client-runtime/src/rpc/session.test.ts index 5112ecdf0fe7..e73de59abff2 100644 --- a/packages/client-runtime/src/rpc/session.test.ts +++ b/packages/client-runtime/src/rpc/session.test.ts @@ -260,12 +260,17 @@ describe("RpcSessionFactory", () => { 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); }), @@ -334,7 +339,7 @@ describe("RpcSessionFactory", () => { 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(firstEvents.map((event) => event.type)).toEqual(["snapshot", "snapshot"]); expect(secondEvents).toEqual(firstEvents); const replay = yield* session.serverConfigEvents!.pipe(Stream.runHead); diff --git a/packages/client-runtime/src/rpc/session.ts b/packages/client-runtime/src/rpc/session.ts index 23c070966be7..d547f99a1b58 100644 --- a/packages/client-runtime/src/rpc/session.ts +++ b/packages/client-runtime/src/rpc/session.ts @@ -70,7 +70,7 @@ interface ServerConfigReplayState { } interface BufferedServerConfigEvent { - readonly event: ServerConfigStreamEvent; + readonly config: ServerConfig; readonly revision: number; } @@ -167,8 +167,9 @@ export const make = Effect.gen(function* () { ); const client = yield* makeWsRpcProtocolClient.pipe(Effect.provide(protocolContext)); const initialConfigDeferred = yield* Deferred.make(); + const serverConfigExit = yield* Deferred.make(); const serverConfigState = yield* Ref.make(undefined); - const serverConfigUpdates = yield* PubSub.bounded(64); + const serverConfigUpdates = yield* PubSub.sliding(64); const serverConfigSource = client[WS_METHODS.subscribeServerConfig]({}).pipe( Stream.runForEach((event) => Effect.gen(function* () { @@ -189,7 +190,7 @@ export const make = Effect.gen(function* () { config, revision: (current?.revision ?? 0) + 1, }; - return [{ event, revision: next.revision }, next] as const; + return [{ config: next.config, revision: next.revision }, next] as const; }); if (buffered !== undefined) { yield* PubSub.publish(serverConfigUpdates, buffered); @@ -197,16 +198,22 @@ export const make = Effect.gen(function* () { }), ), Effect.tapError((error) => - Deferred.fail(initialConfigDeferred, mapSessionRpcError(error)).pipe(Effect.asVoid), + Effect.all([ + Deferred.fail(initialConfigDeferred, mapSessionRpcError(error)), + Deferred.fail(serverConfigExit, error), + ]).pipe(Effect.asVoid), ), Effect.ensuring( - Deferred.fail( - initialConfigDeferred, - new ConnectionTransientErrorClass({ - reason: "remote-unavailable", - detail: `${connection.label} config subscription ended before its initial snapshot.`, - }), - ).pipe(Effect.asVoid), + Effect.all([ + Deferred.fail( + initialConfigDeferred, + new ConnectionTransientErrorClass({ + reason: "remote-unavailable", + detail: `${connection.label} config subscription ended before its initial snapshot.`, + }), + ), + Deferred.succeed(serverConfigExit, undefined), + ]).pipe(Effect.asVoid), ), ); yield* serverConfigSource.pipe(Effect.forkScoped); @@ -221,16 +228,24 @@ export const make = Effect.gen(function* () { if (snapshot === undefined) { return Stream.empty; } + const updates = Stream.fromSubscription(subscription).pipe( + Stream.filter((buffered) => buffered.revision > snapshot.revision), + Stream.map( + (buffered): ServerConfigStreamEvent => ({ + version: 1, + type: "snapshot", + config: buffered.config, + }), + ), + ); + const terminal = Stream.fromEffect(Deferred.await(serverConfigExit)).pipe(Stream.drain); return Stream.concat( Stream.succeed({ version: 1 as const, type: "snapshot" as const, config: snapshot.config, }), - Stream.fromSubscription(subscription).pipe( - Stream.filter((buffered) => buffered.revision > snapshot.revision), - Stream.map((buffered) => buffered.event), - ), + Stream.merge(updates, terminal, { haltStrategy: "either" }), ); }), ); From f169a08ae5c15255c6b34c8d4e1250254fe4e70e Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Thu, 27 Aug 2026 11:01:30 +0200 Subject: [PATCH 4/5] fix(client-runtime): preserve config update events --- .../client-runtime/src/rpc/session.test.ts | 2 +- packages/client-runtime/src/rpc/session.ts | 23 +++++++++++++------ 2 files changed, 17 insertions(+), 8 deletions(-) diff --git a/packages/client-runtime/src/rpc/session.test.ts b/packages/client-runtime/src/rpc/session.test.ts index e73de59abff2..13364bed8426 100644 --- a/packages/client-runtime/src/rpc/session.test.ts +++ b/packages/client-runtime/src/rpc/session.test.ts @@ -339,7 +339,7 @@ describe("RpcSessionFactory", () => { const firstEvents = Array.from(yield* Fiber.join(firstSubscriber)); const secondEvents = Array.from(yield* Fiber.join(secondSubscriber)); - expect(firstEvents.map((event) => event.type)).toEqual(["snapshot", "snapshot"]); + expect(firstEvents.map((event) => event.type)).toEqual(["snapshot", "keybindingsUpdated"]); expect(secondEvents).toEqual(firstEvents); const replay = yield* session.serverConfigEvents!.pipe(Stream.runHead); diff --git a/packages/client-runtime/src/rpc/session.ts b/packages/client-runtime/src/rpc/session.ts index d547f99a1b58..70f299eae364 100644 --- a/packages/client-runtime/src/rpc/session.ts +++ b/packages/client-runtime/src/rpc/session.ts @@ -71,6 +71,7 @@ interface ServerConfigReplayState { interface BufferedServerConfigEvent { readonly config: ServerConfig; + readonly event: ServerConfigStreamEvent; readonly revision: number; } @@ -190,7 +191,7 @@ export const make = Effect.gen(function* () { config, revision: (current?.revision ?? 0) + 1, }; - return [{ config: next.config, revision: next.revision }, next] as const; + return [{ config: next.config, event, revision: next.revision }, next] as const; }); if (buffered !== undefined) { yield* PubSub.publish(serverConfigUpdates, buffered); @@ -230,12 +231,20 @@ export const make = Effect.gen(function* () { } const updates = Stream.fromSubscription(subscription).pipe( Stream.filter((buffered) => buffered.revision > snapshot.revision), - Stream.map( - (buffered): ServerConfigStreamEvent => ({ - version: 1, - type: "snapshot", - config: buffered.config, - }), + Stream.mapAccum( + () => snapshot.revision, + (revision, buffered) => [ + buffered.revision, + [ + buffered.revision === revision + 1 + ? buffered.event + : ({ + version: 1, + type: "snapshot", + config: buffered.config, + } satisfies ServerConfigStreamEvent), + ], + ], ), ); const terminal = Stream.fromEffect(Deferred.await(serverConfigExit)).pipe(Stream.drain); From 66ee0d27e63376f06116f8fe9d83332e6299b729 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Thu, 27 Aug 2026 14:19:51 +0200 Subject: [PATCH 5/5] test(client-runtime): cover config event replay --- .../client-runtime/src/rpc/session.test.ts | 111 ++++++++++++++++++ 1 file changed, 111 insertions(+) diff --git a/packages/client-runtime/src/rpc/session.test.ts b/packages/client-runtime/src/rpc/session.test.ts index 13364bed8426..a6e50fe56d9e 100644 --- a/packages/client-runtime/src/rpc/session.test.ts +++ b/packages/client-runtime/src/rpc/session.test.ts @@ -1,8 +1,12 @@ 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"; @@ -144,6 +148,7 @@ 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, @@ -354,6 +359,112 @@ describe("RpcSessionFactory", () => { ), ); + it.effect.each<{ + readonly event: ServerConfigStreamEventType; + readonly expectedConfig: Partial; + }>([ + { + 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();