diff --git a/apps/mobile/src/connection/background-activity.test.ts b/apps/mobile/src/connection/background-activity.test.ts index 7a2d902557ed..2c5874525fe1 100644 --- a/apps/mobile/src/connection/background-activity.test.ts +++ b/apps/mobile/src/connection/background-activity.test.ts @@ -1,14 +1,139 @@ -import { EnvironmentId, WS_METHODS } from "@t3tools/contracts"; +import { + AVAILABLE_CONNECTION_STATE, + type ConnectionCatalogEntry, + EnvironmentRegistry, + EnvironmentSupervisor, + type NetworkStatus, + type PreparedConnection, + RelayConnectionTarget, + type SupervisorConnectionState, +} from "@t3tools/client-runtime/connection"; +import type { RpcSession, WsRpcProtocolClient } from "@t3tools/client-runtime/rpc"; +import { type ClientActivityReportInput, EnvironmentId, WS_METHODS } from "@t3tools/contracts"; import { describe, expect, it } from "@effect/vitest"; +import * as Clock from "effect/Clock"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Queue from "effect/Queue"; +import * as Stream from "effect/Stream"; +import * as SubscriptionRef from "effect/SubscriptionRef"; +import * as TestClock from "effect/testing/TestClock"; +import { vi } from "vite-plus/test"; +import { MobileStorage } from "../persistence/mobile-storage"; +import { mobileBackgroundActivityReporterLayer } from "./background-activity"; import { onRetainedMobileBackgroundScopesChange, observeMobileBackgroundActivitySubscription, retainedMobileBackgroundScopes, } from "./background-activity-scopes"; +vi.mock("react-native", () => ({ + AppState: { currentState: "active", addEventListener: () => ({ remove: () => {} }) }, +})); +vi.mock("expo-secure-store", () => ({})); + describe("mobile background activity", () => { + it.effect("reports activity after reconnect when the initial report had no session", () => + Effect.gen(function* () { + const environmentId = EnvironmentId.make("reconnecting-environment"); + const target = new RelayConnectionTarget({ environmentId, label: "Test environment" }); + const reports = yield* Queue.unbounded(); + const attempts = yield* Queue.unbounded(); + const debounceArmed = yield* Queue.sliding(1); + const clock = yield* Clock.Clock; + const reporterClock = { + ...clock, + sleep: Effect.fn("TestMobileActivity.sleep")(function* (duration: Duration.Duration) { + if (Duration.toMillis(duration) !== 250) return yield* clock.sleep(duration); + const timer = yield* clock + .sleep(duration) + .pipe(Effect.forkChild({ startImmediately: true })); + yield* Queue.offer(debounceArmed, undefined); + yield* Fiber.join(timer); + }), + }; + const state = yield* SubscriptionRef.make({ + ...AVAILABLE_CONNECTION_STATE, + desired: true, + network: "online", + phase: "connecting", + }); + const session = yield* SubscriptionRef.make>(Option.none()); + const supervisor = EnvironmentSupervisor.of({ + target, + state, + session, + prepared: yield* SubscriptionRef.make>(Option.none()), + connect: Effect.void, + disconnect: Effect.void, + retryNow: Effect.void, + }); + const registryLayer = Layer.mock(EnvironmentRegistry, { + entries: yield* SubscriptionRef.make>( + new Map([[environmentId, { target, profile: Option.none() }]]), + ), + networkStatus: yield* SubscriptionRef.make("online"), + run: (_environmentId, effect) => + Effect.provideService(effect, EnvironmentSupervisor, supervisor).pipe( + Effect.ensuring(Queue.offer(attempts, undefined)), + ), + followStream: (_environmentId, stream) => + Stream.provideService(stream, EnvironmentSupervisor, supervisor), + }); + yield* Layer.build( + mobileBackgroundActivityReporterLayer.pipe( + Layer.provide(registryLayer), + Layer.provide( + Layer.mock(MobileStorage, { + loadOrCreateAgentAwarenessDeviceId: Effect.succeed("test-device"), + }), + ), + ), + ).pipe(Effect.provideService(Clock.Clock, reporterClock)); + + yield* Queue.take(debounceArmed); + yield* TestClock.adjust("250 millis"); + yield* Queue.take(attempts); + yield* Queue.clear(debounceArmed); + expect(yield* Queue.size(reports)).toBe(0); + + const client = { + [WS_METHODS.serverReportClientActivity]: (input: ClientActivityReportInput) => + Queue.offer(reports, input).pipe(Effect.asVoid), + } satisfies Pick; + yield* SubscriptionRef.set( + session, + Option.some({ + client: client as unknown as WsRpcProtocolClient, + initialConfig: Effect.die("Activity reports do not read server config."), + subscribeServerConfig: () => Stream.empty, + ready: Effect.void, + probe: Effect.void, + closed: Effect.never, + }), + ); + yield* SubscriptionRef.update(state, (current) => ({ + ...current, + phase: "connected" as const, + generation: 1, + })); + yield* Queue.take(debounceArmed); + yield* TestClock.adjust("250 millis"); + expect(yield* Queue.take(reports)).toMatchObject({ + environmentId, + clientId: "mobile-test-device", + visible: true, + focused: true, + appState: "active", + scopes: [{ type: "provider-status" }], + }); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("retains VCS demand only while the mobile subscription is active", () => Effect.gen(function* () { const environmentId = EnvironmentId.make("mobile-environment"); diff --git a/apps/mobile/src/connection/background-activity.ts b/apps/mobile/src/connection/background-activity.ts index 3364eed3c2d1..d1cbdfebf39b 100644 --- a/apps/mobile/src/connection/background-activity.ts +++ b/apps/mobile/src/connection/background-activity.ts @@ -1,4 +1,4 @@ -import { EnvironmentRegistry } from "@t3tools/client-runtime/connection"; +import { EnvironmentRegistry, EnvironmentSupervisor } from "@t3tools/client-runtime/connection"; import { EnvironmentRpcSubscriptionObserver, request } from "@t3tools/client-runtime/rpc"; import { type BackgroundScope, @@ -103,6 +103,33 @@ export const mobileBackgroundActivityReporterLayer = Layer.effectDiscard( Stream.runForEach(() => Effect.sync(requestReport)), Effect.forkScoped, ); + // A resume report can fail before the socket reconnects. Report again when + // it connects. Each supervisor has its own generation counter, so restart + // deduplication when the registry replaces the supervisor. + const connectedGenerations = (environmentId: EnvironmentId) => + registry.followStream( + environmentId, + Stream.unwrap( + Effect.map(EnvironmentSupervisor, (supervisor) => + SubscriptionRef.changes(supervisor.state).pipe( + Stream.filter((state) => state.phase === "connected"), + Stream.map((state) => state.generation), + Stream.changes, + ), + ), + ), + ); + yield* SubscriptionRef.changes(registry.entries).pipe( + Stream.map((entries) => [...entries.keys()].sort()), + Stream.changesWith((a, b) => a.length === b.length && a.every((id, i) => id === b[i])), + Stream.switchMap((environmentIds) => + Stream.mergeAll(environmentIds.map(connectedGenerations), { + concurrency: "unbounded", + }), + ), + Stream.runForEach(() => Effect.sync(requestReport)), + Effect.forkScoped, + ); yield* Stream.fromQueue(reportRequests).pipe( Stream.debounce("250 millis"), Stream.runForEach(() => report), diff --git a/apps/mobile/src/connection/network-path-change.test.ts b/apps/mobile/src/connection/network-path-change.test.ts new file mode 100644 index 000000000000..eae446a84445 --- /dev/null +++ b/apps/mobile/src/connection/network-path-change.test.ts @@ -0,0 +1,42 @@ +import { describe, expect, it } from "@effect/vitest"; + +import { + observeNetworkPath, + seedNetworkPathBaseline, + UNKNOWN_NETWORK_PATH, +} from "./network-path-change"; + +describe("network path change detection", () => { + it("does not probe when the listener matches the seeded interface", () => { + const seeded = seedNetworkPathBaseline(UNKNOWN_NETWORK_PATH, "WIFI"); + + expect(observeNetworkPath(seeded, "WIFI")).toEqual({ + baseline: { known: true, type: "WIFI" }, + shouldProbe: false, + }); + }); + + it("probes when a known interface changes", () => { + const seeded = seedNetworkPathBaseline(UNKNOWN_NETWORK_PATH, "WIFI"); + + expect(observeNetworkPath(seeded, "CELLULAR").shouldProbe).toBe(true); + }); + + it("probes when the listener wins the race with the async seed", () => { + const observed = observeNetworkPath(UNKNOWN_NETWORK_PATH, "CELLULAR"); + + expect(observed.shouldProbe).toBe(true); + expect(seedNetworkPathBaseline(observed.baseline, "WIFI")).toEqual({ + known: true, + type: "CELLULAR", + }); + }); + + it("probes on the first known interface after an unknown observation", () => { + const seeded = seedNetworkPathBaseline(UNKNOWN_NETWORK_PATH, "WIFI"); + const unknown = observeNetworkPath(seeded, null); + + expect(unknown.shouldProbe).toBe(false); + expect(observeNetworkPath(unknown.baseline, "CELLULAR").shouldProbe).toBe(true); + }); +}); diff --git a/apps/mobile/src/connection/network-path-change.ts b/apps/mobile/src/connection/network-path-change.ts new file mode 100644 index 000000000000..a189b001d327 --- /dev/null +++ b/apps/mobile/src/connection/network-path-change.ts @@ -0,0 +1,41 @@ +export interface NetworkPathBaseline { + readonly known: boolean; + readonly type: string | null; +} + +export const UNKNOWN_NETWORK_PATH: NetworkPathBaseline = { + known: false, + type: null, +}; + +export function seedNetworkPathBaseline( + baseline: NetworkPathBaseline, + type: string | null, +): NetworkPathBaseline { + if (baseline.known || type === null) { + return baseline; + } + return { known: true, type }; +} + +export function observeNetworkPath( + baseline: NetworkPathBaseline, + type: string | null, +): { + readonly baseline: NetworkPathBaseline; + readonly shouldProbe: boolean; +} { + if (type === null) { + return { + baseline: UNKNOWN_NETWORK_PATH, + shouldProbe: false, + }; + } + return { + baseline: { known: true, type }, + // If the initial async seed has not landed, the first listener event may + // itself be the WiFi/cellular transition. A cheap advisory probe is safer + // than losing that transition and waiting for the socket ping timeout. + shouldProbe: !baseline.known || baseline.type !== type, + }; +} diff --git a/apps/mobile/src/connection/platform.ts b/apps/mobile/src/connection/platform.ts index d6b50e50a6a7..f6494ed8af4e 100644 --- a/apps/mobile/src/connection/platform.ts +++ b/apps/mobile/src/connection/platform.ts @@ -32,6 +32,11 @@ import { appAtomRegistry } from "../state/atom-registry"; import { clearThreadOutboxEnvironment } from "../state/thread-outbox-removal"; import { clearComposerDraftsEnvironment } from "../state/use-composer-drafts"; import { mobileApplicationActiveWakeup } from "./app-state-wakeups"; +import { + observeNetworkPath, + seedNetworkPathBaseline, + UNKNOWN_NETWORK_PATH, +} from "./network-path-change"; import { connectionStorageLayer } from "./storage"; function networkStatus(state: Network.NetworkState): "unknown" | "offline" | "online" { @@ -88,11 +93,13 @@ const connectivityLayer = Connectivity.layer({ const wakeupsLayer = Wakeups.layer({ changes: Stream.merge( - Stream.callback<"application-active-probe" | "application-active-reconnect">((queue) => + Stream.callback< + "application-active-probe" | "application-active-reconnect" | "network-path-changed" + >((queue) => Effect.acquireRelease( Effect.sync(() => { let backgroundedAtMs = AppState.currentState === "background" ? Date.now() : null; - return AppState.addEventListener("change", (state) => { + const appStateSubscription = AppState.addEventListener("change", (state) => { if (state === "background") { backgroundedAtMs = Date.now(); return; @@ -102,6 +109,32 @@ const wakeupsLayer = Wakeups.layer({ backgroundedAtMs = null; } }); + // Wi-Fi/cellular changes can keep isConnected true while the socket + // stops working. Probe active sessions when the interface changes. + // Seed the baseline because the listener only reports changes. + let networkPath = UNKNOWN_NETWORK_PATH; + void Network.getNetworkStateAsync() + .then((current) => { + networkPath = seedNetworkPathBaseline(networkPath, current.type ?? null); + }) + .catch(() => undefined); + const networkSubscription = Network.addNetworkStateListener((state) => { + const observation = observeNetworkPath(networkPath, state.type ?? null); + networkPath = observation.baseline; + if ( + observation.shouldProbe && + state.isConnected === true && + AppState.currentState === "active" + ) { + Queue.offerUnsafe(queue, "network-path-changed"); + } + }); + return { + remove: () => { + appStateSubscription.remove(); + networkSubscription.remove(); + }, + }; }), (subscription) => Effect.sync(() => subscription.remove()), ).pipe(Effect.asVoid), diff --git a/docs/internals/connection-runtime.md b/docs/internals/connection-runtime.md index 1b686c365770..ca565a5525f5 100644 --- a/docs/internals/connection-runtime.md +++ b/docs/internals/connection-runtime.md @@ -80,6 +80,17 @@ Wakeup handling differs by phase, in [supervisor.ts][supervisor]: mobile's `application-active-probe`) rather than reconnecting; a healthy session survives foregrounding. `application-active-reconnect` skips the probe and replaces the lease outright. +- Mobile emits `network-path-changed` when the active network interface changes + while the app is active and online. A connected supervisor probes its current + session with a three-second timeout. Other phases ignore this advisory wakeup, + so repeated interface changes do not shorten backoff. A failed probe uses the + existing immediate reconnect path. It does not request a shell resubscription. + +Mobile sends an activity report after each newly connected session generation. +This restores the server's activity lease if the foreground report failed while +the socket was down, without waiting for the 25-second reporting interval. +Reports use the existing 250-millisecond debounce. Generation deduplication is +per supervisor because replacement supervisors restart their counters. The UI derives `available`, `offline`, `connecting`, `reconnecting`, `connected`, and `error` from supervisor state plus explicit data-sync state. diff --git a/docs/user/remote-access.md b/docs/user/remote-access.md index 40473930c601..f55ef25524cf 100644 --- a/docs/user/remote-access.md +++ b/docs/user/remote-access.md @@ -167,6 +167,14 @@ With mise, asdf, fnm, or nodenv, make sure the tool's shim directory is installe If reconnecting after an app update fails, retry the SSH launch once. The launcher now compares its generated runner script, stops stale launcher-managed remote servers, clears the SSH launch PID/port state, and starts a fresh remote server. You should not normally need to delete `~/.t3/ssh-launch` or kill `t3` processes manually. +## Mobile network changes + +When your phone switches between Wi-Fi and cellular while the app is active, T3 Code checks each +connected environment and reconnects if it does not respond. After reconnecting, the app sends +your current activity state so provider status and repository updates can resume. + +A server at a local network address still requires access to that network. + ## Updating a Remote Server When the T3 Code web or desktop app and a remote server use different versions, a warning appears in diff --git a/packages/client-runtime/src/connection/supervisor.test.ts b/packages/client-runtime/src/connection/supervisor.test.ts index d9f54bb326ca..b24b409f8909 100644 --- a/packages/client-runtime/src/connection/supervisor.test.ts +++ b/packages/client-runtime/src/connection/supervisor.test.ts @@ -5,6 +5,7 @@ import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; import * as SubscriptionRef from "effect/SubscriptionRef"; @@ -122,13 +123,10 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: const prepareCount = yield* Ref.make(0); const sessionCount = yield* Ref.make(0); const releaseCount = yield* Ref.make(0); - const wakeups = yield* SubscriptionRef.make<{ - readonly sequence: number; + const wakeups = yield* Queue.unbounded<{ readonly reason: ConnectionWakeups.ConnectionWakeup; - }>({ - sequence: 0, - reason: "application-active", - }); + readonly consumed: Deferred.Deferred; + }>(); const closedSessions = yield* Ref.make< ReadonlyArray> >([]); @@ -181,8 +179,8 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: Layer.succeed( ConnectionWakeups.ConnectionWakeups, ConnectionWakeups.ConnectionWakeups.of({ - changes: SubscriptionRef.changes(wakeups).pipe( - Stream.drop(1), + changes: Stream.fromQueue(wakeups).pipe( + Stream.tap((event) => Deferred.succeed(event.consumed, undefined)), Stream.map((event) => event.reason), ), }), @@ -199,11 +197,13 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: sessionCount, releaseCount, setNetworkStatus: (status: NetworkStatus) => SubscriptionRef.set(networkStatus, status), - wake: (reason: ConnectionWakeups.ConnectionWakeup) => - SubscriptionRef.update(wakeups, (event) => ({ - sequence: event.sequence + 1, - reason, - })), + wake: Effect.fn("TestConnectionHarness.wake")(function* ( + reason: ConnectionWakeups.ConnectionWakeup, + ) { + const consumed = yield* Deferred.make(); + yield* Queue.offer(wakeups, { reason, consumed }); + yield* Deferred.await(consumed); + }), closeLatestSession: Effect.fn("TestConnectionHarness.closeLatestSession")(function* ( error = transient("Session closed."), ) { @@ -874,6 +874,82 @@ describe("EnvironmentSupervisor", () => { }), ); + it.effect("probes the active session when the network path changes", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const probeCalled = yield* Deferred.make(); + const harness = yield* makeHarness({ + probe: () => + Ref.update(probeCount, (count) => count + 1).pipe( + Effect.andThen(Deferred.succeed(probeCalled, undefined)), + ), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* harness.wake("network-path-changed"); + yield* Deferred.await(probeCalled); + + expect(yield* Ref.get(probeCount)).toBe(1); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + expect(yield* Ref.get(harness.releaseCount)).toBe(0); + expect((yield* SubscriptionRef.get(supervisor.state)).phase).toBe("connected"); + }), + ); + + it.effect("reconnects without backoff when a probe after a path change fails", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + probe: (attempt) => + attempt === 1 + ? Effect.fail(transient("The path changed under the socket.")) + : Effect.void, + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* harness.wake("network-path-changed"); + // The failed wake probe skips the first backoff rung (wakeProbeFailed), + // so the replacement connects without a TestClock advance. + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + expect(yield* Ref.get(harness.releaseCount)).toBe(1); + }), + ); + + it.effect("does not cut backoff short when the network path flaps", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + prepare: () => Effect.fail(transient()), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState( + supervisor.state, + (state) => state.phase === "backoff" && state.attempt === 1, + ); + expect(yield* Ref.get(harness.prepareCount)).toBe(1); + + // Advisory path-change wakeups have no session to probe during backoff + // and must not trigger an early retry. + yield* harness.wake("network-path-changed"); + yield* harness.wake("network-path-changed"); + yield* TestClock.adjust("2999 millis"); + expect(yield* Ref.get(harness.prepareCount)).toBe(1); + expect((yield* SubscriptionRef.get(supervisor.state)).phase).toBe("backoff"); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("immediately replaces a mobile session after a long background resume", () => Effect.gen(function* () { const probeCount = yield* Ref.make(0); @@ -996,8 +1072,12 @@ describe("EnvironmentSupervisor", () => { it.effect("uses the full tolerance window for a stalled desktop foreground probe", () => Effect.gen(function* () { + const probeStarted = yield* Deferred.make(); const harness = yield* makeHarness({ - probe: (attempt) => (attempt === 1 ? Effect.never : Effect.void), + probe: (attempt) => + attempt === 1 + ? Deferred.succeed(probeStarted, undefined).pipe(Effect.andThen(Effect.never)) + : Effect.void, }); const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { initiallyDesired: true, @@ -1005,6 +1085,7 @@ describe("EnvironmentSupervisor", () => { yield* awaitState(supervisor.state, (state) => state.phase === "connected"); yield* harness.wake("application-active"); + yield* Deferred.await(probeStarted); yield* TestClock.adjust("14999 millis"); expect(yield* Ref.get(harness.sessionCount)).toBe(1); yield* TestClock.adjust("1 milli"); @@ -1017,17 +1098,25 @@ describe("EnvironmentSupervisor", () => { }).pipe(Effect.provide(TestClock.layer())), ); - it.effect("quickly times out a stalled mobile foreground liveness probe", () => + it.effect.each([ + { reason: "application-active-probe" as const }, + { reason: "network-path-changed" as const }, + ])("quickly times out a stalled $reason probe", ({ reason }) => Effect.gen(function* () { + const probeStarted = yield* Deferred.make(); const harness = yield* makeHarness({ - probe: (attempt) => (attempt === 1 ? Effect.never : Effect.void), + probe: (attempt) => + attempt === 1 + ? Deferred.succeed(probeStarted, undefined).pipe(Effect.andThen(Effect.never)) + : Effect.void, }); const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { initiallyDesired: true, }).pipe(Effect.provide(harness.dependencies)); yield* awaitState(supervisor.state, (state) => state.phase === "connected"); - yield* harness.wake("application-active-probe"); + yield* harness.wake(reason); + yield* Deferred.await(probeStarted); yield* TestClock.adjust("3 seconds"); // The timed-out wake probe reconnects immediately without a backoff // sleep: no further clock advance is needed. diff --git a/packages/client-runtime/src/connection/supervisor.ts b/packages/client-runtime/src/connection/supervisor.ts index 85fda10ef1a7..0806e2347772 100644 --- a/packages/client-runtime/src/connection/supervisor.ts +++ b/packages/client-runtime/src/connection/supervisor.ts @@ -418,13 +418,17 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( // replaces that lease and starts a fresh attempt without backoff. return true; } - if (next.reason === "application-active" || next.reason === "application-active-probe") { + if ( + next.reason === "application-active" || + next.reason === "application-active-probe" || + next.reason === "network-path-changed" + ) { const probe = yield* lease.session.probe.pipe( Effect.timeoutOrElse({ duration: - next.reason === "application-active-probe" - ? MOBILE_CONNECTION_PROBE_TIMEOUT - : CONNECTION_PROBE_TIMEOUT, + next.reason === "application-active" + ? CONNECTION_PROBE_TIMEOUT + : MOBILE_CONNECTION_PROBE_TIMEOUT, orElse: () => Effect.fail( new ConnectionTransientError({ @@ -622,6 +626,11 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const next = yield* Queue.take(signals); switch (next._tag) { case "Wakeup": + // Path-change wakeups are advisory (probe a connected session) + // and must not cut backoff delays short on a flapping interface. + if (next.reason === "network-path-changed") { + break; + } return ConnectionWakeups.isApplicationActiveWakeup(next.reason); case "ConnectRequested": case "DisconnectRequested": @@ -634,11 +643,16 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( ); }); - const waitForSignal = Queue.take(signals).pipe( - Effect.map( - (next) => next._tag === "Wakeup" && ConnectionWakeups.isApplicationActiveWakeup(next.reason), - ), - ); + const waitForSignal = Effect.gen(function* () { + for (;;) { + const next = yield* Queue.take(signals); + if (next._tag === "Wakeup" && next.reason === "network-path-changed") { + // Advisory only; see waitForRetrySignal. + continue; + } + return next._tag === "Wakeup" && ConnectionWakeups.isApplicationActiveWakeup(next.reason); + } + }); const run = Effect.fnUntraced(function* () { let failureCount = 0; diff --git a/packages/client-runtime/src/connection/wakeups.ts b/packages/client-runtime/src/connection/wakeups.ts index 8573a49c1474..5302fdbd27c5 100644 --- a/packages/client-runtime/src/connection/wakeups.ts +++ b/packages/client-runtime/src/connection/wakeups.ts @@ -6,6 +6,10 @@ export type ConnectionWakeup = | "application-active" | "application-active-probe" | "application-active-reconnect" + // Mobile interface changes can leave the socket open on a dead route. + // Probe connected sessions only. Do not shorten backoff or retry blocked + // connections for this advisory wakeup. + | "network-path-changed" | "credentials-changed"; export function isApplicationActiveWakeup(reason: ConnectionWakeup): boolean { diff --git a/packages/client-runtime/src/rpc/session.ts b/packages/client-runtime/src/rpc/session.ts index 7d975be5c9d3..cbb10962f528 100644 --- a/packages/client-runtime/src/rpc/session.ts +++ b/packages/client-runtime/src/rpc/session.ts @@ -160,6 +160,13 @@ export const make = Effect.fn("RpcSessionFactory.make")(function* ( ), Effect.asVoid, ), + // Record missed pongs separately from ordinary socket closes. + onPingTimeout: Effect.logInfo("Connection ping timed out.").pipe( + Effect.annotateLogs({ + environmentId: connection.environmentId, + connectionLabel: connection.label, + }), + ), }); const socketLayer = Socket.layerWebSocket(connection.socketUrl, { openTimeout: SOCKET_OPEN_TIMEOUT, diff --git a/packages/client-runtime/src/state/shell-sync.test.ts b/packages/client-runtime/src/state/shell-sync.test.ts index 1c0d838026fb..e283b7829be4 100644 --- a/packages/client-runtime/src/state/shell-sync.test.ts +++ b/packages/client-runtime/src/state/shell-sync.test.ts @@ -47,7 +47,7 @@ const LIVE_SHELL_SNAPSHOT: OrchestrationShellSnapshot = { updatedAt: "2026-06-06T00:00:00.000Z", }; -function session(client: WsRpcProtocolClient): RpcSession.RpcSession { +function session(client: WsRpcProtocolClient) { return { client, initialConfig: Effect.succeed({ shellResumeCompletionMarker: true } as never), @@ -55,7 +55,7 @@ function session(client: WsRpcProtocolClient): RpcSession.RpcSession { ready: Effect.void, probe: Effect.void, closed: Effect.never, - }; + } satisfies RpcSession.RpcSession; } describe("environment shell synchronization", () => { @@ -253,18 +253,19 @@ describe("environment shell synchronization", () => { const events = yield* Queue.unbounded(); const wakeups = yield* Queue.unbounded(); const loaderCalls = yield* Ref.make(0); - const capturedAfterSequences = yield* Ref.make>([]); + const subscriptionReceipts = yield* Queue.unbounded(); const client = { [ORCHESTRATION_WS_METHODS.subscribeShell]: (input: { readonly afterSequence?: number }) => Stream.unwrap( - Ref.update(capturedAfterSequences, (captured) => [ - ...captured, - input.afterSequence, - ]).pipe(Effect.as(Stream.fromQueue(events))), + Queue.offer(subscriptionReceipts, input.afterSequence).pipe( + Effect.as(Stream.fromQueue(events)), + ), ), } as unknown as WsRpcProtocolClient; const supervisorState = yield* SubscriptionRef.make(AVAILABLE_CONNECTION_STATE); - const activeSession = yield* SubscriptionRef.make(Option.some(session(client))); + const activeSession = yield* SubscriptionRef.make>( + Option.some(session(client)), + ); const supervisor = EnvironmentSupervisor.EnvironmentSupervisor.of({ target: TARGET, state: supervisorState, @@ -307,11 +308,7 @@ describe("environment shell synchronization", () => { ); // A new session starts from an authoritative HTTP snapshot. - for (let attempt = 0; attempt < 100; attempt += 1) { - if ((yield* Ref.get(capturedAfterSequences)).length >= 1) break; - yield* Effect.yieldNow; - } - expect(yield* Ref.get(capturedAfterSequences)).toEqual([10]); + expect(yield* Queue.take(subscriptionReceipts)).toBe(10); yield* Queue.offer(events, { kind: "synchronized" }); yield* SubscriptionRef.changes(shellState).pipe( Stream.filter((value) => value.status === "live"), @@ -331,35 +328,27 @@ describe("environment shell synchronization", () => { ); yield* Queue.offer(wakeups, "application-active"); - for (let attempt = 0; attempt < 100; attempt += 1) { - if ((yield* Ref.get(capturedAfterSequences)).length >= 2) break; - yield* Effect.yieldNow; - } - expect(yield* Ref.get(capturedAfterSequences)).toEqual([10, 40]); + expect(yield* Queue.take(subscriptionReceipts)).toBe(40); yield* Queue.offer(events, { kind: "synchronized" }); yield* Queue.offer(wakeups, "application-active-probe"); - for (let attempt = 0; attempt < 100; attempt += 1) { - if ((yield* Ref.get(capturedAfterSequences)).length >= 3) break; - yield* Effect.yieldNow; - } - expect(yield* Ref.get(capturedAfterSequences)).toEqual([10, 40, 40]); - - yield* Queue.offer(wakeups, "application-active-reconnect"); - for (let attempt = 0; attempt < 10; attempt += 1) { - yield* Effect.yieldNow; - } - expect((yield* Ref.get(capturedAfterSequences)).length).toBe(3); + expect(yield* Queue.take(subscriptionReceipts)).toBe(40); expect(yield* Ref.get(loaderCalls)).toBe(1); // Replacing the session performs another authoritative refresh. yield* SubscriptionRef.set(activeSession, Option.some(session(client))); - for (let attempt = 0; attempt < 100; attempt += 1) { - if ((yield* Ref.get(capturedAfterSequences)).length >= 4) break; - yield* Effect.yieldNow; - } - expect(yield* Ref.get(capturedAfterSequences)).toEqual([10, 40, 40, 20]); + expect(yield* Queue.take(subscriptionReceipts)).toBe(20); expect(yield* Ref.get(loaderCalls)).toBe(2); }), ); + + it("only resubscribes for foreground wakeups that keep the current session", () => { + expect(ConnectionWakeups.shouldResubscribeAfterWakeup("application-active")).toBe(true); + expect(ConnectionWakeups.shouldResubscribeAfterWakeup("application-active-probe")).toBe(true); + expect(ConnectionWakeups.shouldResubscribeAfterWakeup("application-active-reconnect")).toBe( + false, + ); + expect(ConnectionWakeups.shouldResubscribeAfterWakeup("network-path-changed")).toBe(false); + expect(ConnectionWakeups.shouldResubscribeAfterWakeup("credentials-changed")).toBe(false); + }); });