Skip to content
Open
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
9 changes: 9 additions & 0 deletions .github/workflows/relayflow-pr-proof.yml
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,10 @@ on:
description: Pull request number to prove
required: true
type: number
expected_head_sha:
description: Exact current PR head SHA to prove (required for a fresh manual dispatch)
required: true
type: string

permissions:
actions: read
Expand Down Expand Up @@ -47,9 +51,14 @@ jobs:
id: status
env:
GITHUB_TOKEN: ${{ github.token }}
# Do not interpolate workflow-dispatch input into `run`: Actions
# expands expressions before the shell parses this script. The Node
# guard validates the value after this quoted environment expansion.
EXPECTED_HEAD_SHA: ${{ github.event_name == 'workflow_dispatch' && inputs.expected_head_sha || '' }}
run: >-
node scripts/pr-proof/report-status.mjs start
--event "$GITHUB_EVENT_PATH"
--expected-head-sha "$EXPECTED_HEAD_SHA"
--github-output "$GITHUB_OUTPUT"

- name: Classify PR and validate declared case
Expand Down
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

- Cloud Daytona Fleet provisioning now requires and returns the exact provider sandbox UUID alongside the stable Cloud sandbox ID, enabling ID-bound inspection and cleanup after interrupted launches.

- Node startup and forced shutdown restrict orphan broker cleanup to the selected state directory and exclude PTY workers, preventing an isolated node from stopping other projects' agents.

## [12.0.0] - 2026-09-10

### Added
Expand Down
146 changes: 146 additions & 0 deletions packages/cli/src/cli/commands/core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -980,6 +980,148 @@ describe('registerCoreCommands', () => {
}
);

it('up with an isolated state directory never kills brokers or workers from an ancestor project', async () => {
const fs = createFsMock();
const runningPids = new Set([222, 333, 444, 555, 666, 777, 9001, 4242]);
let now = 0;
const execCommand = vi.fn(async (command: string) => {
if (command === 'ps aux') {
return {
stdout: [
'USER PID %CPU %MEM VSZ RSS TT STAT STARTED TIME COMMAND',
'user 222 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /tmp/project/bin/agent-relay-broker init --name project --persist',
'user 333 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project/peer-state --persist',
'user 444 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /tmp/project/bin/agent-relay-broker pty --agent-name chief -- claude up --state-dir /tmp/project/candidate-state',
'user 555 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /tmp/project/bin/agent-relay node up --state-dir /tmp/project/peer-state',
'user 666 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /bin/zsh -c agent-relay up --state-dir /tmp/project/candidate-state',
'user 777 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project/candidate-state --state-dir',
].join('\n'),
stderr: '',
};
}
return { stdout: 'fcwd\nn/tmp/project\n', stderr: '' };
});
const killImpl = vi.fn((pid: number, signal?: NodeJS.Signals | number) => {
if (signal === 0) {
if (runningPids.has(pid)) return;
throw new Error('not running');
}
runningPids.delete(pid);
});
const { program } = createHarness({
fs,
execCommand,
killImpl,
spawnedProcess: createSpawnedProcessMock({ pid: 9001 }),
nowImpl: vi.fn(() => now),
sleepImpl: vi.fn(async (ms: number) => {
now += ms;
fs.writeFileSync('/tmp/project/candidate-state/connection.json', connectionFile(4242));
}),
});
const exitCode = await runCommand(program, [
'up',
'--background',
'--state-dir',
'/tmp/project/candidate-state',
]);
expect(exitCode).toBe(0);
expect(killImpl.mock.calls.filter(([, signal]) => signal !== 0)).toEqual([]);
expect([...runningPids]).toEqual([222, 333, 444, 555, 666, 777, 9001, 4242]);
});

it('down --force with an isolated state directory kills only its broker and fails closed on ambiguity', async () => {
const runningPids = new Set([222, 333, 444, 555, 666]);
const execCommand = vi.fn(async (command: string) => {
if (command === 'ps aux') {
return {
stdout: [
'USER PID %CPU %MEM VSZ RSS TT STAT STARTED TIME COMMAND',
'user 222 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project/candidate-state --persist',
'user 333 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project/peer-state --persist',
'user 444 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /tmp/project/bin/agent-relay-broker pty --agent-name chief -- claude',
'user 555 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project/candidate-state --state-dir /tmp/project/peer-state',
'user 666 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /bin/zsh -c agent-relay up --state-dir /tmp/project/candidate-state',
].join('\n'),
stderr: '',
};
}
throw new Error(`unexpected command: ${command}`);
});
const killImpl = vi.fn((pid: number, signal?: NodeJS.Signals | number) => {
if (signal === 0) {
if (runningPids.has(pid)) return;
throw new Error('not running');
}
runningPids.delete(pid);
});
let now = 0;
const { program, deps } = createHarness({
execCommand,
killImpl,
nowImpl: vi.fn(() => now),
sleepImpl: vi.fn(async (ms: number) => {
now += ms;
}),
});

const exitCode = await runCommand(program, [
'down',
'--force',
'--state-dir',
'/tmp/project/candidate-state',
]);

expect(exitCode).toBeUndefined();
expect(killImpl).toHaveBeenCalledWith(222, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(333, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(444, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(555, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(666, 'SIGTERM');
expect([...runningPids]).toEqual([333, 444, 555, 666]);
expect(deps.log).toHaveBeenCalledWith('Cleaned up (was not running)');
});

it('down --force resolves a relative state directory from a nested broker cwd', async () => {
const runningPids = new Set([222, 333]);
const execCommand = vi.fn(async (command: string) => {
if (command === 'ps aux') {
return {
stdout: [
'USER PID %CPU %MEM VSZ RSS TT STAT STARTED TIME COMMAND',
'user 222 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir candidate-state --persist',
'user 333 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project/peer-state --persist',
].join('\n'),
stderr: '',
};
}
if (command === 'lsof -nP -a -p 222 -d cwd -Fn') {
return { stdout: 'fcwd\nn/tmp/project/nested\n', stderr: '' };
}
throw new Error(`unexpected command: ${command}`);
});
const killImpl = vi.fn((pid: number, signal?: NodeJS.Signals | number) => {
if (signal === 0) {
if (runningPids.has(pid)) return;
throw new Error('not running');
}
runningPids.delete(pid);
});
const { program } = createHarness({ execCommand, killImpl });

const exitCode = await runCommand(program, [
'down',
'--force',
'--state-dir',
'/tmp/project/nested/candidate-state',
]);

expect(exitCode).toBeUndefined();
expect(killImpl).toHaveBeenCalledWith(222, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(333, 'SIGTERM');
expect([...runningPids]).toEqual([333]);
});

it('down --force only kills actual orphaned broker executables for the project', async () => {
const runningPids = new Set([222, 444, 666]);
const execCommand = vi.fn(async (command: string) => {
Expand All @@ -994,6 +1136,8 @@ describe('registerCoreCommands', () => {
'khaliqgant 555 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /opt/bin/agent-relay-broker init --state-dir /tmp/project-other/.agentworkforce/relay --persist',
'khaliqgant 666 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /Users/test/.agentworkforce/relay/bin/agent-relay up',
'khaliqgant 777 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /Users/test/.agentworkforce/relay/bin/agent-relay status --wait-for=30',
'khaliqgant 888 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /tmp/project/bin/agent-relay-broker pty --agent-name chief -- claude',
'khaliqgant 999 0.0 0.0 1 1 ?? S 1:00PM 0:00.01 /tmp/project/bin/agent-relay-broker init --state-dir /tmp/project/peer-state --persist',
].join('\n'),
stderr: '',
};
Expand Down Expand Up @@ -1036,6 +1180,8 @@ describe('registerCoreCommands', () => {
expect(killImpl).not.toHaveBeenCalledWith(333, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(555, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(777, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(888, 'SIGTERM');
expect(killImpl).not.toHaveBeenCalledWith(999, 'SIGTERM');
expect(deps.warn).toHaveBeenCalledWith('Killing orphaned broker process (pid: 222)');
expect(deps.warn).toHaveBeenCalledWith('Killing orphaned broker process (pid: 444)');
expect(deps.warn).toHaveBeenCalledWith('Killing orphaned broker process (pid: 666)');
Expand Down
89 changes: 58 additions & 31 deletions packages/cli/src/cli/lib/broker-lifecycle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1016,11 +1016,18 @@ function isAttachedBrokerCliCommand(command: string): boolean {
if (!/(?:^|\s)up(?:\s|$)/.test(command) || /(?:^|\s)--background(?:\s|=|$)/.test(command)) {
return false;
}
return /(?:^|\s)(?:\S*agent-relay(?:\.js)?|\S*agent-relay-[^\s]+)(?:\s|$)/.test(command);
const words = command.trim().split(/\s+/);
const executable = commandExecutableBasename(command);
// A shell or worker prompt mentioning `agent-relay up` is not the CLI.
const cli = ['node', 'bun'].includes(executable) ? path.basename(words[1] ?? '') : executable;
return /^(?:agent-relay(?:\.js)?|agent-relay-(?:darwin|linux|win32)[\w.-]*)$/.test(cli);
}

function isBrokerProcessCommand(command: string): boolean {
return isBrokerExecutableCommand(command) || isAttachedBrokerCliCommand(command);
// PTY workers use the same executable as their broker. They are never
// orphan broker candidates, even when their executable lives in this project.
if (isBrokerExecutableCommand(command)) return command.trim().split(/\s+/)[1] === 'init';
return isAttachedBrokerCliCommand(command);
}

function escapeRegExp(value: string): string {
Expand All @@ -1029,27 +1036,36 @@ function escapeRegExp(value: string): string {

function commandHasBrokerName(command: string, brokerName: string): boolean {
const escapedName = escapeRegExp(brokerName);
return new RegExp(`(?:^|\\s)--name(?:\\s+|=)${escapedName}(?:\\s|$)`).test(command);
return new RegExp(`(?:^|\\s)--(?:name|broker-name|instance-name)(?:\\s+|=)${escapedName}(?:\\s|$)`).test(
command
);
}

function commandHasProjectRoot(command: string, projectRoot: string): boolean {
const escapedRoot = escapeRegExp(path.resolve(projectRoot));
return new RegExp(`(?:^|\\s|=|["'])${escapedRoot}(?:$|\\s|["']|${escapeRegExp(path.sep)})`).test(command);
function commandStateDirectory(command: string): string | null | undefined {
const flags = [...command.matchAll(/(?:^|\s)--state-dir(?=\s|=|$)/g)];
if (flags.length === 0) return undefined;
if (flags.length !== 1) return null;
const matches = [...command.matchAll(/(?:^|\s)--state-dir(?:\s+|=)(?:"([^"]+)"|'([^']+)'|(\S+))/g)];
// Ambiguous or incomplete process listings cannot establish ownership.
if (matches.length !== 1) return null;
return matches[0][1] ?? matches[0][2] ?? matches[0][3] ?? null;
}

async function processCwdMatchesProjectRoot(
processInfo: ProcessInfo,
projectRoot: string,
deps: CoreDependencies
): Promise<boolean> {
): Promise<string | null> {
try {
const cwdDetails = await deps.execCommand(`lsof -nP -a -p ${processInfo.pid} -d cwd -Fn`);
return cwdDetails.stdout
const matchingCwd = cwdDetails.stdout
.split('\n')
.filter((line) => line.startsWith('n'))
.some((line) => path.resolve(line.slice(1)) === projectRoot);
.map((line) => path.resolve(line.slice(1)))
.find((cwd) => cwd === projectRoot || cwd.startsWith(`${projectRoot}${path.sep}`));
return matchingCwd ?? null;
} catch {
return false;
return null;
}
}

Expand All @@ -1074,15 +1090,21 @@ async function terminateProcess(pid: number, deps: CoreDependencies, force: bool
}

async function killOrphanedBrokerProcesses(
projectRoot: string,
paths: CoreProjectPaths,
deps: CoreDependencies,
options?: { force?: boolean }
options?: { force?: boolean; brokerName?: string }
): Promise<{ matchedCount: number; killedCount: number }> {
let matchedCount = 0;
let killedCount = 0;
try {
const resolvedProjectRoot = path.resolve(projectRoot);
const brokerName = path.basename(resolvedProjectRoot) || 'project';
const resolvedProjectRoot = path.resolve(paths.projectRoot);
const stateDirectory = path.resolve(paths.dataDir);
const usesDefaultStateDirectory =
stateDirectory === path.join(resolvedProjectRoot, '.agentworkforce/relay');
const brokerName =
options?.brokerName ??
deps.env.AGENT_RELAY_BROKER_NAME ??
(path.basename(resolvedProjectRoot) || 'project');
const candidates: ProcessInfo[] = [];
try {
const processList = await deps.execCommand('ps aux');
Expand All @@ -1092,28 +1114,32 @@ async function killOrphanedBrokerProcesses(
.filter((process): process is ProcessInfo => process !== null)
.filter((process) => isBrokerProcessCommand(process.command));

const matchedPids = new Set<number>();
for (const processInfo of relayProcesses) {
if (commandHasProjectRoot(processInfo.command, resolvedProjectRoot)) {
candidates.push(processInfo);
matchedPids.add(processInfo.pid);
}
}

for (const processInfo of relayProcesses) {
if (matchedPids.has(processInfo.pid)) {
const declaredStateDirectory = commandStateDirectory(processInfo.command);
if (declaredStateDirectory !== undefined) {
if (declaredStateDirectory === null) continue;
if (path.isAbsolute(declaredStateDirectory)) {
if (path.resolve(declaredStateDirectory) === stateDirectory) candidates.push(processInfo);
} else {
const processCwd = await processCwdMatchesProjectRoot(processInfo, resolvedProjectRoot, deps);
if (processCwd && path.resolve(processCwd, declaredStateDirectory) === stateDirectory)
candidates.push(processInfo);
}
continue;
}
const cwdMatches = await processCwdMatchesProjectRoot(processInfo, resolvedProjectRoot, deps);
if (!cwdMatches) continue;
// A custom state directory must never adopt a process using the default
// directory. An executable or arbitrary argument under an ancestor
// project (including the user's home directory) proves no ownership.
if (!usesDefaultStateDirectory) continue;
const processCwd = await processCwdMatchesProjectRoot(processInfo, resolvedProjectRoot, deps);
if (processCwd !== resolvedProjectRoot) continue;
if (
isBrokerExecutableCommand(processInfo.command) &&
!commandHasBrokerName(processInfo.command, brokerName)
) {
continue;
}
candidates.push(processInfo);
matchedPids.add(processInfo.pid);
}
} catch {
// Expected if ps is unavailable; fall through to no matches.
Expand Down Expand Up @@ -1161,7 +1187,8 @@ async function waitForProcessExit(pid: number, timeoutMs: number, deps: CoreDepe

async function recoverHalfStartedBroker(
paths: CoreProjectPaths,
deps: CoreDependencies
deps: CoreDependencies,
brokerName?: string
): Promise<'running' | 'recovered' | 'clear' | 'blocked'> {
deps.fs.mkdirSync(paths.dataDir, { recursive: true });
const readiness = await waitForBrokerReadiness(paths, deps, 0, true);
Expand All @@ -1185,7 +1212,7 @@ async function recoverHalfStartedBroker(
return 'recovered';
}

const orphanCleanup = await killOrphanedBrokerProcesses(paths.projectRoot, deps, { force: true });
const orphanCleanup = await killOrphanedBrokerProcesses(paths, deps, { force: true, brokerName });
if (orphanCleanup.matchedCount > 0) {
if (orphanCleanup.killedCount < orphanCleanup.matchedCount) {
deps.error(
Expand Down Expand Up @@ -1660,7 +1687,7 @@ export async function runUpCommand(options: UpOptions, deps: CoreDependencies):
}

if (options.background) {
const preflight = await recoverHalfStartedBroker(paths, deps);
const preflight = await recoverHalfStartedBroker(paths, deps, options.brokerName);
if (preflight === 'running') {
const pid = readBrokerPid(paths.dataDir, deps);
deps.error(
Expand Down Expand Up @@ -1894,7 +1921,7 @@ export async function runUpCommand(options: UpOptions, deps: CoreDependencies):
// Kill any orphaned broker processes for this project that lost their PID
// files (e.g. user deleted .agentworkforce/relay/ while broker was running).
vlog(deps, options.verbose, 'Checking for orphaned broker processes...');
await killOrphanedBrokerProcesses(paths.projectRoot, deps);
await killOrphanedBrokerProcesses(paths, deps, { brokerName: options.brokerName });

const started = await startBrokerWithPortFallback(
paths,
Expand Down Expand Up @@ -2100,7 +2127,7 @@ export async function runDownCommand(options: DownOptions, deps: CoreDependencie
const conn = readBrokerConnectionFromFs(deps.fs, paths.dataDir);
if (!conn) {
if (options.force) {
await killOrphanedBrokerProcesses(paths.projectRoot, deps, { force: true });
await killOrphanedBrokerProcesses(paths, deps, { force: true });
cleanupBrokerFiles(paths, deps);
deps.log('Cleaned up (was not running)');
} else {
Expand Down
Loading
Loading