diff --git a/apps/server/src/cli/connect.test.ts b/apps/server/src/cli/connect.test.ts index f3cc88d1b58a..317e8ee490e6 100644 --- a/apps/server/src/cli/connect.test.ts +++ b/apps/server/src/cli/connect.test.ts @@ -58,6 +58,7 @@ it.effect("does not install the relay client when the user declines the managed let installCalls = 0; const result = yield* acquireRelayClientForLink( { + prepare: Effect.die("link inspection must not prepare the relay client"), resolve: Effect.succeed({ status: "missing", version: RelayClient.CLOUDFLARED_VERSION, @@ -87,6 +88,7 @@ it.effect("installs the relay client after the user accepts the managed download const progress: Array = []; const result = yield* acquireRelayClientForLink( { + prepare: Effect.die("link inspection must not prepare the relay client"), resolve: Effect.succeed({ status: "missing", version: RelayClient.CLOUDFLARED_VERSION, @@ -125,6 +127,7 @@ it.effect("reuses an available relay client executable without prompting", () => let promptCalls = 0; const result = yield* acquireRelayClientForLink( { + prepare: Effect.die("link inspection must not prepare the relay client"), resolve: Effect.succeed(managedExecutable), install: Effect.die("unexpected install"), installWithProgress: () => Effect.die("unexpected install"), diff --git a/apps/server/src/cloud/ManagedEndpointRuntime.test.ts b/apps/server/src/cloud/ManagedEndpointRuntime.test.ts index 42c5d9b20e1d..747402cce21f 100644 --- a/apps/server/src/cloud/ManagedEndpointRuntime.test.ts +++ b/apps/server/src/cloud/ManagedEndpointRuntime.test.ts @@ -21,7 +21,8 @@ import * as ManagedEndpointRuntime from "./ManagedEndpointRuntime.ts"; const layerRelayClientAvailable = Layer.succeed( RelayClient.RelayClient, RelayClient.RelayClient.of({ - resolve: Effect.succeed({ + resolve: Effect.die("launch must prepare the relay client"), + prepare: Effect.succeed({ status: "available", executablePath: "cloudflared", source: "path", @@ -341,7 +342,20 @@ describe("CloudManagedEndpointRuntime", () => { return handle; }), ); - const runtime = yield* buildCloudManagedEndpointRuntime(spawner); + const relayClient = yield* RelayClient.RelayClient.pipe( + Effect.provide(layerRelayClientAvailable), + ); + let preparations = 0; + const runtime = yield* buildCloudManagedEndpointRuntime( + spawner, + Layer.succeed(RelayClient.RelayClient, { + ...relayClient, + prepare: Effect.sync(() => { + expect(killed).toHaveLength(spawned.length); + preparations += 1; + }).pipe(Effect.andThen(relayClient.prepare)), + }), + ); yield* runtime.applyConfig({ providerKind: "cloudflare_tunnel", @@ -377,6 +391,7 @@ describe("CloudManagedEndpointRuntime", () => { expect(spawned.map((command) => command.options.detached)).toEqual([false, false]); expect(spawned.map((command) => command.options.shell)).toEqual([false, false]); expect(killed).toEqual([100, 101]); + expect(preparations).toBe(2); expect(stopped).toEqual({ status: "disabled" }); }), ); @@ -732,7 +747,8 @@ describe("CloudManagedEndpointRuntime", () => { Layer.succeed( RelayClient.RelayClient, RelayClient.RelayClient.of({ - resolve: Effect.succeed({ + resolve: Effect.die("launch must prepare the relay client"), + prepare: Effect.succeed({ status: "missing", version: RelayClient.CLOUDFLARED_VERSION, }), diff --git a/apps/server/src/cloud/ManagedEndpointRuntime.ts b/apps/server/src/cloud/ManagedEndpointRuntime.ts index 6b54b7ee3bdf..58904f2baa1c 100644 --- a/apps/server/src/cloud/ManagedEndpointRuntime.ts +++ b/apps/server/src/cloud/ManagedEndpointRuntime.ts @@ -298,7 +298,7 @@ export const make = Effect.gen(function* () { yield* stopActive; - const executable = yield* relayClient.resolve; + const executable = yield* relayClient.prepare; if (executable.status !== "available") { return { status: "failed", diff --git a/packages/shared/src/relayClient.test.ts b/packages/shared/src/relayClient.test.ts index f4a6b29cb585..5f871908a914 100644 --- a/packages/shared/src/relayClient.test.ts +++ b/packages/shared/src/relayClient.test.ts @@ -1,11 +1,20 @@ import { sha256 } from "@noble/hashes/sha2"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { describe, expect, it } from "@effect/vitest"; +import * as Clock from "effect/Clock"; +import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; +import * as Scheduler from "effect/Scheduler"; +import * as Scope from "effect/Scope"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import { TestClock } from "effect/testing"; import * as ConfigProvider from "effect/ConfigProvider"; import * as Effect from "effect/Effect"; import * as Hex from "effect/encoding/Hex"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; +import * as Logger from "effect/Logger"; import * as Sink from "effect/Sink"; import * as Stream from "effect/Stream"; import { HttpClient, HttpClientResponse } from "effect/http"; @@ -25,7 +34,7 @@ const layerHostRuntime = (env: Record = {}) => ConfigProvider.layer(ConfigProvider.fromEnv({ env })), ); -function makeHandle(exitCode = 0) { +function makeHandle(exitCode = 0, output = "") { return ChildProcessSpawner.makeHandle({ pid: ChildProcessSpawner.ProcessId(100), exitCode: Effect.succeed(ChildProcessSpawner.ExitCode(exitCode)), @@ -33,7 +42,7 @@ function makeHandle(exitCode = 0) { kill: () => Effect.void, unref: Effect.succeed(Effect.void), stdin: Sink.drain, - stdout: Stream.empty, + stdout: Stream.make(new TextEncoder().encode(output)), stderr: Stream.empty, all: Stream.empty, getInputFd: () => Sink.drain, @@ -60,12 +69,728 @@ const layerSpawner = (commands: Array) => // The pinned Windows executable rejects --version but accepts the version subcommand. return makeHandle( ChildProcess.isStandardCommand(command) && command.args.includes("--version") ? 1 : 0, + `cloudflared version ${RelayClient.CLOUDFLARED_VERSION} (built 2026-05-01-0000 UTC)`, ); }), ), ); +const pinnedBinary = `cloudflared version ${RelayClient.CLOUDFLARED_VERSION} (built 2026-05-01-0000 UTC) +`; + +const forkAutomaticRepair = Effect.fn("test.forkAutomaticRepair")(function* ( + effect: Effect.Effect, +) { + const clock = yield* Clock.Clock; + const deadlineStarted = yield* Deferred.make(); + const fiber = yield* effect.pipe( + Effect.provideService(Clock.Clock, { + ...clock, + sleep: (duration) => + Effect.suspend(() => { + if (Duration.toMillis(duration) === 30_000) { + queueMicrotask(() => Deferred.doneUnsafe(deadlineStarted, Effect.void)); + } + return clock.sleep(duration); + }).pipe( + // TestClock registers sleep before suspending. Defer the signal until that + // registration, without allowing a cooperative yield between the two. + Effect.provideService(Scheduler.PreventSchedulerYield, true), + ), + }), + Effect.forkChild, + ); + yield* Deferred.await(deadlineStarted).pipe(Effect.timeout("10 seconds"), TestClock.withLive); + return fiber; +}); + +const makeManagedFixture = Effect.fn("test.makeManagedFixture")(function* ( + options: { + readonly clientScope?: Scope.Scope; + readonly probeExitCode?: number; + readonly responseBody?: string; + readonly responseStatus?: number; + readonly stallAt?: "response" | "body" | "validation"; + readonly validationExitCode?: number; + } = {}, +) { + const fileSystem = yield* FileSystem.FileSystem; + const baseDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-cloudflared-test-" }); + const managedDirectory = `${baseDir}/tools/cloudflared/${RelayClient.CLOUDFLARED_VERSION}/linux-x64`; + const managedPath = `${managedDirectory}/cloudflared`; + yield* fileSystem.makeDirectory(managedDirectory, { recursive: true }); + const stallStarted = yield* Deferred.make(); + const installCleanedUp = yield* Deferred.make(); + let stallCleanedUp = false; + const stall = Deferred.succeed(stallStarted, undefined).pipe( + Effect.andThen(Effect.never), + Effect.ensuring( + Effect.sync(() => { + stallCleanedUp = true; + }), + ), + ); + const requests: Array = []; + const locksDuringDownload: Array = []; + const commands: Array = []; + const warnings: Array = []; + const writeBinary = Effect.fn("test.writeBinary")(function* (file: string, content: string) { + yield* fileSystem.writeFileString(file, content); + yield* fileSystem.chmod(file, 0o755); + }); + const manager = yield* RelayClient.makeCloudflaredRelayClient({ + baseDir, + releaseAsset: { + url: "https://example.test/cloudflared", + sha256: Hex.encode(sha256(new TextEncoder().encode(pinnedBinary))), + archive: "binary", + }, + }).pipe( + Effect.provideService(Scope.Scope, options.clientScope ?? (yield* Effect.scope)), + Effect.provideService(FileSystem.FileSystem, { + ...fileSystem, + remove: (file, options) => + fileSystem + .remove(file, options) + .pipe( + Effect.tap(() => + file === `${managedPath}.lock` + ? Deferred.succeed(installCleanedUp, undefined) + : Effect.void, + ), + ), + }), + Effect.provideService( + HttpClient.HttpClient, + HttpClient.make((request) => + Effect.gen(function* () { + requests.push(request.url); + locksDuringDownload.push( + yield* fileSystem.exists(`${managedPath}.lock`).pipe(Effect.orDie), + ); + if (options.stallAt === "response" && requests.length === 1) return yield* stall; + const response = HttpClientResponse.fromWeb( + request, + new Response(options.responseBody ?? pinnedBinary, { + status: options.responseStatus ?? 200, + }), + ); + if (options.stallAt === "body" && requests.length === 1) { + Object.defineProperty(response, "arrayBuffer", { value: stall }); + } + return response; + }), + ), + ), + Effect.provideService( + ChildProcessSpawner.ChildProcessSpawner, + ChildProcessSpawner.make((command) => + Effect.gen(function* () { + if (!ChildProcess.isStandardCommand(command)) return yield* Effect.die("Unexpected pipe"); + commands.push(command); + expect(command.args).toEqual(["version"]); + if ( + options.stallAt === "validation" && + command.command !== managedPath && + requests.length === 1 + ) { + return yield* stall; + } + return makeHandle( + command.command === managedPath + ? (options.probeExitCode ?? 0) + : (options.validationExitCode ?? 0), + yield* fileSystem.readFileString(command.command), + ); + }), + ), + ), + ); + const captureWarnings = Logger.layer( + [ + Logger.make(({ logLevel, message }) => { + if (logLevel === "Warn") warnings.push(message); + }), + ], + { mergeWithExisting: false }, + ); + return { + baseDir, + commands, + manager, + stallStarted, + installCleanedUp, + stallCleanedUp: () => stallCleanedUp, + managedPath, + requests, + locksDuringDownload, + warnings, + writeBinary, + captureWarnings, + }; +}); + +const managedTestLayer = Layer.mergeAll(NodeServices.layer, layerHostRuntime()); + describe("RelayClient", () => { + it.effect.skipIf(windowsHost)( + "shares healthy checks across repeated and concurrent preparations", + () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary(fixture.managedPath, pinnedBinary); + const results = yield* Effect.all( + Array.from({ length: 5 }, () => fixture.manager.prepare), + { concurrency: "unbounded" }, + ); + const first = yield* fixture.manager.prepare; + expect(results).toEqual(Array.from({ length: 5 }, () => first)); + expect(fixture.commands).toHaveLength(1); + expect(fixture.requests).toHaveLength(0); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + it.effect.skipIf(windowsHost)("shares failed repairs until the retry cooldown expires", () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture({ responseStatus: 503 }); + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + yield* Effect.gen(function* () { + const results = yield* Effect.all( + Array.from({ length: 5 }, () => fixture.manager.prepare), + { concurrency: "unbounded" }, + ); + expect(results.every((result) => result.status === "available")).toBe(true); + yield* fixture.manager.prepare; + expect(fixture.commands).toHaveLength(2); + expect(fixture.requests).toHaveLength(1); + expect(fixture.warnings).toHaveLength(1); + yield* TestClock.adjust("299 seconds"); + yield* fixture.manager.prepare; + expect(fixture.requests).toHaveLength(1); + yield* TestClock.adjust("1 second"); + yield* fixture.manager.prepare; + expect(fixture.commands).toHaveLength(4); + expect(fixture.requests).toHaveLength(2); + expect(fixture.warnings).toHaveLength(2); + }).pipe(Effect.provide(fixture.captureWarnings)); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + it.effect.skipIf(windowsHost)( + "explicit installation retries during the automatic repair cooldown", + () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture({ responseStatus: 503 }); + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + yield* fixture.manager.prepare.pipe(Effect.provide(fixture.captureWarnings)); + const error = yield* fixture.manager.install.pipe(Effect.flip); + expect(error.reason).toBe("download_failed"); + expect(fixture.requests).toHaveLength(2); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + it.effect.skipIf(windowsHost)("rechecks a replaced file even when its size and mtime match", () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary(fixture.managedPath, pinnedBinary); + yield* fixture.manager.prepare; + const info = yield* fileSystem.stat(fixture.managedPath); + const replacement = `${fixture.managedPath}.replacement`; + yield* fixture.writeBinary(replacement, pinnedBinary.replace("2026.5.2", "2026.9.3")); + if (info.mtime._tag === "Some") + yield* fileSystem.utimes(replacement, info.mtime.value, info.mtime.value); + yield* fileSystem.rename(replacement, fixture.managedPath); + yield* fixture.manager.prepare; + expect(fixture.requests).toHaveLength(1); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + it.effect.skipIf(windowsHost)( + "invalidates a failed repair when the executable changes in place", + () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture({ responseStatus: 503 }); + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + yield* fixture.manager.prepare.pipe(Effect.provide(fixture.captureWarnings)); + yield* fixture.writeBinary(fixture.managedPath, pinnedBinary); + yield* fixture.manager.prepare; + yield* fixture.manager.prepare; + expect(fixture.commands).toHaveLength(3); + expect(fixture.requests).toHaveLength(1); + expect(fixture.warnings).toHaveLength(1); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + it.effect.skipIf(windowsHost)( + "invalidates the managed check when selection switches to an override", + () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary(fixture.managedPath, pinnedBinary); + yield* fixture.manager.prepare; + const overridePath = `${fixture.baseDir}/override`; + yield* fixture.writeBinary(overridePath, "user-selected version"); + expect( + yield* fixture.manager.prepare.pipe( + Effect.provideService( + ConfigProvider.ConfigProvider, + ConfigProvider.fromEnv({ env: { T3CODE_CLOUDFLARED_PATH: overridePath } }), + ), + ), + ).toMatchObject({ source: "override" }); + yield* fixture.manager.prepare; + expect(fixture.commands).toHaveLength(2); + expect(fixture.requests).toHaveLength(0); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + it.effect.skipIf(windowsHost)("inspects a stale managed binary without repairing it", () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(); + const oldBinary = "cloudflared version 2026.9.3"; + yield* fixture.writeBinary(fixture.managedPath, oldBinary); + expect(yield* fixture.manager.resolve).toMatchObject({ + status: "available", + source: "managed", + executablePath: fixture.managedPath, + }); + expect(fixture.commands).toEqual([]); + expect(fixture.requests).toEqual([]); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(oldBinary); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + for (const stallAt of ["response", "body", "validation"] as const) { + it.effect.skipIf(windowsHost)( + `falls back and cleans up a stalled ${stallAt} within 30 seconds`, + () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture({ stallAt }); + const oldBinary = "cloudflared version 2026.9.3"; + yield* fixture.writeBinary(fixture.managedPath, oldBinary); + const resolving = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + forkAutomaticRepair, + ); + yield* Deferred.await(fixture.stallStarted); + yield* TestClock.adjust("30 seconds"); + const result = yield* Fiber.join(resolving).pipe( + Effect.timeout("10 seconds"), + TestClock.withLive, + ); + expect(result).toMatchObject({ + status: "available", + source: "managed", + executablePath: fixture.managedPath, + }); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(oldBinary); + expect((yield* fileSystem.stat(fixture.managedPath)).mode & 0o111).not.toBe(0); + yield* Deferred.await(fixture.installCleanedUp).pipe( + Effect.timeout("10 seconds"), + TestClock.withLive, + ); + expect(fixture.stallCleanedUp()).toBe(true); + expect( + yield* fileSystem.readDirectory( + fixture.managedPath.slice(0, fixture.managedPath.lastIndexOf("/")), + ), + ).toEqual(["cloudflared"]); + expect(fixture.warnings).toHaveLength(1); + // A subsequent explicit install proves the semaphore and lock were released. + yield* fixture.manager.install; + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + for (const invalidation of ["cooldown expiry", "file change"] as const) { + it.effect.skipIf(windowsHost)( + `cools down timeouts behind an explicit install until ${invalidation}`, + () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture({ stallAt: "response" }); + const oldBinary = "cloudflared version 2026.9.3"; + yield* fixture.writeBinary(fixture.managedPath, oldBinary); + const installing = yield* fixture.manager.install.pipe(Effect.forkChild); + yield* Deferred.await(fixture.stallStarted); + const resolving = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + forkAutomaticRepair, + ); + const concurrentResolving = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + forkAutomaticRepair, + ); + yield* TestClock.adjust("29 seconds"); + expect(resolving.pollUnsafe()).toBeUndefined(); + yield* TestClock.adjust("1 second"); + expect( + yield* Fiber.join(resolving).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ + status: "available", + source: "managed", + executablePath: fixture.managedPath, + }); + yield* Fiber.join(concurrentResolving).pipe( + Effect.timeout("10 seconds"), + TestClock.withLive, + ); + // Join without advancing TestClock: cached preparation must not wait for the permit. + const cached = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + Effect.forkChild, + ); + expect( + yield* Fiber.join(cached).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ source: "managed", executablePath: fixture.managedPath }); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(oldBinary); + expect(fixture.warnings).toEqual([ + [ + expect.stringContaining("Keeping the existing binary"), + expect.objectContaining({ reason: "repair_timeout" }), + ], + ]); + expect(installing.pollUnsafe()).toBeUndefined(); + expect(fixture.stallCleanedUp()).toBe(false); + expect(yield* fileSystem.exists(`${fixture.managedPath}.lock`)).toBe(true); + expect(fixture.requests).toHaveLength(1); + expect(fixture.commands).toHaveLength(2); + if (invalidation === "cooldown expiry") { + yield* TestClock.adjust("299 seconds"); + const stillCached = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + Effect.forkChild, + ); + expect( + yield* Fiber.join(stillCached).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ source: "managed" }); + expect(fixture.warnings).toHaveLength(1); + yield* TestClock.adjust("1 second"); + } else { + yield* fixture.writeBinary(fixture.managedPath, `${oldBinary}\nchanged`); + } + const retrying = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + forkAutomaticRepair, + ); + yield* TestClock.adjust("29 seconds"); + expect(retrying.pollUnsafe()).toBeUndefined(); + yield* TestClock.adjust("1 second"); + expect( + yield* Fiber.join(retrying).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ source: "managed", executablePath: fixture.managedPath }); + expect(fixture.warnings).toHaveLength(2); + expect(installing.pollUnsafe()).toBeUndefined(); + expect(fixture.stallCleanedUp()).toBe(false); + expect(yield* fileSystem.exists(`${fixture.managedPath}.lock`)).toBe(true); + expect(fixture.requests).toHaveLength(1); + expect(fixture.commands).toHaveLength(2); + // Only the test owner interrupts the explicit installer, after proving fallback. + yield* Fiber.interrupt(installing); + yield* fixture.manager.install; + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + for (const source of ["override", "path", "missing"] as const) { + it.effect.skipIf(windowsHost)( + `resolves ${source} without waiting for an explicit install`, + () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture({ stallAt: "response" }); + const installing = yield* fixture.manager.install.pipe(Effect.forkChild); + yield* Deferred.await(fixture.stallStarted); + const executablePath = `${fixture.baseDir}/cloudflared`; + if (source !== "missing") + yield* fixture.writeBinary(executablePath, "user-selected version"); + const env = + source === "override" + ? { T3CODE_CLOUDFLARED_PATH: executablePath, PATH: "" } + : { PATH: fixture.baseDir }; + const result = yield* fixture.manager.prepare.pipe( + Effect.provideService(ConfigProvider.ConfigProvider, ConfigProvider.fromEnv({ env })), + Effect.timeout("10 seconds"), + TestClock.withLive, + ); + expect(result).toMatchObject( + source === "missing" + ? { status: "missing" } + : { status: "available", source, executablePath }, + ); + expect(installing.pollUnsafe()).toBeUndefined(); + expect(fixture.stallCleanedUp()).toBe(false); + expect(fixture.commands).toEqual([]); + expect(fixture.requests).toHaveLength(1); + yield* Fiber.interrupt(installing); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + for (const phase of ["staging", "activation"] as const) { + it.effect.skipIf(windowsHost)(`falls back before the stalled ${phase} rename completes`, () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const renameStarted = yield* Deferred.make(); + const finishRename = yield* Deferred.make(); + const fixture = yield* makeManagedFixture().pipe( + Effect.provideService(FileSystem.FileSystem, { + ...fileSystem, + rename: (from, to) => + Effect.gen(function* () { + yield* fileSystem.rename(from, to); + if ((phase === "staging" ? to : from).endsWith(".tmp")) { + // Model a completed filesystem side effect whose callback has not arrived. + yield* Deferred.succeed(renameStarted, undefined); + yield* Deferred.await(finishRename); + } + }), + }), + ); + const oldBinary = "cloudflared version 2026.9.3"; + yield* fixture.writeBinary(fixture.managedPath, oldBinary); + const resolving = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + forkAutomaticRepair, + ); + yield* Effect.gen(function* () { + yield* Deferred.await(renameStarted); + yield* TestClock.adjust("30 seconds"); + expect( + yield* Fiber.join(resolving).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ status: "available", executablePath: fixture.managedPath }); + const cached = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + Effect.forkChild, + ); + expect( + yield* Fiber.join(cached).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ status: "available", executablePath: fixture.managedPath }); + expect(fixture.requests).toHaveLength(1); + expect(yield* fileSystem.exists(`${fixture.managedPath}.lock`)).toBe(true); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe( + phase === "staging" ? oldBinary : pinnedBinary, + ); + }).pipe(Effect.ensuring(Deferred.succeed(finishRename, undefined))); + expect( + yield* Fiber.join(resolving).pipe(Effect.timeout("10 seconds"), TestClock.withLive), + ).toMatchObject({ status: "available", executablePath: fixture.managedPath }); + yield* fixture.manager.install.pipe(Effect.timeout("10 seconds"), TestClock.withLive); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + expect((yield* fileSystem.stat(fixture.managedPath)).mode & 0o111).not.toBe(0); + expect( + yield* fileSystem.readDirectory( + fixture.managedPath.slice(0, fixture.managedPath.lastIndexOf("/")), + ), + ).toEqual(["cloudflared"]); + expect(fixture.warnings).toHaveLength(1); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + it.effect.skipIf(windowsHost)("owns a timed-out activation until the client scope closes", () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const clientScope = yield* Scope.make(); + const renameStarted = yield* Deferred.make(); + const finishRename = yield* Deferred.make(); + const scopeClosing = yield* Deferred.make(); + const fixture = yield* makeManagedFixture({ clientScope }).pipe( + Effect.provideService(FileSystem.FileSystem, { + ...fileSystem, + rename: (from, to) => + Effect.gen(function* () { + if (from.endsWith(".tmp")) { + yield* Deferred.succeed(renameStarted, undefined); + yield* Deferred.await(finishRename); + } + yield* fileSystem.rename(from, to); + }), + }), + ); + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + yield* Effect.gen(function* () { + const preparing = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + forkAutomaticRepair, + ); + yield* Deferred.await(renameStarted); + yield* TestClock.adjust("30 seconds"); + yield* Fiber.join(preparing).pipe(Effect.timeout("10 seconds"), TestClock.withLive); + yield* Scope.addFinalizer(clientScope, Deferred.succeed(scopeClosing, undefined)); + const closing = yield* Scope.close(clientScope, Exit.void).pipe(Effect.forkChild); + yield* Deferred.await(scopeClosing); + expect(closing.pollUnsafe()).toBeUndefined(); + yield* Deferred.succeed(finishRename, undefined); + yield* Fiber.join(closing).pipe(Effect.timeout("10 seconds"), TestClock.withLive); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + expect( + yield* fileSystem.readDirectory( + fixture.managedPath.slice(0, fixture.managedPath.lastIndexOf("/")), + ), + ).toEqual(["cloudflared"]); + }).pipe( + Effect.ensuring( + Deferred.succeed(finishRename, undefined).pipe( + Effect.andThen(Scope.close(clientScope, Exit.void)), + ), + ), + ); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + for (const version of ["2026.9.3", "2026.8.2", "2026.5.20"]) { + it.effect.skipIf(windowsHost)(`repairs a managed ${version} binary before preparing it`, () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary( + fixture.managedPath, + `cloudflared version ${version} (built yesterday)`, + ); + const resolved = yield* fixture.manager.prepare; + expect(resolved).toMatchObject({ + status: "available", + source: "managed", + executablePath: fixture.managedPath, + }); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + expect(fixture.requests).toHaveLength(1); + expect(fixture.locksDuringDownload).toEqual([true]); + expect(yield* fileSystem.exists(`${fixture.managedPath}.lock`)).toBe(false); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + it.effect.skipIf(windowsHost)("uses a matching managed binary without downloading", () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary(fixture.managedPath, pinnedBinary); + expect(yield* fixture.manager.prepare).toMatchObject({ + status: "available", + source: "managed", + }); + expect(fixture.requests).toEqual([]); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + for (const failure of [ + { name: "download failure", responseStatus: 503, reason: "download_failed" }, + { name: "checksum mismatch", responseBody: "tampered", reason: "invalid_checksum" }, + { name: "validation failure", validationExitCode: 1, reason: "validation_failed" }, + ]) { + it.effect.skipIf(windowsHost)( + `keeps the existing managed binary and warns after a ${failure.name}`, + () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(failure); + const staleBinary = "cloudflared version 2026.9.3"; + yield* fixture.writeBinary(fixture.managedPath, staleBinary); + const resolved = yield* fixture.manager.prepare.pipe( + Effect.provide(fixture.captureWarnings), + ); + expect(resolved).toMatchObject({ + status: "available", + source: "managed", + executablePath: fixture.managedPath, + }); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(staleBinary); + expect(fixture.requests).toHaveLength(1); + expect(fixture.locksDuringDownload).toEqual([true]); + expect(fixture.warnings).toEqual([ + [ + expect.stringContaining("Keeping the existing binary"), + expect.objectContaining({ reason: failure.reason }), + ], + ]); + expect(yield* fileSystem.exists(`${fixture.managedPath}.lock`)).toBe(false); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + for (const probe of [ + { name: "unrecognized version output", output: "unknown" }, + { name: "an unsuccessful version command", output: pinnedBinary, probeExitCode: 1 }, + ]) { + it.effect.skipIf(windowsHost)(`repairs a managed binary with ${probe.name}`, () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(probe); + yield* fixture.writeBinary(fixture.managedPath, probe.output); + expect(yield* fixture.manager.prepare).toMatchObject({ source: "managed" }); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + expect(fixture.requests).toHaveLength(1); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + it.effect.skipIf(windowsHost)( + "serializes concurrent repairs before resolving the managed binary", + () => + Effect.gen(function* () { + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + const [first, second] = yield* Effect.all( + [fixture.manager.prepare, fixture.manager.prepare], + { + concurrency: "unbounded", + }, + ); + expect(first).toMatchObject({ status: "available", source: "managed" }); + expect(second).toEqual(first); + expect(fixture.requests).toHaveLength(1); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + + for (const source of ["override", "path"] as const) { + it.effect.skipIf(windowsHost)(`does not validate or replace a ${source} binary`, () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(); + const executablePath = `${fixture.baseDir}/cloudflared`; + yield* fixture.writeBinary(executablePath, "user-selected version"); + if (source === "override") { + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + } + const env = + source === "override" + ? { T3CODE_CLOUDFLARED_PATH: executablePath, PATH: "" } + : { PATH: fixture.baseDir }; + yield* Effect.gen(function* () { + expect(yield* fixture.manager.prepare).toMatchObject({ source, executablePath }); + expect(yield* fixture.manager.install).toMatchObject({ source, executablePath }); + }).pipe( + Effect.provideService(ConfigProvider.ConfigProvider, ConfigProvider.fromEnv({ env })), + ); + expect(yield* fileSystem.readFileString(executablePath)).toBe("user-selected version"); + expect(fixture.commands).toEqual([]); + expect(fixture.requests).toEqual([]); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + } + + it.effect.skipIf(windowsHost)("rechecks the managed version under the installation lock", () => + Effect.gen(function* () { + const fileSystem = yield* FileSystem.FileSystem; + const fixture = yield* makeManagedFixture(); + yield* fixture.writeBinary(fixture.managedPath, "cloudflared version 2026.9.3"); + const installed = yield* fixture.manager.installWithProgress((event) => + event.type === "progress" && event.stage === "waiting_for_lock" + ? fixture.writeBinary(fixture.managedPath, pinnedBinary).pipe(Effect.orDie) + : Effect.void, + ); + expect(installed).toMatchObject({ source: "managed" }); + expect(yield* fileSystem.readFileString(fixture.managedPath)).toBe(pinnedBinary); + expect(fixture.requests).toEqual([]); + }).pipe(Effect.scoped, Effect.provide(managedTestLayer)), + ); + it.effect.skipIf(windowsHost)( "resolves explicit overrides before managed and PATH executables", () => @@ -220,7 +945,7 @@ describe("RelayClient", () => { concurrency: "unbounded", }); expect(second).toEqual(first); - expect(commands).toHaveLength(1); + expect(commands.filter((command) => command.includes(".install-"))).toHaveLength(1); }).pipe( Effect.scoped, Effect.provide( diff --git a/packages/shared/src/relayClient.ts b/packages/shared/src/relayClient.ts index 57576cad32f3..af9d1e050846 100644 --- a/packages/shared/src/relayClient.ts +++ b/packages/shared/src/relayClient.ts @@ -8,13 +8,16 @@ import * as Context from "effect/Context"; import * as Crypto from "effect/Crypto"; import * as Data from "effect/Data"; import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; import * as Hex from "effect/encoding/Hex"; 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 PlatformError from "effect/PlatformError"; +import type * as Scope from "effect/Scope"; import * as Semaphore from "effect/Semaphore"; +import * as Stream from "effect/Stream"; import { HttpClient, HttpClientRequest, HttpClientResponse } from "effect/http"; import { ChildProcess, ChildProcessSpawner } from "effect/process"; import { HostProcessArchitecture, HostProcessPlatform } from "./hostProcess.ts"; @@ -98,6 +101,9 @@ const CLOUDFLARED_RELEASE_ASSETS: Readonly< }, }; +const AUTOMATIC_REPAIR_TIMEOUT = "30 seconds"; +const AUTOMATIC_REPAIR_RETRY_MS = 5 * 60 * 1_000; + const INSTALL_LOCK_RETRY_COUNT = 100; const INSTALL_LOCK_RETRY_DELAY = "100 millis"; const INSTALL_LOCK_STALE_MS = 5 * 60 * 1_000; @@ -125,6 +131,8 @@ export interface CloudflaredRelayClientOptions { export interface RelayClientShape { readonly resolve: Effect.Effect; + /** Check and repair the managed binary before starting a connector. */ + readonly prepare: Effect.Effect; readonly install: Effect.Effect; readonly installWithProgress: ( report: (event: RelayClientInstallProgressEvent) => Effect.Effect, @@ -179,13 +187,16 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function | FileSystem.FileSystem | HttpClient.HttpClient | Path.Path + | Scope.Scope > { + const scope = yield* Effect.scope; const crypto = yield* Crypto.Crypto; const fileSystem = yield* FileSystem.FileSystem; const httpClient = yield* HttpClient.HttpClient; const path = yield* Path.Path; const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; const installSemaphore = yield* Semaphore.make(1); + let managedCheck: { readonly identity: string; readonly retryAfter: number } | undefined; const platform = yield* HostProcessPlatform; const arch = yield* HostProcessArchitecture; const releaseAsset = options.releaseAsset ?? resolveReleaseAsset(platform, arch); @@ -199,6 +210,20 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function executableFileName(platform), ); + const managedIdentity = fileSystem.stat(managedPath).pipe( + Effect.map((info) => + [ + managedPath, + info.dev, + Option.getOrNull(info.ino), + info.size, + info.mode, + Option.map(info.mtime, (time) => time.getTime()).pipe(Option.getOrNull), + ].join(":"), + ), + Effect.orElseSucceed(() => undefined), + ); + const isExecutableFile = Effect.fn("cloudflared.isExecutableFile")(function* ( executablePath: string, ) { @@ -221,7 +246,7 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function return null; }); - const resolve: RelayClientShape["resolve"] = Effect.gen(function* () { + const resolveExecutable: RelayClientShape["resolve"] = Effect.gen(function* () { const config = yield* loadCloudflaredConfig; if (Option.isSome(config.executableOverride)) { return (yield* isExecutableFile(config.executableOverride.value)) @@ -277,6 +302,27 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function } }); + const isPinnedExecutable = Effect.fn("cloudflared.isPinnedExecutable")( + function* (existing: AvailableRelayClient) { + if (existing.source !== "managed") return true; + // The pinned Windows release accepts `version`, but not `--version`. + const child = yield* spawner.spawn( + ChildProcess.make(existing.executablePath, ["version"], { + shell: false, + stdout: "pipe", + stderr: "ignore", + }), + ); + const output = yield* child.stdout.pipe(Stream.decodeText, Stream.mkString); + const exitCode = Number(yield* child.exitCode); + const version = /^cloudflared version (\S+)(?:\s|$)/u.exec(output.trim())?.[1]; + return exitCode === 0 && version === CLOUDFLARED_VERSION; + }, + Effect.scoped, + Effect.timeout("5 seconds"), + Effect.orElseSucceed(() => false), + ); + const downloadAsset = Effect.fn("cloudflared.downloadAsset")(function* ( asset: CloudflaredReleaseAsset, report: (stage: RelayClientInstallProgressStage) => Effect.Effect, @@ -354,8 +400,8 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function report: (stage: RelayClientInstallProgressStage) => Effect.Effect, ) { yield* report("checking"); - const existing = yield* resolve; - if (existing.status === "available") return existing; + const existing = yield* resolveExecutable; + if (existing.status === "available" && (yield* isPinnedExecutable(existing))) return existing; const config = yield* loadCloudflaredConfig; if (Option.isSome(config.executableOverride)) { return yield* new RelayClientInstallError({ @@ -391,8 +437,9 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function }), ); return yield* Effect.gen(function* () { - const afterLock = yield* resolve; - if (afterLock.status === "available") return afterLock; + const afterLock = yield* resolveExecutable; + if (afterLock.status === "available" && (yield* isPinnedExecutable(afterLock))) + return afterLock; const tempDirectory = yield* fileSystem.makeTempDirectoryScoped({ directory: managedDirectory, @@ -426,15 +473,18 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function const stagedPath = `${managedPath}.${yield* crypto.randomUUIDv4}.tmp`; yield* report("activating"); - yield* fileSystem - .rename(executablePath, stagedPath) - .pipe(wrapInstallFailure("write_failed", "Could not stage the relay client.")); - yield* fileSystem - .rename(stagedPath, managedPath) - .pipe( - wrapInstallFailure("write_failed", "Could not activate the relay client."), - Effect.ensuring(fileSystem.remove(stagedPath, { force: true }).pipe(Effect.ignore)), - ); + // Wait for rename callbacks before cleanup, even if automatic repair times out. + yield* Effect.gen(function* () { + yield* fileSystem + .rename(executablePath, stagedPath) + .pipe(wrapInstallFailure("write_failed", "Could not stage the relay client.")); + yield* fileSystem + .rename(stagedPath, managedPath) + .pipe(wrapInstallFailure("write_failed", "Could not activate the relay client.")); + }).pipe( + Effect.ensuring(fileSystem.remove(stagedPath, { force: true }).pipe(Effect.ignore)), + Effect.uninterruptible, + ); return { status: "available", executablePath: managedPath, @@ -459,16 +509,100 @@ export const makeCloudflaredRelayClient = Effect.fn("cloudflared.make")(function }); const installWithProgress: RelayClientShape["installWithProgress"] = (report) => installSemaphore.withPermit( - installUnlocked((stage) => - report({ - type: "progress", - stage, - }), + Effect.sync(() => { + managedCheck = undefined; + }).pipe( + Effect.andThen( + installUnlocked((stage) => + report({ + type: "progress", + stage, + }), + ), + ), ), ); const install = installWithProgress(() => Effect.void); - return RelayClient.of({ resolve, install, installWithProgress }); + const hasManagedCheck = (identity: string | undefined, now: number) => + identity !== undefined && managedCheck?.identity === identity && now < managedCheck.retryAfter; + + const recordManagedFailure = (identity: string | undefined, now: number) => { + if (hasManagedCheck(identity, now)) return false; + managedCheck = + identity === undefined + ? undefined + : { identity, retryAfter: now + AUTOMATIC_REPAIR_RETRY_MS }; + return true; + }; + + const prepare: RelayClientShape["prepare"] = Effect.gen(function* () { + const existing = yield* resolveExecutable; + if (existing.status !== "available" || existing.source !== "managed") { + managedCheck = undefined; + return existing; + } + const selectedIdentity = yield* managedIdentity; + if (hasManagedCheck(selectedIdentity, yield* Clock.currentTimeMillis)) return existing; + let warnOnFailure = false; + const repair = installSemaphore.withPermit( + Effect.gen(function* () { + const identity = yield* managedIdentity; + const now = yield* Clock.currentTimeMillis; + if (hasManagedCheck(identity, now)) { + return existing; + } + return yield* installUnlocked(() => Effect.void).pipe( + // Publish the outcome before releasing the permit, including deadline interruption. + Effect.onExit((exit) => + Effect.gen(function* () { + if (exit._tag === "Success") { + const installedIdentity = yield* managedIdentity; + managedCheck = + installedIdentity === undefined + ? undefined + : { identity: installedIdentity, retryAfter: Infinity }; + } else { + const failedAt = yield* Clock.currentTimeMillis; + warnOnFailure = recordManagedFailure(identity, failedAt); + } + }), + ), + ); + }), + ); + return yield* Effect.acquireUseRelease( + Effect.forkIn(repair, scope), + (fiber) => Fiber.join(fiber).pipe(Effect.timeout(AUTOMATIC_REPAIR_TIMEOUT)), + // Return on deadline while activation finishes atomically in the client's scope. + (fiber) => Fiber.interrupt(fiber).pipe(Effect.forkIn(scope)), + ).pipe( + Effect.catch((error) => + Effect.gen(function* () { + // A queued or activating repair may not have published its outcome yet. + const shouldWarn = + error._tag === "TimeoutError" + ? recordManagedFailure(yield* managedIdentity, yield* Clock.currentTimeMillis) || + warnOnFailure + : warnOnFailure; + if (shouldWarn) { + yield* Effect.logWarning( + "Could not restore the pinned relay client. Keeping the existing binary.", + error._tag === "TimeoutError" + ? { + reason: "repair_timeout", + message: "Automatic relay client repair timed out.", + } + : { reason: error.reason, message: error.message }, + ); + } + return existing; + }), + ), + ); + }); + + return RelayClient.of({ resolve: resolveExecutable, prepare, install, installWithProgress }); }); export const layerCloudflared = (options: CloudflaredRelayClientOptions) =>