From f01a5b71bcfffc13bf020abba33c01f4cbc8b89f Mon Sep 17 00:00:00 2001 From: Mavdol Date: Fri, 17 Apr 2026 14:27:37 +0200 Subject: [PATCH 1/2] support multiple concurrent workers by keying processes by capsule path and working directory --- crates/capsule-sdk/javascript/src/run.ts | 131 ++++++++++-------- .../capsule-sdk/python/src/capsule/worker.py | 8 +- 2 files changed, 77 insertions(+), 62 deletions(-) diff --git a/crates/capsule-sdk/javascript/src/run.ts b/crates/capsule-sdk/javascript/src/run.ts index 030bdae..700f473 100644 --- a/crates/capsule-sdk/javascript/src/run.ts +++ b/crates/capsule-sdk/javascript/src/run.ts @@ -45,72 +45,80 @@ interface PendingRequest { reject: (err: Error) => void; } -let workerProcess: ChildProcess | null = null; -let workerCapsulePath: string | null = null; +// Worker registry keyed by "capsulePath|cwd" to match Python SDK behaviour +const workerRegistry = new Map(); const pending = new Map(); -function getWorker(capsulePath: string): ChildProcess { - if (workerProcess && (workerCapsulePath !== capsulePath || workerProcess.exitCode !== null)) { - workerProcess.kill(); - workerProcess = null; +function workerKey(capsulePath: string, cwd: string): string { + return `${capsulePath}|${cwd}`; +} + +function getWorker(capsulePath: string, cwd: string): ChildProcess { + const key = workerKey(capsulePath, cwd); + const existing = workerRegistry.get(key); + + if (existing && existing.exitCode === null) { + return existing; } - if (!workerProcess) { - const command = getCapsuleCommand(capsulePath); + if (existing) { + existing.kill(); + workerRegistry.delete(key); + } - let child: ChildProcess; - if (process.platform === 'win32') { - const comspec = process.env.comspec || 'cmd.exe'; - child = spawn(comspec, ['/d', '/s', '/c', command, 'worker'], { stdio: ['pipe', 'pipe', 'inherit'] }); - } else { - child = spawn(command, ['worker'], { stdio: ['pipe', 'pipe', 'inherit'] }); - } + const command = getCapsuleCommand(capsulePath); - const rl = createInterface({ input: child.stdout! }); - rl.on('line', (line) => { - let response: { id: string; output?: unknown; error?: string }; - try { - response = JSON.parse(line); - } catch { - return; - } + let child: ChildProcess; + if (process.platform === 'win32') { + const comspec = process.env.comspec || 'cmd.exe'; + child = spawn(comspec, ['/d', '/s', '/c', command, 'worker'], { cwd, stdio: ['pipe', 'pipe', 'inherit'] }); + } else { + child = spawn(command, ['worker'], { cwd, stdio: ['pipe', 'pipe', 'inherit'] }); + } - const request = pending.get(response.id); - if (!request) return; - pending.delete(response.id); + const rl = createInterface({ input: child.stdout! }); + rl.on('line', (line) => { + let response: { id: string; output?: unknown; error?: string }; + try { + response = JSON.parse(line); + } catch { + return; + } - if (response.error) { - request.reject(new Error(response.error)); - } else { - request.resolve(response.output as RunnerResult); - } - }); + const request = pending.get(response.id); + if (!request) return; + pending.delete(response.id); - child.on('exit', () => { - for (const [id, req] of pending) { - req.reject(new Error('Capsule worker process exited unexpectedly')); - pending.delete(id); - } - workerProcess = null; - }); + if (response.error) { + request.reject(new Error(response.error)); + } else { + request.resolve(response.output as RunnerResult); + } + }); - child.on('error', (err) => { - if ((err as NodeJS.ErrnoException).code === 'ENOENT') { - for (const [id, req] of pending) { - req.reject(new Error(`Capsule CLI not found. Use 'npm install -g @capsule-run/cli' to install it.`)); - pending.delete(id); - } - } - workerProcess = null; - }); + child.on('exit', () => { + workerRegistry.delete(key); + for (const [id, req] of pending) { + req.reject(new Error('Capsule worker process exited unexpectedly')); + pending.delete(id); + } + }); - workerProcess = child; - workerCapsulePath = capsulePath; + child.on('error', (err) => { + workerRegistry.delete(key); + const message = (err as NodeJS.ErrnoException).code === 'ENOENT' + ? `Capsule CLI not found. Use 'npm install -g @capsule-run/cli' to install it.` + : err.message; + for (const [id, req] of pending) { + req.reject(new Error(message)); + pending.delete(id); + } + }); - process.once('exit', () => workerProcess?.kill()); - } + workerRegistry.set(key, child); + process.once('exit', () => child.kill()); - return workerProcess; + return child; } function getCapsuleCommand(capsulePath: string): string { @@ -129,11 +137,12 @@ function writeArgsFile(args: string[]): string { // --- run() via persistent worker --- function runViaWorker(options: RunnerOptions): Promise { - const { file, args = [], mounts = [], capsulePath = 'capsule' } = options; + const { file, args = [], mounts = [], cwd, capsulePath = 'capsule' } = options; + const resolvedCwd = cwd || process.cwd(); const id = randomUUID(); - console.time('runViaWorker' + id) + return new Promise((resolve, reject) => { - const worker = getWorker(capsulePath); + const worker = getWorker(capsulePath, resolvedCwd); pending.set(id, { resolve, reject }); @@ -144,7 +153,6 @@ function runViaWorker(options: RunnerOptions): Promise { reject(new Error(`Failed to send task to worker: ${err.message}`)); } }); - console.timeEnd('runViaWorker' + id) }); } @@ -229,7 +237,14 @@ export async function run(options: RunnerOptions): Promise { try { return await runViaWorker(options); - } catch { + } catch (err) { + const msg = (err as Error).message ?? ''; + + const isTransport = + msg.includes('worker process exited') || + msg.includes('CLI not found') || + msg.includes('Failed to send task'); + if (!isTransport) throw err; return runViaSubprocess(options); } } diff --git a/crates/capsule-sdk/python/src/capsule/worker.py b/crates/capsule-sdk/python/src/capsule/worker.py index 8b1ad82..d86d876 100644 --- a/crates/capsule-sdk/python/src/capsule/worker.py +++ b/crates/capsule-sdk/python/src/capsule/worker.py @@ -77,8 +77,7 @@ async def send(self, file: str, args: list[str], mounts: list[str]) -> str: req_id = uuid.uuid4().hex request = json.dumps({"id": req_id, "file": file, "args": args, "mounts": mounts}) - loop = asyncio.get_event_loop() - future: asyncio.Future[str] = loop.create_future() + future: asyncio.Future[str] = asyncio.get_running_loop().create_future() self._pending[req_id] = future assert self._process and self._process.stdin @@ -100,9 +99,7 @@ async def close(self) -> None: _clients: dict[tuple[str, Optional[str]], _WorkerClient] = {} -_clients_lock: Optional[asyncio.Lock] = None -# Tracks capsule_path values where the binary was confirmed missing _unavailable: set[str] = set() @@ -113,9 +110,12 @@ def _mark_unavailable(capsule_path: str) -> None: def _is_unavailable(capsule_path: str) -> bool: return capsule_path in _unavailable +_clients_lock: Optional[asyncio.Lock] = None + def _get_lock() -> asyncio.Lock: global _clients_lock + if _clients_lock is None: _clients_lock = asyncio.Lock() return _clients_lock From a0a534abe6bfeb970a471623bd4d1dce5a544ea4 Mon Sep 17 00:00:00 2001 From: Mavdol Date: Fri, 17 Apr 2026 14:29:33 +0200 Subject: [PATCH 2/2] refactor: remove unused worker registry comment in run.ts --- crates/capsule-sdk/javascript/src/run.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/crates/capsule-sdk/javascript/src/run.ts b/crates/capsule-sdk/javascript/src/run.ts index 700f473..06791a2 100644 --- a/crates/capsule-sdk/javascript/src/run.ts +++ b/crates/capsule-sdk/javascript/src/run.ts @@ -45,7 +45,6 @@ interface PendingRequest { reject: (err: Error) => void; } -// Worker registry keyed by "capsulePath|cwd" to match Python SDK behaviour const workerRegistry = new Map(); const pending = new Map();