Repository navigation
Fix relayflows shared broker startup - #8
Conversation
|
Warning Review limit reached
More reviews will be available in 7 minutes and 45 seconds. Learn how PR review limits work. Your organization has used up its prepaid credits, and credit purchases are no longer available. Enable the review add-on in the billing tab to keep reviews running — you're only billed for reviews past your plan's rate limits ($0.25/file). ⌛ How to resolve this issue?After more reviews become available, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans include higher PR review limits than trial, open-source, and free plans. In all cases, reviews become available again over time. During sustained high-volume PR review activity, CodeRabbit may temporarily slow when the next review becomes available. Please see our Fair Usage Limits Policy for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Plus Run ID: 📒 Files selected for processing (3)
✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Code Review
This pull request introduces shared broker coordination in WorkflowRunner to allow multiple workflow executions to reuse a single healthy broker connection instead of spawning multiple instances. The feedback highlights several robustness issues in the coordination logic, including a bug in isPidRunning where EPERM errors are incorrectly handled as the process not running, potential lock leaks if writing the lock owner file fails, and race conditions from writing lease and owner files directly instead of using atomic renames.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| function isPidRunning(pid: number): boolean { | ||
| try { | ||
| process.kill(pid, 0); | ||
| return true; | ||
| } catch { | ||
| return false; | ||
| } | ||
| } |
There was a problem hiding this comment.
On Unix systems, process.kill(pid, 0) throws an EPERM error if the target process exists but is owned by a different user (or has different privileges). Currently, catching any error and returning false will incorrectly report that such processes are not running. To fix this, explicitly check if the error code is EPERM and return true in that case.
| function isPidRunning(pid: number): boolean { | |
| try { | |
| process.kill(pid, 0); | |
| return true; | |
| } catch { | |
| return false; | |
| } | |
| } | |
| function isPidRunning(pid: number): boolean { | |
| try { | |
| process.kill(pid, 0); | |
| return true; | |
| } catch (err) { | |
| return (err as NodeJS.ErrnoException).code === 'EPERM'; | |
| } | |
| } |
| private async acquireSharedBrokerStartLock( | ||
| stateDir: string, | ||
| startupTimeoutMs: number | ||
| ): Promise<() => void> { | ||
| mkdirSync(stateDir, { recursive: true }); | ||
| const lockDir = path.join(stateDir, SHARED_BROKER_LOCK_DIRNAME); | ||
| const deadline = Date.now() + Math.max(startupTimeoutMs + 5_000, 10_000); | ||
| const staleAfterMs = Math.max(startupTimeoutMs * 2, 30_000); | ||
|
|
||
| for (;;) { | ||
| try { | ||
| mkdirSync(lockDir); | ||
| writeFileSync( | ||
| path.join(lockDir, 'owner.json'), | ||
| JSON.stringify({ pid: process.pid, createdAt: new Date().toISOString() }), | ||
| 'utf-8' | ||
| ); | ||
| return () => { | ||
| rmSync(lockDir, { recursive: true, force: true }); | ||
| }; | ||
| } catch (err) { | ||
| if ((err as NodeJS.ErrnoException).code !== 'EEXIST') { | ||
| throw err; | ||
| } | ||
| } | ||
|
|
||
| try { | ||
| const stat = statSync(lockDir); | ||
| if (Date.now() - stat.mtimeMs > staleAfterMs) { | ||
| rmSync(lockDir, { recursive: true, force: true }); | ||
| continue; | ||
| } | ||
| } catch { | ||
| continue; | ||
| } | ||
|
|
||
| if (Date.now() > deadline) { | ||
| throw new Error(`Timed out waiting for shared broker startup lock at ${lockDir}`); | ||
| } | ||
| await sleepMs(SHARED_BROKER_LOCK_POLL_MS); | ||
| } | ||
| } |
There was a problem hiding this comment.
There are two issues here:
- If
writeFileSyncfails (e.g., due to disk space or permissions), the newly createdlockDiris left behind, causing subsequent runs to block until the lock becomes stale. - The stale lock check only verifies
mtimeMs. If a process is still running but slow, deleting its lock can lead to concurrent broker spawns. Checking if the owner PID is still running before deleting the lock makes this much more robust.
private async acquireSharedBrokerStartLock(
stateDir: string,
startupTimeoutMs: number
): Promise<() => void> {
mkdirSync(stateDir, { recursive: true });
const lockDir = path.join(stateDir, SHARED_BROKER_LOCK_DIRNAME);
const deadline = Date.now() + Math.max(startupTimeoutMs + 5_000, 10_000);
const staleAfterMs = Math.max(startupTimeoutMs * 2, 30_000);
for (;;) {
try {
mkdirSync(lockDir);
try {
writeFileSync(
path.join(lockDir, 'owner.json'),
JSON.stringify({ pid: process.pid, createdAt: new Date().toISOString() }),
'utf-8'
);
} catch (writeErr) {
rmSync(lockDir, { recursive: true, force: true });
throw writeErr;
}
return () => {
rmSync(lockDir, { recursive: true, force: true });
};
} catch (err) {
if ((err as NodeJS.ErrnoException).code !== 'EEXIST') {
throw err;
}
}
try {
const ownerPath = path.join(lockDir, 'owner.json');
let isStale = false;
try {
const owner = JSON.parse(readFileSync(ownerPath, 'utf-8'));
if (typeof owner.pid === 'number' && !isPidRunning(owner.pid)) {
isStale = true;
}
} catch {
const stat = statSync(lockDir);
if (Date.now() - stat.mtimeMs > staleAfterMs) {
isStale = true;
}
}
if (isStale) {
rmSync(lockDir, { recursive: true, force: true });
continue;
}
} catch {
continue;
}
if (Date.now() > deadline) {
throw new Error(`Timed out waiting for shared broker startup lock at ${lockDir}`);
}
await sleepMs(SHARED_BROKER_LOCK_POLL_MS);
}
}| private createSharedBrokerLease( | ||
| stateDir: string, | ||
| connectionPath: string, | ||
| runId: string, | ||
| startedBroker: boolean | ||
| ): SharedBrokerLease { | ||
| const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME); | ||
| const ownerPath = path.join(stateDir, SHARED_BROKER_OWNER_FILENAME); | ||
| mkdirSync(leaseDir, { recursive: true }); | ||
| const leasePath = path.join( | ||
| leaseDir, | ||
| `${process.pid}-${runId}-${randomBytes(4).toString('hex')}.json` | ||
| ); | ||
| writeFileSync( | ||
| leasePath, | ||
| JSON.stringify({ | ||
| pid: process.pid, | ||
| runId, | ||
| startedBroker, | ||
| createdAt: new Date().toISOString(), | ||
| }), | ||
| 'utf-8' | ||
| ); | ||
| return { stateDir, connectionPath, ownerPath, leasePath, startedBroker }; | ||
| } |
There was a problem hiding this comment.
Writing the lease file directly to its final path can cause concurrent processes calling countLiveSharedBrokerLeases to read a partially written or empty file, leading to JSON parsing errors and accidental deletion of the lease. Writing to a temporary file and atomically renaming it using renameSync prevents this race condition.
private createSharedBrokerLease(
stateDir: string,
connectionPath: string,
runId: string,
startedBroker: boolean
): SharedBrokerLease {
const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME);
const ownerPath = path.join(stateDir, SHARED_BROKER_OWNER_FILENAME);
mkdirSync(leaseDir, { recursive: true });
const leasePath = path.join(
leaseDir,
`${process.pid}-${runId}-${randomBytes(4).toString('hex')}.json`
);
const tempPath = `${leasePath}.tmp`;
writeFileSync(
tempPath,
JSON.stringify({
pid: process.pid,
runId,
startedBroker,
createdAt: new Date().toISOString(),
}),
'utf-8'
);
renameSync(tempPath, leasePath);
return { stateDir, connectionPath, ownerPath, leasePath, startedBroker };
}| private writeSharedBrokerOwner(lease: SharedBrokerLease): void { | ||
| const conn = readBrokerConnectionFile(lease.connectionPath); | ||
| writeFileSync( | ||
| lease.ownerPath, | ||
| JSON.stringify({ | ||
| pid: conn?.pid, | ||
| createdByPid: process.pid, | ||
| createdAt: new Date().toISOString(), | ||
| }), | ||
| 'utf-8' | ||
| ); | ||
| } |
There was a problem hiding this comment.
Writing the owner file directly can cause concurrent processes calling isWorkflowOwnedSharedBroker to read a partially written file, leading to JSON parsing errors. Writing to a temporary file and atomically renaming it using renameSync ensures thread-safe/process-safe visibility.
private writeSharedBrokerOwner(lease: SharedBrokerLease): void {
const conn = readBrokerConnectionFile(lease.connectionPath);
const tempPath = `${lease.ownerPath}.tmp`;
writeFileSync(
tempPath,
JSON.stringify({
pid: conn?.pid,
createdByPid: process.pid,
createdAt: new Date().toISOString(),
}),
'utf-8'
);
renameSync(tempPath, lease.ownerPath);
}
Summary
Verification
Summary by cubic
Ensure Relayflows reuses a healthy shared broker and starts it only once across concurrent runs. Prevents duplicate brokers, flaky startups, and accidental shutdowns while other runs are active.
connection.jsonusingHarnessDriverClient.connectand agetStatushealth check from@agent-relay/harness-driver.shutdownRelay()so only the owner stops the broker; others disconnect. CLI signals now route throughshutdownRelay().Written for commit 1227358. Summary will update on new commits.