-
-
Notifications
You must be signed in to change notification settings - Fork 2.4k
test(cloud): verify Freestyle images through the private connection path #12014
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
6 commits
Select commit
Hold shift + click to select a range
9361303
test(cloud): verify Freestyle images through private client connections
austinywang 3a4c85e
test(cloud): reproduce stale socket readiness in image verifier
austinywang 54e83e1
fix(cloud): await client readiness events in image verification
austinywang 19cfbc3
Merge branch 'main' of https://github.com/manaflow-ai/cmux into fix/f…
austinywang 0c84489
test(cloud): cover startup interruption and cleanup failure bounds
austinywang 5473303
fix(cloud): bound verifier cleanup to provider conflict retries
austinywang File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,30 @@ | ||
| // Maintainer-only image verification support; never imported by app/API routes. | ||
| import { Effect, Schedule } from "effect"; | ||
| import { FreestyleApiError } from "freestyle"; | ||
| import { ProviderError } from "../services/vms/drivers/types"; | ||
|
|
||
| /** A provider conflict means deletion is blocked by an attachment still in use. */ | ||
| function isDeletionConflict(error: unknown): boolean { | ||
| const cause = error instanceof ProviderError ? error.cause : error; | ||
| return cause instanceof FreestyleApiError && cause.status === 409; | ||
| } | ||
|
|
||
| /** | ||
| * Only a successful provider delete confirms completion. The SDK exposes no | ||
| * detach-completion event, so retry explicit conflict responses with bounded | ||
| * exponential backoff. Auth/config failures fail immediately. An overall | ||
| * deadline also bounds an unresponsive request, including within scope cleanup. | ||
| * Errors retain only our resource label, never an upstream payload or credential. | ||
| */ | ||
| export function cleanupPrivateLinkResource(label: string, run: (signal: AbortSignal) => Promise<void>) { | ||
| return Effect.tryPromise({ try: run, catch: (error) => error }).pipe( | ||
| Effect.retry({ | ||
| while: isDeletionConflict, | ||
| schedule: Schedule.intersect(Schedule.exponential("100 millis"), Schedule.recurs(7)), | ||
| }), | ||
| Effect.interruptible, | ||
| Effect.timeoutFail({ duration: "30 seconds", onTimeout: () => new Error("Cleanup deadline exceeded") }), | ||
| Effect.mapError(() => new Error(`Cleanup failed: ${label}`)), | ||
| Effect.orDie, | ||
| ); | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,99 @@ | ||
| // Maintainer-only verification subprocesses; these diagnostics are not product CLI copy. | ||
| import { Effect } from "effect"; | ||
| import { spawn, type ChildProcess } from "node:child_process"; | ||
| import { StringDecoder } from "node:string_decoder"; | ||
|
|
||
| type ReadyEvent = { | ||
| event: "hub-ready" | "connection-snapshot"; | ||
| socket: string; | ||
| }; | ||
|
|
||
| /** Wait for the child owner's close event after sending a termination signal. */ | ||
| function terminate(child: ChildProcess, signal: NodeJS.Signals) { | ||
| return Effect.async<void>((resume) => { | ||
| if (child.exitCode !== null || child.signalCode !== null) { | ||
| resume(Effect.void); | ||
| return; | ||
| } | ||
| const closed = () => resume(Effect.void); | ||
| child.once("close", closed); | ||
| child.kill(signal); | ||
| return Effect.sync(() => child.off("close", closed)); | ||
| }); | ||
| } | ||
|
|
||
| /** Close owns completion; the deadline only escalates an unresponsive child. */ | ||
| function stop(child: ChildProcess) { | ||
| return terminate(child, "SIGTERM").pipe( | ||
| Effect.interruptible, | ||
| Effect.timeoutOption("2 seconds"), | ||
| Effect.flatMap((completed) => completed._tag === "Some" ? Effect.void : terminate(child, "SIGKILL")), | ||
| ); | ||
| } | ||
|
|
||
| /** Match only the expected owner's readiness event and the exact requested socket. */ | ||
| function isReady(line: string, expected: ReadyEvent): boolean { | ||
| try { | ||
| const event = JSON.parse(line) as Record<string, unknown>; | ||
| const key = expected.event === "hub-ready" ? "socket" : "local_socket"; | ||
| return event?.event === expected.event && event[key] === expected.socket; | ||
| } catch { | ||
| return false; | ||
| } | ||
| } | ||
|
|
||
| /** | ||
| * Observe stdout from spawn onward, retaining an early ready event until awaited. | ||
| * The readiness promise never rejects unobserved while enrollment is approved. | ||
| * Exit/error before readiness fails immediately; socket file existence is irrelevant. | ||
| */ | ||
| function launch(client: string, args: string[], expected: ReadyEvent) { | ||
| return Effect.async<{ child: ChildProcess; ready: Effect.Effect<void, Error> }, Error>((resume) => { | ||
| const child = spawn(client, args, { stdio: ["ignore", "pipe", "pipe"] }); | ||
| let complete!: (ready: boolean) => void; | ||
| const result = new Promise<boolean>((resolve) => { complete = resolve; }); | ||
| const decoder = new StringDecoder("utf8"); | ||
| let buffered = ""; | ||
| let settled = false; | ||
| const finish = (ready: boolean) => { | ||
| if (settled) return; | ||
| settled = true; | ||
| buffered = ""; | ||
| complete(ready); | ||
| }; | ||
| child.stdout!.on("data", (chunk: Buffer) => { | ||
| if (settled) return; // still drain subsequent connection events | ||
| buffered += decoder.write(chunk); | ||
| if (buffered.length > 1_048_576) { finish(false); return; } | ||
| let boundary: number; | ||
| while (!settled && (boundary = buffered.indexOf("\n")) !== -1) { | ||
| const line = buffered.slice(0, boundary); | ||
| buffered = buffered.slice(boundary + 1); | ||
| if (isReady(line, expected)) finish(true); | ||
| } | ||
| }); | ||
| // Diagnostics can carry invitation material; never echo raw child output. | ||
| child.stderr!.resume(); | ||
| child.once("close", () => finish(false)); | ||
| child.once("error", () => { | ||
| finish(false); | ||
| resume(Effect.fail(new Error("Could not start the verification client"))); | ||
| }); | ||
| child.once("spawn", () => resume(Effect.succeed({ | ||
| child, | ||
| ready: Effect.promise(() => result).pipe( | ||
| Effect.flatMap((ready) => ready ? Effect.void : Effect.fail(new Error("Private connection process exited or returned invalid readiness output"))), | ||
| Effect.timeoutFail({ duration: "30 seconds", onTimeout: () => new Error("Private connection did not announce readiness") }), | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| ), | ||
| }))); | ||
| }); | ||
| } | ||
|
|
||
| /** | ||
| * Own a headless verifier process until scope exit, including failed startup. | ||
| * acquireRelease masks interruption through acquisition and finalizer registration: | ||
| * cancellation between spawn() and its spawn event still acquires, then stops, the child. | ||
| */ | ||
| export function startPrivateLinkClient(client: string, args: string[], ready: ReadyEvent) { | ||
| return Effect.acquireRelease(launch(client, args, ready), ({ child }) => stop(child)); | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,135 @@ | ||
| #!/usr/bin/env bun | ||
| /** | ||
| * Maintainer-only test harness, not a shipped application command. | ||
| * Verify an image through the same private carrier as the Mac app. Local daemon | ||
| * readiness alone misses a stale client, VPC routing, or enrollment failure. | ||
| * | ||
| * bun scripts/verify-devbox-private-link.ts <snapshot-id> <cmux-tui-client> | ||
| * | ||
| * Requires FREESTYLE_API_KEY. Creates an isolated VM, VPC, and temporary tunnel; | ||
| * deletes all three, including on failure. Never modifies existing machines. | ||
| * Uses the supplied client on macOS or Linux, without installing a system VPN. | ||
| */ | ||
| import { Effect } from "effect"; | ||
| import { spawn } from "node:child_process"; | ||
| import { generateKeyPairSync, randomUUID } from "node:crypto"; | ||
| import { mkdtempSync, rmSync, writeFileSync } from "node:fs"; | ||
| import { tmpdir } from "node:os"; | ||
| import path from "node:path"; | ||
| import { cleanupPrivateLinkResource as cleanup } from "./devbox-private-link-cleanup"; | ||
| import { startPrivateLinkClient } from "./devbox-private-link-process"; | ||
| import { FreestyleProvider } from "../services/vms/drivers/freestyle"; | ||
|
|
||
| /** Convert provider failures to local operator stage labels without upstream payloads. */ | ||
| const attempt = <A>(label: string, run: (signal: AbortSignal) => Promise<A>) => | ||
| Effect.tryPromise({ try: run, catch: () => new Error(label) }); | ||
|
|
||
| /** Run one bounded client command and capture its machine-readable response. */ | ||
| function command(client: string, args: string[], label: string) { | ||
| return attempt(label, (signal) => new Promise<string>((resolve, reject) => { | ||
| const child = spawn(client, args, { signal, stdio: ["ignore", "pipe", "pipe"] }); | ||
| let output = ""; | ||
| child.stdout!.on("data", (chunk: Buffer) => { | ||
| output += chunk.toString(); | ||
| if (output.length > 1_048_576) { child.kill(); reject(new Error(label)); } | ||
| }); | ||
| child.stderr!.resume(); | ||
| child.once("error", reject); | ||
| child.once("close", (code) => code === 0 ? resolve(output) : reject(new Error(label))); | ||
| })).pipe(Effect.timeoutFail({ duration: "30 seconds", onTimeout: () => new Error(`${label}: timed out`) })); | ||
| } | ||
|
|
||
| /** Require a usable resource graph from the connected session. */ | ||
| function readSnapshot(client: string, socket: string) { | ||
| return Effect.gen(function* () { | ||
| const stdout = yield* command(client, ["--socket", socket, "--json", "session", "current", "snapshot"], "Session snapshot failed"); | ||
| const snapshot = yield* Effect.try(() => JSON.parse(stdout)); | ||
| if (!Array.isArray(snapshot.workspaces) || !Array.isArray(snapshot.terminals)) { | ||
| return yield* Effect.fail(new Error("Private session returned an invalid snapshot")); | ||
| } | ||
| }); | ||
| } | ||
|
|
||
| /** Verify private enrollment and reconnect using exclusively probe-owned resources. */ | ||
| function verify(image: string, client: string) { | ||
| return Effect.gen(function* () { | ||
| const rawProbe = yield* command(client, ["remote-probe", "--json"], "Client probe failed"); | ||
| const probe = yield* Effect.try(() => JSON.parse(rawProbe)); | ||
| if (probe.app !== "cmux-tui" || !Array.isArray(probe.capabilities) || !probe.capabilities.includes("wireguard-hub")) { | ||
| return yield* Effect.fail(new Error("This client lacks wireguard-hub. Update cmux NIGHTLY before diagnosing or replacing the Freestyle image.")); | ||
| } | ||
| if (!process.env.FREESTYLE_API_KEY) return yield* Effect.fail(new Error("Set FREESTYLE_API_KEY")); | ||
| const provider = new FreestyleProvider(); | ||
| const networking = provider.privateNetworking; | ||
| const root = yield* Effect.acquireRelease( | ||
| Effect.sync(() => mkdtempSync(path.join(tmpdir(), "cmux-image-link-"))), | ||
| (directory) => Effect.sync(() => rmSync(directory, { recursive: true, force: true })), | ||
| ); | ||
| const slug = `cmux-image-check-${randomUUID().slice(0, 12)}`; | ||
| const network = yield* Effect.acquireRelease( | ||
| attempt("Create diagnostic VPC failed", () => networking.ensureNetwork({ slug })), | ||
| (value) => cleanup(`VPC ${value.id}`, () => networking.deleteNetwork(value.id)), | ||
| ); | ||
| const vm = yield* Effect.acquireRelease( | ||
| attempt("Create diagnostic VM failed", () => provider.create({ image, network: { id: network.id }, displayName: slug })), | ||
| (value) => cleanup(`VM ${value.providerVmId}`, () => provider.destroy(value.providerVmId)), | ||
| ); | ||
| const { privateKey, publicKey } = generateKeyPairSync("x25519"); | ||
| const clientPublicKey = publicKey.export({ type: "spki", format: "der" }).subarray(-32).toString("base64"); | ||
| const privateBytes = privateKey.export({ type: "pkcs8", format: "der" }).subarray(-32).toString("base64"); | ||
| const { tunnel } = yield* Effect.acquireRelease( | ||
| attempt("Create diagnostic tunnel failed", () => networking.createTunnel({ slug, networkId: network.id, clientPublicKey })), | ||
| (value) => cleanup(`tunnel ${value.tunnel.id}`, () => networking.deleteTunnel(value.tunnel.id)), | ||
| ); | ||
| if (!/^PrivateKey\s*=/m.test(tunnel.clientConfig)) { | ||
| return yield* Effect.fail(new Error("Tunnel config has no private-key field")); | ||
| } | ||
| const config = tunnel.clientConfig.replace(/^PrivateKey\s*=.*$/m, `PrivateKey = ${privateBytes}`); | ||
| const configPath = path.join(root, "wg.conf"), hubSocket = path.join(root, "wg.sock"); | ||
| yield* Effect.try(() => writeFileSync(configPath, config, { mode: 0o600 })); | ||
| const hub = yield* startPrivateLinkClient(client, ["wg", "hub", "--config", configPath, "--socket", hubSocket], { event: "hub-ready", socket: hubSocket }); | ||
| yield* hub.ready; | ||
| const endpoint = yield* attempt("Read image attach bundle failed", () => provider.openCmuxRemote(vm.providerVmId, { clientCapabilities: probe.capabilities })); | ||
| if (!endpoint.invitation) return yield* Effect.fail(new Error("Fresh VM returned no enrollment invitation")); | ||
| const invitation = endpoint.invitation; | ||
| const invitePath = path.join(root, "invite"); | ||
| yield* Effect.try(() => writeFileSync(invitePath, invitation.uri, { mode: 0o600 })); | ||
| const args = ["remote", "connect", endpoint.route, "--state-dir", path.join(root, "identity"), | ||
| "--device-name", slug, "--headless", "--json", "--wireguard-hub", hubSocket, | ||
| "--connect-timeout-seconds", "30", "--reconnect-attempts", "1"]; | ||
| // Enroll once, then reconnect from the persisted client identity with no | ||
| // control-plane attach/approval call. This also tests older baked daemons. | ||
| yield* Effect.scoped(Effect.gen(function* () { | ||
| const socket = path.join(root, "first.sock"); | ||
| const connection = yield* startPrivateLinkClient(client, [...args, "--local-socket", socket, "--invite-file", invitePath], { event: "connection-snapshot", socket }); | ||
| yield* attempt("Approve image enrollment failed", () => provider.approveCmuxRemoteEnrollment(vm.providerVmId, invitation.invitationId)); | ||
| yield* connection.ready; | ||
| yield* readSnapshot(client, socket); | ||
| })); | ||
| const socket = path.join(root, "reconnect.sock"); | ||
| const connection = yield* startPrivateLinkClient(client, [...args, "--local-socket", socket], { event: "connection-snapshot", socket }); | ||
| yield* connection.ready; | ||
| yield* readSnapshot(client, socket); | ||
| console.log(JSON.stringify({ image, clientCommit: probe.build_identity, | ||
| daemonCommit: endpoint.daemonBuild?.commit, enrollment: "passed", reconnect: "passed", snapshot: "passed" })); | ||
| }); | ||
| } | ||
|
|
||
| const [image, client] = process.argv.slice(2); | ||
| if (!image || !client || !/^sh-[a-zA-Z0-9]+$/.test(image)) { | ||
| console.error("Usage: bun scripts/verify-devbox-private-link.ts <snapshot-id> <cmux-tui-client>"); | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| process.exit(1); | ||
| } | ||
| const controller = new AbortController(); | ||
| const interrupt = () => controller.abort(); | ||
| process.once("SIGINT", interrupt); | ||
| process.once("SIGTERM", interrupt); | ||
| try { | ||
| await Effect.runPromise(Effect.scoped(verify(image, path.resolve(client))), { signal: controller.signal }); | ||
| } catch (error: unknown) { | ||
| console.error(String(error)); | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| process.exitCode = 1; | ||
| } finally { | ||
| process.off("SIGINT", interrupt); | ||
| process.off("SIGTERM", interrupt); | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,73 @@ | ||
| import { expect, test } from "bun:test"; | ||
| import { Effect, Fiber, TestClock, TestContext } from "effect"; | ||
| import { FreestyleApiError } from "freestyle"; | ||
| import { ProviderError } from "../services/vms/drivers/types"; | ||
| import { cleanupPrivateLinkResource } from "../scripts/devbox-private-link-cleanup"; | ||
|
|
||
| const failure = (status: number, code: string) => new ProviderError( | ||
| "freestyle", "deletion failed", new FreestyleApiError(status, { code, message: "provider response" }), | ||
| ); | ||
|
|
||
| test("successful deletion completes without another provider call", async () => { | ||
| let calls = 0; | ||
| await Effect.runPromise(cleanupPrivateLinkResource("VM probe", async () => { calls++; })); | ||
| expect(calls).toBe(1); | ||
| }); | ||
|
|
||
| test("deletion conflicts require a successful provider response before completion", async () => { | ||
| let calls = 0; | ||
| await Effect.runPromise(Effect.gen(function* () { | ||
| const work = yield* Effect.fork(cleanupPrivateLinkResource("VPC probe", async () => { | ||
| if (++calls < 3) throw failure(409, "CONFLICT"); | ||
| })); | ||
| yield* TestClock.adjust("1 second"); | ||
| yield* Fiber.join(work); | ||
| }).pipe(Effect.provide(TestContext.TestContext))); | ||
| expect(calls).toBe(3); | ||
| }); | ||
|
|
||
| test("a permanent provider refusal is not retried", async () => { | ||
| let calls = 0; | ||
| const outcome = await Effect.runPromise(Effect.gen(function* () { | ||
| const work = yield* Effect.fork(cleanupPrivateLinkResource("tunnel probe", async () => { | ||
| calls++; | ||
| throw failure(403, "FORBIDDEN"); | ||
| })); | ||
| yield* TestClock.adjust("15 seconds"); | ||
| return yield* Fiber.await(work); | ||
| }).pipe(Effect.provide(TestContext.TestContext))); | ||
| expect(outcome._tag).toBe("Failure"); | ||
| expect(calls).toBe(1); | ||
| }); | ||
|
|
||
| test("a stalled provider call fails at the cleanup deadline and cancels its request", async () => { | ||
| let aborted = false; | ||
| const outcome = await Effect.runPromise(Effect.gen(function* () { | ||
| const work = yield* Effect.fork(cleanupPrivateLinkResource("VPC probe", (signal) => new Promise(() => { | ||
| signal.addEventListener("abort", () => { aborted = true; }, { once: true }); | ||
| })).pipe(Effect.uninterruptible)); | ||
| yield* TestClock.adjust("30 seconds"); | ||
| const result = yield* Fiber.poll(work); | ||
| yield* Fiber.interrupt(work); | ||
| return result; | ||
| }).pipe(Effect.provide(TestContext.TestContext))); | ||
| expect(outcome._tag).toBe("Some"); | ||
| if (outcome._tag === "Some") expect(outcome.value._tag).toBe("Failure"); | ||
| expect(aborted).toBe(true); | ||
| }); | ||
|
|
||
| test("scope interruption still runs resource cleanup", async () => { | ||
| const { Deferred } = await import("effect"); | ||
| let calls = 0; | ||
| await Effect.runPromise(Effect.gen(function* () { | ||
| const using = yield* Deferred.make<void>(); | ||
| const work = yield* Effect.fork(Effect.scoped(Effect.gen(function* () { | ||
| yield* Effect.acquireRelease(Effect.void, () => cleanupPrivateLinkResource("VM probe", async () => { calls++; })); | ||
| yield* Deferred.succeed(using, undefined); | ||
| yield* Effect.never; | ||
| }))); | ||
| yield* Deferred.await(using); | ||
| yield* Fiber.interrupt(work); | ||
| })); | ||
| expect(calls).toBe(1); | ||
| }); |
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.