merge PR #4839: fix(coding-agent): handle JS eval worker cwd conflicts via worker protocol
This commit is contained in:
@@ -213,8 +213,6 @@ export class WorkerCore {
|
||||
}
|
||||
|
||||
async #runOne(runId: string, code: string, filename: string, snapshot: SessionSnapshot): Promise<void> {
|
||||
const runtime = this.#ensureRuntime(snapshot);
|
||||
runtime.setCwd(snapshot.cwd);
|
||||
const active: ActiveRun = { runId, filename, pendingTools: new Map(), floatingRejections: [] };
|
||||
this.#runs.set(runId, active);
|
||||
const hooks: RuntimeHooks = {
|
||||
@@ -224,6 +222,8 @@ export class WorkerCore {
|
||||
};
|
||||
let result: RunResult;
|
||||
try {
|
||||
const runtime = this.#ensureRuntime(snapshot);
|
||||
runtime.setCwd(snapshot.cwd);
|
||||
const value = await runtime.run(code, filename, hooks, { runId, cwd: snapshot.cwd });
|
||||
runtime.displayValue(value, hooks);
|
||||
result = { type: "result", runId, ok: true };
|
||||
|
||||
@@ -0,0 +1,114 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { WorkerCore } from "@oh-my-pi/pi-coding-agent/eval/js/worker-core";
|
||||
import type {
|
||||
SessionSnapshot,
|
||||
Transport,
|
||||
WorkerInbound,
|
||||
WorkerOutbound,
|
||||
} from "@oh-my-pi/pi-coding-agent/eval/js/worker-protocol";
|
||||
|
||||
interface WorkerHarness {
|
||||
send(message: WorkerInbound): void;
|
||||
onMessage(handler: (message: WorkerOutbound) => void): () => void;
|
||||
}
|
||||
|
||||
function createWorkerHarness(): WorkerHarness {
|
||||
const hostListeners = new Set<(message: WorkerOutbound) => void>();
|
||||
const workerListeners = new Set<(message: WorkerInbound) => void>();
|
||||
const transport: Transport = {
|
||||
send: message => {
|
||||
queueMicrotask(() => {
|
||||
for (const listener of hostListeners) listener(message);
|
||||
});
|
||||
},
|
||||
onMessage: handler => {
|
||||
workerListeners.add(handler);
|
||||
return () => workerListeners.delete(handler);
|
||||
},
|
||||
close: () => {},
|
||||
};
|
||||
new WorkerCore(transport);
|
||||
return {
|
||||
send(message) {
|
||||
queueMicrotask(() => {
|
||||
for (const listener of workerListeners) listener(message);
|
||||
});
|
||||
},
|
||||
onMessage(handler) {
|
||||
hostListeners.add(handler);
|
||||
return () => hostListeners.delete(handler);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function waitForMessage(
|
||||
harness: WorkerHarness,
|
||||
predicate: (message: WorkerOutbound) => boolean,
|
||||
): Promise<WorkerOutbound> {
|
||||
const { promise, resolve } = Promise.withResolvers<WorkerOutbound>();
|
||||
let unsubscribe = (): void => {};
|
||||
unsubscribe = harness.onMessage(message => {
|
||||
if (!predicate(message)) return;
|
||||
unsubscribe();
|
||||
resolve(message);
|
||||
});
|
||||
return promise;
|
||||
}
|
||||
|
||||
async function initializeWorker(harness: WorkerHarness, snapshot: SessionSnapshot): Promise<void> {
|
||||
const ready = waitForMessage(harness, message => message.type === "ready");
|
||||
harness.send({ type: "init", snapshot });
|
||||
expect((await ready).type).toBe("ready");
|
||||
}
|
||||
|
||||
describe("WorkerCore", () => {
|
||||
it("reports same-realm cwd conflicts through the worker protocol", async () => {
|
||||
const first = createWorkerHarness();
|
||||
const second = createWorkerHarness();
|
||||
const cwd = process.cwd();
|
||||
await initializeWorker(first, { cwd, sessionId: "same-realm-first", localRoots: {} });
|
||||
await initializeWorker(second, { cwd, sessionId: "same-realm-second", localRoots: {} });
|
||||
|
||||
const gate = Promise.withResolvers<void>();
|
||||
const entered = Promise.withResolvers<void>();
|
||||
(globalThis as { __omp_worker_core_gate?: { entered(): void; wait: Promise<void> } }).__omp_worker_core_gate = {
|
||||
entered: () => entered.resolve(),
|
||||
wait: gate.promise,
|
||||
};
|
||||
try {
|
||||
first.send({
|
||||
type: "run",
|
||||
runId: "hold-first-runtime",
|
||||
code: "globalThis.__omp_worker_core_gate.entered(); await globalThis.__omp_worker_core_gate.wait;",
|
||||
filename: "[same-realm-first].js",
|
||||
snapshot: { cwd, sessionId: "same-realm-first", localRoots: {} },
|
||||
});
|
||||
await entered.promise;
|
||||
|
||||
const result = waitForMessage(
|
||||
second,
|
||||
message => message.type === "result" && message.runId === "overlap-second-runtime",
|
||||
);
|
||||
second.send({
|
||||
type: "run",
|
||||
runId: "overlap-second-runtime",
|
||||
code: "1 + 1;",
|
||||
filename: "[same-realm-second].js",
|
||||
snapshot: { cwd, sessionId: "same-realm-second", localRoots: {} },
|
||||
});
|
||||
|
||||
expect(await result).toMatchObject({
|
||||
type: "result",
|
||||
runId: "overlap-second-runtime",
|
||||
ok: false,
|
||||
error: { message: "Cannot set cwd while another same-realm JS runtime is running" },
|
||||
});
|
||||
} finally {
|
||||
gate.resolve();
|
||||
delete (globalThis as { __omp_worker_core_gate?: { entered(): void; wait: Promise<void> } })
|
||||
.__omp_worker_core_gate;
|
||||
first.send({ type: "close" });
|
||||
second.send({ type: "close" });
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user