Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 24 additions & 10 deletions src/client/connect.ts
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,8 @@ export interface LinkClientCredential {
}

export interface ClientConnectDeps {
/** Abort enrollment network work and refuse subsequent writes; rollback still drains. */
signal?: AbortSignal;
fetchImpl?: typeof fetch;
now?: () => Date;
lifecycleLockDeps?: ClientLifecycleLockDeps;
Expand Down Expand Up @@ -539,6 +541,18 @@ export async function connectClient(
options: ConnectOptions,
deps: ClientConnectDeps = {},
): Promise<OcxClientConnectionConfig> {
deps.signal?.throwIfAborted();
const rawFetch = deps.fetchImpl ?? fetch;
const fetchImpl: typeof fetch = deps.signal ? Object.assign(async (...[input, init = {}]: Parameters<typeof fetch>) => {
deps.signal!.throwIfAborted();
const signals = [deps.signal, init.signal, input instanceof Request ? input.signal : undefined]
.filter((signal): signal is AbortSignal => signal != null);
return rawFetch(input, { ...init, signal: AbortSignal.any(signals), redirect: "manual" });
}, { preconnect: rawFetch.preconnect }) : rawFetch;
const assertActiveConnectingState = (fingerprint?: string) => {
deps.signal?.throwIfAborted();
assertConnectingState(fingerprint);
};
let serverUrl = "";
let managementUrl = "";
let linkAdmissionToken: string | null = null;
Expand Down Expand Up @@ -586,7 +600,7 @@ export async function connectClient(
throw new Error("link mode requires a valid local config port");
}
withClientLifecycleSync(() => withConfigMutationLockSync(() => {
assertConnectingState();
assertActiveConnectingState();
const externalProvider = currentExternalCodexModelProvider();
if (externalProvider) throw new Error("connect refused: an external Codex provider owns config.toml");
if (linkMode) {
Expand All @@ -603,7 +617,7 @@ export async function connectClient(
}
}), deps.lifecycleLockDeps);

const upstreamFetch = deps.fetchImpl ?? fetch;
const upstreamFetch = fetchImpl;
const readinessFetch = linkMode
? (async (input, init = {}) => {
const url = input instanceof Request ? input.url : String(input);
Expand All @@ -622,18 +636,18 @@ export async function connectClient(
managementUrl,
localGuiOrigin(),
options.credential.value,
{ fetchImpl: deps.fetchImpl },
{ fetchImpl },
);
cleanupCredential = { kind: "gui-session", value: session };
} else if (!linkMode && options.credential.kind === "admin") {
cleanupCredential = { kind: "admin", value: options.credential.value };
}
if (!linkMode) {
issued = await issueClientKey(managementUrl, cleanupCredential!, clientKeyName(), { fetchImpl: deps.fetchImpl });
issued = await issueClientKey(managementUrl, cleanupCredential!, clientKeyName(), { fetchImpl });
}

const initialFiles = withClientLifecycleSync(() => withConfigMutationLockSync(() => {
assertConnectingState(linkMode ? pendingConnectFingerprint! : undefined);
assertActiveConnectingState(linkMode ? pendingConnectFingerprint! : undefined);
const persisted = linkMode
? (() => {
const current = readServiceApiTokenState();
Expand All @@ -656,7 +670,7 @@ export async function connectClient(
if (!admissionToken) throw new Error("client admission credential unavailable");
const apiKeyId = linkMode ? (options.credential as LinkClientCredential).apiKeyId : issued!.id;
const catalog = await downloadClientCatalog(serverUrl, admissionToken, {
fetchImpl: deps.fetchImpl,
fetchImpl,
timeoutMs: options.catalogTimeoutMs,
});
// Fail closed BEFORE the write (#4207). The hub being reachable and the credential working
Expand All @@ -666,7 +680,7 @@ export async function connectClient(
// than writing one and restoring it afterwards.
assertClientCatalogCompatible(catalog.body, deps.catalogCompatibility);
writtenCatalogFingerprint = withClientLifecycleSync(() => withConfigMutationLockSync(() => {
assertConnectingState(persisted.fingerprint);
assertActiveConnectingState(persisted.fingerprint);
atomicWriteFile(DEFAULT_CATALOG_PATH, catalog.body);
return sha256(catalog.body);
}), deps.lifecycleLockDeps);
Expand All @@ -679,7 +693,7 @@ export async function connectClient(
routingTarget: target,
catalogPath: DEFAULT_CATALOG_PATH,
journalOwner: { kind: "client", apiKeyId },
beforeClientWrite: () => assertConnectingState(persisted.fingerprint),
beforeClientWrite: () => assertActiveConnectingState(persisted.fingerprint),
});
if (!preflight.success) throw new Error(preflight.message);

Expand All @@ -688,7 +702,7 @@ export async function connectClient(
routingTarget: target,
catalogPath: DEFAULT_CATALOG_PATH,
journalOwner: { kind: "client", apiKeyId },
beforeClientWrite: () => assertConnectingState(persisted.fingerprint),
beforeClientWrite: () => assertActiveConnectingState(persisted.fingerprint),
});
if (!injected.success || injected.status === "skipped") throw new Error(injected.message);
injectionCommitted = true;
Expand All @@ -715,7 +729,7 @@ export async function connectClient(
...(linkMode ? { transport: "link" as const, link: linkMetadata } : {}),
};
withClientLifecycleSync(() => withConfigMutationLockSync(() => {
assertConnectingState(persisted.fingerprint);
assertActiveConnectingState(persisted.fingerprint);
clearClientConnectPending(persisted.fingerprint);
commitClientConnection(connection);
committed = true;
Expand Down
54 changes: 35 additions & 19 deletions src/client/link-join.ts
Original file line number Diff line number Diff line change
Expand Up @@ -287,7 +287,7 @@ async function waitForReady(
}
const remaining = deadline - now();
if (remaining <= 0) throw new ClientLinkJoinError("join_tunnel_failed");
await sleep(Math.min(JOIN_TUNNEL_POLL_MS, remaining));
await Promise.race([tunnelExited, sleep(Math.min(JOIN_TUNNEL_POLL_MS, remaining))]);
}
}

Expand Down Expand Up @@ -360,29 +360,45 @@ export async function joinHome(deps: ClientLinkJoinDeps, input: { alias: string
throw new ClientLinkJoinError(code);
}

const enrollmentAbort = new AbortController();
let enrollmentFinished = false;
try {
if (!tunnel) throw new ClientLinkJoinError("join_tunnel_failed");
const connect = deps.connect ?? connectClient;
// Keep watching the tunnel until the connection commits: an exited tunnel
// must not let the issued key ride out to whatever next holds the port.
await Promise.race([
tunnel.exited.then(() => { throw new ClientLinkJoinError("join_tunnel_failed"); }),
connect({
serverUrl: `http://127.0.0.1:${tunnelPort}`,
managementUrl: `http://127.0.0.1:${tunnelPort}`,
credential: { kind: "link", apiKeyId: issued.apiKeyId, key: issued.key },
transport: "link",
link: { tunnelPort, linkId: issued.linkId },
selectedClients: deps.selectedClients ?? ["codex", "claude"],
managementTransport: "direct",
}, {
fetchImpl: deps.fetchImpl,
...deps.connectDeps,
}),
]);
// A tunnel exit cancels work; it is NOT a competing terminal result. The connect
// transaction alone decides commit versus rollback, so an exit queued immediately
// after commit cannot revoke a key that a connected client has already retained.
void tunnel.exited.then(() => {
if (!enrollmentFinished) enrollmentAbort.abort(new ClientLinkJoinError("join_tunnel_failed"));
});
const signal = deps.connectDeps?.signal
? AbortSignal.any([enrollmentAbort.signal, deps.connectDeps.signal]) : enrollmentAbort.signal;
// Observe a tunnel that exited after readiness before starting any enrollment write.
await Promise.resolve();
signal.throwIfAborted();
await connect({
serverUrl: `http://127.0.0.1:${tunnelPort}`,
managementUrl: `http://127.0.0.1:${tunnelPort}`,
credential: { kind: "link", apiKeyId: issued.apiKeyId, key: issued.key },
transport: "link",
link: { tunnelPort, linkId: issued.linkId },
selectedClients: deps.selectedClients ?? ["codex", "claude"],
managementTransport: "direct",
}, {
fetchImpl: deps.fetchImpl,
...deps.connectDeps,
signal,
});
} catch (error) {
// connectClient has drained its local rollback before rejecting. Only then can
// the tunnel and the remote key be compensated without racing a late writer.
const tunnelAborted = enrollmentAbort.signal.aborted;
enrollmentFinished = true;
await rollback(deps, issued.linkId, tunnel);
throw new ClientLinkJoinError(error instanceof ClientLinkJoinError ? error.code : "join_connect_failed");
throw new ClientLinkJoinError(tunnelAborted ? "join_tunnel_failed"
: error instanceof ClientLinkJoinError ? error.code : "join_connect_failed");
} finally {
enrollmentFinished = true;
}

await stopTunnel(tunnel);
Expand Down
2 changes: 2 additions & 0 deletions structure/remote-link.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,3 +63,5 @@ Codex keeps the standalone loopback routing: `routingTarget` in `src/client/conn
> Decision record: [ADR-6032](decisions/ADR-6032-link-relay-credential-boundary.md)

Regression coverage lives in `tests/clients/link-ssh-argv.test.ts`, `tests/clients/link-ssh-config.test.ts`, `tests/clients/link-tunnel-state.test.ts`, `tests/clients/link-store.test.ts`, `tests/clients/link-boundary.test.ts`, `tests/clients/link-routes.test.ts`, `tests/clients/client-link-connect.test.ts`, `tests/clients/client-link-relay.test.ts`, `tests/clients/client-machine-listener.test.ts`, `tests/clients/client-link-status.test.ts`, `tests/clients/client-link-runtime.test.ts`, `tests/codex-integration/injection-link-websocket.test.ts`, `tests/clients/link-supervisor.test.ts`, `tests/clients/link-status-projection.test.ts`, `tests/clients/link-admission-wait.test.ts`, `tests/clients/link-fingerprint.test.ts`, `tests/cli/cli-link.test.ts`, `tests/server/link-management-routes.test.ts`, `tests/server/link-join-route.test.ts`, `tests/server/port-reclaim.test.ts`, `tests/server/link-listener-lifecycle.test.ts`, `tests/clients/client-link-teardown.test.ts` and `gui/tests/remote-link.test.tsx`.

Enrollment cancellation in `src/client/connect.ts` reaches actual network requests and every subsequent write boundary. `src/client/link-join.ts` awaits that transaction's commit or completed rollback instead of racing a separate failure against its terminal result. An exit observed before commit aborts enrollment, drains local rollback, and only then revokes the issued key; an exit queued after the synchronous commit retains the committed link for restart. Readiness polling also observes tunnel exit while sleeping, rather than waiting out its deadline. `tests/server/link-join-route.test.ts` and `tests/clients/client-link-connect.test.ts` pin those boundaries.
77 changes: 77 additions & 0 deletions tests/clients/client-link-connect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import { getDefaultConfig, saveConfig } from "../../src/config";
import { clientConnectionSchema } from "../../src/config/schema/leaf-validators";
import { handleConnectCommand, handleDisconnectCommand } from "../../src/cli/connect";
import { readSecretBytes } from "../../src/cli/runtime-api";
import { joinHome } from "../../src/client/link-join";
import { quoteRemote, remoteOcxArgv } from "../../src/link/ssh-argv";
import { connectClient, routingTarget } from "../../src/client/connect";
import { readServiceApiTokenState } from "../../src/lib/service-secrets";
import { isLinkConnection, readClientConnectionState } from "../../src/client/state";
Expand Down Expand Up @@ -130,6 +132,81 @@ describe("client link connection contracts", () => {
});
});

test("cancellation during catalog download prevents late enrollment writes and drains token rollback", async () => {
await withLinkHome(async home => {
const prior = '{"models":[{"id":"prior"}]}\n';
writeFileSync(DEFAULT_CATALOG_PATH, prior);
const abort = new AbortController();
let cancelledFetch = false;
await expect(connectClient(linkOptions(), {
signal: abort.signal,
fetchImpl: async (input, init) => {
if (String(input).endsWith("/readyz")) return Response.json({
service: "opencodex", version: "0.0.0", uptime: 1, pid: 1, port: 34567,
status: "ready", protocol: 1, minimumClientProtocol: 1,
managementUrl: "http://127.0.0.1:34567",
});
abort.abort(new Error("fixture enrollment cancelled"));
cancelledFetch = init?.signal?.aborted === true;
// Even a fetch implementation returning after abort cannot authorize a write.
return Response.json({ models: [] });
},
lifecycleLockDeps: { lockPath: join(home, "lifecycle.sqlite") },
})).rejects.toThrow("fixture enrollment cancelled");
expect(cancelledFetch).toBe(true);
expect(readFileSync(DEFAULT_CATALOG_PATH, "utf8")).toBe(prior);
expect(readServiceApiTokenState()).toEqual({ kind: "absent" });
expect(readClientConnectionState()).toEqual({ kind: "disconnected" });
});
});

test("a real enrollment commit survives a tunnel exit queued before join completion", async () => {
await withLinkHome(async home => {
let exit!: (code: number) => void;
const exited = new Promise<number>(resolve => { exit = resolve; });
let revoked = 0, restarted = 0;
const result = await joinHome({
runner: {
run: async argv => {
if (argv.at(-1) === quoteRemote(remoteOcxArgv(["link", "revoke", "--link-id", linkId]))) {
revoked += 1;
return { code: 0, stdout: "", stderr: "" };
}
return { code: 0, stdout: JSON.stringify({ linkId, apiKeyId: "key-1", key, listenerPort: 45678 }), stderr: "" };
},
spawnTunnel: () => { throw new Error("unexpected real tunnel"); },
},
knownHostsFile: join(home, "known-hosts"),
confirmedHost: { alias: "home", fingerprint: "SHA256:fixture", probedAt: 1 },
now: () => 1, choosePort: async () => 34567,
writeState: () => {}, clearState: () => {}, readSidecar: () => null,
spawnTunnel: () => ({ pid: 123, exited, stop: async () => {} }),
scanListenPids: () => ({ ok: true, pids: [123] }),
selectedClients: ["claude"],
fetchImpl: async (input, init) => {
if (!String(input).endsWith("/readyz")) return Response.json({ models: [] });
if (!new Headers(init?.headers).has("x-opencodex-api-key")) return new Response(null, { status: 401 });
return Response.json({ service: "opencodex", version: "0.0.0", uptime: 1, pid: 123,
port: 34567, status: "ready", protocol: 1, minimumClientProtocol: 1,
managementUrl: "http://127.0.0.1:34567" });
},
connectDeps: { lifecycleLockDeps: { lockPath: join(home, "lifecycle.sqlite") },
catalogCompatibility: { supportedEfforts: () => new Set() } },
connect: async (options, deps) => {
const committed = await connectClient(options, deps);
exit(255);
return committed;
},
scheduleRestart: () => { restarted += 1; },
}, { alias: "home" });
expect(result).toEqual({ linkId, apiKeyId: "key-1" });
expect(readClientConnectionState()).toMatchObject({ kind: "connected", value: { link: { linkId } } });
expect(readServiceApiTokenState()).toMatchObject({ kind: "present", token: key });
expect(revoked).toBe(0);
expect(restarted).toBe(1);
});
});

test("catalog failure removes the pending link token and leaves config.client unset", async () => {
await withLinkHome(async home => {
await expect(connectClient(linkOptions(), {
Expand Down
Loading
Loading