diff --git a/apps/web/src/components/ChatView.logic.test.ts b/apps/web/src/components/ChatView.logic.test.ts index 57c12959ffb9..f57a497a3485 100644 --- a/apps/web/src/components/ChatView.logic.test.ts +++ b/apps/web/src/components/ChatView.logic.test.ts @@ -25,6 +25,7 @@ import { reconcileRetainedMountedThreadIds, resolveThreadMetadataUpdateForNextTurn, resolveSendEnvMode, + shouldQuietlyRecoverManagedPrimaryEnvironment, shouldShowBranchMismatchBanner, shouldWriteThreadErrorToCurrentServerThread, } from "./ChatView.logic"; @@ -34,6 +35,49 @@ const projectId = ProjectId.make("project-1"); const threadId = ThreadId.make("thread-1"); const now = "2026-03-29T00:00:00.000Z"; +describe("managed primary recovery", () => { + it.each(["connecting", "reconnecting"])( + "keeps the managed primary surface interactive while %s", + (connectionPhase) => { + expect( + shouldQuietlyRecoverManagedPrimaryEnvironment({ + managed: true, + activeEnvironmentId: environmentId, + primaryEnvironmentId: environmentId, + connectionPhase, + }), + ).toBe(true); + }, + ); + + it("does not hide unavailable state for unmanaged, secondary, or failed environments", () => { + expect( + shouldQuietlyRecoverManagedPrimaryEnvironment({ + managed: false, + activeEnvironmentId: environmentId, + primaryEnvironmentId: environmentId, + connectionPhase: "reconnecting", + }), + ).toBe(false); + expect( + shouldQuietlyRecoverManagedPrimaryEnvironment({ + managed: true, + activeEnvironmentId: environmentId, + primaryEnvironmentId: EnvironmentId.make("environment-primary"), + connectionPhase: "reconnecting", + }), + ).toBe(false); + expect( + shouldQuietlyRecoverManagedPrimaryEnvironment({ + managed: true, + activeEnvironmentId: environmentId, + primaryEnvironmentId: environmentId, + connectionPhase: "error", + }), + ).toBe(false); + }); +}); + function makeThread(overrides: Partial = {}): Thread { return { id: threadId, diff --git a/apps/web/src/components/ChatView.logic.ts b/apps/web/src/components/ChatView.logic.ts index c0e967914bd9..3fe2b81c9d51 100644 --- a/apps/web/src/components/ChatView.logic.ts +++ b/apps/web/src/components/ChatView.logic.ts @@ -27,6 +27,20 @@ export const MAX_HIDDEN_MOUNTED_PREVIEW_THREADS = 3; export const LastInvokedScriptByProjectSchema = Schema.Record(ProjectId, Schema.String); +export function shouldQuietlyRecoverManagedPrimaryEnvironment(input: { + readonly managed: boolean; + readonly activeEnvironmentId: EnvironmentId | null; + readonly primaryEnvironmentId: EnvironmentId | null; + readonly connectionPhase: string; +}): boolean { + return Boolean( + input.managed && + input.activeEnvironmentId !== null && + input.activeEnvironmentId === input.primaryEnvironmentId && + ["connecting", "reconnecting"].includes(input.connectionPhase), + ); +} + export function resolveThreadMetadataUpdateForNextTurn(input: { currentModelSelection: ModelSelection; nextModelSelection?: ModelSelection; diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 79afad5f9bb6..fc105ebb458d 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -130,7 +130,7 @@ import { setActivePreviewTab, useThreadPreviewState, } from "../previewStateStore"; -import { managedWorkspaceBrowserUrl } from "~/managedDevPc"; +import { isManagedDevPc, managedWorkspaceBrowserUrl } from "~/managedDevPc"; import { addBrowserSurface } from "./preview/addBrowserSurface"; import { closePreviewSession } from "./preview/closePreviewSession"; import { subscribePreviewAction } from "./preview/previewActionBus"; @@ -269,6 +269,7 @@ import { resolveSendEnvMode, revokeBlobPreviewUrl, revokeUserMessagePreviewUrls, + shouldQuietlyRecoverManagedPrimaryEnvironment, waitForStartedServerThread, } from "./ChatView.logic"; import { useLocalStorage } from "~/hooks/useLocalStorage"; @@ -1625,6 +1626,17 @@ function ChatViewContent(props: ChatViewProps) { connection: activeEnvironment.connection, }; }, [activeEnvironment, activeEnvironmentUnavailable, activeEnvironmentUnavailableLabel]); + const quietlyRecoveringManagedPrimary = shouldQuietlyRecoverManagedPrimaryEnvironment({ + managed: isManagedDevPc, + activeEnvironmentId: activeEnvironment?.environmentId ?? null, + primaryEnvironmentId, + connectionPhase: activeEnvironmentConnectionPhase, + }); + const activeEnvironmentActionUnavailable = + activeEnvironmentUnavailable && !quietlyRecoveringManagedPrimary; + const activeEnvironmentActionUnavailableState = quietlyRecoveringManagedPrimary + ? null + : activeEnvironmentUnavailableState; const handleReconnectActiveEnvironment = useCallback( async (environmentId: EnvironmentId) => { const result = await retryEnvironment(environmentId); @@ -1838,13 +1850,13 @@ function ChatViewContent(props: ChatViewProps) { const versionMismatchSelfUpdate = resolveServerSelfUpdateCapability(serverConfig); const systemComposerBannerItems = useMemo(() => { const items: ComposerBannerStackItem[] = []; - if (activeEnvironmentUnavailableState) { - const connection = activeEnvironmentUnavailableState.connection; + if (activeEnvironmentActionUnavailableState) { + const connection = activeEnvironmentActionUnavailableState.connection; const isReconnecting = connection.phase === "connecting" || connection.phase === "reconnecting"; if (isReconnecting) { items.push({ - id: `environment-unavailable:${activeEnvironmentUnavailableState.environmentId}`, + id: `environment-unavailable:${activeEnvironmentActionUnavailableState.environmentId}`, variant: "info", icon: , title: connection.phase === "connecting" ? "Connecting…" : "Reconnecting…", @@ -1853,10 +1865,10 @@ function ChatViewContent(props: ChatViewProps) { }); } else { items.push({ - id: `environment-unavailable:${activeEnvironmentUnavailableState.environmentId}`, + id: `environment-unavailable:${activeEnvironmentActionUnavailableState.environmentId}`, variant: connection.phase === "error" ? "error" : "warning", icon: , - title: `${activeEnvironmentUnavailableState.label}: ${connectionStatusTitle(connection)}`, + title: `${activeEnvironmentActionUnavailableState.label}: ${connectionStatusTitle(connection)}`, description: connection.error ?? "Reconnect this environment before sending messages or running actions.", @@ -1866,7 +1878,7 @@ function ChatViewContent(props: ChatViewProps) { size="xs" onClick={() => void handleReconnectActiveEnvironment( - activeEnvironmentUnavailableState.environmentId, + activeEnvironmentActionUnavailableState.environmentId, ) } > @@ -1922,7 +1934,7 @@ function ChatViewContent(props: ChatViewProps) { } return items; }, [ - activeEnvironmentUnavailableState, + activeEnvironmentActionUnavailableState, handleReconnectActiveEnvironment, navigate, setDismissedVersionMismatchKey, @@ -4413,7 +4425,7 @@ function ChatViewContent(props: ChatViewProps) { const localApi = readLocalApi(); if (!localApi || !activeThread || isRevertingCheckpoint) return; - if (activeEnvironmentUnavailable && activeEnvironmentUnavailableLabel) { + if (activeEnvironmentActionUnavailable && activeEnvironmentUnavailableLabel) { setThreadError( activeThread.id, `Reconnect ${activeEnvironmentUnavailableLabel} before reverting checkpoints.`, @@ -4455,7 +4467,7 @@ function ChatViewContent(props: ChatViewProps) { }, [ activeThread, - activeEnvironmentUnavailable, + activeEnvironmentActionUnavailable, activeEnvironmentUnavailableLabel, environmentId, isConnecting, @@ -4473,7 +4485,7 @@ function ChatViewContent(props: ChatViewProps) { !activeThread || isSendBusy || isConnecting || - activeEnvironmentUnavailable || + activeEnvironmentActionUnavailable || sendInFlightRef.current ) return; @@ -5060,6 +5072,7 @@ function ChatViewContent(props: ChatViewProps) { !isServerThread || isSendBusy || isConnecting || + activeEnvironmentActionUnavailable || sendInFlightRef.current ) { return; @@ -5202,6 +5215,7 @@ function ChatViewContent(props: ChatViewProps) { [ activeThread, activeProposedPlan, + activeEnvironmentActionUnavailable, beginLocalDispatch, isConnecting, isSendBusy, @@ -5227,7 +5241,7 @@ function ChatViewContent(props: ChatViewProps) { !isServerThread || isSendBusy || isConnecting || - activeEnvironmentUnavailable || + activeEnvironmentActionUnavailable || sendInFlightRef.current ) { return; @@ -5364,7 +5378,7 @@ function ChatViewContent(props: ChatViewProps) { activeThreadBranch, activeThread, beginLocalDispatch, - activeEnvironmentUnavailable, + activeEnvironmentActionUnavailable, createThread, deleteThread, isConnecting, @@ -5848,7 +5862,7 @@ function ChatViewContent(props: ChatViewProps) { isConnecting={isConnecting} isSendBusy={isSendBusy} isPreparingWorktree={isPreparingWorktree} - environmentUnavailable={activeEnvironmentUnavailableState} + environmentUnavailable={activeEnvironmentActionUnavailableState} activePendingApproval={activePendingApproval} pendingApprovals={pendingApprovals} pendingUserInputs={pendingUserInputs} diff --git a/packages/client-runtime/src/operations/commands.ts b/packages/client-runtime/src/operations/commands.ts index ad25d6544dc1..ca6ddc718531 100644 --- a/packages/client-runtime/src/operations/commands.ts +++ b/packages/client-runtime/src/operations/commands.ts @@ -12,7 +12,7 @@ import { type EnvironmentRpcFailure, type EnvironmentRpcSuccess, type EnvironmentRpcUnavailableError, - request, + requestWhenConnected, } from "../rpc/client.ts"; type CommandType = ClientOrchestrationCommand["type"]; @@ -80,7 +80,7 @@ function timestampedCommandMetadata(input: { } function dispatch(command: ClientOrchestrationCommand) { - return request(ORCHESTRATION_WS_METHODS.dispatchCommand, command); + return requestWhenConnected(ORCHESTRATION_WS_METHODS.dispatchCommand, command); } export const createProject: (input: CreateProjectInput) => CommandEffect = Effect.fn( diff --git a/packages/client-runtime/src/rpc/client.test.ts b/packages/client-runtime/src/rpc/client.test.ts index 507d137caccb..fa75bf573354 100644 --- a/packages/client-runtime/src/rpc/client.test.ts +++ b/packages/client-runtime/src/rpc/client.test.ts @@ -25,7 +25,13 @@ import { import * as EnvironmentSupervisor from "../connection/supervisor.ts"; import * as RpcSession from "../rpc/session.ts"; import type { WsRpcProtocolClient } from "../rpc/protocol.ts"; -import { EnvironmentRpcRequestObserver, request, runStream, subscribe } from "./client.ts"; +import { + EnvironmentRpcRequestObserver, + request, + requestWhenConnected, + runStream, + subscribe, +} from "./client.ts"; const TARGET = new PrimaryConnectionTarget({ environmentId: EnvironmentId.make("environment-1"), @@ -77,6 +83,43 @@ const makeHarness = Effect.fn("TestEnvironmentRpc.makeHarness")(function* () { }); describe("environment RPC", () => { + it.effect("defers a primary command until the reconnecting session is available", () => + Effect.gen(function* () { + const requests: string[] = []; + const client = { + [WS_METHODS.cloudGetRelayClientStatus]: () => + Effect.sync(() => { + requests.push("dispatched"); + return { status: "available" as const, version: "2026.8.0" }; + }), + } as unknown as WsRpcProtocolClient; + const { activeSession, supervisor } = yield* makeHarness(); + const reconnectingState: SupervisorConnectionState = { + ...AVAILABLE_CONNECTION_STATE, + desired: true, + phase: "backoff", + }; + yield* SubscriptionRef.set(supervisor.state, reconnectingState); + + const requestFiber = yield* requestWhenConnected( + WS_METHODS.cloudGetRelayClientStatus, + {}, + ).pipe( + Effect.provideService(EnvironmentSupervisor.EnvironmentSupervisor, supervisor), + Effect.forkChild, + ); + yield* Effect.yieldNow; + expect(requests).toEqual([]); + + yield* SubscriptionRef.set(activeSession, Option.some(session(client))); + expect(yield* Fiber.join(requestFiber)).toEqual({ + status: "available", + version: "2026.8.0", + }); + expect(requests).toEqual(["dispatched"]); + }), + ); + 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 c7c928b3c954..c4f52e093b76 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -86,24 +86,63 @@ export type EnvironmentRpcStreamFailure = ? E : never; -const currentSession = Effect.fn("EnvironmentRpc.currentSession")(function* () { - const supervisor = yield* EnvironmentSupervisor; +function unavailableError(supervisor: EnvironmentSupervisor["Service"]) { + return new EnvironmentRpcUnavailableError({ + environmentId: supervisor.target.environmentId, + message: `${supervisor.target.label} is not connected.`, + }); +} + +const currentSessionFor = Effect.fn("EnvironmentRpc.currentSessionFor")(function* ( + supervisor: EnvironmentSupervisor["Service"], +) { return yield* SubscriptionRef.get(supervisor.session).pipe( Effect.flatMap( Option.match({ - onNone: () => - Effect.fail( - new EnvironmentRpcUnavailableError({ - environmentId: supervisor.target.environmentId, - message: `${supervisor.target.label} is not connected.`, - }), - ), + onNone: () => Effect.fail(unavailableError(supervisor)), onSome: Effect.succeed, }), ), ); }); +const currentSession = Effect.fn("EnvironmentRpc.currentSession")(function* () { + return yield* currentSessionFor(yield* EnvironmentSupervisor); +}); + +const dispatchSession = Effect.fn("EnvironmentRpc.dispatchSession")(function* ( + supervisor: EnvironmentSupervisor["Service"], +) { + const activeSession = yield* SubscriptionRef.get(supervisor.session); + if (Option.isSome(activeSession)) return activeSession.value; + + const state = yield* SubscriptionRef.get(supervisor.state); + const canReconnectPrimary = + supervisor.target._tag === "PrimaryConnectionTarget" && + ["connecting", "backoff", "connected"].includes(state.phase); + if (!canReconnectPrimary) return yield* unavailableError(supervisor); + + type DispatchSessionEvent = + | { readonly _tag: "Session"; readonly session: RpcSession } + | { readonly _tag: "Unavailable" }; + + const next = yield* Stream.merge( + SubscriptionRef.changes(supervisor.session).pipe( + Stream.filter(Option.isSome), + Stream.map((session): DispatchSessionEvent => ({ _tag: "Session", session: session.value })), + ), + SubscriptionRef.changes(supervisor.state).pipe( + Stream.filter((connection) => ["available", "offline", "blocked"].includes(connection.phase)), + Stream.map((): DispatchSessionEvent => ({ _tag: "Unavailable" })), + ), + ).pipe(Stream.runHead); + + if (Option.isNone(next) || next.value._tag === "Unavailable") { + return yield* unavailableError(supervisor); + } + return next.value.session; +}); + export const request = Effect.fn("EnvironmentRpc.request")(function* < TTag extends EnvironmentUnaryRpcTag, >(tag: TTag, input: EnvironmentRpcInput) { @@ -124,6 +163,26 @@ export const request = Effect.fn("EnvironmentRpc.request")(function* < return yield* method(input).pipe(Effect.ensuring(completeObservation)); }); +export const requestWhenConnected = Effect.fn("EnvironmentRpc.requestWhenConnected")(function* < + TTag extends EnvironmentUnaryRpcTag, +>(tag: TTag, input: EnvironmentRpcInput) { + const supervisor = yield* EnvironmentSupervisor; + yield* Effect.annotateCurrentSpan({ + "environment.id": supervisor.target.environmentId, + "rpc.method": tag, + }); + const session = yield* dispatchSession(supervisor); + const observer = yield* EnvironmentRpcRequestObserver; + const method = session.client[tag] as ( + input: EnvironmentRpcInput, + ) => Effect.Effect, EnvironmentRpcFailure>; + const completeObservation = yield* observer.observe({ + environmentId: supervisor.target.environmentId, + method: tag, + }); + return yield* method(input).pipe(Effect.ensuring(completeObservation)); +}); + export function runStream( tag: TTag, input: EnvironmentRpcInput,