From cdf28bef875aceb9633114ed194b0998f6378558 Mon Sep 17 00:00:00 2001 From: Bear Huddleston Date: Sat, 26 Sep 2026 10:57:41 -0500 Subject: [PATCH] fix(client-runtime): connection status follows a re-registered environment Registering an environment again replaces its supervisor even when the catalog entry is identical, which is what a re-pair does: same target and profile, new credential. followStream only switched supervisors when the entry changed, so status and durable streams stayed on the closed supervisor and showed its last state (often a stale reconnect error) until the environment was switched off and on. It now follows the environment's current supervisor. --- .../src/connection/registry.test.ts | 56 +++++++++++++++++++ .../client-runtime/src/connection/registry.ts | 28 ++++++++-- 2 files changed, 79 insertions(+), 5 deletions(-) diff --git a/packages/client-runtime/src/connection/registry.test.ts b/packages/client-runtime/src/connection/registry.test.ts index 8a71344ad173..4b81fd272253 100644 --- a/packages/client-runtime/src/connection/registry.test.ts +++ b/packages/client-runtime/src/connection/registry.test.ts @@ -1145,6 +1145,62 @@ describe("EnvironmentRegistry", () => { }), ); + it.effect("moves durable streams to the supervisor an identical re-registration installs", () => + Effect.gen(function* () { + const harness = yield* makeHarness([RELAY_TARGET]); + + yield* Effect.gen(function* () { + const registry = yield* EnvironmentRegistry.EnvironmentRegistry; + const firstObserved = yield* Deferred.make(); + const secondObserved = yield* Deferred.make(); + const supervisors = yield* Ref.make>([]); + yield* registry.start; + yield* awaitConnectionState( + registry, + RELAY_TARGET.environmentId, + (state) => state.phase === "connected", + ); + + const subscription = yield* Effect.forkChild( + registry + .followStream( + RELAY_TARGET.environmentId, + Stream.unwrap( + EnvironmentSupervisor.EnvironmentSupervisor.pipe( + Effect.map((supervisor) => + Stream.concat(Stream.succeed(supervisor), Stream.never), + ), + ), + ), + ) + .pipe( + Stream.tap((supervisor) => + Ref.updateAndGet(supervisors, (current) => [...current, supervisor]).pipe( + Effect.flatMap((current) => + current.length === 1 + ? Deferred.succeed(firstObserved, undefined) + : Deferred.succeed(secondObserved, undefined), + ), + ), + ), + Stream.runDrain, + ), + ); + + yield* Deferred.await(firstObserved).pipe(Effect.timeout("1 second")); + // A re-pair re-registers the same target and profile; only the stored + // credential changes, so the catalog entry is identical. + yield* registry.register(new RelayConnectionRegistration({ target: RELAY_TARGET })); + yield* Deferred.await(secondObserved).pipe(Effect.timeout("1 second")); + yield* Fiber.interrupt(subscription); + + const [first, second] = yield* Ref.get(supervisors); + expect(second).toBeDefined(); + expect(second).not.toBe(first); + }).pipe(Effect.provide(harness.layer), Effect.scoped); + }), + ); + it.effect("ignores retry signals for environments that are no longer registered", () => Effect.gen(function* () { const harness = yield* makeHarness([]); diff --git a/packages/client-runtime/src/connection/registry.ts b/packages/client-runtime/src/connection/registry.ts index af1cc46faa9b..5bdc95f7b30b 100644 --- a/packages/client-runtime/src/connection/registry.ts +++ b/packages/client-runtime/src/connection/registry.ts @@ -5,6 +5,7 @@ import * as Equal from "effect/Equal"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Predicate from "effect/Predicate"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; @@ -378,11 +379,18 @@ export const make = Effect.gen(function* () { acquireSupervisor(environmentId).pipe( Effect.match({ onFailure: () => Stream.empty, - onSuccess: (supervisor) => - Stream.provideService( - stream, - EnvironmentSupervisor.EnvironmentSupervisor, - supervisor, + // Reinstalling an entry replaces its supervisor even when the + // entry is unchanged (a re-pair keeps the same target and + // profile), so follow the supervisor itself. + onSuccess: () => + supervisorChanges(environmentId).pipe( + Stream.switchMap((supervisor) => + Stream.provideService( + stream, + EnvironmentSupervisor.EnvironmentSupervisor, + supervisor, + ), + ), ), }), ), @@ -391,6 +399,16 @@ export const make = Effect.gen(function* () { ), ); + const supervisorChanges = (environmentId: EnvironmentId) => + Stream.concat( + Stream.fromEffect(SubscriptionRef.get(serviceScopes)), + SubscriptionRef.changes(serviceScopes), + ).pipe( + Stream.map((current) => current.get(environmentId)?.supervisor), + Stream.filter(Predicate.isNotUndefined), + Stream.changesWith((previous, next) => previous === next), + ); + const start = Effect.gen(function* () { if (yield* Ref.getAndSet(started, true)) { return;