Skip to content
Merged
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
180 changes: 180 additions & 0 deletions src/lib/onboard/forward-start.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -601,9 +601,189 @@ describe("runDetachedForwardStartWithDiagnostics", () => {
expect(result.reason).toBe("timeout");
expect(isPortListening).not.toHaveBeenCalled();
});

it("rejects an untracked live-port fallback while the expected row is dead (#7140)", () => {
let now = 0;
vi.spyOn(Date, "now").mockImplementation(() => now);
const fetchList = vi
.fn()
.mockReturnValue(forwardListWith([{ sandbox: "my-sandbox", port: 18789, status: "dead" }]));
const spawn = vi.fn().mockImplementation(({ stderr }: { stderr: number }) => {
fs.writeSync(
stderr,
"Could not discover backgrounded SSH process; forward may be running but is not tracked\n",
);
return { pid: 788 };
});
const isPortListening = vi.fn().mockReturnValue(true);
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => true);

const result = runDetachedForwardStartWithDiagnostics(
spawn,
fetchList,
{ port: 18789, sandboxName: "my-sandbox" },
{
overallTimeoutMs: 10_000,
pollIntervalMs: 500,
sleepMs: (ms) => {
now += ms;
},
isPortListening,
},
);

expect(result.ok).toBe(false);
expect(result.reason).toBe("dead-forward");
expect(isPortListening).not.toHaveBeenCalled();
expect(killSpy).toHaveBeenCalledWith(788, "SIGTERM");
});
});

describe("runDetachedForwardStartWithRetries", () => {
it("stops and retries an expected forward whose ANSI status stays dead (#7140)", () => {
let now = 0;
vi.spyOn(Date, "now").mockImplementation(() => now);
const spawn = vi.fn().mockReturnValueOnce({ pid: 41 }).mockReturnValueOnce({ pid: 42 });
const fetchList = vi
.fn()
.mockImplementation(() =>
spawn.mock.calls.length === 1
? forwardListWith([
{ sandbox: "my-sandbox", port: 18789, status: "\u001B[31mdead\u001B[0m" },
])
: forwardListWith([{ sandbox: "my-sandbox", port: 18789 }]),
);
const beforeRetry = vi.fn();
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => true);

const result = runDetachedForwardStartWithRetries(
spawn,
fetchList,
{ port: 18789, sandboxName: "my-sandbox" },
beforeRetry,
{
overallTimeoutMs: 10_000,
pollIntervalMs: 500,
sleepMs: (ms) => {
now += ms;
},
isPortListening: vi.fn().mockReturnValue(false),
},
);

expect(result.ok).toBe(true);
expect(result.reason).toBe("ok");
expect(beforeRetry).toHaveBeenCalledOnce();
expect(spawn).toHaveBeenCalledTimes(2);
expect(killSpy).toHaveBeenCalledWith(41, "SIGTERM");
expect(killSpy.mock.invocationCallOrder[0]).toBeLessThan(
beforeRetry.mock.invocationCallOrder[0],
);
expect(killSpy).not.toHaveBeenCalledWith(42, expect.anything());
});

it("allows only one cleanup when the replacement forward also stays dead (#7140)", () => {
let now = 0;
vi.spyOn(Date, "now").mockImplementation(() => now);
const spawn = vi.fn().mockReturnValueOnce({ pid: 41 }).mockReturnValueOnce({ pid: 42 });
const fetchList = vi
.fn()
.mockReturnValue(forwardListWith([{ sandbox: "my-sandbox", port: 18789, status: "dead" }]));
const beforeRetry = vi.fn();
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => true);

const result = runDetachedForwardStartWithRetries(
spawn,
fetchList,
{ port: 18789, sandboxName: "my-sandbox" },
beforeRetry,
{
overallTimeoutMs: 10_000,
pollIntervalMs: 500,
sleepMs: (ms) => {
now += ms;
},
isPortListening: vi.fn().mockReturnValue(false),
},
);

expect(result.ok).toBe(false);
expect(result.reason).toBe("dead-forward");
expect(beforeRetry).toHaveBeenCalledOnce();
expect(spawn).toHaveBeenCalledTimes(2);
expect(killSpy).toHaveBeenCalledTimes(2);
expect(killSpy).toHaveBeenNthCalledWith(1, 41, "SIGTERM");
expect(killSpy).toHaveBeenNthCalledWith(2, 42, "SIGTERM");
});

it("accepts an expected forward that recovers during the dead-state grace (#7140)", () => {
let now = 0;
vi.spyOn(Date, "now").mockImplementation(() => now);
const spawn = vi.fn().mockReturnValue({ pid: 41 });
const fetchList = vi
.fn()
.mockReturnValueOnce(
forwardListWith([{ sandbox: "my-sandbox", port: 18789, status: "dead" }]),
)
.mockReturnValue(forwardListWith([{ sandbox: "my-sandbox", port: 18789 }]));
const beforeRetry = vi.fn();
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => true);

const result = runDetachedForwardStartWithRetries(
spawn,
fetchList,
{ port: 18789, sandboxName: "my-sandbox" },
beforeRetry,
{
overallTimeoutMs: 10_000,
pollIntervalMs: 500,
sleepMs: (ms) => {
now += ms;
},
isPortListening: vi.fn().mockReturnValue(false),
},
);

expect(result.ok).toBe(true);
expect(result.reason).toBe("ok");
expect(beforeRetry).not.toHaveBeenCalled();
expect(spawn).toHaveBeenCalledOnce();
expect(killSpy).not.toHaveBeenCalled();
});

it("does not clean up dead rows for another sandbox or port (#7140)", () => {
let now = 0;
vi.spyOn(Date, "now").mockImplementation(() => now);
const spawn = vi.fn().mockReturnValue({ pid: 41 });
const fetchList = vi.fn().mockReturnValue(
forwardListWith([
{ sandbox: "my-sandbox", port: 18790, status: "dead" },
{ sandbox: "other-sandbox", port: 18789, status: "dead" },
]),
);
const beforeRetry = vi.fn();

const result = runDetachedForwardStartWithRetries(
spawn,
fetchList,
{ port: 18789, sandboxName: "my-sandbox" },
beforeRetry,
{
overallTimeoutMs: 3_000,
pollIntervalMs: 500,
sleepMs: (ms) => {
now += ms;
},
isPortListening: vi.fn().mockReturnValue(false),
},
);

expect(result.ok).toBe(false);
expect(result.reason).toBe("timeout");
expect(beforeRetry).not.toHaveBeenCalled();
expect(spawn).toHaveBeenCalledOnce();
});

it("retries after a port-conflict diagnostic, then succeeds", () => {
const fetchList = vi
.fn()
Expand Down
129 changes: 85 additions & 44 deletions src/lib/onboard/forward-start.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import path from "node:path";

import { compactText } from "../core/url-utils";
import { redact } from "../security/redact";
import { parseForwardList } from "../state/sandbox-session";
import { getOccupiedPorts } from "./dashboard-port";
import { cleanupTempDir, secureTempFile } from "./temp-files";

Expand Down Expand Up @@ -40,6 +41,7 @@ export interface DetachedForwardStartOutcome {
| "spawn-error"
| "timeout"
| "spawn-conflict"
| "dead-forward"
| "listener-ownership-conflict"
| "listener-start-failure";
}
Expand Down Expand Up @@ -163,6 +165,15 @@ function blockingSleepMs(ms: number): void {
});
}

// OpenShell's external background-forward registry can briefly report a newly
// registered row as `dead` while its SSH process settles. NemoClaw cannot fix
// that producer in this repository, so require the exact sandbox+port row to
// stay dead across several normal polls before using the existing
// sandbox-scoped stop/retry boundary. Remove this reconciliation once every
// supported OpenShell version either stops retaining persistent dead rows or
// exposes an atomic recovery operation.
const DEAD_FORWARD_GRACE_MS = 2_000;

/**
* Build a `DetachedForwardSpawnRunner` that spawns the given argv as a
* detached child, writing stdio to the file descriptors supplied by
Expand Down Expand Up @@ -363,6 +374,7 @@ export function runDetachedForwardStartWithDiagnostics(
const portProbeIntervalMs = Math.max(pollIntervalMs, 5_000);
let nextPortProbeAt = start;
let lastListSnapshot = "";
let deadForwardObservedAt: number | null = null;
while (Date.now() < deadline) {
let list = "";
try {
Expand All @@ -380,39 +392,62 @@ export function runDetachedForwardStartWithDiagnostics(
return { ok: true, diagnostic: readDiag(), pid, reason: "ok" };
}
const listedForAnotherSandbox = Boolean(listedOwner);
const expectedForwardIsDead = parseForwardList(list).some(
(entry) =>
entry.sandboxName === expect.sandboxName &&
entry.port === String(expect.port) &&
entry.status === "dead",
);
const diagSoFar = readDiag();
if (looksLikeForwardPortConflict(diagSoFar)) {
terminateDetachedForwardChild(pid);
return { ok: false, diagnostic: diagSoFar, pid, reason: "spawn-conflict" };
}
// A completed listener-start failure cannot recover through more list
// polling. Preserve the established ControlMaster exception from #6099
// only when forward-list ownership enumeration succeeded and openshell
// also emitted that narrower untracked-forward diagnostic. A live TCP
// port alone is not evidence that this attempt owns the listener.
const listenerOutcome = classifyListenerStartDiagnostic({
diagnostic: diagSoFar,
pid,
port: expect.port,
ownerLookupSucceeded: lastFetchError === null,
listedForAnotherSandbox,
isPortListening,
});
if (listenerOutcome) return listenerOutcome;
// Preserve the established "untracked forward" compatibility path
// (GitHub #6099). It requires openshell's narrow diagnostic, a successful
// list query with no foreign sandbox row, and a live local port. This is
// intentionally not widened to other diagnostics because the probe does
// not establish process identity.
if (
lastFetchError === null &&
!listedForAnotherSandbox &&
looksLikeUntrackedForward(diagSoFar) &&
Date.now() >= nextPortProbeAt
) {
nextPortProbeAt = Date.now() + portProbeIntervalMs;
if (isPortListening(expect.port)) {
return { ok: true, diagnostic: readDiag(), pid, reason: "ok-port-live" };
if (expectedForwardIsDead) {
deadForwardObservedAt ??= Date.now();
if (Date.now() - deadForwardObservedAt >= DEAD_FORWARD_GRACE_MS) {
terminateDetachedForwardChild(pid);
const deadSummary =
`forward on port ${expect.port} remained dead for ` +
`${DEAD_FORWARD_GRACE_MS}ms after startup`;
return {
ok: false,
diagnostic: diagSoFar ? `${deadSummary} ${diagSoFar}` : deadSummary,
pid,
reason: "dead-forward",
};
}
} else {
deadForwardObservedAt = null;
// A completed listener-start failure cannot recover through more list
// polling. Preserve the established ControlMaster exception from #6099
// only when forward-list ownership enumeration succeeded and openshell
// also emitted that narrower untracked-forward diagnostic. A live TCP
// port alone is not evidence that this attempt owns the listener.
const listenerOutcome = classifyListenerStartDiagnostic({
diagnostic: diagSoFar,
pid,
port: expect.port,
ownerLookupSucceeded: lastFetchError === null,
listedForAnotherSandbox,
isPortListening,
});
if (listenerOutcome) return listenerOutcome;
// Preserve the established "untracked forward" compatibility path
// (GitHub #6099). It requires openshell's narrow diagnostic, a successful
// list query with no foreign sandbox row, and a live local port. This is
// intentionally not widened to other diagnostics because the probe does
// not establish process identity.
if (
lastFetchError === null &&
!listedForAnotherSandbox &&
looksLikeUntrackedForward(diagSoFar) &&
Date.now() >= nextPortProbeAt
) {
nextPortProbeAt = Date.now() + portProbeIntervalMs;
if (isPortListening(expect.port)) {
return { ok: true, diagnostic: readDiag(), pid, reason: "ok-port-live" };
}
}
}
if (onProgress && Date.now() >= nextProgressAt) {
Expand Down Expand Up @@ -444,19 +479,20 @@ export function runDetachedForwardStartWithDiagnostics(

/**
* Retry the detached forward-start after an EADDRINUSE-style port conflict or
* a definitive listener-start failure. `beforePortConflictRetry` preserves the
* established conflict-recovery behavior. Listener-start failures retry
* without sandbox/port cleanup because OpenShell does not expose immutable
* attempt identity.
* a definitive listener-start failure. A persistently dead exact sandbox+port
* row receives one independent recovery after the sandbox-scoped cleanup
* callback. Listener-start failures retry without sandbox/port cleanup because
* OpenShell does not expose immutable attempt identity.
*/
export function runDetachedForwardStartWithRetries(
runDetachedSpawn: DetachedForwardSpawnRunner,
fetchForwardList: ForwardListFetcher,
expect: { port: number; sandboxName: string },
beforePortConflictRetry: () => void,
beforeRetryCleanup: () => void,
options: DetachedForwardStartOptions = {},
): DetachedForwardStartOutcome {
const maxRetries = options.maxRetries ?? 3;
let deadForwardRecoveryAvailable = true;
const isPortListening = options.isPortListening ?? probeLocalPortListening;
const runAttempt = (): DetachedForwardStartOutcome =>
isPortListening(expect.port)
Expand All @@ -467,17 +503,22 @@ export function runDetachedForwardStartWithRetries(
}
: runDetachedForwardStartWithDiagnostics(runDetachedSpawn, fetchForwardList, expect, options);
let attempt = runAttempt();
for (
let retries = 0;
!attempt.ok &&
((attempt.reason !== "listener-ownership-conflict" &&
looksLikeForwardPortConflict(attempt.diagnostic)) ||
attempt.reason === "listener-start-failure") &&
retries < maxRetries;
retries++
) {
if (looksLikeForwardPortConflict(attempt.diagnostic)) {
beforePortConflictRetry();
let standardRetries = 0;
while (!attempt.ok) {
if (attempt.reason === "dead-forward") {
if (!deadForwardRecoveryAvailable) break;
deadForwardRecoveryAvailable = false;
beforeRetryCleanup();
} else {
const isRetryableStandardFailure =
(attempt.reason !== "listener-ownership-conflict" &&
looksLikeForwardPortConflict(attempt.diagnostic)) ||
attempt.reason === "listener-start-failure";
if (!isRetryableStandardFailure || standardRetries >= maxRetries) break;
if (looksLikeForwardPortConflict(attempt.diagnostic)) {
beforeRetryCleanup();
}
standardRetries++;
}
attempt = runAttempt();
}
Expand Down
Loading