From 52e1c9d178e9dc3476a90467488d4658a2de366f Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 2 Oct 2026 21:41:33 -0700 Subject: [PATCH 01/12] fix(server): subagent threads stop publishing tombstones to the relay (#15016) Co-authored-by: Claude Opus 5.5 (1M context) (cherry picked from commit d668ffcd64d74daba16520cbbfb0fdb431983eac) Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/relay/AgentAwarenessRelay.test.ts | 24 +++++++++++++++++++ apps/server/src/relay/AgentAwarenessRelay.ts | 10 ++++++++ 2 files changed, 34 insertions(+) diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index 00681f223..8b2ee4653 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -715,6 +715,30 @@ describe("AgentAwarenessRelay", () => { }), ); + it.effect.each([ + { label: "live", archived: false }, + { label: "archived", archived: true }, + ])("never publishes tombstones for $label subagent threads", ({ archived }) => + Effect.gen(function* () { + const { relay, currentShell, publications } = yield* makeTestRelay(); + yield* Ref.set( + currentShell, + shell({ + lineage: { + rootThreadId: THREAD_ID, + parentThreadId: THREAD_ID, + relationshipToParent: "subagent", + }, + ...(archived ? { archivedAt: yield* DateTime.now } : {}), + }), + ); + yield* relay.publishThread(THREAD_ID); + yield* TestClock.adjust("5 seconds"); + yield* relay.drain; + assert.equal(publications.length, 0); + }), + ); + it.effect("confirms a first completed state and respects disabling during confirmation", () => Effect.gen(function* () { const { relay, secrets, currentShell, publications } = yield* makeTestRelay(); diff --git a/apps/server/src/relay/AgentAwarenessRelay.ts b/apps/server/src/relay/AgentAwarenessRelay.ts index c708377e2..e9ce611c7 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.ts @@ -551,6 +551,16 @@ export const make = Effect.gen(function* () { // domain event, so materializing the full shell here would make the cost // of one thread's activity proportional to how many threads exist. const threadShell = yield* threads.getThreadShell(threadId); + if ( + threadShell?.lineage.relationshipToParent === "subagent" && + !(yield* Ref.get(publishedStateByThreadRef)).has(threadId) + ) { + // Subagents never project activity, so the relay holds no row to clear. + // Their events would otherwise publish a tombstone each, and every + // publish re-delivers the user's aggregate. Checked before the archive + // filter so archiving one stays quiet too. + return; + } const thread = threadShell === null || threadShell.archivedAt !== null ? Option.none() From 0f9d3e45f95442bef809bb7b571b2dd58baa5206 Mon Sep 17 00:00:00 2001 From: Dara Adedeji <76637177+SunkenInTime@users.noreply.github.com> Date: Fri, 2 Oct 2026 17:07:18 -0400 Subject: [PATCH 02/12] fix(server): idle shells stop reloading every project once a minute (#14893) Co-authored-by: Claude Opus 5.5 (1M context) (cherry picked from commit 8bc40b4e07bb7b4b0f71876d59520360c9bf958c) Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/ws.ts | 37 ++++++++++++++++++++++++++++--------- 1 file changed, 28 insertions(+), 9 deletions(-) diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index fae3441f6..dfc7a801a 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -971,16 +971,35 @@ export const subscribeOrchestrationV2Shell = Effect.fn("ws.orchestrationV2.subsc const enrichmentRefreshes = Stream.fromSubscription(enrichmentChanges).pipe( Stream.filter((change) => change.repositoryIdentityResolved), Stream.groupedWithin(64, Duration.millis(25)), + // Build the refresh from the identities the changes carry. Re-enriching + // every project here re-requested each expired root, whose resolution + // published again, so one expiry kept every subscriber reloading every + // project's metadata once a minute. Stream.mapEffect((changes) => - applicationEvents.latestApplicationSequence.pipe( - Effect.flatMap(loadProjectMetadataSnapshot), - Effect.map(({ snapshot }) => - shellStreamItemFromEnrichmentRefresh({ - snapshot, - changes: Array.from(changes), - }), - ), - ), + Effect.gen(function* () { + const identities = new Map( + Array.from(changes, (change) => [ + change.workspaceRoot, + change.enrichment.repositoryIdentity, + ]), + ); + const snapshotSequence = yield* applicationEvents.latestApplicationSequence; + const changedProjects = (yield* projects.listShells()).flatMap((project) => + identities.has(project.workspaceRoot) + ? [{ ...project, repositoryIdentity: identities.get(project.workspaceRoot) ?? null }] + : [], + ); + return shellStreamItemFromEnrichmentRefresh({ + snapshot: { + schemaVersion: ORCHESTRATION_V2_PROJECTION_SCHEMA_VERSION, + snapshotSequence, + projects: changedProjects, + threads: [], + archivedThreads: [], + } as OrchestrationV2ShellSnapshot, + changes: Array.from(changes), + }); + }), ), ); From 564d451e1922cf111172cb524f46d119c45a21f3 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 2 Oct 2026 15:29:45 -0700 Subject: [PATCH 03/12] fix(settings): Pi and ACP Registry no longer show an Early Access badge (#14915) Co-authored-by: Claude Opus 5.5 (1M context) (cherry picked from commit a7b3ce8c0896d123a7c3f02c586ff8193841cc09) Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/provider/Layers/PiProvider.ts | 1 - apps/web/src/components/settings/providerDriverMeta.ts | 2 -- packages/contracts/src/settings.ts | 2 +- 3 files changed, 1 insertion(+), 4 deletions(-) diff --git a/apps/server/src/provider/Layers/PiProvider.ts b/apps/server/src/provider/Layers/PiProvider.ts index 3f87f63f8..f4a8954e6 100644 --- a/apps/server/src/provider/Layers/PiProvider.ts +++ b/apps/server/src/provider/Layers/PiProvider.ts @@ -59,7 +59,6 @@ import { const PI_PRESENTATION = { displayName: "Pi", - badgeLabel: "Early Access", showInteractionModeToggle: false, supportedRuntimeModes: ["approval-required", "auto-accept-edits", "full-access"], // The adapter reports context usage from Pi's streaming usage while a diff --git a/apps/web/src/components/settings/providerDriverMeta.ts b/apps/web/src/components/settings/providerDriverMeta.ts index 316a66e0a..7496e5c2d 100644 --- a/apps/web/src/components/settings/providerDriverMeta.ts +++ b/apps/web/src/components/settings/providerDriverMeta.ts @@ -99,13 +99,11 @@ const PROVIDER_CLIENT_DEFINITIONS: readonly ProviderClientDefinition[] = [ { value: ProviderDriverKind.make("pi"), label: "Pi", - badgeLabel: "Early Access", settingsSchema: PiSettings, }, { value: ProviderDriverKind.make("acpRegistry"), label: "ACP Registry", - badgeLabel: "Early Access", settingsSchema: AcpRegistrySettings, hasDefaultInstance: false, }, diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index 6adc163bc..481284b25 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -990,7 +990,7 @@ export type AntigravitySettings = typeof AntigravitySettings.Type; export const PiSettings = makeProviderSettingsSchema( { - // Disabled by default while Pi support is Early Access. + // Off by default like Cursor and Grok. Users opt in from Settings. enabled: Schema.Boolean.pipe( Schema.withDecodingDefault(Effect.succeed(false)), Schema.annotateKey({ providerSettingsForm: { hidden: true } }), From 78b7421403bca1b0ac85de28ed98c9b54761ede4 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Fri, 2 Oct 2026 21:25:44 -0700 Subject: [PATCH 04/12] fix(web): outdated servers can be updated even when the client can't connect (#15002) Co-authored-by: Claude Opus 5.5 (1M context) (cherry picked from commit 22e9d35613305a04a28e96f83abc10c1bcf4aa8a) Co-Authored-By: Claude Opus 5.5 (1M context) --- .../web/src/components/ServerUpdateAction.tsx | 55 +++- .../settings/ConnectionsSettings.tsx | 13 + apps/web/src/state/server.ts | 3 + packages/client-runtime/package.json | 4 + .../client-runtime/src/connection/catalog.ts | 2 + .../src/connection/compatibility.test.ts | 24 ++ .../src/connection/compatibility.ts | 26 +- .../client-runtime/src/connection/index.ts | 2 + .../client-runtime/src/connection/layer.ts | 2 + .../client-runtime/src/connection/model.ts | 2 + .../src/connection/onboarding.test.ts | 52 +++- .../src/connection/onboarding.ts | 5 +- .../src/connection/outdatedHostUpdate.test.ts | 278 ++++++++++++++++++ .../src/connection/outdatedHostUpdate.ts | 190 ++++++++++++ .../client-runtime/src/connection/registry.ts | 36 ++- .../client-runtime/src/connection/resolver.ts | 27 +- .../src/state/outdatedServerUpdate.ts | 78 +++++ packages/client-runtime/src/state/server.ts | 5 +- 18 files changed, 784 insertions(+), 20 deletions(-) create mode 100644 packages/client-runtime/src/connection/outdatedHostUpdate.test.ts create mode 100644 packages/client-runtime/src/connection/outdatedHostUpdate.ts create mode 100644 packages/client-runtime/src/state/outdatedServerUpdate.ts diff --git a/apps/web/src/components/ServerUpdateAction.tsx b/apps/web/src/components/ServerUpdateAction.tsx index fa7a9a2c9..f4ae391d7 100644 --- a/apps/web/src/components/ServerUpdateAction.tsx +++ b/apps/web/src/components/ServerUpdateAction.tsx @@ -10,7 +10,7 @@ import { type ComponentProps, useRef, useState } from "react"; import { requestConfirmDialog } from "~/confirmDialog"; import { useCopyToClipboard } from "~/hooks/useCopyToClipboard"; import { useEnvironmentSettings } from "~/hooks/useSettings"; -import { serverEnvironment } from "~/state/server"; +import { serverEnvironment, updateOutdatedServer } from "~/state/server"; import { useAtomCommand } from "~/state/use-atom-command"; import { manualServerUpdateCommand } from "~/versionSkew"; import { Button } from "./ui/button"; @@ -292,3 +292,56 @@ export function ServerUpdateAction({ ); } + +/** + * Updates a host too old for this client to connect to. Its version comes + * from the host descriptor because the host never delivers a server config. + */ +export function OutdatedServerUpdateAction({ + environmentId, + serverLabel, + fromVersion, + targetVersion, + label = "Update", +}: { + readonly environmentId: EnvironmentId; + readonly serverLabel: string; + readonly fromVersion: string | undefined; + readonly targetVersion: string; + readonly label?: string; +}) { + const update = useAtomCommand(updateOutdatedServer, { reportFailure: false }); + const handleUpdate = async () => { + if (pendingUpdateEnvironmentIds.has(environmentId)) return; + pendingUpdateEnvironmentIds.add(environmentId); + try { + const result = await update({ + environmentId, + input: { targetVersion }, + ...(fromVersion === undefined ? {} : { fromVersion }), + }); + if (result._tag === "Failure") { + if (isAtomCommandInterrupted(result)) return; + throw squashAtomCommandFailure(result); + } + toastManager.add({ + type: "success", + title: `${serverLabel} updated`, + description: `Reconnected on t3@${result.value.targetVersion}.`, + }); + } catch (error) { + toastManager.add({ + type: "error", + title: "Server update failed", + description: updateFailureMessage(error), + }); + } finally { + pendingUpdateEnvironmentIds.delete(environmentId); + } + }; + return ( + + ); +} diff --git a/apps/web/src/components/settings/ConnectionsSettings.tsx b/apps/web/src/components/settings/ConnectionsSettings.tsx index 95c4e93b6..83f98b293 100644 --- a/apps/web/src/components/settings/ConnectionsSettings.tsx +++ b/apps/web/src/components/settings/ConnectionsSettings.tsx @@ -156,11 +156,13 @@ import { usePrimaryEnvironment, useRelayEnvironmentDiscovery, } from "~/state/environments"; +import { APP_VERSION } from "~/branding"; import { requestConfirmDialog } from "~/confirmDialog"; import { useAtomCommand } from "../../state/use-atom-command"; import { primaryServerKeybindingsAtom, serverEnvironment } from "~/state/server"; import { ConnectionStatusDot } from "../ConnectionStatusDot"; import { + OutdatedServerUpdateAction, ServerUpdateAction, ServerUpdateProgress, ServerUpdatesAction, @@ -1609,6 +1611,17 @@ function SavedBackendListRow({ Check again ) : null} + {unsupported && + environment.entry.serverUpdateRequired === true && + serverUpdateState.status !== "running" ? ( + + ) : null} {showUpdateAction ? ( ()( diff --git a/packages/client-runtime/src/connection/compatibility.test.ts b/packages/client-runtime/src/connection/compatibility.test.ts index e6a2cf3ca..3f9dc8f99 100644 --- a/packages/client-runtime/src/connection/compatibility.test.ts +++ b/packages/client-runtime/src/connection/compatibility.test.ts @@ -50,5 +50,29 @@ describe("orchestration protocol compatibility", () => { ); expect(error).toMatchObject({ reason: "unsupported" }); expect(error?.message).toContain("This client is not supported"); + expect(error).not.toHaveProperty("serverUpdateRequired"); + }); + + it("offers a remote update only for an older host that can update itself", () => { + const older = descriptor(ORCHESTRATION_PROTOCOL_VERSION - 1); + const withCapabilities = (capabilities: ExecutionEnvironmentDescriptor["capabilities"]) => + orchestrationProtocolCompatibilityError({ ...older, capabilities }); + + expect( + withCapabilities({ repositoryIdentity: true, serverSelfUpdate: "boot-service" }), + ).toMatchObject({ serverUpdateRequired: true }); + expect(withCapabilities({ repositoryIdentity: true })).not.toHaveProperty( + "serverUpdateRequired", + ); + expect( + withCapabilities({ repositoryIdentity: true, serverSelfUpdate: "desktop-managed" }), + ).not.toHaveProperty("serverUpdateRequired"); + expect( + withCapabilities({ + repositoryIdentity: true, + serverSelfUpdate: "desktop-managed", + desktopAppUpdate: true, + }), + ).toMatchObject({ serverUpdateRequired: true }); }); }); diff --git a/packages/client-runtime/src/connection/compatibility.ts b/packages/client-runtime/src/connection/compatibility.ts index a99320bde..e9c66b642 100644 --- a/packages/client-runtime/src/connection/compatibility.ts +++ b/packages/client-runtime/src/connection/compatibility.ts @@ -14,13 +14,25 @@ export function orchestrationProtocolCompatibilityError( if (serverProtocolVersion === ORCHESTRATION_PROTOCOL_VERSION) { return null; } - return new ConnectionBlockedError({ - reason: "unsupported", - detail: - serverProtocolVersion > ORCHESTRATION_PROTOCOL_VERSION - ? `This client is not supported by this server. Update your app or use a compatible release to connect to ${descriptor.label}.` - : `This client requires a newer server. Update Pylon on ${descriptor.label} to connect.`, - }); + return serverProtocolVersion > ORCHESTRATION_PROTOCOL_VERSION + ? new ConnectionBlockedError({ + reason: "unsupported", + detail: `This client is not supported by this server. Update your app or use a compatible release to connect to ${descriptor.label}.`, + }) + : new ConnectionBlockedError({ + reason: "unsupported", + detail: `This client requires a newer server. Update Pylon on ${descriptor.label} to connect.`, + ...(canSelfUpdate(descriptor) ? { serverUpdateRequired: true } : {}), + }); +} + +/** Whether this client can drive the host's update remotely. */ +function canSelfUpdate(descriptor: ExecutionEnvironmentDescriptor): boolean { + const { serverSelfUpdate, desktopAppUpdate } = descriptor.capabilities; + return ( + serverSelfUpdate !== undefined && + (serverSelfUpdate !== "desktop-managed" || desktopAppUpdate === true) + ); } export function appendOrchestrationProtocol(socketUrl: string): string { diff --git a/packages/client-runtime/src/connection/index.ts b/packages/client-runtime/src/connection/index.ts index 14ba95763..c3c10f006 100644 --- a/packages/client-runtime/src/connection/index.ts +++ b/packages/client-runtime/src/connection/index.ts @@ -15,3 +15,5 @@ export * as EnvironmentSupervisor from "./supervisor.ts"; export * as Wakeups from "./wakeups.ts"; export { orchestrationProtocolCompatibilityError } from "./compatibility.ts"; +// Flat so consumers' inferred command types can name it. +export { OutdatedHostUpdateError } from "./outdatedHostUpdate.ts"; diff --git a/packages/client-runtime/src/connection/layer.ts b/packages/client-runtime/src/connection/layer.ts index 25d834484..8abeca62d 100644 --- a/packages/client-runtime/src/connection/layer.ts +++ b/packages/client-runtime/src/connection/layer.ts @@ -79,6 +79,8 @@ export function layerWithOptions(options: RpcSession.RpcSessionOptions) { registryLayer, RelayEnvironmentDiscovery.layer, onboardingLayer, + // Exposed for updating hosts too old to connect through the driver. + ConnectionResolver.layer, ); const connectionStartupLayer = Layer.effectDiscard( Effect.gen(function* () { diff --git a/packages/client-runtime/src/connection/model.ts b/packages/client-runtime/src/connection/model.ts index 2d5dcbba9..e9145ffba 100644 --- a/packages/client-runtime/src/connection/model.ts +++ b/packages/client-runtime/src/connection/model.ts @@ -94,6 +94,8 @@ export class ConnectionBlockedError extends Schema.TaggedError, - options?: { readonly failDescriptor?: boolean; readonly protocolVersion?: number }, + options?: { + readonly failDescriptor?: boolean; + readonly protocolVersion?: number; + readonly selfUpdate?: boolean; + }, ) { const fetchFn = ((input, init = {}) => { const url = String(input); @@ -56,6 +60,7 @@ function pairingHttpLayer( orchestrationProtocolVersion: options?.protocolVersion ?? ORCHESTRATION_PROTOCOL_VERSION, capabilities: { repositoryIdentity: true, + ...(options?.selfUpdate === true ? { serverSelfUpdate: "boot-service" } : {}), }, }), ); @@ -145,6 +150,51 @@ describe("connection onboarding", () => { }), ); + it.effect("pairs an outdated server so it can be updated from this client", () => + Effect.gen(function* () { + const calls: Array<{ readonly url: string; readonly init: RequestInit }> = []; + const registration = yield* preparePairingRegistration({ + host: "remote.example.test", + pairingCode: "pairing-token", + }).pipe( + Effect.provide( + Layer.mergeAll( + CLIENT_PRESENTATION_LAYER, + pairingHttpLayer(calls, { + protocolVersion: ORCHESTRATION_PROTOCOL_VERSION - 1, + selfUpdate: true, + }), + ), + ), + ); + expect(registration.target.environmentId).toBe("environment-paired"); + expect(calls.map((call) => call.url)).toContain("https://remote.example.test/oauth/token"); + }), + ); + + it.effect("refuses an outdated server that cannot update itself", () => + Effect.gen(function* () { + const calls: Array<{ readonly url: string; readonly init: RequestInit }> = []; + const error = yield* preparePairingRegistration({ + host: "remote.example.test", + pairingCode: "pairing-token", + }).pipe( + Effect.provide( + Layer.mergeAll( + CLIENT_PRESENTATION_LAYER, + pairingHttpLayer(calls, { protocolVersion: ORCHESTRATION_PROTOCOL_VERSION - 1 }), + ), + ), + Effect.flip, + ); + expect(error).toMatchObject({ reason: "unsupported" }); + expect(error).not.toHaveProperty("serverUpdateRequired"); + expect(calls.map((call) => call.url)).toEqual([ + "https://remote.example.test/.well-known/t3/environment", + ]); + }), + ); + it.effect("does not consume a pairing credential when descriptor discovery fails", () => Effect.gen(function* () { const calls: Array<{ readonly url: string; readonly init: RequestInit }> = []; diff --git a/packages/client-runtime/src/connection/onboarding.ts b/packages/client-runtime/src/connection/onboarding.ts index 24c03addf..fec7847b5 100644 --- a/packages/client-runtime/src/connection/onboarding.ts +++ b/packages/client-runtime/src/connection/onboarding.ts @@ -93,7 +93,10 @@ export const preparePairingRegistration = Effect.fn( httpBaseUrl: target.httpBaseUrl, }).pipe(Effect.mapError(mapRemoteEnvironmentError)); const compatibilityError = orchestrationProtocolCompatibilityError(descriptor); - if (compatibilityError !== null) return yield* compatibilityError; + // An outdated server is still saved so it can be updated from this client. + if (compatibilityError !== null && compatibilityError.serverUpdateRequired !== true) { + return yield* compatibilityError; + } const access = yield* bootstrapRemoteBearerSession({ httpBaseUrl: target.httpBaseUrl, credential: target.credential, diff --git a/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts b/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts new file mode 100644 index 000000000..e63d095fa --- /dev/null +++ b/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts @@ -0,0 +1,278 @@ +import { + EnvironmentId, + ORCHESTRATION_PROTOCOL_VERSION, + ExecutionEnvironmentDescriptor, + WS_METHODS, +} from "@t3tools/contracts"; +import { describe, expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; +import * as SubscriptionRef from "effect/SubscriptionRef"; +import * as HttpClient from "effect/unstable/http/HttpClient"; +import * as HttpClientResponse from "effect/unstable/http/HttpClientResponse"; +import * as Socket from "effect/unstable/socket/Socket"; + +import type { ConnectionCatalogEntry } from "./catalog.ts"; +import { orchestrationProtocolCompatibilityError } from "./compatibility.ts"; +import { PrimaryConnectionTarget } from "./model.ts"; +import { updateOutdatedHost } from "./outdatedHostUpdate.ts"; +import * as EnvironmentRegistry from "./registry.ts"; +import * as ConnectionResolver from "./resolver.ts"; +import * as RelayEnvironmentDiscovery from "../relay/discovery.ts"; + +const TARGET = new PrimaryConnectionTarget({ + environmentId: EnvironmentId.make("environment-old"), + label: "Build Mac", + httpBaseUrl: "https://build.example.test", + wsBaseUrl: "wss://build.example.test", +}); + +const descriptor = (protocol: number | undefined, serverVersion: string) => + ({ + environmentId: TARGET.environmentId, + label: TARGET.label, + platform: { os: "darwin", arch: "arm64" }, + serverVersion, + ...(protocol === undefined ? {} : { orchestrationProtocolVersion: protocol }), + capabilities: { repositoryIdentity: true, serverSelfUpdate: "boot-service" }, + }) satisfies ExecutionEnvironmentDescriptor; + +const RpcRequest = Schema.TaggedStruct("Request", { + id: Schema.Union([Schema.String, Schema.Number]), + payload: Schema.Unknown, + tag: Schema.String, +}); +const isRpcRequest = Schema.is(RpcRequest); +const decodeJson = Schema.decodeUnknownSync(Schema.fromJsonString(Schema.Unknown)); +const encodeJson = Schema.encodeUnknownSync(Schema.fromJsonString(Schema.Unknown)); +const encodeDescriptor = Schema.encodeSync(Schema.fromJsonString(ExecutionEnvironmentDescriptor)); + +type Listener = (event: { readonly type: string; readonly data?: unknown }) => void; + +/** Answers update RPCs the way a protocol-1 server does. */ +class OutdatedHostSocket { + static readonly OPEN = 1; + readyState = 0; + readonly requests: Array = []; + private readonly listeners = new Map>(); + + readonly url: string; + private readonly onUpdate: () => void; + + constructor(url: string, onUpdate: () => void) { + this.url = url; + this.onUpdate = onUpdate; + queueMicrotask(() => { + this.readyState = OutdatedHostSocket.OPEN; + this.emit({ type: "open" }); + }); + } + + addEventListener(type: string, listener: Listener) { + const listeners = this.listeners.get(type) ?? new Set(); + listeners.add(listener); + this.listeners.set(type, listeners); + } + + removeEventListener(type: string, listener: Listener) { + this.listeners.get(type)?.delete(listener); + } + + send(data: string) { + const message = decodeJson(data); + if (!isRpcRequest(message)) return; + this.requests.push(message); + if (message.tag !== WS_METHODS.serverUpdateServer) return; + this.onUpdate(); + queueMicrotask(() => + this.emit({ + type: "message", + data: encodeJson({ + _tag: "Exit", + requestId: message.id, + exit: { + _tag: "Success", + value: { targetVersion: "0.0.46", method: "boot-service", updateId: "update-1" }, + }, + }), + }), + ); + } + + close() { + this.readyState = 3; + this.emit({ type: "close" }); + } + + private emit(event: { readonly type: string; readonly data?: unknown }) { + for (const listener of this.listeners.get(event.type) ?? []) listener(event); + } +} + +describe("updateOutdatedHost", () => { + it.effect("updates a protocol-1 host over a bare socket, then switches it back on", () => + Effect.gen(function* () { + const blocked = orchestrationProtocolCompatibilityError(descriptor(undefined, "0.0.45")); + expect(blocked).toMatchObject({ serverUpdateRequired: true }); + + const sockets: Array = []; + let served: ExecutionEnvironmentDescriptor = descriptor(undefined, "0.0.45"); + const entries = yield* SubscriptionRef.make< + ReadonlyMap + >( + new Map([ + [ + TARGET.environmentId, + { + target: TARGET, + profile: Option.none(), + enabled: false, + unsupportedReason: blocked?.message ?? "", + serverUpdateRequired: true, + }, + ], + ]), + ); + const calls: Array = []; + const registry = EnvironmentRegistry.EnvironmentRegistry.of({ + entries, + setCompatibility: (_environmentId: EnvironmentId, error: unknown) => + Effect.sync(() => calls.push(`compatibility:${error === null ? "clear" : "block"}`)), + setEnabled: (_environmentId: EnvironmentId, enabled: boolean) => + Effect.sync(() => calls.push(`enabled:${enabled}`)), + } as unknown as EnvironmentRegistry.EnvironmentRegistry["Service"]); + const resolver = ConnectionResolver.ConnectionResolver.of({ + prepare: () => Effect.die(new Error("The update must bypass the protocol gate.")), + prepareForUpdate: () => + Effect.sync(() => served).pipe( + Effect.map((current) => ({ + descriptor: current, + prepared: { + environmentId: TARGET.environmentId, + label: TARGET.label, + httpBaseUrl: TARGET.httpBaseUrl, + socketUrl: "wss://build.example.test/ws?wsTicket=ticket", + httpAuthorization: null, + target: TARGET, + }, + })), + ), + }); + const httpClient = HttpClient.make((request) => + Effect.sync(() => + HttpClientResponse.fromWeb(request, new Response(encodeDescriptor(served))), + ), + ); + + const result = yield* updateOutdatedHost( + TARGET.environmentId, + { targetVersion: "0.0.46" }, + () => Effect.void, + ).pipe( + Effect.provide( + Layer.mergeAll( + Layer.succeed(EnvironmentRegistry.EnvironmentRegistry, registry), + Layer.succeed(ConnectionResolver.ConnectionResolver, resolver), + Layer.succeed( + RelayEnvironmentDiscovery.RelayEnvironmentDiscovery, + RelayEnvironmentDiscovery.RelayEnvironmentDiscovery.of({ + state: yield* SubscriptionRef.make( + RelayEnvironmentDiscovery.EMPTY_RELAY_ENVIRONMENT_DISCOVERY_STATE, + ), + refresh: Effect.void, + }), + ), + Layer.succeed(HttpClient.HttpClient, httpClient), + Layer.succeed(Socket.WebSocketConstructor, (url) => { + // The host relaunches on a compatible protocol once the update lands. + const socket = new OutdatedHostSocket(url, () => { + served = descriptor(ORCHESTRATION_PROTOCOL_VERSION, "0.0.46"); + }); + sockets.push(socket); + return socket as unknown as globalThis.WebSocket; + }), + ), + ), + ); + + expect(sockets[0]?.url).not.toContain("orchestrationProtocol"); + expect(sockets[0]?.requests.map((request) => request.tag)).toEqual([ + WS_METHODS.serverUpdateServer, + ]); + expect(result.targetVersion).toBe("0.0.46"); + expect(calls).toEqual(["compatibility:clear", "enabled:true"]); + }), + ); + + it.effect("refuses a host that cannot update itself without opening a socket", () => + Effect.gen(function* () { + const manual = { + ...descriptor(undefined, "0.0.45"), + capabilities: { repositoryIdentity: true }, + }; + let opened = false; + const error = yield* Effect.flip( + updateOutdatedHost(TARGET.environmentId, { targetVersion: "0.0.46" }, () => Effect.void), + ).pipe( + Effect.provide( + Layer.mergeAll( + Layer.succeed( + EnvironmentRegistry.EnvironmentRegistry, + EnvironmentRegistry.EnvironmentRegistry.of({ + entries: yield* SubscriptionRef.make< + ReadonlyMap + >( + new Map([ + [ + TARGET.environmentId, + { target: TARGET, profile: Option.none(), enabled: false }, + ], + ]), + ), + } as unknown as EnvironmentRegistry.EnvironmentRegistry["Service"]), + ), + Layer.succeed( + ConnectionResolver.ConnectionResolver, + ConnectionResolver.ConnectionResolver.of({ + prepare: () => Effect.die(new Error("unused")), + prepareForUpdate: () => + Effect.succeed({ + descriptor: manual, + prepared: { + environmentId: TARGET.environmentId, + label: TARGET.label, + httpBaseUrl: TARGET.httpBaseUrl, + socketUrl: "wss://build.example.test/ws", + httpAuthorization: null, + target: TARGET, + }, + }), + }), + ), + Layer.succeed( + RelayEnvironmentDiscovery.RelayEnvironmentDiscovery, + RelayEnvironmentDiscovery.RelayEnvironmentDiscovery.of({ + state: yield* SubscriptionRef.make( + RelayEnvironmentDiscovery.EMPTY_RELAY_ENVIRONMENT_DISCOVERY_STATE, + ), + refresh: Effect.void, + }), + ), + Layer.succeed( + HttpClient.HttpClient, + HttpClient.make(() => Effect.die(new Error("unused"))), + ), + Layer.succeed(Socket.WebSocketConstructor, () => { + opened = true; + throw new Error("unused"); + }), + ), + ), + ); + expect(error).toMatchObject({ _tag: "OutdatedHostUpdateError" }); + expect(opened).toBe(false); + }), + ); +}); diff --git a/packages/client-runtime/src/connection/outdatedHostUpdate.ts b/packages/client-runtime/src/connection/outdatedHostUpdate.ts new file mode 100644 index 000000000..de2e496d1 --- /dev/null +++ b/packages/client-runtime/src/connection/outdatedHostUpdate.ts @@ -0,0 +1,190 @@ +import { + ORCHESTRATION_PROTOCOL_VERSION, + type EnvironmentId, + type ExecutionEnvironmentDescriptor, + type ServerSelfUpdateInput, + type ServerSelfUpdateResult, + WS_METHODS, +} from "@t3tools/contracts"; +import * as Duration from "effect/Duration"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as Ref from "effect/Ref"; +import * as Schedule from "effect/Schedule"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import * as SubscriptionRef from "effect/SubscriptionRef"; +import * as HttpClient from "effect/unstable/http/HttpClient"; +import * as RpcClient from "effect/unstable/rpc/RpcClient"; +import * as RpcSerialization from "effect/unstable/rpc/RpcSerialization"; +import * as Socket from "effect/unstable/socket/Socket"; + +import { fetchRemoteEnvironmentDescriptor } from "../environment/descriptor.ts"; +import { makeWsRpcProtocolClient } from "../rpc/protocol.ts"; +import { isLegacyUpdateHandoffLoss, resolveServerUpdateProgressResult } from "../state/server.ts"; +import * as RelayEnvironmentDiscovery from "../relay/discovery.ts"; +import * as ConnectionResolver from "./resolver.ts"; +import * as EnvironmentRegistry from "./registry.ts"; + +// A v1 host restarting into v2 runs migrations before its descriptor answers again. +const OUTDATED_HOST_RESTART_TIMEOUT = Duration.minutes(4); +const SOCKET_OPEN_TIMEOUT = "15 seconds"; + +export class OutdatedHostUpdateError extends Schema.TaggedError()( + "OutdatedHostUpdateError", + { + environmentId: Schema.String, + message: Schema.String, + }, +) {} + +export type OutdatedHostUpdateStage = "downloading" | "installing" | "resuming"; + +/** + * Updates a host whose orchestration protocol is too old for this client. + * + * The normal session refuses to open against such a host, so this opens a + * bare socket and calls only the self-update RPCs, whose wire shape has not + * changed across protocol versions. Once the host relaunches on a compatible + * protocol the environment is switched back on and connects normally. + */ +export const updateOutdatedHost = Effect.fn("clientRuntime.connection.updateOutdatedHost")( + function* ( + environmentId: EnvironmentId, + input: ServerSelfUpdateInput, + onStage: (stage: OutdatedHostUpdateStage) => Effect.Effect, + ) { + const registry = yield* EnvironmentRegistry.EnvironmentRegistry; + const resolver = yield* ConnectionResolver.ConnectionResolver; + const webSocketConstructor = yield* Socket.WebSocketConstructor; + const httpClient = yield* HttpClient.HttpClient; + const entry = (yield* SubscriptionRef.get(registry.entries)).get(environmentId); + if (entry === undefined) { + return yield* new EnvironmentRegistry.EnvironmentNotRegisteredError({ environmentId }); + } + const { prepared, descriptor } = yield* resolver.prepareForUpdate(entry); + const capabilities = descriptor.capabilities; + if ( + capabilities.serverSelfUpdate === undefined || + (capabilities.serverSelfUpdate === "desktop-managed" && + capabilities.desktopAppUpdate !== true) + ) { + return yield* new OutdatedHostUpdateError({ + environmentId, + message: `Update Pylon on ${descriptor.label} manually; it cannot update itself.`, + }); + } + + const result = yield* Effect.scoped( + Effect.gen(function* () { + const protocolContext = yield* Layer.build( + Layer.effect( + RpcClient.Protocol, + RpcClient.makeProtocolSocket({ + retryTransientErrors: false, + retryPolicy: Schedule.recurs(0), + }), + ).pipe( + Layer.provide( + Layer.mergeAll( + Socket.layerWebSocket(prepared.socketUrl, { + openTimeout: SOCKET_OPEN_TIMEOUT, + }).pipe( + Layer.provide(Layer.succeed(Socket.WebSocketConstructor, webSocketConstructor)), + ), + RpcSerialization.layerJson, + ), + ), + ), + ); + const client = yield* makeWsRpcProtocolClient.pipe(Effect.provide(protocolContext)); + + const updateResult: ServerSelfUpdateResult = + capabilities.serverSelfUpdateProgress === true + ? yield* Effect.gen(function* () { + const terminal = yield* Ref.make(Option.none()); + const streamExit = yield* client[WS_METHODS.serverUpdateServerWithProgress]( + input, + ).pipe( + Stream.runForEach((event) => + event.type === "complete" + ? Ref.set(terminal, Option.some(event.result)) + : onStage(event.stage), + ), + Effect.exit, + ); + return yield* resolveServerUpdateProgressResult( + input.targetVersion, + yield* Ref.get(terminal), + streamExit, + ); + }) + : yield* client[WS_METHODS.serverUpdateServer](input).pipe( + // Older servers can drop the socket before acknowledging the restart. + Effect.catchCauseIf( + (cause) => + (capabilities.serverSelfUpdate === "boot-service" || + capabilities.serverSelfUpdate === "respawn") && + isLegacyUpdateHandoffLoss(cause), + () => + Effect.succeed({ + targetVersion: input.targetVersion, + method: capabilities.serverSelfUpdate as "boot-service" | "respawn", + } satisfies ServerSelfUpdateResult), + ), + ); + + if ( + updateResult.method === "desktop-app" && + updateResult.desktopUpdateToken !== undefined + ) { + // The commit relaunches the desktop app, so a dropped socket is success. + yield* client[WS_METHODS.serverCommitDesktopUpdate]({ + requestId: updateResult.desktopUpdateToken, + }).pipe( + Effect.catchCauseIf( + (cause) => isLegacyUpdateHandoffLoss(cause), + () => Effect.void, + ), + ); + } + return updateResult; + }), + ); + + yield* onStage("resuming"); + const resumed = yield* fetchRemoteEnvironmentDescriptor({ + httpBaseUrl: prepared.httpBaseUrl, + }).pipe( + Effect.provideService(HttpClient.HttpClient, httpClient), + Effect.option, + Effect.repeat({ + schedule: Schedule.spaced(Duration.seconds(1)), + until: (current) => Option.exists(current, isCompatibleDescriptor), + }), + Effect.timeoutOption(OUTDATED_HOST_RESTART_TIMEOUT), + Effect.map(Option.flatten), + ); + if (Option.isNone(resumed)) { + return yield* new OutdatedHostUpdateError({ + environmentId, + message: `${descriptor.label} did not come back on a compatible Pylon version.`, + }); + } + + // Discovery still holds the old relay descriptor and would re-block the + // environment from it, so replace that before clearing the block. + if (entry.target._tag === "RelayConnectionTarget") { + const discovery = yield* RelayEnvironmentDiscovery.RelayEnvironmentDiscovery; + yield* discovery.refresh; + } + yield* registry.setCompatibility(environmentId, null); + yield* registry.setEnabled(environmentId, true); + return { ...result, targetVersion: resumed.value.serverVersion }; + }, +); + +function isCompatibleDescriptor(descriptor: ExecutionEnvironmentDescriptor): boolean { + return (descriptor.orchestrationProtocolVersion ?? 1) === ORCHESTRATION_PROTOCOL_VERSION; +} diff --git a/packages/client-runtime/src/connection/registry.ts b/packages/client-runtime/src/connection/registry.ts index 7724a12dc..07e54f1db 100644 --- a/packages/client-runtime/src/connection/registry.ts +++ b/packages/client-runtime/src/connection/registry.ts @@ -39,6 +39,17 @@ import * as ConnectionWakeups from "./wakeups.ts"; const isSshConnectionProfile = Schema.is(SshConnectionProfile); +function unsupportedState( + entry: ConnectionCatalogEntry, +): Pick { + return { + ...(entry.unsupportedReason === undefined + ? {} + : { unsupportedReason: entry.unsupportedReason }), + ...(entry.serverUpdateRequired === true ? { serverUpdateRequired: true } : {}), + }; +} + export class EnvironmentNotRegisteredError extends Schema.TaggedError()( "EnvironmentNotRegisteredError", { @@ -454,7 +465,7 @@ export const make = Effect.gen(function* () { ...(previous.unsupportedReason !== undefined && Equal.equals(previous.target, registered.target) && Equal.equals(previous.profile, registered.profile) - ? { unsupportedReason: previous.unsupportedReason } + ? unsupportedState(previous) : {}), }; yield* registrations.register(registration); @@ -480,7 +491,7 @@ export const make = Effect.gen(function* () { previous?.unsupportedReason !== undefined && Equal.equals(previous.target, registered.target) && Equal.equals(previous.profile, registered.profile) - ? { ...registered, enabled: false, unsupportedReason: previous.unsupportedReason } + ? { ...registered, enabled: false, ...unsupportedState(previous) } : registered; yield* Ref.update(platformEnvironmentIds, (current) => { const next = new Set(current); @@ -814,11 +825,26 @@ export const make = Effect.gen(function* () { // Discovery can race with a replacement registration for this ID. // Match the observed relay target under the lock that protects the write. if (options?.expectedRelayTarget && entry?.target !== options.expectedRelayTarget) return; - if (entry === undefined || entry.unsupportedReason === (error?.message ?? undefined)) + if ( + entry === undefined || + (entry.unsupportedReason === (error?.message ?? undefined) && + entry.serverUpdateRequired === (error?.serverUpdateRequired ?? undefined)) + ) return; - const { unsupportedReason: _previousReason, ...rest } = entry; + const { + unsupportedReason: _previousReason, + serverUpdateRequired: _previousUpdateRequired, + ...rest + } = entry; const next: ConnectionCatalogEntry = - error === null ? rest : { ...rest, enabled: false, unsupportedReason: error.message }; + error === null + ? rest + : { + ...rest, + enabled: false, + unsupportedReason: error.message, + ...(error.serverUpdateRequired === true ? { serverUpdateRequired: true } : {}), + }; if ( error !== null && entry.enabled && diff --git a/packages/client-runtime/src/connection/resolver.ts b/packages/client-runtime/src/connection/resolver.ts index 66d99fa60..0a9ab3359 100644 --- a/packages/client-runtime/src/connection/resolver.ts +++ b/packages/client-runtime/src/connection/resolver.ts @@ -1,4 +1,7 @@ -import type { AuthClientPresentationMetadata } from "@t3tools/contracts"; +import type { + AuthClientPresentationMetadata, + ExecutionEnvironmentDescriptor, +} from "@t3tools/contracts"; import { withRelayClientTracing } from "@t3tools/shared/relayTracing"; import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; @@ -45,6 +48,17 @@ export class ConnectionResolver extends Context.Service< readonly prepare: ( entry: ConnectionCatalogEntry, ) => Effect.Effect; + /** + * Authorizes a socket without the orchestration protocol gate, for hosts + * too old to connect normally. Only update RPCs may run over it. + */ + readonly prepareForUpdate: (entry: ConnectionCatalogEntry) => Effect.Effect< + { + readonly prepared: PreparedConnection; + readonly descriptor: ExecutionEnvironmentDescriptor; + }, + ConnectionAttemptError + >; } >()("@t3tools/client-runtime/connection/resolver/ConnectionResolver") {} @@ -232,7 +246,7 @@ export const make = Effect.gen(function* () { const ssh = yield* makeSshBroker(); const httpClient = yield* HttpClient.HttpClient; - const prepare = Effect.fn("clientRuntime.connection.broker.prepare")(function* ( + const authorize = Effect.fn("clientRuntime.connection.broker.authorize")(function* ( entry: ConnectionCatalogEntry, ) { const target: ConnectionTarget = entry.target; @@ -264,6 +278,13 @@ export const make = Effect.gen(function* () { actual: descriptor.environmentId, }); } + return { prepared, descriptor }; + }); + + const prepare = Effect.fn("clientRuntime.connection.broker.prepare")(function* ( + entry: ConnectionCatalogEntry, + ) { + const { prepared, descriptor } = yield* authorize(entry); const compatibilityError = orchestrationProtocolCompatibilityError(descriptor); if (compatibilityError !== null) { return yield* compatibilityError; @@ -274,7 +295,7 @@ export const make = Effect.gen(function* () { }; }); - return ConnectionResolver.of({ prepare }); + return ConnectionResolver.of({ prepare, prepareForUpdate: authorize }); }); export const layer = Layer.effect(ConnectionResolver, make); diff --git a/packages/client-runtime/src/state/outdatedServerUpdate.ts b/packages/client-runtime/src/state/outdatedServerUpdate.ts new file mode 100644 index 000000000..36d7ba1f6 --- /dev/null +++ b/packages/client-runtime/src/state/outdatedServerUpdate.ts @@ -0,0 +1,78 @@ +import type { EnvironmentId, ServerSelfUpdateInput } from "@t3tools/contracts"; +import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import type * as HttpClient from "effect/unstable/http/HttpClient"; +import type { Atom } from "effect/unstable/reactivity"; +import type * as Socket from "effect/unstable/socket/Socket"; + +import { updateOutdatedHost } from "../connection/outdatedHostUpdate.ts"; +import type * as ConnectionResolver from "../connection/resolver.ts"; +import type * as EnvironmentRegistry from "../connection/registry.ts"; +import type * as RelayEnvironmentDiscovery from "../relay/discovery.ts"; +import { createAtomCommandScheduler, createRuntimeCommand } from "./runtime.ts"; +import { + serverUpdateFailureMessage, + serverUpdateStateAtom, + type ServerUpdateStage, +} from "./server.ts"; + +export interface OutdatedServerUpdateTarget { + readonly environmentId: EnvironmentId; + readonly input: ServerSelfUpdateInput; + /** From the host descriptor when known; an outdated host never delivers a server config. */ + readonly fromVersion?: string; +} + +/** + * Updates a host too old for this client to connect to. Progress lands in the + * same per-environment update state as a normal server update. + */ +export function createOutdatedServerUpdateCommand( + runtime: Atom.AtomRuntime< + | EnvironmentRegistry.EnvironmentRegistry + | ConnectionResolver.ConnectionResolver + | RelayEnvironmentDiscovery.RelayEnvironmentDiscovery + | Socket.WebSocketConstructor + | HttpClient.HttpClient, + E + >, +) { + return createRuntimeCommand(runtime, { + label: "environment-data:server:update-outdated-server", + scheduler: createAtomCommandScheduler(), + concurrency: { + mode: "singleFlight", + key: ({ environmentId }: OutdatedServerUpdateTarget) => environmentId, + }, + execute: (target: OutdatedServerUpdateTarget, atomRegistry) => { + const stateAtom = serverUpdateStateAtom(target.environmentId); + const targetVersion = target.input.targetVersion; + const fromVersion = target.fromVersion ?? targetVersion; + let currentStage: ServerUpdateStage = "downloading"; + const setStage = (stage: ServerUpdateStage) => + Effect.sync(() => { + currentStage = stage; + atomRegistry.set(stateAtom, { status: "running", stage, fromVersion, targetVersion }); + }); + return setStage(currentStage).pipe( + Effect.andThen(updateOutdatedHost(target.environmentId, target.input, setStage)), + Effect.onExit((exit) => + Effect.sync(() => { + if (Exit.isSuccess(exit) || Cause.hasInterruptsOnly(exit.cause)) { + atomRegistry.set(stateAtom, { status: "idle" }); + return; + } + atomRegistry.set(stateAtom, { + status: "failed", + stage: currentStage, + fromVersion, + targetVersion, + message: serverUpdateFailureMessage(Cause.squash(exit.cause)), + }); + }), + ), + ); + }, + }); +} diff --git a/packages/client-runtime/src/state/server.ts b/packages/client-runtime/src/state/server.ts index 51200366f..6f16be9b9 100644 --- a/packages/client-runtime/src/state/server.ts +++ b/packages/client-runtime/src/state/server.ts @@ -85,7 +85,8 @@ const IDLE_SERVER_UPDATE_STATE: ServerUpdateState = { status: "idle" }; const EMPTY_SERVER_UPDATE_STATE_ATOM = Atom.make(IDLE_SERVER_UPDATE_STATE).pipe( Atom.withLabel("environment-data:server:update-state:empty"), ); -const serverUpdateStateAtom = Atom.family((environmentId: EnvironmentId) => +/** Shared with the outdated-host update, which reports through the same state. */ +export const serverUpdateStateAtom = Atom.family((environmentId: EnvironmentId) => Atom.make(IDLE_SERVER_UPDATE_STATE).pipe( Atom.withLabel(`environment-data:server:update-state:${environmentId}`), ), @@ -308,7 +309,7 @@ export function serverUpdateStateForServerVersion( : IDLE_SERVER_UPDATE_STATE; } -function serverUpdateFailureMessage(error: unknown): string { +export function serverUpdateFailureMessage(error: unknown): string { return error instanceof Error ? error.message : "Server update failed."; } From b09e8d8dffde8f82b57f75495218a9d2627a2af2 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 2 Oct 2026 21:42:36 -0700 Subject: [PATCH 05/12] fix(client-runtime): reconnects back off with jitter and keep healthy sockets (#14897) Co-authored-by: Claude Opus 5.5 (1M context) (cherry picked from commit 9333509c918083cbfb7e66759ddbcbfa71a3a452) Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/telemetry/AnalyticsService.test.ts | 84 ++++++ apps/server/src/telemetry/AnalyticsService.ts | 124 +++++++-- docs/internals/connection-runtime.md | 22 +- docs/internals/product-analytics.md | 7 + .../src/connection/registry.test.ts | 40 +++ .../src/connection/supervisor.test.ts | 245 ++++++++++++++++- .../src/connection/supervisor.ts | 251 +++++++++++------- .../client-runtime/src/connection/wakeups.ts | 1 + packages/client-runtime/src/state/server.ts | 6 +- 9 files changed, 632 insertions(+), 148 deletions(-) diff --git a/apps/server/src/telemetry/AnalyticsService.test.ts b/apps/server/src/telemetry/AnalyticsService.test.ts index 8f4306039..b32897ef6 100644 --- a/apps/server/src/telemetry/AnalyticsService.test.ts +++ b/apps/server/src/telemetry/AnalyticsService.test.ts @@ -4,9 +4,14 @@ import { assert, it } from "@effect/vitest"; import * as ConfigProvider from "effect/ConfigProvider"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Schema from "effect/Schema"; +import * as TestClock from "effect/testing/TestClock"; +import * as HttpClient from "effect/unstable/http/HttpClient"; +import * as HttpClientError from "effect/unstable/http/HttpClientError"; import * as HttpServer from "effect/unstable/http/HttpServer"; import * as HttpServerRequest from "effect/unstable/http/HttpServerRequest"; import * as HttpServerResponse from "effect/unstable/http/HttpServerResponse"; +import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess"; import * as ServerConfig from "../config.ts"; import { getTelemetryIdentifier } from "./Identify.ts"; @@ -35,6 +40,42 @@ interface RecordedBatchBody { }>; } +const SentBatch = Schema.fromJsonString( + Schema.Struct({ + batch: Schema.Array(Schema.Struct({ uuid: Schema.String })), + }), +); + +/** + * HTTP client that reads each batch, then fails as if the connection dropped + * before the response arrived. PostHog stores these batches, so the server + * must not send them forever. + */ +const acceptThenFailClient = (batches: Array>) => + Layer.succeed( + HttpClient.HttpClient, + HttpClient.make((request) => + Effect.gen(function* () { + if (request.body._tag === "Uint8Array") { + const body = yield* Schema.decodeEffect(SentBatch)( + new TextDecoder().decode(request.body.body), + ).pipe(Effect.orDie); + batches.push(body.batch); + } + return yield* new HttpClientError.HttpClientError({ + reason: new HttpClientError.TransportError({ request, cause: "connection reset" }), + }); + }), + ), + ); + +it("retryDelayMs doubles from 2s and stays under the 5 minute cap", () => { + assert.equal(AnalyticsService.retryDelayMs(1, 0), 1_000); + assert.equal(AnalyticsService.retryDelayMs(2, 0.999_999), 4_000); + assert.equal(AnalyticsService.retryDelayMs(30, 0), 150_000); + assert.equal(AnalyticsService.retryDelayMs(30, 0.999_999), 300_000); +}); + it.layer(NodeServices.layer)("AnalyticsService test", (it) => { it.effect("sends nothing when no PostHog key is configured", () => Effect.gen(function* () { @@ -81,6 +122,49 @@ it.layer(NodeServices.layer)("AnalyticsService test", (it) => { }), ); + it.effect("a batch that keeps failing is retried with backoff, then dropped", () => + Effect.gen(function* () { + const batches: Array> = []; + const runtimeLayer = AnalyticsService.layer.pipe( + Layer.provideMerge( + ServerConfig.ServerConfig.layerTest(process.cwd(), { prefix: "t3-telemetry-retry-" }), + ), + Layer.provide( + ConfigProvider.layer( + ConfigProvider.fromUnknown({ + T3CODE_TELEMETRY_ENABLED: true, + T3CODE_POSTHOG_KEY: "phc_test_key", + T3CODE_POSTHOG_HOST: "http://localhost", + }), + ), + ), + Layer.provide( + Layer.mergeAll( + Layer.succeed(HostProcessPlatform, "win32"), + Layer.succeed(HostProcessArchitecture, "x64"), + acceptThenFailClient(batches), + ), + ), + ); + + yield* Effect.gen(function* () { + const analytics = yield* AnalyticsService.AnalyticsService; + for (let index = 0; index < 20; index += 1) { + yield* analytics.record("test.retry", { index }); + } + // Before the fix this loop sent the batch about once a second. + for (let second = 0; second < 600; second += 1) { + yield* TestClock.adjust("1 second"); + } + }).pipe(Effect.provide(runtimeLayer)); + + assert.equal(batches.length, 5); + const uuids = batches.map((batch) => batch.map((event) => event.uuid).join(",")); + assert.equal(new Set(uuids).size, 1, "every retry carries the same uuids"); + assert.equal(new Set(batches[0]?.map((event) => event.uuid)).size, 20); + }), + ); + it.effect("flush drains all buffered events across multiple batches", () => Effect.gen(function* () { const capturedRequests: Array = []; diff --git a/apps/server/src/telemetry/AnalyticsService.ts b/apps/server/src/telemetry/AnalyticsService.ts index 2bfb77c94..4726ffd59 100644 --- a/apps/server/src/telemetry/AnalyticsService.ts +++ b/apps/server/src/telemetry/AnalyticsService.ts @@ -2,18 +2,26 @@ * Anonymous PostHog telemetry service. * * Persists an installation-scoped anonymous identifier, buffers events in - * memory, and flushes batches over Effect's HTTP client. + * memory, and flushes batches over Effect's HTTP client. A failed batch is + * retried with backoff and dropped after a few tries. Each event carries a + * uuid, so PostHog can tell a retried copy from a new event. * * @module AnalyticsService */ import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess"; +import * as Clock from "effect/Clock"; import * as Config from "effect/Config"; import * as Context from "effect/Context"; +import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; +import * as Random from "effect/Random"; import * as Ref from "effect/Ref"; +import * as Result from "effect/Result"; +import * as Semaphore from "effect/Semaphore"; import * as HttpClient from "effect/unstable/http/HttpClient"; import * as HttpClientRequest from "effect/unstable/http/HttpClientRequest"; import * as HttpClientResponse from "effect/unstable/http/HttpClientResponse"; @@ -23,11 +31,40 @@ import * as ServerConfig from "../config.ts"; import { getTelemetryIdentifier } from "./Identify.ts"; interface BufferedAnalyticsEvent { + readonly uuid: string; readonly event: string; readonly properties?: Readonly>; readonly capturedAt: string; } +interface DeliveryState { + /** Batch that failed last. It is sent again before newer events. */ + readonly failedBatch: ReadonlyArray; + /** Failed sends of `failedBatch`. */ + readonly batchAttempts: number; + /** Failed sends since the last success. Sets the backoff delay. */ + readonly failures: number; + /** The background flush does not send before this time (epoch ms). */ + readonly retryAt: number; +} + +const FLUSH_INTERVAL_MS = 1_000; +// A hung send would hold the flush lock, and with it the shutdown flush. +const SEND_TIMEOUT = "10 seconds"; +const MAX_BATCH_ATTEMPTS = 5; +const RETRY_BASE_DELAY_MS = 2_000; +const RETRY_MAX_DELAY_MS = 300_000; + +/** + * Delay before the next send after `failures` consecutive failed sends. The + * ceiling doubles from 2s up to 5 minutes, and the delay is a random point in + * its upper half. `random` is in [0, 1). + */ +export function retryDelayMs(failures: number, random: number): number { + const ceiling = Math.min(RETRY_MAX_DELAY_MS, RETRY_BASE_DELAY_MS * 2 ** (failures - 1)); + return Math.round(ceiling / 2 + (ceiling / 2) * random); +} + const TelemetryEnvConfig = Config.all({ // No baked-in key. Pylon inherited upstream's, which meant every install // reported into T3 Code's PostHog project — data Pylon cannot read and did not @@ -74,17 +111,31 @@ export const make = Effect.gen(function* () { const httpClient = yield* HttpClient.HttpClient; const serverConfig = yield* ServerConfig.ServerConfig; const identifier = yield* getTelemetryIdentifier; + const crypto = yield* Crypto.Crypto; const bufferRef = yield* Ref.make>([]); + const deliveryRef = yield* Ref.make({ + failedBatch: [], + batchAttempts: 0, + failures: 0, + retryAt: 0, + }); + // The background flush and the shutdown flush must not send the same batch at once. + const flushLock = yield* Semaphore.make(1); const clientType = serverConfig.mode === "desktop" ? "desktop-app" : "cli-web-client"; const hostPlatform = yield* HostProcessPlatform; const hostArchitecture = yield* HostProcessArchitecture; - const enqueueBufferedEvent = (event: string, properties?: Readonly>) => + const enqueueBufferedEvent = ( + uuid: string, + event: string, + properties?: Readonly>, + ) => Effect.flatMap(DateTime.now, (now) => Ref.modify(bufferRef, (current) => { const appended = [ ...current, { + uuid, event, ...(properties ? { properties } : {}), capturedAt: DateTime.formatIso(now), @@ -114,6 +165,7 @@ export const make = Effect.gen(function* () { const payload = { api_key: telemetryConfig.posthogKey, batch: events.map((event) => ({ + uuid: event.uuid, event: event.event, distinct_id: identifier, properties: { @@ -133,39 +185,69 @@ export const make = Effect.gen(function* () { HttpClientRequest.bodyJson(payload), Effect.flatMap(httpClient.execute), Effect.flatMap(HttpClientResponse.filterStatusOk), + Effect.timeout(SEND_TIMEOUT), ); }); + const takeBatch = Ref.modify(bufferRef, (current) => { + const nextBatch = current.slice(0, telemetryConfig.flushBatchSize); + return [nextBatch, current.slice(nextBatch.length)] as const; + }); + + // Sends batches until the buffer is empty or a send fails. A failed batch is + // kept for the next flush, and dropped after MAX_BATCH_ATTEMPTS failed sends. const flush: AnalyticsService["Service"]["flush"] = Effect.gen(function* () { while (true) { - const batch = yield* Ref.modify(bufferRef, (current) => { - if (current.length === 0) { - return [[] as ReadonlyArray, current] as const; - } - const nextBatch = current.slice(0, telemetryConfig.flushBatchSize); - const remaining = current.slice(nextBatch.length); - return [nextBatch, remaining] as const; - }); - + const delivery = yield* Ref.get(deliveryRef); + const batch = delivery.failedBatch.length > 0 ? delivery.failedBatch : yield* takeBatch; if (batch.length === 0) { return; } - yield* sendBatch(batch).pipe( - Effect.catch((error) => - Ref.update(bufferRef, (current) => [...batch, ...current]).pipe( - Effect.flatMap(() => Effect.fail(error)), - ), - ), - ); + const sent = yield* Effect.result(sendBatch(batch)); + if (Result.isSuccess(sent)) { + yield* Ref.set(deliveryRef, { failedBatch: [], batchAttempts: 0, failures: 0, retryAt: 0 }); + continue; + } + + const failures = delivery.failures + 1; + const batchAttempts = delivery.batchAttempts + 1; + const retryAt = (yield* Clock.currentTimeMillis) + retryDelayMs(failures, yield* Random.next); + if (batchAttempts < MAX_BATCH_ATTEMPTS) { + yield* Ref.set(deliveryRef, { failedBatch: batch, batchAttempts, failures, retryAt }); + yield* Effect.logDebug("Failed to send telemetry batch; will retry", { + attempt: batchAttempts, + cause: sent.failure, + }); + return; + } + yield* Ref.set(deliveryRef, { failedBatch: [], batchAttempts: 0, failures, retryAt }); + yield* Effect.logWarning("Dropped telemetry batch after repeated send failures", { + events: batch.length, + attempts: batchAttempts, + cause: sent.failure, + }); + return; } - }).pipe(Effect.catch((cause) => Effect.logError("Failed to flush telemetry", { cause }))); + }).pipe(flushLock.withPermit); + + const flushWhenDue = Effect.gen(function* () { + const { retryAt } = yield* Ref.get(deliveryRef); + if ((yield* Clock.currentTimeMillis) >= retryAt) { + yield* flush; + } + }); const record: AnalyticsService["Service"]["record"] = Effect.fn("AnalyticsService.record")( function* (event, properties) { if (!telemetryConfig.enabled || !telemetryConfig.posthogKey || !identifier) return; - const enqueueResult = yield* enqueueBufferedEvent(event, properties); + // Telemetry is best effort: an event without a uuid is not sent. The + // Node implementation throws (a defect) rather than failing, so catch both. + const uuid = yield* Effect.exit(crypto.randomUUIDv7); + if (Exit.isFailure(uuid)) return; + + const enqueueResult = yield* enqueueBufferedEvent(uuid.value, event, properties); if (enqueueResult.dropped) { yield* Effect.logDebug("analytics buffer full; dropping oldest event", { size: enqueueResult.size, @@ -175,7 +257,7 @@ export const make = Effect.gen(function* () { }, ); - yield* Effect.forever(Effect.sleep(1000).pipe(Effect.flatMap(() => flush)), { + yield* Effect.forever(Effect.sleep(FLUSH_INTERVAL_MS).pipe(Effect.flatMap(() => flushWhenDue)), { disableYield: true, }).pipe(Effect.forkScoped); diff --git a/docs/internals/connection-runtime.md b/docs/internals/connection-runtime.md index 052e841c4..27c8ccb7d 100644 --- a/docs/internals/connection-runtime.md +++ b/docs/internals/connection-runtime.md @@ -10,16 +10,20 @@ several views need the same environment. The [supervisor](../../packages/client-runtime/src/connection/supervisor.ts) owns retry policy; resolving an endpoint and opening an RPC session are single -attempts. Transient failures retry with capped backoff. Offline states and -authentication failures wait for a wakeup instead of spending attempts on -unchanged conditions. +attempts. Transient failures retry with jittered exponential backoff, capped at +five minutes, that resets only after a connection stays up. Without jitter, every +client of a restarted server reconnects in the same second; with a short cap, a +client that can never connect retries all day. Offline states and authentication +failures wait for a wakeup instead of spending attempts on unchanged conditions. -Foregrounding needs different treatment depending on the connection's state. -It wakes a retry immediately, leaves an ordinary in-flight attempt alone, and -probes an established session before replacing it. A long mobile background -suspension forces replacement because the OS can kill a socket without reporting -closure. Treating every foreground event as a reconnect delays healthy attempts; -treating every resume as harmless leaves suspended sockets stuck. +Foregrounding, an explicit retry, and an offline report probe the established +session, and only a failed probe reconnects. Offline reports are often wrong, for +example for a loopback server. A long mobile background suspension is the one +exception: it replaces the session at once, because the OS can kill a socket +without reporting closure, and a probe would hold a dead socket in "Resuming" +until it times out. That fresh attempt runs even while the network reports +offline. Foregrounding also wakes a pending retry immediately and +leaves an ordinary in-flight attempt alone. The [registry](../../packages/client-runtime/src/connection/registry.ts) scopes connections by environment. An involuntary disconnect retains the registration diff --git a/docs/internals/product-analytics.md b/docs/internals/product-analytics.md index 3aa0eb1eb..d6c789ad0 100644 --- a/docs/internals/product-analytics.md +++ b/docs/internals/product-analytics.md @@ -44,6 +44,13 @@ complete usage, no observed subagents, and no mixed models; compare matching mod interaction mode, and terminal status. Aggregate output/input ratios should divide the summed totals. Averaging per-turn ratios lets small-input turns dominate. +## Delivery + +A send can fail after PostHog has stored the batch, so every retry is a copy. +[Delivery](../../apps/server/src/telemetry/AnalyticsService.ts) gives each event a +uuid when it is recorded, backs off after a failed send, and drops a batch after a +few tries. Without these limits, one stuck batch was sent every second for days. + ## Collection boundary Keep analytics payloads to product metadata and normalized measurements. Do not send prompts, diff --git a/packages/client-runtime/src/connection/registry.test.ts b/packages/client-runtime/src/connection/registry.test.ts index 7e2e9ef85..88783dbef 100644 --- a/packages/client-runtime/src/connection/registry.test.ts +++ b/packages/client-runtime/src/connection/registry.test.ts @@ -1570,6 +1570,46 @@ describe("EnvironmentRegistry", () => { }), ); + it.effect("keeps one session per environment across concurrent registrations and retries", () => + Effect.gen(function* () { + const harness = yield* makeHarness([]); + + yield* Effect.gen(function* () { + const registry = yield* EnvironmentRegistry.EnvironmentRegistry; + const registration = new PrimaryConnectionRegistration({ target: TARGET }); + yield* Effect.all( + Array.from({ length: 5 }, () => registry.registerPlatform(registration)), + { concurrency: "unbounded", discard: true }, + ); + yield* awaitConnectionState( + registry, + TARGET.environmentId, + (state) => state.phase === "connected", + ); + + // Platform polls and explicit retries reach a healthy connection at once. + yield* Effect.all( + [ + ...Array.from({ length: 5 }, () => registry.registerPlatform(registration)), + ...Array.from({ length: 5 }, () => registry.retryNow(TARGET.environmentId)), + registry.reconcilePlatform([registration]), + ], + { concurrency: "unbounded", discard: true }, + ); + for (let attempt = 0; attempt < 100; attempt += 1) { + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(harness.sessions)).toHaveLength(1); + expect(yield* Ref.get(harness.releasedSessions)).toBe(0); + expect(yield* registry.state(TARGET.environmentId)).toMatchObject({ + phase: "connected", + generation: 1, + }); + }).pipe(Effect.provide(harness.layer), Effect.scoped); + }), + ); + it.effect("retains a healthy runtime when the platform repeats an identical registration", () => Effect.gen(function* () { const harness = yield* makeHarness([]); diff --git a/packages/client-runtime/src/connection/supervisor.test.ts b/packages/client-runtime/src/connection/supervisor.test.ts index d64863b9f..1d518f017 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 Random from "effect/Random"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; import * as SubscriptionRef from "effect/SubscriptionRef"; @@ -186,6 +187,11 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: }); const dependencies = Layer.mergeAll( + // Jitter at its maximum, so each retry waits exactly its ceiling: 2s, 4s, 8s... + Layer.succeed(Random.Random, { + nextDoubleUnsafe: () => 1 - Number.EPSILON, + nextIntUnsafe: () => 0, + }), Layer.succeed(Connectivity.Connectivity, connectivity), Layer.succeed( ConnectionWakeups.ConnectionWakeups, @@ -225,6 +231,17 @@ const makeHarness = Effect.fn("TestConnectionHarness.make")(function* (options?: }; }); +describe("retryDelayMs", () => { + it("doubles from 2 seconds to a 5 minute cap, jittered within the upper half of each step", () => { + const ceilings = [2_000, 4_000, 8_000, 16_000, 32_000, 64_000, 128_000, 256_000, 300_000]; + for (const [failureCount, ceiling] of [...ceilings, 300_000].entries()) { + expect(EnvironmentSupervisor.retryDelayMs(failureCount, 0)).toBe(ceiling / 2); + expect(EnvironmentSupervisor.retryDelayMs(failureCount, 0.5)).toBe((ceiling * 3) / 4); + expect(EnvironmentSupervisor.retryDelayMs(failureCount, 1 - Number.EPSILON)).toBe(ceiling); + } + }); +}); + describe("EnvironmentSupervisor", () => { it.effect("exports each relay setup as a standalone linked trace that ends at readiness", () => Effect.gen(function* () { @@ -353,7 +370,7 @@ describe("EnvironmentSupervisor", () => { }), ); - it.effect("retries forever with exponential backoff capped at sixteen seconds", () => + it.effect("retries forever with exponential backoff capped at five minutes", () => Effect.gen(function* () { const harness = yield* makeHarness({ prepare: () => Effect.fail(transient()), @@ -368,15 +385,20 @@ describe("EnvironmentSupervisor", () => { ); expect(yield* Ref.get(harness.prepareCount)).toBe(1); - for (const [index, delay] of [3_000, 4_000, 8_000, 16_000, 16_000, 16_000].entries()) { - yield* TestClock.adjust(delay); + const delays = [ + 2_000, 4_000, 8_000, 16_000, 32_000, 64_000, 128_000, 256_000, 300_000, 300_000, + ]; + for (const [index, delay] of delays.entries()) { + yield* TestClock.adjust(delay - 1); + expect(yield* Ref.get(harness.prepareCount)).toBe(index + 1); + yield* TestClock.adjust(1); yield* eventuallyState( supervisor.state, (state) => state.phase === "backoff" && state.attempt === index + 2, ); } - expect(yield* Ref.get(harness.prepareCount)).toBe(7); + expect(yield* Ref.get(harness.prepareCount)).toBe(delays.length + 1); }).pipe(Effect.provide(TestClock.layer())), ); @@ -578,7 +600,7 @@ describe("EnvironmentSupervisor", () => { ); expect(yield* Ref.get(harness.prepareCount)).toBe(3); - yield* TestClock.adjust("2999 millis"); + yield* TestClock.adjust("1999 millis"); expect(yield* Ref.get(harness.prepareCount)).toBe(3); yield* TestClock.adjust("1 milli"); yield* eventuallyState( @@ -641,9 +663,12 @@ describe("EnvironmentSupervisor", () => { }).pipe(Effect.provide(TestClock.layer())), ); - it.effect("releases a live session while offline and starts a new generation when online", () => + it.effect("keeps a session that still answers when the network reports offline", () => Effect.gen(function* () { - const harness = yield* makeHarness(); + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { initiallyDesired: true, }).pipe(Effect.provide(harness.dependencies)); @@ -652,21 +677,160 @@ describe("EnvironmentSupervisor", () => { supervisor.state, (state) => state.phase === "connected" && state.generation === 1, ); + // A loopback server, or a flap shorter than the probe, keeps working. yield* harness.setNetworkStatus("offline"); - yield* awaitState(supervisor.state, (state) => state.phase === "offline"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + yield* harness.setNetworkStatus("online"); + yield* Effect.yieldNow; - expect(yield* Ref.get(harness.releaseCount)).toBe(1); - expect(Option.isNone(yield* SubscriptionRef.get(supervisor.session))).toBe(true); + 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)).toMatchObject({ + phase: "connected", + generation: 1, + }); + }), + ); + + it.effect("replaces the session on a long resume while the network reports offline", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); - yield* harness.setNetworkStatus("online"); yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 1, + ); + // A wrong offline report: the probe answers, so the session stays. + yield* harness.setNetworkStatus("offline"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + + // The replacement connects although the network still reports offline. + yield* harness.wake("application-active-reconnect"); + const replaced = yield* awaitState( supervisor.state, (state) => state.phase === "connected" && state.generation === 2, ); + + expect(replaced.attempt).toBe(1); + expect(yield* Ref.get(probeCount)).toBe(1); expect(yield* Ref.get(harness.sessionCount)).toBe(2); + expect(yield* Ref.get(harness.releaseCount)).toBe(1); + }), + ); + + it.effect( + "releases a session that stops answering while offline and reconnects when online", + () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + probe: (attempt) => (attempt === 1 ? 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" && state.generation === 1, + ); + yield* harness.setNetworkStatus("offline"); + yield* TestClock.adjust("3 seconds"); + yield* awaitState(supervisor.state, (state) => state.phase === "offline"); + + expect(yield* Ref.get(harness.releaseCount)).toBe(1); + expect(Option.isNone(yield* SubscriptionRef.get(supervisor.session))).toBe(true); + + yield* harness.setNetworkStatus("online"); + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + }).pipe(Effect.provide(TestClock.layer())), + ); + + it.effect("probes instead of replacing a healthy session on an explicit retry", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* supervisor.retryNow; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(probeCount)).toBe(1); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + expect(yield* Ref.get(harness.releaseCount)).toBe(0); }), ); + it.effect("keeps the backoff ladder after an explicit retry finds a healthy session", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(TARGET_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState(supervisor.state, (state) => state.phase === "connected"); + yield* harness.closeLatestSession(); + yield* awaitState( + supervisor.state, + (state) => state.phase === "backoff" && state.attempt === 1, + ); + yield* TestClock.adjust("2 seconds"); + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + + yield* supervisor.retryNow; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + expect(yield* Ref.get(probeCount)).toBe(1); + + // The flapping session keeps climbing the ladder: the answered retry does + // not reset it after this unrelated close. + yield* harness.closeLatestSession(); + yield* awaitState( + supervisor.state, + (state) => state.phase === "backoff" && state.attempt === 2, + ); + yield* TestClock.adjust("4 seconds"); + const reconnected = yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 3, + ); + expect(reconnected.attempt).toBe(3); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("retries a blocked connection when platform credentials change", () => Effect.gen(function* () { const harness = yield* makeHarness({ @@ -844,7 +1008,7 @@ describe("EnvironmentSupervisor", () => { supervisor.state, (state) => state.phase === "backoff" && state.attempt === 1, ); - yield* TestClock.adjust("3 seconds"); + yield* TestClock.adjust("2 seconds"); yield* awaitState( supervisor.state, (state) => state.phase === "connected" && state.generation === 2 && state.attempt === 2, @@ -964,6 +1128,35 @@ describe("EnvironmentSupervisor", () => { }), ); + it.effect("reconnects immediately when the session closes during a resume probe", () => + Effect.gen(function* () { + const probeStarted = yield* Deferred.make(); + const harness = yield* makeHarness({ + 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* Deferred.await(probeStarted); + // The OS reports the suspended socket's close before the probe answers. + yield* harness.closeLatestSession(); + + // No TestClock advance: the unanswered probe skips the first backoff rung. + const reconnected = yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2, + ); + expect(reconnected.attempt).toBe(1); + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("reconnects immediately when the foreground liveness probe fails", () => Effect.gen(function* () { const allowReconnect = yield* Deferred.make(); @@ -1020,7 +1213,7 @@ describe("EnvironmentSupervisor", () => { supervisor.state, (state) => state.phase === "backoff" && state.attempt === 1, ); - yield* TestClock.adjust("2999 millis"); + yield* TestClock.adjust("1999 millis"); expect(yield* Ref.get(harness.prepareCount)).toBe(2); yield* TestClock.adjust("1 milli"); yield* eventuallyState( @@ -1055,6 +1248,32 @@ describe("EnvironmentSupervisor", () => { }).pipe(Effect.provide(TestClock.layer())), ); + it.effect("an explicit retry shortens a stalled desktop foreground probe", () => + Effect.gen(function* () { + const harness = yield* makeHarness({ + probe: (attempt) => (attempt === 1 ? 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"); + yield* TestClock.adjust("5 seconds"); + yield* supervisor.retryNow; + // The retry's 3 second limit applies, not the 10 seconds left of the 15. + yield* TestClock.adjust("2999 millis"); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + yield* TestClock.adjust("1 milli"); + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 2 && state.attempt === 1, + ); + + expect(yield* Ref.get(harness.sessionCount)).toBe(2); + }).pipe(Effect.provide(TestClock.layer())), + ); + it.effect("quickly times out a stalled mobile foreground liveness probe", () => Effect.gen(function* () { const harness = yield* makeHarness({ diff --git a/packages/client-runtime/src/connection/supervisor.ts b/packages/client-runtime/src/connection/supervisor.ts index 36956987b..4e1d4c11c 100644 --- a/packages/client-runtime/src/connection/supervisor.ts +++ b/packages/client-runtime/src/connection/supervisor.ts @@ -2,11 +2,13 @@ import { withRelayClientTracing } from "@t3tools/shared/relayTracing"; import * as Cause from "effect/Cause"; import * as Clock from "effect/Clock"; import * as Context from "effect/Context"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; +import * as Random from "effect/Random"; import * as Ref from "effect/Ref"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; @@ -30,10 +32,13 @@ import { safeErrorLogAttributes } from "../errors/safeLog.ts"; import { NETWORK_BLOCKING_HINT } from "../errors/network.ts"; import * as ConnectionWakeups from "./wakeups.ts"; -const RETRY_DELAYS_MS = [3_000, 4_000, 8_000, 16_000] as const; +const RETRY_BASE_DELAY_MS = 1_000; +const RETRY_MAX_DELAY_MS = 300_000; const CONNECTION_ESTABLISHMENT_TIMEOUT = "15 seconds"; const CONNECTION_PROBE_TIMEOUT = "15 seconds"; -const MOBILE_CONNECTION_PROBE_TIMEOUT = "3 seconds"; +// Mobile resumes, explicit retries, and offline events want a fast answer: +// the user is waiting, or the network may be gone. +const QUICK_CONNECTION_PROBE_TIMEOUT = "3 seconds"; const BACKOFF_RESET_AFTER_MS = 30_000; interface SupervisorIntent { @@ -102,8 +107,19 @@ export interface EnvironmentSupervisorOptions { readonly initiallyDesired?: boolean; } -function retryDelayMs(failureCount: number): number { - return RETRY_DELAYS_MS[Math.min(failureCount, RETRY_DELAYS_MS.length - 1)] ?? 16_000; +/** + * Delay before the next attempt after `failureCount` consecutive failures + * (0 for the first retry). The ceiling doubles from 2s up to 5 minutes, and + * the delay is a random point in its upper half: never quicker than half the + * ceiling, and spread out so clients that lost the same server do not all + * reconnect in the same second. `random` is in [0, 1). + * + * The long cap only applies to a connection that keeps failing. Returning to + * the app, the network coming back, and an explicit retry all skip the wait. + */ +export function retryDelayMs(failureCount: number, random: number): number { + const ceiling = Math.min(RETRY_MAX_DELAY_MS, RETRY_BASE_DELAY_MS * 2 ** (failureCount + 1)); + return Math.round(ceiling / 2 + (ceiling / 2) * random); } function annotateTarget(target: ConnectionTarget) { @@ -236,10 +252,11 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const intent = yield* Ref.make(initialIntent); const signals = yield* Queue.unbounded(); const resetRetryState = yield* Ref.make(false); - // Set when a foreground wake probe fails or times out: the user is actively - // returning to the app on a dead transport, so the follow-up reconnect skips - // the first backoff rung instead of sleeping. - const wakeProbeFailed = yield* Ref.make(false); + // Set while a probe of the live session is running, and kept when it fails + // or times out: something asked whether the connection still works and it + // closed or failed before answering, so the follow-up reconnect skips the + // first backoff rung instead of sleeping. + const probeUnanswered = yield* Ref.make(false); const state = yield* SubscriptionRef.make( !initialIntent.desired ? availableState(initialIntent, 0) @@ -393,96 +410,118 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( } }); + // Signals that end a connected lease whatever its health: "reset" ends it + // and restarts the retry ladder, "end" ends it, undefined keeps it. + const connectedLeaseEnd = Effect.fnUntraced(function* (next: SupervisorSignal) { + if (next._tag === "DisconnectRequested") { + return "end" as const; + } + if (next._tag !== "Wakeup") { + return undefined; + } + if (next.reason === "application-active-reconnect") { + // Mobile operating systems often kill a suspended socket without a close + // event. A probe would show a dead socket as "Resuming" until it times + // out, so a long background resume replaces the session at once. + return "reset" as const; + } + if (next.reason === "credentials-changed" && target._tag === "RelayConnectionTarget") { + yield* logManagedRelayAccountChange; + return "end" as const; + } + return undefined; + }); + + // How long a signal waits for the live session to answer a probe, or + // undefined when the signal does not question the connection. + const probeTimeoutFor = (next: SupervisorSignal): Duration.Input | undefined => { + switch (next._tag) { + case "RetryRequested": + return QUICK_CONNECTION_PROBE_TIMEOUT; + case "NetworkChanged": + return next.network === "offline" ? QUICK_CONNECTION_PROBE_TIMEOUT : undefined; + case "Wakeup": + if (next.reason === "application-active") { + return CONNECTION_PROBE_TIMEOUT; + } + return next.reason === "application-active-probe" + ? QUICK_CONNECTION_PROBE_TIMEOUT + : undefined; + case "ConnectRequested": + case "DisconnectRequested": + return undefined; + } + }; + + // Holds a connected lease until it must end, and returns whether to restart + // the retry ladder. Returning to the app, an explicit retry, and the network + // reporting offline all probe the live session instead of replacing it, so a + // healthy socket is not torn down (the offline report is often wrong, for + // example for a loopback server). Only a long mobile resume replaces the + // session without a probe. A failed probe fails this effect, and the + // supervisor reconnects. const monitorConnectedLease = Effect.fnUntraced(function* ( lease: ConnectionDriver.EnvironmentConnectionLease, ) { + // A probe answers an explicit retry here, so the retry must not also reset + // the backoff of a later, unrelated failure. + const takeSignal = Queue.take(signals).pipe( + Effect.tap((next) => + next._tag === "RetryRequested" ? Ref.set(resetRetryState, false) : Effect.void, + ), + ); for (;;) { - const next = yield* Queue.take(signals); - switch (next._tag) { - case "DisconnectRequested": - case "RetryRequested": - return false; - case "NetworkChanged": - if (next.network === "offline") { - return false; - } - break; - case "Wakeup": - if (next.reason === "credentials-changed" && target._tag === "RelayConnectionTarget") { - yield* logManagedRelayAccountChange; - return false; - } - if (next.reason === "application-active-reconnect") { - // Mobile operating systems commonly suspend sockets without - // delivering a close event. A long background resume deliberately - // replaces that lease and starts a fresh attempt without backoff. - return true; - } - if (next.reason === "application-active" || next.reason === "application-active-probe") { - const probe = yield* lease.session.probe.pipe( - Effect.timeoutOrElse({ - duration: - next.reason === "application-active-probe" - ? MOBILE_CONNECTION_PROBE_TIMEOUT - : CONNECTION_PROBE_TIMEOUT, - orElse: () => - Effect.fail( - new ConnectionTransientError({ - reason: "timeout", - detail: `${target.label} did not respond to a connection health check.`, - }), - ), - }), - Effect.forkChild, - ); - for (;;) { - const probeEvent = yield* Effect.raceFirst( - Fiber.await(probe).pipe( - Effect.map((exit) => ({ _tag: "ProbeCompleted" as const, exit })), - ), - Queue.take(signals).pipe( - Effect.map((signal) => ({ _tag: "Signal" as const, signal })), - ), - ); - if (probeEvent._tag === "ProbeCompleted") { - if (Exit.isFailure(probeEvent.exit)) { - yield* Ref.set(wakeProbeFailed, true); - } - yield* probeEvent.exit; - break; - } - switch (probeEvent.signal._tag) { - case "DisconnectRequested": - case "RetryRequested": - yield* Fiber.interrupt(probe); - return false; - case "NetworkChanged": - if (probeEvent.signal.network === "offline") { - yield* Fiber.interrupt(probe); - return false; - } - break; - case "Wakeup": - if (probeEvent.signal.reason === "application-active-reconnect") { - yield* Fiber.interrupt(probe); - return true; - } - if ( - probeEvent.signal.reason === "credentials-changed" && - target._tag === "RelayConnectionTarget" - ) { - yield* Fiber.interrupt(probe); - return false; - } - break; - case "ConnectRequested": - break; - } - } + const next = yield* takeSignal; + const end = yield* connectedLeaseEnd(next); + if (end !== undefined) { + return end === "reset"; + } + const probeTimeout = probeTimeoutFor(next); + if (probeTimeout === undefined) { + continue; + } + yield* Ref.set(probeUnanswered, true); + const probe = yield* Effect.forkChild(lease.session.probe); + // Monotonic nanoseconds, so a wall-clock correction cannot move the deadline. + let deadline = (yield* Clock.monotonicTimeNanos) + Duration.toNanosUnsafe(probeTimeout); + for (;;) { + const remaining = deadline - (yield* Clock.monotonicTimeNanos); + const probeEvent = yield* Effect.raceAllFirst([ + Fiber.await(probe).pipe( + Effect.map((exit) => ({ _tag: "ProbeCompleted" as const, exit })), + ), + takeSignal.pipe(Effect.map((signal) => ({ _tag: "Signal" as const, signal }))), + Effect.sleep(Duration.nanos(remaining > 0n ? remaining : 0n)).pipe( + Effect.as({ _tag: "TimedOut" as const }), + ), + ]); + if (probeEvent._tag === "TimedOut") { + yield* Fiber.interrupt(probe); + return yield* new ConnectionTransientError({ + reason: "timeout", + detail: `${target.label} did not respond to a connection health check.`, + }); + } + if (probeEvent._tag === "ProbeCompleted") { + if (Exit.isSuccess(probeEvent.exit)) { + yield* Ref.set(probeUnanswered, false); } + yield* probeEvent.exit; break; - case "ConnectRequested": - break; + } + const endDuringProbe = yield* connectedLeaseEnd(probeEvent.signal); + if (endDuringProbe !== undefined) { + yield* Fiber.interrupt(probe); + return endDuringProbe === "reset"; + } + // A retry or an offline report during a desktop foreground probe wants + // its quicker answer, so it shortens the running probe. + const signalTimeout = probeTimeoutFor(probeEvent.signal); + if (signalTimeout !== undefined) { + const signalDeadline = + (yield* Clock.monotonicTimeNanos) + Duration.toNanosUnsafe(signalTimeout); + if (signalDeadline < deadline) deadline = signalDeadline; + } } } }); @@ -492,6 +531,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( generation: number, lastFailure: ConnectionAttemptError | null, pendingRetry: Option.Option, + ignoreOffline: boolean, ) { yield* SubscriptionRef.set(prepared, Option.none()); const establishment = yield* Effect.raceAllFirst([ @@ -557,7 +597,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const active = establishment.exit.value; const currentIntent = yield* Ref.get(intent); - if (!currentIntent.desired || currentIntent.network === "offline") { + if (!currentIntent.desired || (currentIntent.network === "offline" && !ignoreOffline)) { return { _tag: "Interrupted", established: false, @@ -643,6 +683,10 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( failureCount = 0; pendingRetry = Option.none(); }; + // Set after a long resume ends an attempt or a session. The fresh attempt + // runs even while the network reports offline: the report is often wrong, + // and the replaced session must not leave the client offline. + let replacing = false; for (;;) { if (yield* Ref.getAndSet(resetRetryState, false)) { @@ -659,7 +703,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( yield* waitForSignal; continue; } - if (currentIntent.network === "offline") { + if (currentIntent.network === "offline" && !replacing) { yield* clearLease; yield* setState(offlineState(currentIntent, generation, failureCount + 1, latestFailure)); const applicationActivated = yield* waitForSignal; @@ -672,11 +716,12 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( const attempt = failureCount + 1; const nextGeneration = generation + 1; const outcome: AttemptOutcome = yield* Effect.scoped( - runAttempt(attempt, nextGeneration, latestFailure, pendingRetry), + runAttempt(attempt, nextGeneration, latestFailure, pendingRetry, replacing), ); + replacing = false; // Consumed on every iteration so a stale marker can never leak into a // later, unrelated failure. - const failedWakeProbe = yield* Ref.getAndSet(wakeProbeFailed, false); + const failedProbe = yield* Ref.getAndSet(probeUnanswered, false); if (outcome.established) { generation = nextGeneration; if (outcome.stable) { @@ -687,6 +732,7 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( if (outcome._tag === "Interrupted") { if (outcome.resetRetry) { resetRetryLadder(); + replacing = true; } continue; } @@ -713,18 +759,19 @@ export const make = Effect.fn("EnvironmentSupervisor.make")(function* ( continue; } - if (failedWakeProbe) { - // The wake probe found a dead transport while the user is returning to - // the app, so reconnect immediately instead of sleeping the first - // backoff rung. Only this first attempt skips the ladder; if it fails - // too, normal backoff resumes. + if (failedProbe) { + // A probe found a dead transport, or the transport closed while a probe + // waited for an answer (the user returned to the app, asked to retry, + // or the network changed), so reconnect immediately instead + // of sleeping the first backoff rung. Only this first attempt skips the + // ladder; if it fails too, normal backoff resumes. resetRetryLadder(); yield* setState(connectingState(yield* Ref.get(intent), generation, 1, error)); continue; } failureCount += 1; - const delayMs = retryDelayMs(failureCount - 1); + const delayMs = retryDelayMs(failureCount - 1, yield* Random.next); pendingRetry = Option.map(attemptSpan, (previousAttempt) => ({ previousAttempt, failureCount, diff --git a/packages/client-runtime/src/connection/wakeups.ts b/packages/client-runtime/src/connection/wakeups.ts index 721b94164..5f19924c4 100644 --- a/packages/client-runtime/src/connection/wakeups.ts +++ b/packages/client-runtime/src/connection/wakeups.ts @@ -16,6 +16,7 @@ export function isApplicationActiveWakeup(reason: ConnectionWakeup): boolean { ); } +// A long resume replaces the session, and the new session subscribes on its own. export function shouldResubscribeAfterWakeup(reason: ConnectionWakeup): boolean { return reason === "application-active" || reason === "application-active-probe"; } diff --git a/packages/client-runtime/src/state/server.ts b/packages/client-runtime/src/state/server.ts index 6f16be9b9..921e8e2ef 100644 --- a/packages/client-runtime/src/state/server.ts +++ b/packages/client-runtime/src/state/server.ts @@ -180,9 +180,9 @@ export function validateServerUpdateReadyEvent( * Keeps reconnect attempts ~1s apart for the whole update restart. * * A restart takes the server down for ~15 seconds, but the supervisor's normal - * backoff ladder (1/2/4/8/16s) assumes an unexpected failure and lands attempts - * at ~3, 5, 9, 17 and 33 seconds — so a 15-second restart is observed as a - * 33-second "Resuming". Nudging on every backoff entry (not just the first) + * backoff assumes an unexpected failure and doubles its delay after each failed + * attempt, so a 15-second restart can be observed as a ~30-second "Resuming". + * Nudging on every backoff entry (not just the first) * holds the retry cadence flat until the server answers again. The sleep before * each nudge is the pacer: a connection that fails instantly re-enters backoff * immediately and would otherwise spin a tight retry loop. From 0d785e216ee3e77a08a3cec1d2ee4e4af72c2532 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Fri, 2 Oct 2026 17:45:42 -0400 Subject: [PATCH 06/12] fix(server): API-key Codex installs no longer warn on every start (#14903) Signed-off-by: Yordis Prieto (cherry picked from commit bf7121d7a25fe84d76561c3142b9433615a703a1) Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/telemetry/Identify.test.ts | 43 ++++++++++++++++++++++ apps/server/src/telemetry/Identify.ts | 20 ++++++++-- 2 files changed, 59 insertions(+), 4 deletions(-) diff --git a/apps/server/src/telemetry/Identify.test.ts b/apps/server/src/telemetry/Identify.test.ts index 92d322326..a8c107b57 100644 --- a/apps/server/src/telemetry/Identify.test.ts +++ b/apps/server/src/telemetry/Identify.test.ts @@ -57,6 +57,49 @@ it.layer(NodeServices.layer)("telemetry identity", (it) => { ), ); + it.effect("falls back quietly when Codex authenticates with an API key", () => { + const logs: CapturedLog[] = []; + const logger = makeCaptureLogger(logs); + + return Effect.gen(function* () { + const config = yield* ServerConfig.ServerConfig; + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const homeDirectory = path.join(config.baseDir, "home"); + const codexAuthPath = path.join(homeDirectory, ".codex", "auth.json"); + const anonymousId = "api-key-fallback-anonymous-id"; + const privateApiKey = "sk-private-openai-api-key"; + + yield* fileSystem.makeDirectory(path.dirname(codexAuthPath), { recursive: true }); + yield* fileSystem.writeFileString( + codexAuthPath, + `{"auth_mode":"apikey","OPENAI_API_KEY":"${privateApiKey}"}`, + ); + yield* fileSystem.writeFileString(config.anonymousIdPath, anonymousId); + + const identifier = yield* Identify.getTelemetryIdentifierForHome(homeDirectory); + + assert.equal(identifier, sha256(anonymousId)); + assert.isUndefined(findIdentityLog(logs, "codex", "TelemetryIdentityDecodeError")); + assert.isUndefined(findIdentityLog(logs, "codex", "TelemetryIdentityReadError")); + const allLogs = logs + .map((log) => + [String(log.message), ...Object.values(log.annotations).map(String)].join("\n"), + ) + .join("\n"); + assert.notInclude(allLogs, privateApiKey); + }).pipe( + Effect.provide( + Layer.merge( + ServerConfig.layerTest(process.cwd(), { + prefix: "t3-telemetry-identify-apikey-", + }), + Logger.layer([logger], { mergeWithExisting: false }), + ), + ), + ); + }); + it.effect("logs structured decode context and falls back from malformed Codex auth", () => { const logs: CapturedLog[] = []; const logger = makeCaptureLogger(logs); diff --git a/apps/server/src/telemetry/Identify.ts b/apps/server/src/telemetry/Identify.ts index 50f1687bb..fcccfa550 100644 --- a/apps/server/src/telemetry/Identify.ts +++ b/apps/server/src/telemetry/Identify.ts @@ -10,10 +10,17 @@ import * as Schema from "effect/Schema"; import * as ServerConfig from "../config.ts"; +/** + * Codex omits `tokens` entirely when the install authenticates with an API key + * rather than a ChatGPT account, so an absent `tokens` is a supported install + * and not a malformed file. + */ const CodexAuthJsonSchema = Schema.Struct({ - tokens: Schema.Struct({ - account_id: Schema.String, - }), + tokens: Schema.optional( + Schema.Struct({ + account_id: Schema.String, + }), + ), }); const ClaudeJsonSchema = Schema.Struct({ @@ -183,7 +190,9 @@ const getCodexAccountId = Effect.fn("TelemetryIdentity.getCodexAccountId")(funct ), ); - return Option.some(authJson.tokens.account_id); + return authJson.tokens === undefined + ? Option.none() + : Option.some(authJson.tokens.account_id); }); const getClaudeUserId = Effect.fn("TelemetryIdentity.getClaudeUserId")(function* ( @@ -250,6 +259,9 @@ const upsertAnonymousId = Effect.gen(function* () { * 1. ~/.codex/auth.json tokens.account_id * 2. ~/.claude.json userID * 3. ~/.pylon-code/telemetry/anonymous-id + * + * A missing file or an API-key-only Codex auth.json falls through quietly. Only + * unreadable or malformed files warn. */ export const getTelemetryIdentifierForHome = Effect.fn("getTelemetryIdentifierForHome")( function* (homeDirectory: string) { From 579581d2a095e38525407f3e3b985158dba9efec Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Fri, 2 Oct 2026 18:50:01 -0400 Subject: [PATCH 07/12] fix(server): Codex auth tokens are an omitted key, never undefined (#14908) Signed-off-by: Yordis Prieto (cherry picked from commit a8927712f8b3338fa3c47a0c657b6e11d7e9aaec) Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/telemetry/Identify.test.ts | 31 +++++++++++++++------- apps/server/src/telemetry/Identify.ts | 8 +++--- 2 files changed, 26 insertions(+), 13 deletions(-) diff --git a/apps/server/src/telemetry/Identify.test.ts b/apps/server/src/telemetry/Identify.test.ts index a8c107b57..c4240be18 100644 --- a/apps/server/src/telemetry/Identify.test.ts +++ b/apps/server/src/telemetry/Identify.test.ts @@ -57,7 +57,24 @@ it.layer(NodeServices.layer)("telemetry identity", (it) => { ), ); - it.effect("falls back quietly when Codex authenticates with an API key", () => { + it.effect.each([ + { + login: "an API key", + secret: "sk-private-openai-api-key", + authJson: (secret: string) => `{"auth_mode":"apikey","OPENAI_API_KEY":"${secret}"}`, + }, + { + login: "an agent identity", + secret: "private-agent-identity-jwt", + authJson: (secret: string) => + `{"auth_mode":"agentIdentity","OPENAI_API_KEY":null,"agent_identity":"${secret}"}`, + }, + { + login: "a personal access token", + secret: "private-personal-access-token", + authJson: (secret: string) => `{"OPENAI_API_KEY":null,"personal_access_token":"${secret}"}`, + }, + ])("falls back quietly when Codex authenticates with $login", ({ secret, authJson }) => { const logs: CapturedLog[] = []; const logger = makeCaptureLogger(logs); @@ -67,14 +84,10 @@ it.layer(NodeServices.layer)("telemetry identity", (it) => { const path = yield* Path.Path; const homeDirectory = path.join(config.baseDir, "home"); const codexAuthPath = path.join(homeDirectory, ".codex", "auth.json"); - const anonymousId = "api-key-fallback-anonymous-id"; - const privateApiKey = "sk-private-openai-api-key"; + const anonymousId = "tokenless-codex-anonymous-id"; yield* fileSystem.makeDirectory(path.dirname(codexAuthPath), { recursive: true }); - yield* fileSystem.writeFileString( - codexAuthPath, - `{"auth_mode":"apikey","OPENAI_API_KEY":"${privateApiKey}"}`, - ); + yield* fileSystem.writeFileString(codexAuthPath, authJson(secret)); yield* fileSystem.writeFileString(config.anonymousIdPath, anonymousId); const identifier = yield* Identify.getTelemetryIdentifierForHome(homeDirectory); @@ -87,12 +100,12 @@ it.layer(NodeServices.layer)("telemetry identity", (it) => { [String(log.message), ...Object.values(log.annotations).map(String)].join("\n"), ) .join("\n"); - assert.notInclude(allLogs, privateApiKey); + assert.notInclude(allLogs, secret); }).pipe( Effect.provide( Layer.merge( ServerConfig.layerTest(process.cwd(), { - prefix: "t3-telemetry-identify-apikey-", + prefix: "t3-telemetry-identify-tokenless-", }), Logger.layer([logger], { mergeWithExisting: false }), ), diff --git a/apps/server/src/telemetry/Identify.ts b/apps/server/src/telemetry/Identify.ts index fcccfa550..46fd519d4 100644 --- a/apps/server/src/telemetry/Identify.ts +++ b/apps/server/src/telemetry/Identify.ts @@ -11,12 +11,12 @@ import * as Schema from "effect/Schema"; import * as ServerConfig from "../config.ts"; /** - * Codex omits `tokens` entirely when the install authenticates with an API key - * rather than a ChatGPT account, so an absent `tokens` is a supported install - * and not a malformed file. + * Codex writes `tokens` only for ChatGPT logins and omits the key for API-key, + * agent-identity, and personal-access-token logins, so its absence is a + * supported install and not a malformed file. */ const CodexAuthJsonSchema = Schema.Struct({ - tokens: Schema.optional( + tokens: Schema.optionalKey( Schema.Struct({ account_id: Schema.String, }), From f95e10fb6c346d1fd8f063f60f70351a28ae4dec Mon Sep 17 00:00:00 2001 From: Bob Fowler Date: Sat, 3 Oct 2026 11:45:10 -0400 Subject: [PATCH 08/12] fix(server): editors appear once a slow discovery scan finishes (#13917) (cherry picked from commit 0080e80c00a3e333e3dcafd85b2615d56f84e3b8) Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/process/externalLauncher.test.ts | 41 ++-- apps/server/src/process/externalLauncher.ts | 79 +++++--- apps/server/src/ws.test.ts | 188 +++++++++++++++++- apps/server/src/ws.ts | 110 ++++++++-- 4 files changed, 359 insertions(+), 59 deletions(-) diff --git a/apps/server/src/process/externalLauncher.test.ts b/apps/server/src/process/externalLauncher.test.ts index 967389ade..1470c6827 100644 --- a/apps/server/src/process/externalLauncher.test.ts +++ b/apps/server/src/process/externalLauncher.test.ts @@ -7,6 +7,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; import * as Clock from "effect/Clock"; import * as ConfigProvider from "effect/ConfigProvider"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; import * as FileSystem from "effect/FileSystem"; @@ -1268,26 +1269,24 @@ it.effect("memoizes editor discovery and refreshes after the cache window", () = ); }); -// A client that disconnects mid-scan interrupts the shared discovery effect on -// the connection fiber. The cache must not retain that interrupt: doing so -// replayed it to every later connect for the whole TTL, so `server.getConfig` -// failed and no client could reconnect until the server restarted. -it.effect("rescans after an interrupted discovery instead of caching the interrupt", () => { +// Connects run discovery under a timeout and may disconnect mid-scan. Neither +// may cancel the scan: on a busy host every connect would time out partway +// through, cache nothing, and leave every client without editors. +it.effect("keeps scanning after the caller is interrupted and shares that scan", () => { const fileInfo = { type: "File" } as FileSystem.File.Info; - let blockFirstScan = true; - let scans = 0; + const release = Deferred.makeUnsafe(); + let parkedStats = 0; const launcherLayer = ExternalLauncher.layer.pipe( Layer.provide( Layer.mergeAll( FileSystem.layerNoop({ - // The first scan parks inside `stat` so the interrupt lands while - // discovery is in flight, which is what a client disconnecting - // mid-connect does to the shared effect. + // Scans park inside `stat` until released, so the interrupt lands + // while discovery is in flight. stat: () => Effect.gen(function* () { - scans += 1; - if (blockFirstScan) { - return yield* Effect.never; + if (!Deferred.isDoneUnsafe(release)) { + parkedStats += 1; + yield* Deferred.await(release); } return fileInfo; }), @@ -1304,16 +1303,18 @@ it.effect("rescans after an interrupted discovery instead of caching the interru return Effect.gen(function* () { const launcher = yield* ExternalLauncher.ExternalLauncher; - const fiber = yield* Effect.forkChild(launcher.resolveAvailableEditors()); + const interrupted = yield* Effect.forkChild(launcher.resolveAvailableEditors()); yield* Effect.yieldNow; - yield* Fiber.interrupt(fiber); + yield* Fiber.interrupt(interrupted); - // The next connect must still get a real answer well inside the TTL. - blockFirstScan = false; - scans = 0; - const editors = yield* launcher.resolveAvailableEditors(); + // The next connect joins the running scan instead of starting its own. + const next = yield* Effect.forkChild(launcher.resolveAvailableEditors()); + yield* Effect.yieldNow; + assert.equal(parkedStats, 1); + + yield* Deferred.succeed(release, undefined); + const editors = yield* Fiber.join(next); assert.equal(editors.includes("vscode"), true); - assert.isAbove(scans, 0); }).pipe( Effect.provide( Layer.mergeAll( diff --git a/apps/server/src/process/externalLauncher.ts b/apps/server/src/process/externalLauncher.ts index e8f076443..1ae3e16f9 100644 --- a/apps/server/src/process/externalLauncher.ts +++ b/apps/server/src/process/externalLauncher.ts @@ -28,13 +28,16 @@ import { import * as Clock from "effect/Clock"; import * as Config from "effect/Config"; import * as Context from "effect/Context"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Encoding from "effect/Encoding"; +import * as Exit from "effect/Exit"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as Ref from "effect/Ref"; +import * as Scope from "effect/Scope"; import * as ChildProcess from "effect/unstable/process/ChildProcess"; import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; @@ -481,21 +484,21 @@ const resolveFileManagerRevealKind = Effect.fn("externalLauncher.resolveFileMana // the discovered set for a bounded window so repeat connects skip even the // per-command cache lookups in @t3tools/shared/shell. // -// This deliberately does not use `Effect.cachedWithTTL`: that memoizes the -// first caller's Exit whatever it is, including an interrupt. Callers run this -// on the connection fiber under a timeout (`resolveAvailableEditorsForConfig`), -// so one client disconnecting mid-scan would cache the interrupt and replay it -// to every later connect for the whole TTL, breaking `server.getConfig` -// permanently. Storing only on success means an interrupted scan leaves the -// cache untouched and the next connect simply rescans. +// The scan runs on its own fiber in the service scope, and every caller awaits +// that one scan. Callers apply a timeout (`resolveAvailableEditorsForConfig`) +// and disconnect mid-connect; neither may cancel a scan other connects are +// waiting on, or throw away work a slow host (a busy server at startup, a +// long PATH) needs more than one connect to finish. A failed scan clears the +// entry so the next caller starts over rather than replaying the failure. // Expiry uses the monotonic clock (Clock.currentTimeNanos), matching the // command-resolution cache in @t3tools/shared/shell, so a backward wall-clock // adjustment cannot keep an expired entry alive. const EDITOR_DISCOVERY_CACHE_TTL_NANOS = 60_000_000_000n; interface EditorDiscoveryCacheEntry { - readonly editors: ReadonlyArray; - readonly expiresAtNanos: bigint; + readonly scan: Deferred.Deferred>; + /** Undefined while the scan is still running. */ + readonly expiresAtNanos: bigint | undefined; } /** @@ -779,27 +782,57 @@ export const make = Effect.gen(function* () { Effect.provideService(Path.Path, path), ); + const scope = yield* Scope.Scope; const editorDiscoveryCache = yield* Ref.make>( Option.none(), ); - const cachedAvailableEditors = Effect.gen(function* () { - const nowNanos = yield* Clock.currentTimeNanos; - const entry = yield* Ref.get(editorDiscoveryCache); - if (Option.isSome(entry) && entry.value.expiresAtNanos > nowNanos) { - return entry.value.editors; - } - const editors = yield* provideCommandResolutionServices(resolveAvailableEditors()).pipe( + const runEditorDiscovery = (scan: Deferred.Deferred>) => + provideCommandResolutionServices(resolveAvailableEditors()).pipe( Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner), + Effect.onExit((exit) => + Effect.gen(function* () { + const expiresAtNanos = (yield* Clock.currentTimeNanos) + EDITOR_DISCOVERY_CACHE_TTL_NANOS; + yield* Ref.update(editorDiscoveryCache, (current) => + Option.isNone(current) || current.value.scan !== scan + ? current + : Exit.isSuccess(exit) + ? Option.some({ scan, expiresAtNanos }) + : Option.none(), + ); + yield* Deferred.done(scan, exit); + }), + ), + Effect.interruptible, + Effect.forkIn(scope), ); - yield* Ref.set( + // Claiming the cache entry and starting its scan must not be split by an + // interrupt, or the entry would wait on a scan that never runs. + const acquireEditorDiscovery = Effect.gen(function* () { + const nowNanos = yield* Clock.currentTimeNanos; + const [scan, isNewScan] = yield* Ref.modify( editorDiscoveryCache, - Option.some({ - editors, - expiresAtNanos: nowNanos + EDITOR_DISCOVERY_CACHE_TTL_NANOS, - }), + ( + current, + ): [ + [EditorDiscoveryCacheEntry["scan"], boolean], + Option.Option, + ] => { + if ( + Option.isSome(current) && + (current.value.expiresAtNanos === undefined || current.value.expiresAtNanos > nowNanos) + ) { + return [[current.value.scan, false], current]; + } + const scan = Deferred.makeUnsafe>(); + return [[scan, true], Option.some({ scan, expiresAtNanos: undefined })]; + }, ); - return editors; - }); + if (isNewScan) { + yield* runEditorDiscovery(scan); + } + return scan; + }).pipe(Effect.uninterruptible); + const cachedAvailableEditors = Effect.flatMap(acquireEditorDiscovery, Deferred.await); return ExternalLauncher.of({ resolveAvailableEditors: () => cachedAvailableEditors, diff --git a/apps/server/src/ws.test.ts b/apps/server/src/ws.test.ts index e280f0750..630d64c49 100644 --- a/apps/server/src/ws.test.ts +++ b/apps/server/src/ws.test.ts @@ -1,15 +1,29 @@ import { assert, it } from "@effect/vitest"; -import { ORCHESTRATION_PROTOCOL_VERSION } from "@t3tools/contracts"; +import { + ORCHESTRATION_PROTOCOL_VERSION, + type ServerConfig, + type ServerConfigStreamEvent, +} from "@t3tools/contracts"; +import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; +import * as ConfigProvider from "effect/ConfigProvider"; import * as Deferred from "effect/Deferred"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import * as Queue from "effect/Queue"; +import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; +import { ChildProcessSpawner } from "effect/unstable/process"; +import * as ExternalLauncher from "./process/externalLauncher.ts"; import { hasCompatibleOrchestrationProtocol, resolveAvailableEditorsForConfig, shouldUseBoundedThreadSnapshot, + withLateEditorConfig, } from "./ws.ts"; it("accepts only the current orchestration protocol before websocket RPC setup", () => { @@ -48,3 +62,175 @@ it.effect("does not block server config when editor discovery never resolves", ( assert.deepEqual(availableEditors, []); }), ); + +// Only the fields the late-editor fold reads or rewrites. +const snapshotConfig = (fields: Partial) => + ({ availableEditors: [], settings: {}, ...fields }) as unknown as ServerConfig; + +const settingsUpdated = (settings: object): ServerConfigStreamEvent => ({ + version: 1, + type: "settingsUpdated", + payload: { settings: settings as ServerConfig["settings"] }, +}); + +it.effect("resends late editors without rolling back updates already sent", () => + Effect.gen(function* () { + const settingsSent = yield* Deferred.make(); + const events = yield* withLateEditorConfig( + snapshotConfig({ settings: { enableProviderUpdateChecks: true } as never }), + Stream.make(settingsUpdated({ enableProviderUpdateChecks: false })), + { + resolveAvailableEditors: () => Effect.succeed(["file-manager"]), + // Holds the late snapshot until the settings change has gone out. + resolveFileManagerRevealKind: () => + Deferred.await(settingsSent).pipe(Effect.as("file-explorer" as const)), + }, + ).pipe( + Stream.tap((event) => + event.type === "settingsUpdated" ? Deferred.succeed(settingsSent, undefined) : Effect.void, + ), + Stream.runCollect, + ); + + const [first, second] = Array.from(events); + assert.equal(events.length, 2); + assert.equal(first?.type, "settingsUpdated"); + assert.equal(second?.type, "snapshot"); + if (second?.type === "snapshot") { + assert.deepEqual(second.config.availableEditors, ["file-manager"]); + assert.equal(second.config.shellRevealInFileManagerKind, "file-explorer"); + assert.deepEqual(second.config.settings, { enableProviderUpdateChecks: false } as never); + } + }), +); + +it.effect("sends no late snapshot when the scan matches the snapshot", () => + Effect.gen(function* () { + const events = yield* withLateEditorConfig( + snapshotConfig({ availableEditors: ["vscode"] }), + Stream.empty, + { + resolveAvailableEditors: () => Effect.succeed(["vscode"]), + resolveFileManagerRevealKind: () => Effect.succeed(undefined), + }, + ).pipe(Stream.runCollect); + + assert.equal(events.length, 0); + }), +); + +it.effect("resends a file manager reveal kind that missed the snapshot", () => + Effect.gen(function* () { + const events = yield* withLateEditorConfig( + snapshotConfig({ availableEditors: ["file-manager"] }), + Stream.empty, + { + resolveAvailableEditors: () => Effect.succeed(["file-manager"]), + resolveFileManagerRevealKind: () => Effect.succeed("file-explorer"), + }, + ).pipe(Stream.runCollect); + + const [late] = Array.from(events); + assert.equal(events.length, 1); + assert.equal(late?.type, "snapshot"); + if (late?.type === "snapshot") { + assert.equal(late.config.shellRevealInFileManagerKind, "file-explorer"); + } + }), +); + +// The real launcher on Windows over a filesystem whose probes park until +// released, like a host too busy to finish discovery inside the snapshot timeout. +const makeParkedWindowsLauncher = Effect.gen(function* () { + const parkedProbes = yield* Queue.unbounded(); + const release = yield* Deferred.make(); + const launcher = yield* ExternalLauncher.make.pipe( + Effect.provide( + Layer.mergeAll( + FileSystem.layerNoop({ + stat: () => + Queue.offer(parkedProbes, undefined).pipe( + Effect.andThen(Deferred.await(release)), + Effect.as({ type: "File" } as FileSystem.File.Info), + ), + }), + Path.layer, + Layer.succeed( + ChildProcessSpawner.ChildProcessSpawner, + ChildProcessSpawner.make(() => Effect.die("unexpected spawn")), + ), + ), + ), + ); + const onWindows = (effect: Effect.Effect) => + effect.pipe( + Effect.provideService(HostProcessPlatform, "win32"), + Effect.provide( + ConfigProvider.layer( + ConfigProvider.fromEnv({ + env: { PATH: "C:\\t3-late-editors-test", PATHEXT: ".EXE" }, + }), + ), + ), + ); + return { + editors: { + resolveAvailableEditors: () => onWindows(launcher.resolveAvailableEditors()), + resolveFileManagerRevealKind: () => onWindows(launcher.resolveFileManagerRevealKind()), + }, + probeParked: Queue.take(parkedProbes), + releaseProbes: Deferred.succeed(release, undefined), + }; +}); + +it.effect("recovers editors after a real scan outlasts the config timeout", () => + Effect.gen(function* () { + const { editors, probeParked, releaseProbes } = yield* makeParkedWindowsLauncher; + + const snapshotFiber = yield* resolveAvailableEditorsForConfig( + editors.resolveAvailableEditors(), + ).pipe(Effect.forkChild); + yield* probeParked; + yield* TestClock.adjust(Duration.seconds(5)); + const snapshotEditors = yield* Fiber.join(snapshotFiber); + assert.deepEqual(snapshotEditors, []); + + const lateFiber = yield* withLateEditorConfig( + snapshotConfig({ availableEditors: snapshotEditors }), + Stream.empty, + editors, + ).pipe(Stream.runCollect, Effect.forkChild); + yield* releaseProbes; + + const [late] = Array.from(yield* Fiber.join(lateFiber)); + assert.equal(late?.type, "snapshot"); + if (late?.type === "snapshot") { + assert.equal(late.config.availableEditors.includes("vscode"), true); + } + }).pipe(Effect.scoped), +); + +it.effect("recovers a reveal kind whose real probe outlasts the config timeout", () => + Effect.gen(function* () { + const { editors, probeParked, releaseProbes } = yield* makeParkedWindowsLauncher; + + // The snapshot's bounded probe timed out: file manager, but no reveal kind. + const lateFiber = yield* withLateEditorConfig( + snapshotConfig({ availableEditors: ["file-manager"] }), + Stream.empty, + { + resolveAvailableEditors: () => Effect.succeed(["file-manager"]), + resolveFileManagerRevealKind: editors.resolveFileManagerRevealKind, + }, + ).pipe(Stream.runCollect, Effect.forkChild); + yield* probeParked; + yield* TestClock.adjust(Duration.seconds(6)); + yield* releaseProbes; + + const [late] = Array.from(yield* Fiber.join(lateFiber)); + assert.equal(late?.type, "snapshot"); + if (late?.type === "snapshot") { + assert.equal(late.config.shellRevealInFileManagerKind, "file-explorer"); + } + }).pipe(Effect.scoped), +); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index dfc7a801a..f44b9e43e 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -97,6 +97,8 @@ import { type RelayClientInstallProgressEvent, type ServerSelfUpdateError, type ServerSelfUpdateProgressEvent, + type ServerConfig as ClientServerConfig, + type ServerConfigStreamEvent, type ServerLifecycleStreamEvent, type FilesystemBrowseFailure, FilesystemBrowseError, @@ -284,6 +286,94 @@ const resolveFileManagerRevealKindForConfig = ( discovery: Effect.Effect, ) => resolveDiscoveryForConfig(discovery, () => undefined); +type EditorDiscovery = Pick< + ExternalLauncher.ExternalLauncher["Service"], + "resolveAvailableEditors" | "resolveFileManagerRevealKind" +>; + +// The config fields that follow from which editors are installed. +const resolveEditorConfig = ( + availableEditors: ReadonlyArray, + revealKind: Effect.Effect, +) => + Effect.gen(function* () { + const fileManagerRevealKind = availableEditors.includes("file-manager") + ? yield* revealKind + : undefined; + return { + availableEditors, + ...(fileManagerRevealKind === undefined + ? {} + : { + shellRevealInFileManager: true, + shellRevealInFileManagerKind: fileManagerRevealKind, + }), + }; + }); + +/** + * Live config updates that follow a snapshot of `config`. A busy host can + * outlast the snapshot's discovery timeouts, which send no editors, or no + * reveal kind for the file manager. The scan keeps running, so once it lands + * this resends the config: clients replace theirs on any snapshot. The resent + * config is folded from the live updates already sent, so it cannot roll back + * a change that landed while the scan ran. + */ +export const withLateEditorConfig = ( + config: ClientServerConfig, + liveUpdates: Stream.Stream, + launcher: EditorDiscovery, +) => { + const lateEditorConfig = Stream.fromEffect(launcher.resolveAvailableEditors()).pipe( + Stream.filter( + (editors) => + editors.join() !== config.availableEditors.join() || + (editors.includes("file-manager") && config.shellRevealInFileManagerKind === undefined), + ), + // Unbounded, unlike the snapshot: the reveal-kind probe is not shared, so + // a timeout here would cancel a probe that outlasts it every time. + Stream.mapEffect((editors) => + resolveEditorConfig(editors, launcher.resolveFileManagerRevealKind()), + ), + Stream.filter( + (editorConfig) => + editorConfig.availableEditors.join() !== config.availableEditors.join() || + editorConfig.shellRevealInFileManagerKind !== config.shellRevealInFileManagerKind, + ), + Stream.map((editorConfig) => ({ type: "editorsResolved" as const, editorConfig })), + ); + + return Stream.merge(liveUpdates, lateEditorConfig).pipe( + Stream.mapAccum( + (): ClientServerConfig => config, + (current, event): readonly [ClientServerConfig, ReadonlyArray] => { + switch (event.type) { + case "editorsResolved": { + const { + availableEditors: _editors, + shellRevealInFileManager: _reveal, + shellRevealInFileManagerKind: _revealKind, + ...rest + } = current; + const next = { ...rest, ...event.editorConfig }; + return [next, [{ version: 1, type: "snapshot", config: next }]]; + } + case "keybindingsUpdated": + return [{ ...current, ...event.payload }, [event]]; + case "providerStatuses": + return [{ ...current, providers: event.payload.providers }, [event]]; + case "settingsUpdated": + return [{ ...current, settings: event.payload.settings }, [event]]; + // Themes and usage-limit sources never ride in a snapshot; clients + // carry their projected values across one. + default: + return [current, [event]]; + } + }, + ), + ); +}; + function unexpectedCompatibilityError(error: never): never { throw new Error(`Unhandled compatibility error: ${String(error)}`); } @@ -1655,14 +1745,10 @@ const makeWsRpcLayer = ( const environment = yield* serverEnvironment.getDescriptor; const auth = yield* serverAuth.getDescriptor(); const scratchWorkspaceRoot = yield* managedFolders.scratchRoot; - const availableEditors: ReadonlyArray = yield* resolveAvailableEditorsForConfig( - externalLauncher.resolveAvailableEditors(), + const editorConfig = yield* resolveEditorConfig( + yield* resolveAvailableEditorsForConfig(externalLauncher.resolveAvailableEditors()), + resolveFileManagerRevealKindForConfig(externalLauncher.resolveFileManagerRevealKind()), ); - const fileManagerRevealKind = availableEditors.includes("file-manager") - ? yield* resolveFileManagerRevealKindForConfig( - externalLauncher.resolveFileManagerRevealKind(), - ) - : undefined; return { environment, @@ -1672,7 +1758,7 @@ const makeWsRpcLayer = ( keybindings: keybindingsConfig.keybindings, issues: keybindingsConfig.issues, providers, - availableEditors, + ...editorConfig, // Same discovery-with-timeout treatment as editors: a slow probe // must not stall server.getConfig, so it degrades to no targets. remoteOpenTargets: yield* resolveAvailableEditorsForConfig( @@ -1694,12 +1780,6 @@ const makeWsRpcLayer = ( }, settings, shellResumeCompletionMarker: true, - ...(fileManagerRevealKind === undefined - ? {} - : { - shellRevealInFileManager: true, - shellRevealInFileManagerKind: fileManagerRevealKind, - }), threadResumeCompletionMarker: true, threadSnapshotPagination: true, ...Option.match(scratchWorkspaceRoot, { @@ -3869,7 +3949,7 @@ const makeWsRpcLayer = ( return Stream.concat( rpcInitialItems([{ version: 1 as const, type: "snapshot" as const, config }]), - liveUpdates, + withLateEditorConfig(config, liveUpdates, externalLauncher), ); }), { "rpc.aggregate": "server" }, From d825275c32c34e16598f1138b57a062978895c99 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sat, 3 Oct 2026 13:28:32 -0600 Subject: [PATCH 09/12] test: cover shell enrichment refreshes built from resolved identities Extract the upstream 8bc40b4e07 refresh mapping into ShellStream so the no-re-enrichment behavior has focused regression coverage. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/orchestration-v2/ShellStream.test.ts | 35 +++++++++++++++++++ .../src/orchestration-v2/ShellStream.ts | 25 +++++++++++++ apps/server/src/ws.ts | 14 +++----- 3 files changed, 64 insertions(+), 10 deletions(-) diff --git a/apps/server/src/orchestration-v2/ShellStream.test.ts b/apps/server/src/orchestration-v2/ShellStream.test.ts index 3435468a4..64995e546 100644 --- a/apps/server/src/orchestration-v2/ShellStream.test.ts +++ b/apps/server/src/orchestration-v2/ShellStream.test.ts @@ -4,6 +4,8 @@ import type { OrchestrationV2ShellSnapshot, OrchestrationV2StoredEvent, OrchestrationV2ThreadShell, + OrchestrationProjectShell, + RepositoryIdentity, } from "@t3tools/contracts"; import { ProjectId, ThreadId } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; @@ -16,6 +18,7 @@ import { coalesceStoredThreadEvents, composeShellStreamWithEnrichment, dedupeShellEnrichment, + projectsWithResolvedRepositoryIdentities, shellStreamItemFromEnrichmentRefresh, shellStreamItemFromThreadShell, shellStreamItemsFromInitialSnapshot, @@ -211,6 +214,38 @@ describe("archivedShellStreamItemFromThreadShell", () => { }); }); +describe("projectsWithResolvedRepositoryIdentities", () => { + const identity = (canonicalKey: string) => + ({ canonicalKey, locator: { source: "git-remote" } }) as unknown as RepositoryIdentity; + const projects = [ + { id: "project-a", workspaceRoot: "/workspace/a", repositoryIdentity: null }, + { id: "project-b", workspaceRoot: "/workspace/b", repositoryIdentity: identity("stale-b") }, + { id: "project-c", workspaceRoot: "/workspace/c", repositoryIdentity: identity("kept-c") }, + ] as unknown as ReadonlyArray; + + it("refreshes only changed roots with the identity each resolution carried", () => { + const refreshed = projectsWithResolvedRepositoryIdentities(projects, [ + { workspaceRoot: "/workspace/a", enrichment: { repositoryIdentity: identity("first-a") } }, + { workspaceRoot: "/workspace/b", enrichment: { repositoryIdentity: null } }, + { workspaceRoot: "/workspace/a", enrichment: { repositoryIdentity: identity("latest-a") } }, + { workspaceRoot: "/workspace/unknown", enrichment: { repositoryIdentity: null } }, + ]); + + expect(refreshed.map((project) => [project.id, project.repositoryIdentity])).toEqual([ + ["project-a", identity("latest-a")], + ["project-b", null], + ]); + }); + + it("emits no projects when no stored project matches a change", () => { + expect( + projectsWithResolvedRepositoryIdentities(projects, [ + { workspaceRoot: "/workspace/other", enrichment: { repositoryIdentity: null } }, + ]), + ).toEqual([]); + }); +}); + describe("shellStreamItemFromEnrichmentRefresh", () => { it("batches nearby completion roots onto one snapshot item", () => { expect( diff --git a/apps/server/src/orchestration-v2/ShellStream.ts b/apps/server/src/orchestration-v2/ShellStream.ts index 31f556ce6..c6cfdf46b 100644 --- a/apps/server/src/orchestration-v2/ShellStream.ts +++ b/apps/server/src/orchestration-v2/ShellStream.ts @@ -7,6 +7,7 @@ import type { OrchestrationV2ThreadShellSnapshot, OrchestrationV2ShellStreamItem, OrchestrationV2StoredEvent, + RepositoryIdentity, } from "@t3tools/contracts"; import { OrchestrationProjectShell as ProjectShellSchema } from "@t3tools/contracts"; import * as Schema from "effect/Schema"; @@ -125,6 +126,30 @@ export function dedupeShellEnrichment( } /** Build a shell snapshot stream item for a batched enrichment completion. */ +/** + * Projects whose repository identity just resolved, carrying the identity from + * the resolution itself. Re-enriching the projects here would re-request every + * expired root, whose resolution publishes again, so one expiry would keep + * every shell subscriber reloading every project's metadata once a minute. + * The latest change for a root wins. + */ +export function projectsWithResolvedRepositoryIdentities( + projects: ReadonlyArray, + changes: ReadonlyArray<{ + readonly workspaceRoot: string; + readonly enrichment: { readonly repositoryIdentity: RepositoryIdentity | null }; + }>, +): ReadonlyArray { + const identities = new Map( + changes.map((change) => [change.workspaceRoot, change.enrichment.repositoryIdentity]), + ); + return projects.flatMap((project) => + identities.has(project.workspaceRoot) + ? [{ ...project, repositoryIdentity: identities.get(project.workspaceRoot) ?? null }] + : [], + ); +} + export function shellStreamItemFromEnrichmentRefresh(input: { readonly snapshot: OrchestrationV2ShellSnapshot; readonly changes: ReadonlyArray<{ readonly workspaceRoot: string }>; diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index f44b9e43e..f8598de66 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -147,6 +147,7 @@ import { coalesceStoredThreadEvents, composeShellStreamWithEnrichment, dedupeShellEnrichment, + projectsWithResolvedRepositoryIdentities, shellStreamItemFromEnrichmentRefresh, shellStreamItemFromThreadShell, shellStreamItemsFromInitialSnapshot, @@ -1067,17 +1068,10 @@ export const subscribeOrchestrationV2Shell = Effect.fn("ws.orchestrationV2.subsc // project's metadata once a minute. Stream.mapEffect((changes) => Effect.gen(function* () { - const identities = new Map( - Array.from(changes, (change) => [ - change.workspaceRoot, - change.enrichment.repositoryIdentity, - ]), - ); const snapshotSequence = yield* applicationEvents.latestApplicationSequence; - const changedProjects = (yield* projects.listShells()).flatMap((project) => - identities.has(project.workspaceRoot) - ? [{ ...project, repositoryIdentity: identities.get(project.workspaceRoot) ?? null }] - : [], + const changedProjects = projectsWithResolvedRepositoryIdentities( + yield* projects.listShells(), + Array.from(changes), ); return shellStreamItemFromEnrichmentRefresh({ snapshot: { From d1e2c0974ac71e0ead8fc20fd76ec24e75b97705 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sat, 3 Oct 2026 13:28:33 -0600 Subject: [PATCH 10/12] test: cover relay probe reconnects and clear adopted lint findings Add relay-target supervisor coverage for probe-instead-of-replace and credential changes during a probe, hoist the analytics test decoder, and use Effect.undefined in the late editor config test. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/telemetry/AnalyticsService.test.ts | 7 +- apps/server/src/ws.test.ts | 2 +- .../src/connection/supervisor.test.ts | 68 +++++++++++++++++++ 3 files changed, 73 insertions(+), 4 deletions(-) diff --git a/apps/server/src/telemetry/AnalyticsService.test.ts b/apps/server/src/telemetry/AnalyticsService.test.ts index b32897ef6..1141d609f 100644 --- a/apps/server/src/telemetry/AnalyticsService.test.ts +++ b/apps/server/src/telemetry/AnalyticsService.test.ts @@ -45,6 +45,7 @@ const SentBatch = Schema.fromJsonString( batch: Schema.Array(Schema.Struct({ uuid: Schema.String })), }), ); +const decodeSentBatch = Schema.decodeEffect(SentBatch); /** * HTTP client that reads each batch, then fails as if the connection dropped @@ -57,9 +58,9 @@ const acceptThenFailClient = (batches: Array Effect.gen(function* () { if (request.body._tag === "Uint8Array") { - const body = yield* Schema.decodeEffect(SentBatch)( - new TextDecoder().decode(request.body.body), - ).pipe(Effect.orDie); + const body = yield* decodeSentBatch(new TextDecoder().decode(request.body.body)).pipe( + Effect.orDie, + ); batches.push(body.batch); } return yield* new HttpClientError.HttpClientError({ diff --git a/apps/server/src/ws.test.ts b/apps/server/src/ws.test.ts index 630d64c49..ab0794f7c 100644 --- a/apps/server/src/ws.test.ts +++ b/apps/server/src/ws.test.ts @@ -111,7 +111,7 @@ it.effect("sends no late snapshot when the scan matches the snapshot", () => Stream.empty, { resolveAvailableEditors: () => Effect.succeed(["vscode"]), - resolveFileManagerRevealKind: () => Effect.succeed(undefined), + resolveFileManagerRevealKind: () => Effect.undefined, }, ).pipe(Stream.runCollect); diff --git a/packages/client-runtime/src/connection/supervisor.test.ts b/packages/client-runtime/src/connection/supervisor.test.ts index 1d518f017..ca41ebda9 100644 --- a/packages/client-runtime/src/connection/supervisor.test.ts +++ b/packages/client-runtime/src/connection/supervisor.test.ts @@ -1353,6 +1353,74 @@ describe("EnvironmentSupervisor", () => { }), ); + it.effect("keeps a relay session that answers a probe on an offline report or retry", () => + Effect.gen(function* () { + const probeCount = yield* Ref.make(0); + const harness = yield* makeHarness({ + probe: () => Ref.update(probeCount, (count) => count + 1), + }); + const supervisor = yield* EnvironmentSupervisor.make(RELAY_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 1, + ); + yield* harness.setNetworkStatus("offline"); + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 0) break; + yield* Effect.yieldNow; + } + yield* harness.setNetworkStatus("online"); + yield* supervisor.retryNow; + for (let attempt = 0; attempt < 100; attempt += 1) { + if ((yield* Ref.get(probeCount)) > 1) break; + yield* Effect.yieldNow; + } + + expect(yield* Ref.get(probeCount)).toBe(2); + expect(yield* Ref.get(harness.sessionCount)).toBe(1); + expect(yield* Ref.get(harness.releaseCount)).toBe(0); + expect(yield* SubscriptionRef.get(supervisor.state)).toMatchObject({ + phase: "connected", + generation: 1, + }); + }), + ); + + it.effect("ends a probing relay session when relay credentials change", () => + Effect.gen(function* () { + const probeStarted = yield* Deferred.make(); + const harness = yield* makeHarness({ + probe: (attempt) => + attempt === 1 + ? Deferred.succeed(probeStarted, undefined).pipe(Effect.andThen(Effect.never)) + : Effect.void, + }); + const supervisor = yield* EnvironmentSupervisor.make(RELAY_ENTRY, { + initiallyDesired: true, + }).pipe(Effect.provide(harness.dependencies)); + + yield* awaitState( + supervisor.state, + (state) => state.phase === "connected" && state.generation === 1, + ); + yield* supervisor.retryNow; + yield* Deferred.await(probeStarted); + // The account behind the relay changed: the stalled probe must not hold + // the old session until its timeout. + yield* harness.wake("credentials-changed"); + 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("keeps a healthy relay session when its HTTP access token expires", () => Effect.gen(function* () { const harness = yield* makeHarness({ From 86459581ef5f3007f8a9807297f9b37a1cac6640 Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sat, 3 Oct 2026 13:44:26 -0600 Subject: [PATCH 11/12] fix(web): confirm before updating a host too old to connect The outdated-host update restarted the remote server, or relaunched a desktop-managed app, on a single click. The runtime now reads the host descriptor, asks for confirmation naming the method and the agent sessions that will stop, and sends nothing when declined. The success toast reports a desktop relaunch, the block message tells clients without an update action that desktop or web can start it, and the connection-runtime doc records the bare-socket path and its four-minute no-recommit timeout. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/orchestration-v2/ShellStream.ts | 2 +- .../web/src/components/ServerUpdateAction.tsx | 11 +- docs/internals/connection-runtime.md | 16 ++ .../src/connection/compatibility.test.ts | 14 ++ .../src/connection/compatibility.ts | 16 +- .../client-runtime/src/connection/index.ts | 2 +- .../src/connection/outdatedHostUpdate.test.ts | 153 +++++++++++++++++- .../src/connection/outdatedHostUpdate.ts | 34 ++++ .../src/state/outdatedServerUpdate.ts | 28 +++- 9 files changed, 260 insertions(+), 16 deletions(-) diff --git a/apps/server/src/orchestration-v2/ShellStream.ts b/apps/server/src/orchestration-v2/ShellStream.ts index c6cfdf46b..15e870b8b 100644 --- a/apps/server/src/orchestration-v2/ShellStream.ts +++ b/apps/server/src/orchestration-v2/ShellStream.ts @@ -125,7 +125,6 @@ export function dedupeShellEnrichment( }); } -/** Build a shell snapshot stream item for a batched enrichment completion. */ /** * Projects whose repository identity just resolved, carrying the identity from * the resolution itself. Re-enriching the projects here would re-request every @@ -150,6 +149,7 @@ export function projectsWithResolvedRepositoryIdentities( ); } +/** Build a shell snapshot stream item for a batched enrichment completion. */ export function shellStreamItemFromEnrichmentRefresh(input: { readonly snapshot: OrchestrationV2ShellSnapshot; readonly changes: ReadonlyArray<{ readonly workspaceRoot: string }>; diff --git a/apps/web/src/components/ServerUpdateAction.tsx b/apps/web/src/components/ServerUpdateAction.tsx index f4ae391d7..9a7a36484 100644 --- a/apps/web/src/components/ServerUpdateAction.tsx +++ b/apps/web/src/components/ServerUpdateAction.tsx @@ -4,6 +4,7 @@ import { isAtomCommandInterrupted, squashAtomCommandFailure, } from "@t3tools/client-runtime/state/runtime"; +import { outdatedHostUpdateConfirmation } from "@t3tools/client-runtime/connection"; import { CircleArrowUpIcon } from "lucide-react"; import { type ComponentProps, useRef, useState } from "react"; @@ -319,6 +320,11 @@ export function OutdatedServerUpdateAction({ environmentId, input: { targetVersion }, ...(fromVersion === undefined ? {} : { fromVersion }), + // The host restarts without asking anyone there, and this client cannot + // see its threads, so this is the only confirmation in the flow. No + // themed host mounted (undefined) means proceed: the click was the request. + confirm: async (plan) => + (await requestConfirmDialog(outdatedHostUpdateConfirmation(plan))) ?? true, }); if (result._tag === "Failure") { if (isAtomCommandInterrupted(result)) return; @@ -327,7 +333,10 @@ export function OutdatedServerUpdateAction({ toastManager.add({ type: "success", title: `${serverLabel} updated`, - description: `Reconnected on t3@${result.value.targetVersion}.`, + description: + result.value.method === "desktop-app" + ? `Desktop app relaunched on ${result.value.targetVersion}.` + : `Reconnected on t3@${result.value.targetVersion}.`, }); } catch (error) { toastManager.add({ diff --git a/docs/internals/connection-runtime.md b/docs/internals/connection-runtime.md index 27c8ccb7d..ee1f4c35c 100644 --- a/docs/internals/connection-runtime.md +++ b/docs/internals/connection-runtime.md @@ -31,6 +31,22 @@ and cached data. Explicit removal closes the scope and clears credentials, projections, and platform-owned state such as drafts. Cloud-account changes apply to relay registrations; they must not discard directly paired environments. +## Updating a host too old to connect + +A host on an older orchestration protocol is blocked before a session opens, so +the normal update path, which runs over the session, cannot reach it. The +[outdated-host update](../../packages/client-runtime/src/connection/outdatedHostUpdate.ts) +authorizes a bare socket without the protocol gate and calls only the self-update +RPCs, whose wire shape has not changed across protocol versions. It confirms with +the user first, with the method read from the host's descriptor: nobody at the +host is asked, and this client cannot see the threads a restart stops. + +It then polls the descriptor for up to four minutes until the host reports a +compatible protocol. Unlike the session path, it does not re-commit a desktop +update when the app relaunches on the old version; that case fails after the +timeout and the user retries. Mobile has no update action, so the block message +names desktop and web as the place to start one. + ## HTTP authorization RPC sessions authenticate at socket upgrade, while HTTP snapshot loaders need current diff --git a/packages/client-runtime/src/connection/compatibility.test.ts b/packages/client-runtime/src/connection/compatibility.test.ts index 3f9dc8f99..5e6108b11 100644 --- a/packages/client-runtime/src/connection/compatibility.test.ts +++ b/packages/client-runtime/src/connection/compatibility.test.ts @@ -53,6 +53,20 @@ describe("orchestration protocol compatibility", () => { expect(error).not.toHaveProperty("serverUpdateRequired"); }); + it("tells clients without an update action where an outdated host can be updated", () => { + const older = descriptor(ORCHESTRATION_PROTOCOL_VERSION - 1); + const updatable = orchestrationProtocolCompatibilityError({ + ...older, + capabilities: { repositoryIdentity: true, serverSelfUpdate: "respawn" }, + }); + expect(updatable?.message).toContain("Update Pylon on Build Mac to connect."); + expect(updatable?.message).toContain("Pylon on desktop or web can start the update"); + + const manual = orchestrationProtocolCompatibilityError(older); + expect(manual?.message).toContain("Update Pylon on Build Mac to connect."); + expect(manual?.message).not.toContain("desktop or web"); + }); + it("offers a remote update only for an older host that can update itself", () => { const older = descriptor(ORCHESTRATION_PROTOCOL_VERSION - 1); const withCapabilities = (capabilities: ExecutionEnvironmentDescriptor["capabilities"]) => diff --git a/packages/client-runtime/src/connection/compatibility.ts b/packages/client-runtime/src/connection/compatibility.ts index e9c66b642..683428ec6 100644 --- a/packages/client-runtime/src/connection/compatibility.ts +++ b/packages/client-runtime/src/connection/compatibility.ts @@ -19,11 +19,17 @@ export function orchestrationProtocolCompatibilityError( reason: "unsupported", detail: `This client is not supported by this server. Update your app or use a compatible release to connect to ${descriptor.label}.`, }) - : new ConnectionBlockedError({ - reason: "unsupported", - detail: `This client requires a newer server. Update Pylon on ${descriptor.label} to connect.`, - ...(canSelfUpdate(descriptor) ? { serverUpdateRequired: true } : {}), - }); + : canSelfUpdate(descriptor) + ? new ConnectionBlockedError({ + reason: "unsupported", + // Mobile has no update action, so name where the update can start. + detail: `This client requires a newer server. Update Pylon on ${descriptor.label} to connect. Pylon on desktop or web can start the update from Settings → Connections.`, + serverUpdateRequired: true, + }) + : new ConnectionBlockedError({ + reason: "unsupported", + detail: `This client requires a newer server. Update Pylon on ${descriptor.label} to connect.`, + }); } /** Whether this client can drive the host's update remotely. */ diff --git a/packages/client-runtime/src/connection/index.ts b/packages/client-runtime/src/connection/index.ts index c3c10f006..3a06502ec 100644 --- a/packages/client-runtime/src/connection/index.ts +++ b/packages/client-runtime/src/connection/index.ts @@ -16,4 +16,4 @@ export * as Wakeups from "./wakeups.ts"; export { orchestrationProtocolCompatibilityError } from "./compatibility.ts"; // Flat so consumers' inferred command types can name it. -export { OutdatedHostUpdateError } from "./outdatedHostUpdate.ts"; +export { OutdatedHostUpdateError, outdatedHostUpdateConfirmation } from "./outdatedHostUpdate.ts"; diff --git a/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts b/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts index e63d095fa..e88e7ab25 100644 --- a/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts +++ b/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts @@ -5,7 +5,9 @@ import { WS_METHODS, } from "@t3tools/contracts"; import { describe, expect, it } from "@effect/vitest"; +import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; @@ -17,7 +19,11 @@ import * as Socket from "effect/unstable/socket/Socket"; import type { ConnectionCatalogEntry } from "./catalog.ts"; import { orchestrationProtocolCompatibilityError } from "./compatibility.ts"; import { PrimaryConnectionTarget } from "./model.ts"; -import { updateOutdatedHost } from "./outdatedHostUpdate.ts"; +import { + type OutdatedHostUpdatePlan, + outdatedHostUpdateConfirmation, + updateOutdatedHost, +} from "./outdatedHostUpdate.ts"; import * as EnvironmentRegistry from "./registry.ts"; import * as ConnectionResolver from "./resolver.ts"; import * as RelayEnvironmentDiscovery from "../relay/discovery.ts"; @@ -166,10 +172,17 @@ describe("updateOutdatedHost", () => { ), ); + const events: Array = []; const result = yield* updateOutdatedHost( TARGET.environmentId, { targetVersion: "0.0.46" }, - () => Effect.void, + (stage) => Effect.sync(() => events.push(`stage:${stage}`)), + (plan) => + Effect.sync(() => { + events.push(`confirm:${plan.method}:${plan.fromVersion}->${plan.targetVersion}`); + expect(sockets).toHaveLength(0); + return true; + }), ).pipe( Effect.provide( Layer.mergeAll( @@ -203,6 +216,12 @@ describe("updateOutdatedHost", () => { ]); expect(result.targetVersion).toBe("0.0.46"); expect(calls).toEqual(["compatibility:clear", "enabled:true"]); + // Confirmation comes before any progress or socket. + expect(events).toEqual([ + "confirm:boot-service:0.0.45->0.0.46", + "stage:downloading", + "stage:resuming", + ]); }), ); @@ -214,7 +233,12 @@ describe("updateOutdatedHost", () => { }; let opened = false; const error = yield* Effect.flip( - updateOutdatedHost(TARGET.environmentId, { targetVersion: "0.0.46" }, () => Effect.void), + updateOutdatedHost( + TARGET.environmentId, + { targetVersion: "0.0.46" }, + () => Effect.void, + () => Effect.die(new Error("A host that cannot update itself is never confirmed.")), + ), ).pipe( Effect.provide( Layer.mergeAll( @@ -275,4 +299,127 @@ describe("updateOutdatedHost", () => { expect(opened).toBe(false); }), ); + it.effect("asks before relaunching a desktop-managed host and sends nothing when declined", () => + Effect.gen(function* () { + const desktopManaged = { + ...descriptor(undefined, "0.0.45"), + capabilities: { + repositoryIdentity: true, + serverSelfUpdate: "desktop-managed", + desktopAppUpdate: true, + }, + } satisfies ExecutionEnvironmentDescriptor; + const plans: Array = []; + const stages: Array = []; + const registryCalls: Array = []; + let opened = false; + const exit = yield* Effect.exit( + updateOutdatedHost( + TARGET.environmentId, + { targetVersion: "0.0.46" }, + (stage) => Effect.sync(() => stages.push(stage)), + (plan) => + Effect.sync(() => { + plans.push(plan); + return false; + }), + ), + ).pipe( + Effect.provide( + Layer.mergeAll( + Layer.succeed( + EnvironmentRegistry.EnvironmentRegistry, + EnvironmentRegistry.EnvironmentRegistry.of({ + entries: yield* SubscriptionRef.make< + ReadonlyMap + >( + new Map([ + [ + TARGET.environmentId, + { target: TARGET, profile: Option.none(), enabled: false }, + ], + ]), + ), + setCompatibility: () => Effect.sync(() => registryCalls.push("compatibility")), + setEnabled: () => Effect.sync(() => registryCalls.push("enabled")), + } as unknown as EnvironmentRegistry.EnvironmentRegistry["Service"]), + ), + Layer.succeed( + ConnectionResolver.ConnectionResolver, + ConnectionResolver.ConnectionResolver.of({ + prepare: () => Effect.die(new Error("unused")), + prepareForUpdate: () => + Effect.succeed({ + descriptor: desktopManaged, + prepared: { + environmentId: TARGET.environmentId, + label: TARGET.label, + httpBaseUrl: TARGET.httpBaseUrl, + socketUrl: "wss://build.example.test/ws", + httpAuthorization: null, + target: TARGET, + }, + }), + }), + ), + Layer.succeed( + RelayEnvironmentDiscovery.RelayEnvironmentDiscovery, + RelayEnvironmentDiscovery.RelayEnvironmentDiscovery.of({ + state: yield* SubscriptionRef.make( + RelayEnvironmentDiscovery.EMPTY_RELAY_ENVIRONMENT_DISCOVERY_STATE, + ), + refresh: Effect.void, + }), + ), + Layer.succeed( + HttpClient.HttpClient, + HttpClient.make(() => Effect.die(new Error("unused"))), + ), + Layer.succeed(Socket.WebSocketConstructor, () => { + opened = true; + throw new Error("A declined update must not open a socket."); + }), + ), + ), + ); + + expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true); + expect(plans).toEqual([ + { + hostLabel: "Build Mac", + method: "desktop-managed", + fromVersion: "0.0.45", + targetVersion: "0.0.46", + }, + ]); + expect(stages).toEqual([]); + expect(registryCalls).toEqual([]); + expect(opened).toBe(false); + }), + ); +}); + +describe("outdatedHostUpdateConfirmation", () => { + const plan = { + hostLabel: "Build Mac", + fromVersion: "0.0.45", + targetVersion: "0.0.46", + } as const; + + it("names the desktop app relaunch for a desktop-managed host", () => { + const message = outdatedHostUpdateConfirmation({ ...plan, method: "desktop-managed" }); + expect(message).toContain("Pylon desktop app on Build Mac"); + expect(message).toContain("close and relaunch"); + expect(message).toContain("agent sessions running there will stop"); + }); + + it.each(["boot-service", "respawn"] as const)( + "warns that running sessions stop for a %s restart", + (method) => { + const message = outdatedHostUpdateConfirmation({ ...plan, method }); + expect(message).toContain("Update Pylon on Build Mac to 0.0.46?"); + expect(message).toContain("agent sessions running there will stop"); + expect(message).not.toContain("desktop app"); + }, + ); }); diff --git a/packages/client-runtime/src/connection/outdatedHostUpdate.ts b/packages/client-runtime/src/connection/outdatedHostUpdate.ts index de2e496d1..d1fcc2d58 100644 --- a/packages/client-runtime/src/connection/outdatedHostUpdate.ts +++ b/packages/client-runtime/src/connection/outdatedHostUpdate.ts @@ -2,6 +2,7 @@ import { ORCHESTRATION_PROTOCOL_VERSION, type EnvironmentId, type ExecutionEnvironmentDescriptor, + type ServerSelfUpdateCapability, type ServerSelfUpdateInput, type ServerSelfUpdateResult, WS_METHODS, @@ -41,6 +42,26 @@ export class OutdatedHostUpdateError extends Schema.TaggedError Effect.Effect, + /** Asked once the method is known; declining interrupts before anything is sent. */ + confirm: (plan: OutdatedHostUpdatePlan) => Effect.Effect, ) { const registry = yield* EnvironmentRegistry.EnvironmentRegistry; const resolver = yield* ConnectionResolver.ConnectionResolver; @@ -76,6 +99,17 @@ export const updateOutdatedHost = Effect.fn("clientRuntime.connection.updateOutd }); } + const confirmed = yield* confirm({ + hostLabel: descriptor.label, + method: capabilities.serverSelfUpdate, + fromVersion: descriptor.serverVersion, + targetVersion: input.targetVersion, + }); + if (!confirmed) { + return yield* Effect.interrupt; + } + yield* onStage("downloading"); + const result = yield* Effect.scoped( Effect.gen(function* () { const protocolContext = yield* Layer.build( diff --git a/packages/client-runtime/src/state/outdatedServerUpdate.ts b/packages/client-runtime/src/state/outdatedServerUpdate.ts index 36d7ba1f6..2b0c1bf17 100644 --- a/packages/client-runtime/src/state/outdatedServerUpdate.ts +++ b/packages/client-runtime/src/state/outdatedServerUpdate.ts @@ -6,7 +6,10 @@ import type * as HttpClient from "effect/unstable/http/HttpClient"; import type { Atom } from "effect/unstable/reactivity"; import type * as Socket from "effect/unstable/socket/Socket"; -import { updateOutdatedHost } from "../connection/outdatedHostUpdate.ts"; +import { + type OutdatedHostUpdatePlan, + updateOutdatedHost, +} from "../connection/outdatedHostUpdate.ts"; import type * as ConnectionResolver from "../connection/resolver.ts"; import type * as EnvironmentRegistry from "../connection/registry.ts"; import type * as RelayEnvironmentDiscovery from "../relay/discovery.ts"; @@ -22,6 +25,11 @@ export interface OutdatedServerUpdateTarget { readonly input: ServerSelfUpdateInput; /** From the host descriptor when known; an outdated host never delivers a server config. */ readonly fromVersion?: string; + /** + * Asked before the host restarts, with the method its descriptor reports. + * Resolving false cancels the update without touching the host. + */ + readonly confirm: (plan: OutdatedHostUpdatePlan) => Promise; } /** @@ -48,19 +56,29 @@ export function createOutdatedServerUpdateCommand( execute: (target: OutdatedServerUpdateTarget, atomRegistry) => { const stateAtom = serverUpdateStateAtom(target.environmentId); const targetVersion = target.input.targetVersion; - const fromVersion = target.fromVersion ?? targetVersion; + // The descriptor read before confirmation is the authoritative version. + let fromVersion = target.fromVersion ?? targetVersion; let currentStage: ServerUpdateStage = "downloading"; + let started = false; const setStage = (stage: ServerUpdateStage) => Effect.sync(() => { + started = true; currentStage = stage; atomRegistry.set(stateAtom, { status: "running", stage, fromVersion, targetVersion }); }); - return setStage(currentStage).pipe( - Effect.andThen(updateOutdatedHost(target.environmentId, target.input, setStage)), + // Progress starts once the user confirms; the host descriptor is read first. + return updateOutdatedHost(target.environmentId, target.input, setStage, (plan) => + Effect.sync(() => { + fromVersion = plan.fromVersion; + }).pipe(Effect.andThen(Effect.promise(() => target.confirm(plan)))), + ).pipe( Effect.onExit((exit) => Effect.sync(() => { if (Exit.isSuccess(exit) || Cause.hasInterruptsOnly(exit.cause)) { - atomRegistry.set(stateAtom, { status: "idle" }); + // A declined confirmation leaves any earlier failure visible. + if (started || Exit.isSuccess(exit)) { + atomRegistry.set(stateAtom, { status: "idle" }); + } return; } atomRegistry.set(stateAtom, { From 8df76f3d407af888f6d0aaf07f506830b8be957b Mon Sep 17 00:00:00 2001 From: Trevor Walker Date: Sat, 3 Oct 2026 14:03:02 -0600 Subject: [PATCH 12/12] test: build a typed registry mock for outdated-host update tests Replace the EnvironmentRegistry service casts with a helper that builds the full service and dies on unexpected calls. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/connection/outdatedHostUpdate.test.ts | 109 ++++++++++++------ 1 file changed, 71 insertions(+), 38 deletions(-) diff --git a/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts b/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts index e88e7ab25..b69dd4d44 100644 --- a/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts +++ b/packages/client-runtime/src/connection/outdatedHostUpdate.test.ts @@ -11,6 +11,7 @@ import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; import * as SubscriptionRef from "effect/SubscriptionRef"; import * as HttpClient from "effect/unstable/http/HttpClient"; import * as HttpClientResponse from "effect/unstable/http/HttpClientResponse"; @@ -18,7 +19,7 @@ import * as Socket from "effect/unstable/socket/Socket"; import type { ConnectionCatalogEntry } from "./catalog.ts"; import { orchestrationProtocolCompatibilityError } from "./compatibility.ts"; -import { PrimaryConnectionTarget } from "./model.ts"; +import { type NetworkStatus, PrimaryConnectionTarget } from "./model.ts"; import { type OutdatedHostUpdatePlan, outdatedHostUpdateConfirmation, @@ -117,6 +118,37 @@ class OutdatedHostSocket { } } +type RegistryService = EnvironmentRegistry.EnvironmentRegistry["Service"]; + +/** + * A full registry whose unused methods die, so a test only supplies the + * entries and the calls the update is expected to make. + */ +const makeTestRegistry = Effect.fn("makeTestRegistry")(function* ( + entries: ReadonlyMap, + overrides: Partial> = {}, +) { + const unused = (method: string) => Effect.die(new Error(`Unexpected registry.${method}.`)); + return EnvironmentRegistry.EnvironmentRegistry.of({ + entries: yield* SubscriptionRef.make(entries), + networkStatus: yield* SubscriptionRef.make("online"), + start: unused("start"), + register: () => unused("register"), + registerPlatform: () => unused("registerPlatform"), + reconcilePlatform: () => unused("reconcilePlatform"), + remove: () => unused("remove"), + removeRelayEnvironments: () => unused("removeRelayEnvironments"), + retryNow: () => unused("retryNow"), + setEnabled: overrides.setEnabled ?? (() => unused("setEnabled")), + setCompatibility: overrides.setCompatibility ?? (() => unused("setCompatibility")), + state: () => unused("state"), + stateChanges: () => Stream.die(new Error("Unexpected registry.stateChanges.")), + run: () => unused("run"), + runStream: () => Stream.die(new Error("Unexpected registry.runStream.")), + followStream: () => Stream.die(new Error("Unexpected registry.followStream.")), + }); +}); + describe("updateOutdatedHost", () => { it.effect("updates a protocol-1 host over a bare socket, then switches it back on", () => Effect.gen(function* () { @@ -125,9 +157,8 @@ describe("updateOutdatedHost", () => { const sockets: Array = []; let served: ExecutionEnvironmentDescriptor = descriptor(undefined, "0.0.45"); - const entries = yield* SubscriptionRef.make< - ReadonlyMap - >( + const calls: Array = []; + const registry = yield* makeTestRegistry( new Map([ [ TARGET.environmentId, @@ -140,15 +171,17 @@ describe("updateOutdatedHost", () => { }, ], ]), + { + setCompatibility: (_environmentId, error) => + Effect.sync(() => { + calls.push(`compatibility:${error === null ? "clear" : "block"}`); + }), + setEnabled: (_environmentId, enabled) => + Effect.sync(() => { + calls.push(`enabled:${enabled}`); + }), + }, ); - const calls: Array = []; - const registry = EnvironmentRegistry.EnvironmentRegistry.of({ - entries, - setCompatibility: (_environmentId: EnvironmentId, error: unknown) => - Effect.sync(() => calls.push(`compatibility:${error === null ? "clear" : "block"}`)), - setEnabled: (_environmentId: EnvironmentId, enabled: boolean) => - Effect.sync(() => calls.push(`enabled:${enabled}`)), - } as unknown as EnvironmentRegistry.EnvironmentRegistry["Service"]); const resolver = ConnectionResolver.ConnectionResolver.of({ prepare: () => Effect.die(new Error("The update must bypass the protocol gate.")), prepareForUpdate: () => @@ -244,18 +277,14 @@ describe("updateOutdatedHost", () => { Layer.mergeAll( Layer.succeed( EnvironmentRegistry.EnvironmentRegistry, - EnvironmentRegistry.EnvironmentRegistry.of({ - entries: yield* SubscriptionRef.make< - ReadonlyMap - >( - new Map([ - [ - TARGET.environmentId, - { target: TARGET, profile: Option.none(), enabled: false }, - ], - ]), - ), - } as unknown as EnvironmentRegistry.EnvironmentRegistry["Service"]), + yield* makeTestRegistry( + new Map([ + [ + TARGET.environmentId, + { target: TARGET, profile: Option.none(), enabled: false }, + ], + ]), + ), ), Layer.succeed( ConnectionResolver.ConnectionResolver, @@ -329,20 +358,24 @@ describe("updateOutdatedHost", () => { Layer.mergeAll( Layer.succeed( EnvironmentRegistry.EnvironmentRegistry, - EnvironmentRegistry.EnvironmentRegistry.of({ - entries: yield* SubscriptionRef.make< - ReadonlyMap - >( - new Map([ - [ - TARGET.environmentId, - { target: TARGET, profile: Option.none(), enabled: false }, - ], - ]), - ), - setCompatibility: () => Effect.sync(() => registryCalls.push("compatibility")), - setEnabled: () => Effect.sync(() => registryCalls.push("enabled")), - } as unknown as EnvironmentRegistry.EnvironmentRegistry["Service"]), + yield* makeTestRegistry( + new Map([ + [ + TARGET.environmentId, + { target: TARGET, profile: Option.none(), enabled: false }, + ], + ]), + { + setCompatibility: () => + Effect.sync(() => { + registryCalls.push("compatibility"); + }), + setEnabled: () => + Effect.sync(() => { + registryCalls.push("enabled"); + }), + }, + ), ), Layer.succeed( ConnectionResolver.ConnectionResolver,