diff --git a/apps/server/src/environment/ServerEnvironment.test.ts b/apps/server/src/environment/ServerEnvironment.test.ts index 6b3290246fe..61892d53d63 100644 --- a/apps/server/src/environment/ServerEnvironment.test.ts +++ b/apps/server/src/environment/ServerEnvironment.test.ts @@ -67,6 +67,7 @@ it.layer(NodeServices.layer)("ServerEnvironmentLive", (it) => { expect(first.environmentId).toBe(second.environmentId); expect(second.capabilities.repositoryIdentity).toBe(true); + expect(second.capabilities.connectionProbe).toBe(true); }), ); diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index b5fbd8e1088..1c0d34ea5bc 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -135,6 +135,7 @@ export const make = Effect.gen(function* () { serverVersion: packageJson.version, capabilities: { repositoryIdentity: true, + connectionProbe: true, }, }; diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 32dcb8a13d6..c7223d09dc5 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -284,6 +284,7 @@ const RPC_REQUIRED_SCOPE = new Map([ [ORCHESTRATION_WS_METHODS.subscribeShell, AuthOrchestrationReadScope], [ORCHESTRATION_WS_METHODS.getArchivedShellSnapshot, AuthOrchestrationReadScope], [ORCHESTRATION_WS_METHODS.subscribeThread, AuthOrchestrationReadScope], + [WS_METHODS.serverProbe, AuthOrchestrationReadScope], [WS_METHODS.serverGetConfig, AuthOrchestrationReadScope], [WS_METHODS.serverRefreshProviders, AuthOrchestrationOperateScope], [WS_METHODS.serverUpdateProvider, AuthOrchestrationOperateScope], @@ -1238,6 +1239,10 @@ const makeWsRpcLayer = ( }), { "rpc.aggregate": "orchestration" }, ), + [WS_METHODS.serverProbe]: (_input) => + observeRpcEffect(WS_METHODS.serverProbe, Effect.succeed({}), { + "rpc.aggregate": "server", + }), [WS_METHODS.serverGetConfig]: (_input) => observeRpcEffect(WS_METHODS.serverGetConfig, loadServerConfig, { "rpc.aggregate": "server", diff --git a/packages/client-runtime/src/connection/supervisor.test.ts b/packages/client-runtime/src/connection/supervisor.test.ts index eadeceacc2c..f3901e42251 100644 --- a/packages/client-runtime/src/connection/supervisor.test.ts +++ b/packages/client-runtime/src/connection/supervisor.test.ts @@ -652,7 +652,7 @@ describe("EnvironmentSupervisor", () => { }).pipe(Effect.provide(TestClock.layer())), ); - it.effect("keeps a healthy session when the application becomes active", () => + it.effect("probes the active session without reconnecting on application activation", () => Effect.gen(function* () { const probeCount = yield* Ref.make(0); const probeCalled = yield* Deferred.make(); diff --git a/packages/client-runtime/src/rpc/session.test.ts b/packages/client-runtime/src/rpc/session.test.ts index 7820c93a935..402faf4f0e2 100644 --- a/packages/client-runtime/src/rpc/session.test.ts +++ b/packages/client-runtime/src/rpc/session.test.ts @@ -109,6 +109,7 @@ const SERVER_CONFIG: ServerConfigType = { serverVersion: "0.0.0-test", capabilities: { repositoryIdentity: true, + connectionProbe: true, }, }, auth: { @@ -141,6 +142,16 @@ const decodeJson = Schema.decodeUnknownSync(Schema.UnknownFromJsonString); const decodeRpcRequest = Schema.decodeUnknownSync(RpcRequest); const encodeJson = Schema.encodeUnknownSync(Schema.UnknownFromJsonString); const encodeServerConfig = Schema.encodeSync(ServerConfig); +const ENCODED_SERVER_CONFIG = encodeServerConfig(SERVER_CONFIG); +const LEGACY_SERVER_CONFIG = { + ...ENCODED_SERVER_CONFIG, + environment: { + ...ENCODED_SERVER_CONFIG.environment, + capabilities: { + repositoryIdentity: true, + }, + }, +}; const makeFactory = Effect.fn("TestRpcSessionFactory.make")(function* () { const sockets: TestWebSocket[] = []; @@ -169,9 +180,10 @@ const awaitSocket = Effect.fn("TestRpcSessionFactory.awaitSocket")(function* ( const awaitRequest = Effect.fn("TestRpcSessionFactory.awaitRequest")(function* ( socket: TestWebSocket, + index = 0, ) { for (let attempt = 0; attempt < 100; attempt += 1) { - const request = socket.sent[0]; + const request = socket.sent[index]; if (request) { return decodeRpcRequest(decodeJson(request)); } @@ -182,6 +194,7 @@ const awaitRequest = Effect.fn("TestRpcSessionFactory.awaitRequest")(function* ( const completeInitialConfig = Effect.fn("TestRpcSessionFactory.completeInitialConfig")(function* ( socket: TestWebSocket, + config: unknown = ENCODED_SERVER_CONFIG, ) { const request = yield* awaitRequest(socket); expect(request).toMatchObject({ @@ -195,7 +208,7 @@ const completeInitialConfig = Effect.fn("TestRpcSessionFactory.completeInitialCo requestId: request.id, exit: { _tag: "Success", - value: encodeServerConfig(SERVER_CONFIG), + value: config, }, }), ); @@ -218,6 +231,30 @@ describe("RpcSessionFactory", () => { expect(config).toEqual(SERVER_CONFIG); expect(socket.sent).toHaveLength(1); + const probeFiber = yield* Effect.forkChild(session.probe); + const probeRequest = yield* awaitRequest(socket, 1); + expect(probeRequest).toMatchObject({ + _tag: "Request", + tag: WS_METHODS.serverProbe, + payload: {}, + }); + socket.serverMessage( + encodeJson({ + _tag: "Exit", + requestId: probeRequest.id, + exit: { + _tag: "Success", + value: {}, + }, + }), + ); + yield* Fiber.join(probeFiber); + + expect(socket.sent.map((request) => decodeRpcRequest(decodeJson(request)).tag)).toEqual([ + WS_METHODS.serverGetConfig, + WS_METHODS.serverProbe, + ]); + socket.close(1012, "service restart"); const error = yield* Effect.flip(session.closed); @@ -250,6 +287,45 @@ describe("RpcSessionFactory", () => { }), ); + it.effect("uses the legacy config RPC for probes when the server lacks the capability", () => + 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, LEGACY_SERVER_CONFIG); + yield* Fiber.join(readyFiber); + + const probeFiber = yield* Effect.forkChild(session.probe); + const probeRequest = yield* awaitRequest(socket, 1); + expect(probeRequest).toMatchObject({ + _tag: "Request", + tag: WS_METHODS.serverGetConfig, + payload: {}, + }); + socket.serverMessage( + encodeJson({ + _tag: "Exit", + requestId: probeRequest.id, + exit: { + _tag: "Success", + value: LEGACY_SERVER_CONFIG, + }, + }), + ); + yield* Fiber.join(probeFiber); + + expect(socket.sent.map((request) => decodeRpcRequest(decodeJson(request)).tag)).toEqual([ + WS_METHODS.serverGetConfig, + WS_METHODS.serverGetConfig, + ]); + }), + ), + ); + it.effect("fails readiness when the websocket never opens", () => 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 f9594c1b7ca..9625effa406 100644 --- a/packages/client-runtime/src/rpc/session.ts +++ b/packages/client-runtime/src/rpc/session.ts @@ -42,8 +42,9 @@ export class RpcSessionFactory extends Context.Service< type InitialConfigError = Effect.Error< ReturnType >; +type ProbeError = Effect.Error>; -function mapInitialConfigError(error: InitialConfigError): ConnectionAttemptError { +function mapSessionRpcError(error: InitialConfigError | ProbeError): ConnectionAttemptError { switch (error._tag) { case "EnvironmentAuthorizationError": return new ConnectionBlockedError({ @@ -115,12 +116,17 @@ export const make = Effect.gen(function* () { const client = yield* makeWsRpcProtocolClient.pipe(Effect.provide(protocolContext)); const initialConfig = yield* Effect.cached( client[WS_METHODS.serverGetConfig]({}).pipe( - Effect.mapError(mapInitialConfigError), + Effect.mapError(mapSessionRpcError), Effect.withSpan("environment.initialSync"), ), ); - const probe = client[WS_METHODS.serverGetConfig]({}).pipe( - Effect.mapError(mapInitialConfigError), + const probe = initialConfig.pipe( + Effect.flatMap((config) => + (config.environment.capabilities.connectionProbe === true + ? client[WS_METHODS.serverProbe]({}) + : client[WS_METHODS.serverGetConfig]({}) + ).pipe(Effect.mapError(mapSessionRpcError)), + ), Effect.asVoid, Effect.withSpan("clientRuntime.connection.rpcSession.probe"), ); diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index fb52972d5aa..bc8db25a95a 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -22,6 +22,7 @@ export type ExecutionEnvironmentPlatform = typeof ExecutionEnvironmentPlatform.T export const ExecutionEnvironmentCapabilities = Schema.Struct({ repositoryIdentity: Schema.Boolean.pipe(Schema.withDecodingDefault(Effect.succeed(false))), + connectionProbe: Schema.optionalKey(Schema.Boolean), }); export type ExecutionEnvironmentCapabilities = typeof ExecutionEnvironmentCapabilities.Type; diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index 48c5d9a774d..0356aa1807b 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -201,6 +201,7 @@ export const WS_METHODS = { previewAutomationFocusHost: "previewAutomation.focusHost", // Server meta + serverProbe: "server.probe", serverGetConfig: "server.getConfig", serverRefreshProviders: "server.refreshProviders", serverUpdateProvider: "server.updateProvider", @@ -246,6 +247,12 @@ export const WsServerRemoveKeybindingRpc = Rpc.make(WS_METHODS.serverRemoveKeybi error: Schema.Union([KeybindingsConfigError, EnvironmentAuthorizationError]), }); +export const WsServerProbeRpc = Rpc.make(WS_METHODS.serverProbe, { + payload: Schema.Struct({}), + success: Schema.Struct({}), + error: EnvironmentAuthorizationError, +}); + export const WsServerGetConfigRpc = Rpc.make(WS_METHODS.serverGetConfig, { payload: Schema.Struct({}), success: ServerConfig, @@ -682,6 +689,7 @@ export const WsSubscribeAuthAccessRpc = Rpc.make(WS_METHODS.subscribeAuthAccess, }); export const WsRpcGroup = RpcGroup.make( + WsServerProbeRpc, WsServerGetConfigRpc, WsServerRefreshProvidersRpc, WsServerUpdateProviderRpc,