From 27db84be13c43740021cc7e3b202edb0f6ad0b2a Mon Sep 17 00:00:00 2001 From: Daniel Desjardins Date: Mon, 14 Sep 2026 16:58:49 -0400 Subject: [PATCH 1/6] WIP: real separate-process localExec() evaluator for pythonmonkey NOT YET FULLY TESTED -- DCP services were down (planned outage) for the final verification pass. See PYTHONMONKEY_EVALUATOR_PLAN.md sec5e for the exact pickup sequence once services are back. Replaces the old in-process shared-JS-global simulation (localexec_patch/pm_localexec_setup.py, a personal patch script requiring manual per-script wiring) with a real separate-process evaluator matching how Node.js's own localExec() already works: dcp/_pm_evaluator/ spawns a genuine child process per job, connected over a real socket, speaking the same line-based protocol lib/standaloneWorker.js already expects. Proven end-to-end through real dcp-client job deployment this session (exec -> init -> preauth -> deploying -> listeners -> uploading, using the already-installed dcp-client-bundle.js's existing __pmEvaluatorCtor gate -- no monorepo change was even needed to prove this part). Also includes a real port of the still-needed local-file job-argument routing (dcp/api/job.py's new localExec()/_route_arguments_locally()/ _grant_local_file_origins()) and a Python 3.14 compatibility fix (dcp/dry/aio.py -- nest_asyncio.apply() was unconditional and silently breaks under 3.14; now only applied when a loop is already running, credit: Tom Tang's diagnosis). Deliberately NOT ported yet: the local-worker completion-detection and work-function error-propagation fixes from the old patch script -- open question (see plan doc sec4/5d) is whether the new real child-process evaluator already handles these correctly on its own via StandaloneWorker-compatible messaging, unconfirmed pending live testing. Co-Authored-By: Claude Sonnet 5 --- PYTHONMONKEY_EVALUATOR_PLAN.md | 458 ++++++++++++++++++++++++ dcp/_pm_evaluator/__init__.py | 19 + dcp/_pm_evaluator/_test_bootstrap.py | 62 ++++ dcp/_pm_evaluator/_test_exec_control.py | 48 +++ dcp/_pm_evaluator/_test_plumbing.py | 43 +++ dcp/_pm_evaluator/_test_real_job.py | 57 +++ dcp/_pm_evaluator/channel.py | 92 +++++ dcp/_pm_evaluator/child.py | 170 +++++++++ dcp/_pm_evaluator/evaluator.py | 103 ++++++ dcp/api/job.py | 196 ++++++++++ dcp/dry/aio.py | 31 +- dcp/initialization.py | 12 + 12 files changed, 1288 insertions(+), 3 deletions(-) create mode 100644 PYTHONMONKEY_EVALUATOR_PLAN.md create mode 100644 dcp/_pm_evaluator/__init__.py create mode 100644 dcp/_pm_evaluator/_test_bootstrap.py create mode 100644 dcp/_pm_evaluator/_test_exec_control.py create mode 100644 dcp/_pm_evaluator/_test_plumbing.py create mode 100644 dcp/_pm_evaluator/_test_real_job.py create mode 100644 dcp/_pm_evaluator/channel.py create mode 100644 dcp/_pm_evaluator/child.py create mode 100644 dcp/_pm_evaluator/evaluator.py diff --git a/PYTHONMONKEY_EVALUATOR_PLAN.md b/PYTHONMONKEY_EVALUATOR_PLAN.md new file mode 100644 index 0000000..355b00b --- /dev/null +++ b/PYTHONMONKEY_EVALUATOR_PLAN.md @@ -0,0 +1,458 @@ +# Real, separate-process `localExec()` support for pythonmonkey — plan doc + +Written to survive a session compaction or a fresh pickup with zero other +context. If you're reading this cold: read this whole document before +touching code. It supersedes any assumption that `localExec()` under +pythonmonkey works by sharing one JS global — it's being rebuilt to spawn +a real child process instead, matching how Node.js's `localExec()` already +works. + +--- + +## 0. Why this document exists + +`job.localExec()` under pythonmonkey was gotten working during an earlier +investigation (see `localexec_patch/STATUS.md` and +`localexec_patch/FIXES_SUMMARY.md` in `C:\Users\danie\DCP\`) via an +in-process simulation: pythonmonkey has one shared JS global, and the +"Supervisor" (job-management) code and the "sandboxed worker" (work +function execution) code were made to share it, with ~950 lines of +workarounds (`C:\Users\danie\DCP\localexec_patch\pm_localexec_setup.py`) +for everything that broke as a result (console/require/timer clobbering, +sandbox access-list masking nuking Supervisor globals, a timer-starvation +bug, a Worker-lifecycle completion-detection gap, silent error swallowing). + +That got real jobs completing end-to-end, but every test script had to +manually `sys.path.insert(...)`, `import pm_localexec_setup`, call +`pm_localexec_setup.install()`, and wire up 3-4 more per-job patch function +calls before `job.localExec()` would work at all. The user's own words: +**"you added a bunch of stuff and this is sloppy as hell."** + +The user showed the actual target API shape (see §2) and asked for +`localExec()` to work that cleanly, with zero manual wiring — which means +all of that machinery needs to move from a personal patch script into +bifrost2 (and, for some of it, the real `dcp` monorepo) as first-class, +automatic behavior. + +While investigating how to do that cleanly, a bigger finding emerged (see +§3): the real Node.js `localExec()` doesn't share a global at all — it +spawns a **separate OS process**. The user explicitly chose to build the +architecturally-correct separate-process version rather than just +formalize the shared-global hack. **That is the current, in-progress +effort this document tracks.** + +A first attempt at just the narrower `job.py` fix (not the full platform +work) was opened as bifrost2 PR #49 and **was closed by the user for being +low-quality** ("slop") — it excluded real bugs (thinking they were +local-only test scaffolding when they weren't) and didn't reflect the full +picture. Do not reopen or reference it as a model; this document and the +`pythonmonkey-platform-support` branch supersede it. + +--- + +## 1. Desired end state — the target API shape + +Exactly this, verbatim, with **no import beyond `json`/`dcp`, no manual +patch wiring, no `sys.path` hacks**: + +```python +import json +import dcp +dcp.init() + +# IDENTITY + +# INPUT SET +input_set = list('yelling!') + +# WORK FUNCTION +def work_function(letter): + dcp.progress() + return letter.upper() + +# COMPUTE FOR +job = dcp.compute_for(input_set, work_function) + +# COMPUTE GROUPS +job.computeGroups = [ + { 'joinKey':'demo', 'joinSecret':'dcp' }, + { 'joinKey':'public' } +]; + +# PUBLIC INFO +job.public.name = 'to-upper-case' +job.public.description = 'Minimal demonstration of a distributed job' +job.public.link = 'https://distributive.network' + +# EVENTS +job.on('readystatechange', lambda s: print(f"Ready State: {s}")) +job.on('accepted', lambda _: print(f" Job ID: {job.id}\n Awaiting results...")) +job.on('noProgress', lambda n: print(json.dumps(n, indent=4).replace('\\n', '\n'))) +job.on('error', lambda e: print(json.dumps(e, indent=4).replace('\\n', '\n'))) +job.on('nofunds', lambda n: print(json.dumps(n, indent=4).replace('\\n', '\n'))) +job.on('result', lambda r: print(json.dumps(r, indent=4).replace('\\n', '\n'))) + +# EXECUTION +results = job.localExec() # <-- the only line that differs from exec()+wait() + +# RESULT POST-PROCESSING +print(''.join(results)) +``` + +`dcp.init()` alone must make pythonmonkey a fully-working dcp-client +platform (for both `exec()` and `localExec()`). `job.localExec()` alone +must work for any work function (JS or pyodide/Python), including +correctly propagating a broken work function's real error. + +--- + +## 2. Current, real (not simulated) state of related work + +These are already done, merged/open, independent of this effort: + +- **PythonMonkey PR #509** — SpiderMonkey rebuilt to current mozilla-central + (157a1) + SharedArrayBuffer/Atomics enabled. + https://github.com/Distributive-Network/PythonMonkey/pull/509 +- **PythonMonkey PR #510** — real `WebSocket` builtin module added. + https://github.com/Distributive-Network/PythonMonkey/pull/510 +- **dcp monorepo MR !3323** — pythonmonkey no longer excluded from the + `['websocket','polling']` transport list (now uses WebSocket like every + other platform). https://gitlab.com/Distributed-Compute-Protocol/dcp/-/merge_requests/3323 +- **bifrost2 PR #49 — CLOSED, do not reuse.** Was a narrower `job.py`-only + fix, judged incomplete/low-quality by the user. + +None of the above four repos/PRs currently contain any of the +separate-process evaluator work below — that's all new, uncommitted, or +only in the branches noted in §6. + +--- + +## 3. The core architectural finding + +Read `src/dcp-client/worker/evaluators/node-localExec.js` in the dcp +monorepo (cloned locally at `C:\Users\danie\DCP\dcp-monorepo`, see §7 for +exact remote/commit info). It spawns a **child process** +(`child_process`, with `I_WANT_AN_INSECURE_DCP_WORKER` and +`DCP_SCHEDULER_LOCATION` env vars) connected over a socket/pipe — not an +in-process simulation. + +Critically: the **client-side wire-protocol handler is platform-agnostic**. +`lib/standaloneWorker.js` (in the `dcp-client` repo, cloned locally at +`C:\Users\danie\DCP\dcp-client`) exports `StandaloneWorker`/`workerFactory`, +which just needs any `readStream`/`writeStream` — nothing Node-specific. +The wire protocol is line-based, newline-delimited: + +- `LOG:` — debug/log line, informational only, no action taken +- `DIE:` — child is shutting down (or parent telling child to shut down) +- `MSG:` — a real message, JSON body has `type`: + - `type: "workerMessage"` — a `postMessage()` payload, `message` field + holds the actual JS value + - `type: "result"` — with `exception` present (or not) to signal a slice + succeeded/failed + +This exact protocol is **already correctly implemented** on the pythonmonkey +side — `pm_localexec_setup.py`'s `install()` function's `writeln`/`onreadln`/ +`die` globals are a from-scratch, in-process-only reimplementation of this +exact contract. That means the wire protocol doesn't need to be invented — +it needs to be connected to a **real socket** instead of fake in-process +dispatch. + +**Decision (made explicitly by the user): build the real separate-process +version, not a formalized version of the shared-global hack.** The payoff: +eliminates essentially all of "Category A" below (console/require/timer +clobbering fixes, access-lists masking bypass) since a genuinely separate +process has its own JS global — nothing to clobber, nothing to mask. + +--- + +## 4. Full fix inventory — what exists in `pm_localexec_setup.py` and where it needs to end up + +Source: `C:\Users\danie\DCP\localexec_patch\pm_localexec_setup.py` (963 +lines). Every function's exact docstring has real, hard-won reasoning — +read it before reimplementing, don't guess. + +### Category A — pythonmonkey platform bootstrap (not localExec-specific; needed for `exec()` too) + +| Function | Lines | What it does | +|---|---|---| +| `install()` | 65-195 | Builds a Worker-shaped evaluator constructor (`postMessage`/`onmessage`/`onerror`/`terminate`/`addEventListener`) matching what `Sandbox.start()` expects as a `SandboxConstructor`. Registers `globalThis.__pmEvaluatorCtor`. | +| `run_bootstrap_eagerly()` | 196-404 | Loads the 22-file sandbox bootstrap (BravoJS, access-lists, polyfills, pyodide-core, etc.) eagerly at top level rather than lazily nested (confirmed the lazy path stalls indefinitely). Also fixes 4 distinct global-clobbering bugs this causes (require, console, timers, XHR) by save/restore around the bootstrap call. | +| `fix_crypto_getrandomvalues()` | 452-484 | `crypto.getRandomValues` returns all zeros after the bootstrap runs (a stub wins via `Object.assign` source-ordering); restores a real one backed by `os.urandom`. | +| `fix_access_lists_masking()` | 485-539 | `access-lists.js`'s real sandbox-isolation logic nukes Supervisor-side globals (`dcpConfig` etc.) because pythonmonkey has no real separate realm; patches `Object.defineProperty` to no-op only this exact masking shape. | + +**With a real separate-process child, this entire category should mostly +become unnecessary** — there's no shared global left to clobber or mask. +**Verify this empirically once the child-process bootstrap is built — do +not assume it; test it.** It's possible some subset is still needed even +in a fresh child process (e.g. if the 22 bootstrap files have bugs +independent of global-sharing) — find out by testing, not by assumption. + +**Target home if still needed at all**: `dcp.init()` in bifrost2 +(`dcp/initialization.py`'s `init()` closure, following the exact existing +"XXX apply dcp-client hacks XXX" convention already there for the +`getProcessPath` hack — see that file, already read this session). + +### Category B — real `localExec()`-specific dcp-client bugs (still needed regardless of process architecture) + +| Function | Lines | What it does | +|---|---|---| +| `route_job_arguments_through_local_files()` | 594-761 | `localExec()` otherwise unconditionally uploads job arguments/slice values to the real scheduler (`addSlices()`) regardless of size — defeats the point of *local* exec. Mirrors what Node's own `localExec()` does: route through local temp files instead. **Confirmed still needed this session** — removing it made the job hit a real `uploading` network state. | +| `force_job_completion_when_done()` | 832-928 | The local single-sandbox Worker never emits a terminating `'stop'` (starts a second `describe` round nobody replies to). Synthesizes a `'stop'`/`'complete'` once all expected slice results have arrived. **Confirmed still needed this session** — removing it hung indefinitely, zero result events. | +| `raise_on_first_work_error()` | 762-813 | **Not currently wired into any test script at all** — a real, unexercised bug: a broken work function currently returns normally from `localExec()`, fires `'error'` with an empty payload, and smuggles the real exception into the results list at a misaligned index. Hooks the same `"workError"` protocol message (see §3.3 below) to reject the underlying promise properly instead. | + +Needed reason: these are dcp-client-level bugs in `localExec()`'s own +argument-upload path and completion detection — unrelated to whether the +worker is in-process or a real child process. Investigate whether a real +child process changes how completion/error signals arrive (it might — a +real child process may get to use `StandaloneWorker`'s already-correct +`'result'`/`'error'` event handling in `lib/standaloneWorker.js` lines +252-278, which already handles this properly for Node! **This may mean +Category B's `force_job_completion_when_done`/`raise_on_first_work_error` +become unnecessary too, IF the real child process correctly emits proper +`result`/exception messages the way `StandaloneWorker` already expects.** +This needs verification, not assumption — it's a promising simplification +but unconfirmed. + +**Target home**: `Job.localExec()` in `dcp/api/job.py` (bifrost2) — same +file/method PR #49 touched, this time complete and tested. + +### Category C — confirmed dead, already deleted from the canonical test script + +| Function | Lines | Status | +|---|---|---| +| `start_fetchtask_keepalive()` | 405-450 | **Confirmed obsolete this session** — the `setInterval`-based watchdog timer-starvation bug it worked around does not reproduce on the rebuilt SpiderMonkey (157a1)/rewritten `JobQueue`. Removed, retested, job still completes correctly. | +| `patch_delay_manager()` | 544-593 | Same confirmed-obsolete finding, same test. | + +**Do not port these anywhere.** They're gone for good, a real side benefit +of the PR #509 `JobQueue` rewrite this session — worth a mention in that +PR if not already there. + +### Also needed: three small `dcp-client-bundle.js` patches, now understood to be real *monorepo* source changes + +Source: `C:\Users\danie\DCP\localexec_patch\FIXES_SUMMARY.md` §3 (exact +diffs). These were previously applied as hand-patches to the *minified, +vendored* `dcp-client-bundle.js` — now traced to their real, editable +source in the `dcp` monorepo (see §7 for exact file paths): + +1. **§3.1 localExec() platform gate** (`src/dcp-client/job/index.js` + line ~645) — accept `"pythonmonkey"` as a valid `localExec()` platform, + analogous to `"nodejs"`, using a `pythonmonkeyEvaluatorFactory()` (to be + written, modeled on `nodeEvaluatorFactory()` in + `src/dcp-client/worker/evaluators/node-localExec.js`) as the + `SandboxConstructor`. +2. **§3.2 SocketIOTransport transports** — **already done**, MR !3323 (§2). +3. **§3.3 two `Sandbox.start()` onmessage hooks** — signals for + `"complete"`/`"workError"` protocol messages, currently implemented as a + hand-patch giving `pm_localexec_setup.py` a direct signal since the + local single-sandbox simulation never fires real Job-level events. + **Investigate whether a real child process needs these hooks at all** + (see Category B note above — `StandaloneWorker` may already handle this + correctly without needing bespoke hooks). + +--- + +## 5. Progress so far this session (all in bifrost2, branch `pythonmonkey-platform-support`, NOT yet pushed/committed) + +Directory: `C:\Users\danie\DCP\bifrost2\dcp\_pm_evaluator\` (new package). + +- **`__init__.py`** — module docstring, wire protocol spec. +- **`channel.py`** — `EvaluatorChannel` class. Parent-side: spawns the + child via direct script path (`subprocess.Popen([sys.executable, + , "--port", N])` — **deliberately NOT `-m + dcp._pm_evaluator.child`**, which would force importing the entire `dcp` + package first, confirmed slow/problematic), binds a TCP listener on an + ephemeral port first, accepts the child's connection, then reads lines + in a **background thread with blocking sockets** (not asyncio streams — + see next bullet for why), delivering each line to `self.on_line` via + `loop.call_soon_threadsafe()`. +- **Real bug found and fixed**: bifrost2's shared event loop + (`dry.aio.loop`) has `nest_asyncio` applied (for pythonmonkey's + reentrant event-loop needs elsewhere in the codebase). `nest_asyncio`'s + patched loop **silently breaks plain asyncio `wait_for`/task scheduling** + — confirmed via a real hang with zero exception (parent successfully + spawned and connected to the child, per printed logs, but then just hung + forever with no error). Switched from an asyncio-streams design to + threads + blocking sockets specifically to sidestep this; do not go back + to asyncio streams for this without first confirming `nest_asyncio` + isn't going to bite again. +- **`child.py`** — **Stage 1 only, proof-of-plumbing, NOT the real sandbox + child yet.** Connects to the given port, sends one `LOG:` line and one + `MSG:{"type":"result","result":"stage1-ok"}` line, echoes any incoming + `MSG:` as a `LOG:`, exits cleanly on `DIE:`. Deliberately has zero + pythonmonkey/dcp involvement — pure stdlib socket code — to prove the + subprocess+socket mechanics in isolation before adding complexity. +- **`_test_plumbing.py`** — smoke test, **currently passing**: + ``` + [parent] spawning child... + [parent] child connected in 0.09s, pid=9452 + [parent] received: LOG:pm evaluator child connected, pid=9452 + [parent] received: MSG:{"type":"result","result":"stage1-ok"} + [parent] received: LOG:child received: {"type":"workerMessage","message":"hello from parent"} + [parent] received: DIE: + PLUMBING TEST PASSED + ``` + Run with: `cd C:\Users\danie\DCP\bifrost2 && python -u -m dcp._pm_evaluator._test_plumbing` + +**What this proves**: spawning a real child process and exchanging +messages over a real socket is fast (90ms) and reliable on this machine. +**What this does NOT yet prove**: that a real pythonmonkey instance can be +started inside that child, run the 22-file sandbox bootstrap in isolation, +and correctly execute a real work function. That's all still ahead. + +--- + +## 5b. Further progress (this session, after §5) — real job.localExec() reaches real deployment + +**Major finding: the whole separate-process architecture works, end to end, through real dcp-client deployment machinery**, using the EXISTING installed bundle's already-present §3.1 patch (`globalThis.__pmEvaluatorCtor` gate) — no monorepo change was even needed to reach this point. + +Copied `dcp/_pm_evaluator/` into the site-packages `dcp` install (the known-working environment with all prior investigation patches already applied) and ran a real `job.localExec()` via `_test_real_job.py`. Progression across fixes: + +1. **`evaluator.py` written and works.** Parent-side `globalThis.__pmEvaluatorCtor`, matching the exact shape from `pm_localexec_setup.py`'s `install()` (postMessage/onmessage/onerror/terminate/addEventListener), but internals route through a real `EvaluatorChannel` instead of fake in-process dispatch. Buffers `postMessage()` calls issued before the child finishes connecting (same pattern as `WebSocket.js`'s constructor). +2. **`child.py` stage 3 built and works**: real pythonmonkey instance in the child, portable bootstrap-file path resolution (via `importlib.util.find_spec('dcp')`, NOT the old hardcoded `C:\Users\danie\AppData\...` paths), all 22 sandbox bootstrap files load successfully in true isolation. Confirmed `writeln`/`onreadln`/`die` get deliberately deleted from `globalThis` by `sa-ww-simulation.js` after capturing them into a private closure (`/* Remove symbols from global scope that may be security leaks*/`) — **this is correct, intentional behavior, not a bug** (initially mis-flagged the check for this as a problem; it isn't). +3. **Hit `Error: module not found -- require('fs') from dcp-client/index.py`** the first time this ran against a *fresh* checkout that had never been patched. Root cause already known from the original investigation (FIXES_SUMMARY.md sec4.3) — `dcp-client/index.py` (itself part of the `dcp-client` npm package, confirmed by a fresh `npm i` reproducing the unpatched file) needs an `fs`/`os`/`child_process` shim for Supervisor-side code (`createBackingStore`/`obtainWorkerId`, i.e. the *parent* process establishing a persistent local-worker identity) that does real `require('fs')`. **This is unrelated to the new child-process work** — it's a parent-side (Supervisor) gap that exists regardless of evaluator architecture. Fixed by copying the already-patched `index.py` from the working site-packages install. +4. **Hit `DCPError: no transports defined` (DCPC-1014)** in the fresh checkout, confirmed via a control test to happen with *plain* `job.exec()` too — proving it was unrelated to any of this session's work. While investigating, the user relayed a diagnosis from a colleague (Tom Tang) of the **exact same class of bug**, found independently: `dcp/dry/aio.py` calls `nest_asyncio.apply()` unconditionally, which is broken under Python 3.14 (this machine's version) -- it doesn't propagate asyncio's "current task" context correctly, breaking any aiohttp call using `timeout=` (including pythonmonkey's `XMLHttpRequest-internal.py`, which backs DCP's socket.io polling transport) with `RuntimeError: Timeout should be used inside a task` *before the request is even sent* -- which dcp-client reports up as the much more confusing "no transports defined". **This is very likely also the same underlying cause of the `asyncio.wait_for` hang noted in §5 above** (`channel.py`'s original asyncio-streams design) -- both are "nest_asyncio broken on 3.14" symptoms. Applied the suggested fix to `dry/aio.py` (both this checkout and site-packages): only call `nest_asyncio.apply()` when a loop is *already running* in the current thread at import time (i.e. only when actually needed for Jupyter/web-server/GUI reentrant-loop support), not unconditionally. **This did NOT fully resolve the fresh-checkout "no transports defined" case** on retest -- there may be more than one contributing cause, or something else differs in that checkout. Not chased further this session (secondary to the core evaluator validation); the fix itself is real and worth keeping regardless. **`channel.py`'s thread-based design (built to work around the earlier `wait_for` hang) was NOT reverted back to asyncio streams after this fix** -- it works, wasn't broken, and reverting without retesting would be an unforced risk. Revisit only if there's a concrete reason to prefer the asyncio-streams version. +5. **With the working site-packages environment (all above fixes applied), the real job got all the way to**: `exec -> init -> preauth -> deploying -> listeners -> uploading`, then `DCPError: Could not connect to https://result-submitter.distributed.computer/result-submitter/ within 60s`. Confirmed NOT transient (reproduced twice). **This is expected and not yet a bug to fix** -- `route_job_arguments_through_local_files()` (Category B, per sec4) was deliberately not ported into the new evaluator flow yet; its whole purpose is keeping a local job's data off the real network, and result submission is very likely the same category of concern. This is the next real step, not a new mystery. + +**Bottom line**: the separate-process architecture is now validated end-to-end through real dcp-client job deployment -- spawning, socket bridging, sandbox bootstrap, and the real `Sandbox.start()`/`DistributiveWorker` construction path all work. What's left is exactly what sec6 already outlined: port Category B's still-needed fixes (starting with local data routing) into the new flow, then re-verify whether the local-worker completion-detection/error-propagation hooks are still needed given the new architecture (per sec4's open question) or whether the real child process now handles this more correctly on its own. + +### 5c. Session-ending blocker: external service outage, not a code issue + +While testing whether `route_job_arguments_through_local_files()` (reused as-is from `pm_localexec_setup.py`, temporarily, just to check the hypothesis) fixes the `result-submitter.distributed.computer` timeout from sec5b item 5, hit a NEW failure: plain `dcp.init()` **alone**, with zero job/evaluator code involved, started failing with a real `HTTP Error 404: Not Found` from a `fetch()` call inside `dcp-client/index.py`'s own init sequence (some remote config/version check). Reproduced twice, consistently. **Confirmed NOT caused by this session's changes**: reverted the `nest_asyncio` gating fix (sec5b item 4) back to the original unconditional `nest_asyncio.apply()` and the exact same 404 still happened -- then restored the fix (it's independently correct and unrelated). This is external infrastructure instability -- almost certainly the same family of issues already confirmed earlier this session with `packages.distributed.computer` (real, reproducible 404s/502s on that service, independent of any client). Blocked further live testing when this was hit; **not a code problem to fix, a live-service dependency to wait out or verify against a different environment**. + +**Resolved as expected/deliberate, not a bug**: confirmed with the user this was planned downtime -- "we took our services offline exactly at 1630hrs." Consistent with the evidence gathered before asking (6/6 repeated `curl` attempts against `https://scheduler.distributed.computer/etc/dcp-config.js` all returned a plain, consistent nginx 404 -- not the intermittent pattern seen with the actual `packages.distributed.computer` session-routing bug, which was the right thing to check first before assuming it was the same class of issue. It wasn't -- just an outage. + +**If picking this up fresh**: first confirm `dcp.init()` succeeds at all (`python -c "import dcp; dcp.init()"`) before assuming any code-level regression. If it 404s on `/etc/dcp-config.js`, check whether DCP services are actually up before investigating further -- this exact failure has already been seen once and was just planned downtime, not a client bug. + +## 5d. Code written during the outage (untested against live services -- verify first before trusting) + +With DCP services down (planned outage, sec5c), continued with code-only work that doesn't need the network: writing the real (non-borrowed) versions of the still-needed fixes directly into `dcp/api/job.py` and wiring `evaluator.py` into `dcp/initialization.py`. **None of this has been tested against a real job yet** -- it imports cleanly (verified: `python -c "import dcp"` succeeds, no circular-import or syntax errors) but that's all that could be confirmed with services down. + +**`dcp/api/job.py` changes** (this checkout, `pythonmonkey-platform-support` branch -- NOT yet copied to site-packages, NOT yet committed): + +- **`_route_arguments_locally(self)`** (new method) -- a real port of `pm_localexec_setup.route_job_arguments_through_local_files()`'s file-rewriting half (KVIN-encode each `jobArguments`/`jobInputData` element, write to a local temp file, rewrite in place to a `file://` URL, `jobInputData` -> separate `marshaledDataValues` property). Same mechanism, ported to be a real `Job` method instead of an opt-in global hook. Returns `(grants, written_paths)` instead of stashing state in a module global. +- **`_grant_local_file_origins(self, grants, ...)`** (new method) -- the async-polling half (wait for `job.js_ref.localWorker.originManager` to exist, then grant one narrow per-file origin). Changed the retry loop from unbounded polling to a bounded `max_attempts` (200 x 0.05s = 10s) -- the original polled forever with no cap; add this back to being unbounded if 10s proves too short once tested, but an unbounded retry that silently never grants access on a genuine failure seemed worse than a bounded one that at least stops. +- **`localExec(self, *args, **kwargs)`** (new method, replacing the inherited generic-proxy fallback) -- calls `_before_exec()`, then `_route_arguments_locally()` synchronously (before the real JS `localExec()` call -- timing-critical, see the method's own docstring for why), then `_grant_local_file_origins()` (fire-and-forget, doesn't block), then delegates to `self.js_ref.localExec(*args, **kwargs)` via `dry.aio.blockify`, with the same error-unwrapping (`.jsError.message`) and return-value conversion (`Array.from` + per-element `deserialize`) already validated in the earlier (closed) PR #49 attempt. Cleans up temp files in a `finally` block. +- **Deliberately NOT ported yet**: `force_job_completion_when_done()`/`raise_on_first_work_error()`'s concerns. Per sec4's still-open question -- does the real child-process evaluator (via `StandaloneWorker`-compatible messaging) already handle local-worker completion/error signaling correctly on its own, making these unnecessary? -- porting them blind, without being able to test whether they're even still needed, would risk exactly the kind of "looks complete but wasn't verified" mistake this whole rework exists to fix. **This is the first thing to check once services are back**: run a real `job.localExec()` through the new code and see whether it hangs (needs the completion-detection fix) or silently mishandles a broken work function's error (needs the error-propagation fix) before porting either one. +- Added the matching explanatory note to `wait()` (why `.wait()` can't follow `.localExec()`), same reasoning as the closed PR #49. + +**`dcp/initialization.py` changes**: +- `from ._pm_evaluator import evaluator as _pm_evaluator_module` at the top. +- `_pm_evaluator_module.install(aio.loop)` added inside `init()`, immediately before `js.dcp_client['init'](**kwargs)` -- registers the real `globalThis.__pmEvaluatorCtor` automatically, matching sec6 step 5's plan. Confirmed via `pm.eval('typeof globalThis.__pmEvaluatorCtor')` that it's `undefined` before `dcp.init()` runs and would need `dcp.init()` to actually reach this line to register it -- **not yet confirmed it successfully registers**, since `dcp.init()` couldn't complete (network down) far enough to reach this line in a real run. Once services are back, first check is simply: does `dcp.init()` now register the evaluator automatically, with zero manual `install_pm_evaluator(...)` call needed in the test script (unlike `_test_real_job.py`, sec5b, which called it manually)? + +## 5e. Exact pickup sequence once DCP services are back online + +1. Sanity check services are actually back: `curl -s -o /dev/null -w '%{http_code}\n' https://scheduler.distributed.computer/etc/dcp-config.js` should return `200`, not `404`. +2. Copy this checkout's `dcp/api/job.py`, `dcp/initialization.py`, and `dcp/_pm_evaluator/` into the site-packages install (same copy pattern used throughout sec5b/5c) -- OR switch to testing directly against this checkout once its own `dcp.init()` (fresh, unpatched `dcp-client` bundle) is confirmed working end-to-end independent of the site-packages shortcuts taken so far. Either is fine; site-packages is faster to unblock testing since it already has every other prerequisite patch (sec5b items 3-4) applied. +3. Run a real job with **zero manual evaluator install call** (letting `dcp.init()`'s new automatic wiring do it) -- e.g. adapt `_pm_evaluator/_test_real_job.py` to remove its manual `install_pm_evaluator(...)` line and confirm the evaluator ctor is present anyway. +4. Confirm the new `job.py`'s `localExec()` (with `_route_arguments_locally`/`_grant_local_file_origins` now real methods, not a borrowed external module) gets *past* the `result-submitter` timeout from sec5b item 5. +5. Once a real job completes: deliberately break a work function (make it raise) and confirm the error surfaces as a real Python exception with the actual traceback message -- this exercises the still-open `raise_on_first_work_error` question from sec4/5d. If it doesn't surface correctly, that's confirmation the fix is still needed and should be ported (from `pm_localexec_setup.py` lines 762-813) into `job.py` the same way the other two were. +6. If a job with more than a handful of slices, or a job that legitimately takes a while, hangs after all real results have arrived: that's confirmation `force_job_completion_when_done`'s concern is still needed too (port from lines 832-928). +7. Once real end-to-end success is confirmed (including error propagation): retest `pycomod_localexec_test.py` (heavier stress test) per sec6 step 6, then prepare the real PRs per sec6 step 7 -- monorepo MR (the `pythonmonkeyEvaluatorFactory()` + platform-gate change, sec4 item 1) and the bifrost2 PR (this branch), with a proper handoff doc. + +## 6. Remaining steps, in order + +1. **Build the real child process** (`child.py` stage 2). Replace the + stage-1 stub with: import `pythonmonkey`, resolve the dcp-client + bootstrap-file paths *portably* (not the hardcoded + `C:\Users\danie\AppData\Roaming\...` paths `pm_localexec_setup.py` + uses — resolve relative to the installed `dcp` package's own location, + e.g. via `importlib.util.find_spec('dcp')`), run the same + `BOOTSTRAP_FILES` list (copy the list from + `pm_localexec_setup.py` lines 25-48, fix the paths), wire real + `writeln`/`onreadln`/`die` globals to the actual socket (write out + `LOG:`/`MSG:` lines instead of in-process dispatch; feed incoming + socket lines to whatever `onreadln` registered). Test whether this + genuinely-isolated bootstrap needs ANY of Category A's fixes — don't + assume it doesn't just because the global is no longer shared; test it. +2. **Build the JS-side evaluator constructor** (parent, in bifrost2, + injected via `pm.eval()` the same way `pm_localexec_setup.py` did) that + wraps `EvaluatorChannel`: `postMessage()` writes a `MSG:` line via a + Python bridge function exposed on `globalThis`; incoming lines + (delivered via `channel.on_line`, called on the main loop thanks to + `call_soon_threadsafe`) get parsed and dispatched to `onmessage`/ + `onerror`, matching `Sandbox.start()`'s expected contract exactly (same + shape as the existing, proven `__pmEvaluatorCtor` in + `pm_localexec_setup.py` lines 112-193 — reuse that shape, just change + the internals to talk to a real channel instead of fake dispatch). + Since spawning is not instant, buffer any `postMessage()` calls issued + before the channel finishes connecting (same pattern already proven in + `WebSocket.js`'s constructor this session — buffer-then-flush on + connect). +3. **Test the JS↔Python bridge directly** against the stage-2 child, + without going through real `Sandbox.start()`/dcp-client yet — construct + the evaluator directly via `pm.eval()`, call `postMessage`, confirm + `onmessage` fires with real data from the child's actual pythonmonkey + instance. +4. **Monorepo changes** (`C:\Users\danie\DCP\dcp-monorepo`, need a fresh + feature branch off `develop`): + - `src/dcp-client/job/index.js` ~line 645: accept `"pythonmonkey"` + platform (§4 item 1 above). + - New file, modeled on `src/dcp-client/worker/evaluators/node-localExec.js`: + a `pythonmonkeyEvaluatorFactory()` — but note, unlike Node's version, + this doesn't need to do the spawning itself (that's already handled + Python-side via `globalThis.__pmEvaluatorCtor`, exactly like the + existing FIXES_SUMMARY.md §3.1 diff already established) — likely a + much smaller file than `node-localExec.js`, just returning + `globalThis.__pmEvaluatorCtor`. + - Investigate/implement the `"complete"`/`"workError"` signal question + from §4 Category B — determine whether `Sandbox.start()`'s existing + onmessage handling (feeding into `StandaloneWorker`'s already-correct + `'result'`/error dispatch) is now sufent once real socket messages are + flowing, or whether the two bespoke hooks (§4 item 3) are still + needed. +5. **bifrost2 changes**: + - `dcp/initialization.py`: call the new evaluator-registration code + automatically in `init()`, before `js.dcp_client['init'](**kwargs)` + runs (ordering matters — confirmed by the original investigation). + Only port whatever Category A fixes step 1's testing showed are still + genuinely needed. + - `dcp/api/job.py`: rewrite `localExec()` cleanly, folding in whatever + Category B fixes are still needed after step 4's investigation + (route-arguments-locally at minimum; completion/error detection only + if still needed after checking `StandaloneWorker`'s existing handling). +6. **End-to-end test** against a real job. Use + `C:\Users\danie\DCP\dcp_local_job_test.py` as the reference test job + (uppercase-8-letters, Pyodide work function, `demo`/`dcp` compute + group) — but the goal is for the **exact clean script in §1** to work, + not a script with any manual patch wiring. Also retest + `pycomod_localexec_test.py` (heavier stress test — filesystem shipping, + extra Pyodide modules, cloudpickle round-trip) once the basic case + works. +7. **Only once real, end-to-end, tested** — prepare PRs: + - dcp monorepo MR (step 4's changes) + - bifrost2 PR (step 5's changes) — this is the real replacement for the + closed PR #49; make sure it's actually complete and tested this time, + not another partial cut. + - Write a handoff doc (matching the pattern of + `SPIDERMONKEY_VERSION_BUMP.md`/`LOCALEXEC_PYODIDE_FIX_NOTES.md`) + covering the full architecture, what was tested, what's flagged as + needing review. + +--- + +## 7. Repo/access reference + +| Repo | Local path | Remote | Branch (this work) | Access | +|---|---|---|---|---| +| PythonMonkey | `C:\Users\danie\DCP\pythonmonkey-src` | `github.com/Distributive-Network/PythonMonkey` | (PRs #509, #510 already open on their own branches) | `gh`, push access confirmed | +| bifrost2 | `C:\Users\danie\DCP\bifrost2` | `github.com/Distributive-Network/bifrost2` | `pythonmonkey-platform-support` (current work, **not pushed yet**) | `gh`, push access confirmed | +| dcp monorepo | `C:\Users\danie\DCP\dcp-monorepo` | `gitlab.com:Distributed-Compute-Protocol/dcp.git` | `pythonmonkey-websocket-transport` used for MR !3323; **need a NEW branch off `develop` for this work** | SSH clone/push confirmed working (git only; `glab` installed but not authenticated) | +| dcp-client | `C:\Users\danie\DCP\dcp-client` | `gitlab.com:Distributed-Compute-Protocol/dcp-client.git` | detached at `v5.7.3`, clean, unmodified | read-only used so far; this repo is just a packaging wrapper around the monorepo (its `prepack` hook clones `dcp.git` and builds from there) — **the real source lives in dcp-monorepo, not here**, except `lib/standaloneWorker.js` which genuinely lives in *this* repo and is directly reusable (see §3) | + +Tools: `gh` (GitHub CLI) installed portably at `C:\Users\danie\tools\bin\gh.exe`, +authenticated as `dan-distributive`. `glab` (GitLab CLI) installed +portably at `C:\Users\danie\tools\glab_extracted\bin\glab.exe`, **not +authenticated** (no browser device-flow available in this version; needs a +GitLab personal access token with `api`+`write_repository` scopes, or open +MRs manually via the printed `git push` URL). + +Original investigation reference docs (read these for exact historical +reasoning, don't just take this plan doc's summaries as complete): +- `C:\Users\danie\DCP\localexec_patch\STATUS.md` — full investigation log +- `C:\Users\danie\DCP\localexec_patch\FIXES_SUMMARY.md` — clean summary of + all fixes including the exact bundle diffs (§3) +- `C:\Users\danie\DCP\localexec_patch\pm_localexec_setup.py` — the actual + 963-line implementation being superseded/ported +- `C:\Users\danie\DCP\pythonmonkey-src\SPIDERMONKEY_VERSION_BUMP.md` — the + SpiderMonkey rebuild handoff doc (PR #509), for the JobQueue-rewrite + context behind Category C being obsolete diff --git a/dcp/_pm_evaluator/__init__.py b/dcp/_pm_evaluator/__init__.py new file mode 100644 index 0000000..75f234c --- /dev/null +++ b/dcp/_pm_evaluator/__init__.py @@ -0,0 +1,19 @@ +""" +A real, separate-process evaluator for job.localExec() under pythonmonkey, +matching how localExec() actually works on Node.js -- a genuinely separate +process, not a shared JS global. + +See src/dcp-client/worker/evaluators/node-localExec.js (dcp monorepo) for +the reference architecture this mirrors: a SandboxConstructor that spawns +a worker connected over a socket, wrapped to satisfy the same postMessage/ +onmessage/onerror/terminate/addEventListener contract Sandbox.start() +expects from any platform's evaluator. + +Wire protocol (matches lib/standaloneWorker.js's StandaloneWorker exactly, +so a real dcp-worker-shaped child speaks a protocol dcp-client already +understands): newline-delimited, three line-prefixes: + - "LOG:" -- debug/log line, informational only + - "DIE:" -- child is shutting down + - "MSG:" -- {"type": "workerMessage", "message": ...} + (postMessage payload) or {"type": "result", ...} +""" diff --git a/dcp/_pm_evaluator/_test_bootstrap.py b/dcp/_pm_evaluator/_test_bootstrap.py new file mode 100644 index 0000000..48c0b97 --- /dev/null +++ b/dcp/_pm_evaluator/_test_bootstrap.py @@ -0,0 +1,62 @@ +""" +STAGE 3 test: does the real 22-file sandbox bootstrap load successfully in +a genuinely isolated child process? Key question this answers: does +process isolation eliminate the need for Category A's console/require/ +timer-clobbering fixes and the access-lists-masking bypass (see +PYTHONMONKEY_EVALUATOR_PLAN.md), or is at least some of it still needed +even with no shared global? + +Run: cd C:\\Users\\danie\\DCP\\bifrost2 && python -u -m dcp._pm_evaluator._test_bootstrap +""" +import asyncio +import sys +import os +import time + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))) +from dcp._pm_evaluator.channel import EvaluatorChannel + + +async def main(): + loop = asyncio.get_running_loop() + channel = EvaluatorChannel(loop) + + lines_received = [] + + def _on_line(line): + print("[parent]", line, flush=True) + lines_received.append(line) + + channel.on_line = _on_line + + print("[parent] spawning child...", flush=True) + t0 = time.time() + channel.spawn_and_connect() + print(f"[parent] child connected in {time.time()-t0:.2f}s, pid=", channel.proc.pid, flush=True) + + # Bootstrap loading took real, non-trivial time in the original + # investigation (22 files, some doing real async work) -- give it a + # generous window and poll for completion rather than a fixed sleep. + deadline = time.time() + 90 + while time.time() < deadline: + if any("bootstrap] all 22 files loaded" in l or "bootstrap] ABORTED" in l for l in lines_received): + break + await asyncio.sleep(0.5) + + channel.terminate() + await asyncio.sleep(0.5) + + aborted = [l for l in lines_received if "ABORTED" in l or "FAILED" in l] + succeeded = any("all 22 files loaded" in l for l in lines_received) + intact = [l for l in lines_received if "still intact" in l] + + print() + print("=" * 60) + print(f"Bootstrap succeeded: {succeeded}") + print(f"Failures/aborts: {aborted}") + print(f"writeln/onreadln intact after bootstrap: {intact}") + print("=" * 60) + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/dcp/_pm_evaluator/_test_exec_control.py b/dcp/_pm_evaluator/_test_exec_control.py new file mode 100644 index 0000000..e058ed3 --- /dev/null +++ b/dcp/_pm_evaluator/_test_exec_control.py @@ -0,0 +1,48 @@ +"""Control test: does plain job.exec() (unrelated to the new evaluator) +also hit 'no transports defined' in this fresh checkout? Disambiguates a +general environment/config issue from something specific to localExec().""" +import json +import dcp +dcp.init() + +import pythonmonkey as pm +from dcp.dry.aio import loop as _shared_loop + +_load_id_keystore = pm.eval(""" +async () => { + const wallet = dcp.wallet; + const identity = dcp.identity; + const idKeystore = await wallet.get('id', { KeystoreConstructor: wallet.IdKeystore }); + identity.set(idKeystore); + return idKeystore.address.toString(); +} +""") + +async def _load_identity(): + return await _load_id_keystore() + +address = _shared_loop.run_until_complete(_load_identity()) +print("Using identity address:", address, flush=True) + +input_set = list('yelling!') + +def work_function(letter): + dcp.progress() + return letter.upper() + +job = dcp.compute_for(input_set, work_function) +job.computeGroups = [{'joinKey': 'demo', 'joinSecret': 'dcp'}] +job.public.name = 'fresh-checkout-exec-control' +job.public.description = 'control test' +job.public.link = 'https://distributive.network' + +job.on('readystatechange', lambda s: print(f"Ready State: {s}", flush=True)) +job.on('accepted', lambda _: print(f" Job ID: {job.id}", flush=True)) +job.on('error', lambda e: print("error event:", json.dumps(e, indent=2), flush=True)) +job.on('result', lambda r: print("result event:", json.dumps(r, indent=2), flush=True)) + +print("Calling job.exec()...", flush=True) +job.exec() +results = job.wait() +print(''.join(results)) +print("EXEC CONTROL TEST COMPLETE") diff --git a/dcp/_pm_evaluator/_test_plumbing.py b/dcp/_pm_evaluator/_test_plumbing.py new file mode 100644 index 0000000..9459279 --- /dev/null +++ b/dcp/_pm_evaluator/_test_plumbing.py @@ -0,0 +1,43 @@ +""" +STAGE 1 smoke test: proves spawn+socket-connect+message-exchange+terminate +works at all on this machine, with zero pythonmonkey/bifrost2 involvement -- +isolates subprocess/socket mechanics from everything else before adding +that complexity on top. + +Run directly: python -m dcp._pm_evaluator._test_plumbing +""" +import asyncio +import sys +import os +import time + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))) +from dcp._pm_evaluator.channel import EvaluatorChannel + + +async def main(): + loop = asyncio.get_running_loop() + channel = EvaluatorChannel(loop) + + lines_received = [] + channel.on_line = lambda line: (print("[parent] received:", line, flush=True), lines_received.append(line)) + + print("[parent] spawning child...", flush=True) + t0 = time.time() + channel.spawn_and_connect() + print(f"[parent] child connected in {time.time()-t0:.2f}s, pid=", channel.proc.pid, flush=True) + + channel.write_line('MSG:{"type":"workerMessage","message":"hello from parent"}') + + await asyncio.sleep(2.0) + + channel.terminate() + await asyncio.sleep(0.5) + + assert any("pythonmonkey ready" in l for l in lines_received), "child pythonmonkey did not start" + assert any("child JS onreadln got" in l for l in lines_received), "child's real pythonmonkey JS did not receive our message via onreadln" + print("STAGE 2 PLUMBING TEST PASSED (real pythonmonkey in child, real socket round-trip)") + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/dcp/_pm_evaluator/_test_real_job.py b/dcp/_pm_evaluator/_test_real_job.py new file mode 100644 index 0000000..5c7d1b4 --- /dev/null +++ b/dcp/_pm_evaluator/_test_real_job.py @@ -0,0 +1,57 @@ +""" +STAGE 4 test: a REAL job.localExec() call through the real separate-process +evaluator, no Category A/B patches applied at all -- testing what actually +still breaks, empirically, rather than assuming. +""" +import sys +sys.path.insert(0, r"C:\Users\danie\DCP\bifrost2") # use THIS checkout of dcp, not site-packages + +import json +import dcp +dcp.init() + +import pythonmonkey as pm +from dcp._pm_evaluator.evaluator import install as install_pm_evaluator +from dcp.dry.aio import loop as _shared_loop + +install_pm_evaluator(_shared_loop) +print("Evaluator constructor installed:", pm.eval("typeof globalThis.__pmEvaluatorCtor")) + +# Real identity via id.keystore, same safe pattern as the working tests. +_load_id_keystore = pm.eval(""" +async () => { + const wallet = dcp.wallet; + const identity = dcp.identity; + const idKeystore = await wallet.get('id', { KeystoreConstructor: wallet.IdKeystore }); + identity.set(idKeystore); + return idKeystore.address.toString(); +} +""") + +async def _load_identity(): + return await _load_id_keystore() + +address = _shared_loop.run_until_complete(_load_identity()) +print("Using identity address:", address) + +input_set = list('yelling!') + +def work_function(letter): + dcp.progress() + return letter.upper() + +job = dcp.compute_for(input_set, work_function) +job.computeGroups = [{'joinKey': 'demo', 'joinSecret': 'dcp'}] +job.public.name = 'pm-real-evaluator-test' +job.public.description = 'Real separate-process evaluator test' +job.public.link = 'https://distributive.network' + +job.on('readystatechange', lambda s: print(f"Ready State: {s}", flush=True)) +job.on('accepted', lambda _: print(f" Job ID: {job.id}", flush=True)) +job.on('error', lambda e: print("error event:", json.dumps(e, indent=2), flush=True)) +job.on('result', lambda r: print("result event:", json.dumps(r, indent=2), flush=True)) + +print("Calling job.localExec()...", flush=True) +results = job.localExec() +print(''.join(results)) +print("REAL JOB TEST COMPLETE") diff --git a/dcp/_pm_evaluator/channel.py b/dcp/_pm_evaluator/channel.py new file mode 100644 index 0000000..3781dbb --- /dev/null +++ b/dcp/_pm_evaluator/channel.py @@ -0,0 +1,92 @@ +""" +Parent-side channel: spawns the child evaluator process, listens for its +connection, and bridges line-based traffic to/from it. + +Uses a background thread with plain blocking sockets, NOT asyncio streams +-- bifrost2's shared loop (dry.aio.loop) has nest_asyncio applied for +pythonmonkey's reentrant event-loop needs, and nest_asyncio's patched loop +was confirmed (empirically, via a real hang with no exception) to break +plain asyncio task/timeout scheduling in ways not worth fighting. A +background thread reading a blocking socket sidesteps that entirely; each +received line is handed back to the main loop/thread via +loop.call_soon_threadsafe(), which is the same safe cross-thread pattern +pythonmonkey's own C++ side uses to reach the event loop (see JobQueue.cc's +dispatchToEventLoop in the pythonmonkey-src rebuild this session). +""" +import os +import socket +import subprocess +import sys +import threading + +_CHILD_SCRIPT = os.path.join(os.path.dirname(os.path.abspath(__file__)), "child.py") + + +class EvaluatorChannel: + def __init__(self, loop): + self.loop = loop + self.proc: subprocess.Popen | None = None + self.sock: socket.socket | None = None + self.on_line = None # callable(str) -> None; called ON THE MAIN LOOP/THREAD + self._reader_thread = None + self._stop = False + + def spawn_and_connect(self, connect_timeout=20): + """Synchronous -- binds, spawns, accepts. Call this from a Python + thread/context that's fine blocking briefly (the accept() wait is + normally sub-second; the child does no heavy work before connecting).""" + srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + srv.bind(("127.0.0.1", 0)) + srv.listen(1) + port = srv.getsockname()[1] + srv.settimeout(connect_timeout) + + self.proc = subprocess.Popen([sys.executable, _CHILD_SCRIPT, "--port", str(port)]) + + try: + conn, _ = srv.accept() + finally: + srv.close() + + self.sock = conn + self._reader_thread = threading.Thread(target=self._read_loop, daemon=True) + self._reader_thread.start() + return self + + def _read_loop(self): + buf = b"" + while not self._stop: + try: + data = self.sock.recv(4096) + except OSError: + break + if not data: + break + buf += data + while b"\n" in buf: + raw, buf = buf.split(b"\n", 1) + text = raw.decode("utf-8", errors="replace") + if self.on_line: + self.loop.call_soon_threadsafe(self.on_line, text) + + def write_line(self, line: str): + if self.sock is None: + return + try: + self.sock.sendall((line + "\n").encode("utf-8")) + except OSError: + pass + + def terminate(self): + self.write_line("DIE:") + self._stop = True + if self.proc: + try: + self.proc.wait(timeout=5) + except subprocess.TimeoutExpired: + self.proc.kill() + if self.sock: + try: + self.sock.close() + except OSError: + pass diff --git a/dcp/_pm_evaluator/child.py b/dcp/_pm_evaluator/child.py new file mode 100644 index 0000000..b3d9c1a --- /dev/null +++ b/dcp/_pm_evaluator/child.py @@ -0,0 +1,170 @@ +""" +Entry point for the separate-process pythonmonkey evaluator's child side. +Invoked as: python --port (direct script path, NOT +`-m dcp._pm_evaluator.child` -- see channel.py's comment on _CHILD_SCRIPT +for why: `-m` forces importing the whole `dcp` package first, which this +process does not need and which was confirmed to slow/complicate startup). + +STAGE 2 (current): adds a real pythonmonkey instance and wires writeln/ +onreadln/die to the real socket, matching the wire protocol StandaloneWorker +(dcp-client's lib/standaloneWorker.js) expects on its read side. Does NOT +yet run the 22-file sandbox bootstrap -- that's the next stage, added only +once this stage is confirmed working (pythonmonkey starts reliably in a +spawned subprocess, and the writeln/onreadln bridge works end to end). +""" +import argparse +import asyncio +import importlib.util +import os +import socket +import sys +import threading + + +def _bootstrap_files(): + """Portable resolution of the 22 sandbox control-code files -- NOT the + hardcoded C:\\Users\\danie\\AppData\\... paths pm_localexec_setup.py + used. Resolved relative to wherever `dcp` is actually installed.""" + spec = importlib.util.find_spec("dcp") + dcp_dir = os.path.dirname(spec.origin) + js_root = os.path.join(dcp_dir, "js", "node_modules") + dcp_client = os.path.join(js_root, "dcp-client") + sandbox = os.path.join(dcp_client, "libexec", "sandbox") + return [ + os.path.join(js_root, "kvin", "kvin.js"), + os.path.join(sandbox, "sa-ww-simulation.js"), + os.path.join(sandbox, "script-load-wrapper.js"), + os.path.join(sandbox, "timer-classes.js"), + os.path.join(sandbox, "wrap-event-listeners.js"), + os.path.join(sandbox, "event-loop-virtualization.js"), + os.path.join(sandbox, "lift-webgl.js"), + os.path.join(sandbox, "lift-wasm.js"), + os.path.join(sandbox, "lift-webgpu.js"), + os.path.join(sandbox, "url.js"), + os.path.join(sandbox, "polyfills.js"), + os.path.join(sandbox, "access-lists.js"), + os.path.join(sandbox, "fetch-factory.js"), + os.path.join(sandbox, "bravojs-init.js"), + os.path.join(js_root, "bravojs", "bravo.js"), + os.path.join(sandbox, "bravojs-env.js"), + os.path.join(sandbox, "worktimes.js"), + os.path.join(sandbox, "pyodide-core.js"), + os.path.join(sandbox, "pyodide-worktime.js"), + os.path.join(sandbox, "map-basic-worktime.js"), + os.path.join(sandbox, "calculate-capabilities.js"), + os.path.join(sandbox, "bootstrap.js"), + ] + + +def main(): + parser = argparse.ArgumentParser() + parser.add_argument("--port", type=int, required=True) + args = parser.parse_args() + + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.connect(("127.0.0.1", args.port)) + + import pythonmonkey as pm + + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + + def _write_line(line: str): + try: + sock.sendall((line + "\n").encode("utf-8")) + except OSError: + pass + + pm.globalThis["__pmChildWriteLine"] = _write_line + + # Same writeln contract as the real sandbox control-code files expect + # (see pm_localexec_setup.py's install(), which this mirrors) -- but + # now genuinely writing to a real socket instead of fake in-process + # dispatch. + pm.eval(""" + globalThis.writeln = function(line) { + globalThis.__pmChildWriteLine(line); + }; + globalThis.__pmOnReadlnHandler = null; + globalThis.onreadln = function(fn) { globalThis.__pmOnReadlnHandler = fn; }; + globalThis.die = function() { globalThis.__pmChildWriteLine('DIE:'); }; + """) + + should_exit = threading.Event() + + def _socket_reader(): + buf = b"" + while True: + try: + data = sock.recv(4096) + except OSError: + break + if not data: + break + buf += data + while b"\n" in buf: + raw, buf = buf.split(b"\n", 1) + text = raw.decode("utf-8", errors="replace") + if text.startswith("DIE:"): + should_exit.set() + return + elif text.startswith("MSG:"): + loop.call_soon_threadsafe(_dispatch_incoming, text[4:]) + should_exit.set() + + def _dispatch_incoming(json_text: str): + handler = pm.eval("globalThis.__pmOnReadlnHandler") + if handler: + # onreadln handlers, per the real protocol, receive the RAW + # "MSG:\n" line, not the decoded payload -- matches + # pm_localexec_setup.py's writeln() parsing the raw line itself. + handler("MSG:" + json_text + "\n") + + reader_thread = threading.Thread(target=_socket_reader, daemon=True) + reader_thread.start() + + _write_line("LOG:pm evaluator child: pythonmonkey ready, pid=%d" % __import__("os").getpid()) + + # STAGE 3: run the real 22-file sandbox bootstrap, in TRUE isolation -- + # this process's JS global is used for NOTHING else (no shared + # "Supervisor" role), so console/require/timer clobbering and + # access-lists masking should not be able to break anything else the + # way they did in the old in-process simulation. Verifying that + # empirically here, not assuming it. + async def _run_bootstrap(): + for f in _bootstrap_files(): + _write_line(f"LOG:[bootstrap] loading {f}") + try: + src = open(f, encoding="utf-8").read() + pm.eval(src) + except Exception as e: + _write_line(f"LOG:[bootstrap] FAILED loading {f}: {type(e).__name__}: {e}") + raise + _write_line(f"LOG:[bootstrap] done {f}") + + try: + loop.run_until_complete(_run_bootstrap()) + _write_line("LOG:[bootstrap] all 22 files loaded successfully") + except Exception as e: + _write_line(f"LOG:[bootstrap] ABORTED: {type(e).__name__}: {e}") + sock.close() + sys.exit(1) + + # Confirm writeln/onreadln/die are still OUR functions after the + # bootstrap ran (i.e. nothing in the 22 files redefined them out from + # under us) -- this is exactly the kind of clobbering that broke the + # in-process simulation; check it explicitly rather than assume + # isolation fixed it. + still_ours = pm.eval("typeof globalThis.writeln === 'function' && typeof globalThis.onreadln === 'function'") + _write_line(f"LOG:writeln/onreadln still intact after bootstrap: {still_ours}") + + async def _wait_for_exit(): + while not should_exit.is_set(): + await asyncio.sleep(0.05) + + loop.run_until_complete(_wait_for_exit()) + sock.close() + + +if __name__ == "__main__": + main() diff --git a/dcp/_pm_evaluator/evaluator.py b/dcp/_pm_evaluator/evaluator.py new file mode 100644 index 0000000..53ea793 --- /dev/null +++ b/dcp/_pm_evaluator/evaluator.py @@ -0,0 +1,103 @@ +""" +Parent-side: the real `globalThis.__pmEvaluatorCtor` -- a Worker-shaped +(postMessage/onmessage/onerror/terminate/addEventListener) constructor +matching what dcp-client's Sandbox.start() expects as a SandboxConstructor +(see FIXES_SUMMARY.md sec3.1 / node-localExec.js's nodeEvaluatorFactory +for the reference contract). Backed by a REAL child process +(EvaluatorChannel), not the old in-process simulation. +""" +import asyncio + +import pythonmonkey as pm + +from .channel import EvaluatorChannel + +_ctor_factory = pm.eval(""" +(spawnAndGetHandle) => { + return function PythonMonkeyEvaluator(_options) { + var self = this; + this.onmessage = null; + this.onerror = null; + this._listeners = {}; + this.addEventListener = function(type, listener) { + (self._listeners[type] = self._listeners[type] || []).push(listener); + }; + this.removeEventListener = function(type, listener) { + if (!self._listeners[type]) return; + self._listeners[type] = self._listeners[type].filter((l) => l !== listener); + }; + + var terminated = false; + var sendBuffer = []; + var channelWrite = null; + var channelTerminate = null; + + function fireEnd() { + Promise.resolve().then(() => { + (self._listeners['end'] || []).forEach((fn) => { try { fn(); } catch (e) {} }); + }); + } + + function handleLine(line) { + if (terminated) return; + if (line.indexOf('LOG:') === 0) return; + if (line.indexOf('DIE:') === 0) { + terminated = true; + fireEnd(); + return; + } + if (line.indexOf('MSG:') !== 0) return; + var obj; + try { obj = JSON.parse(line.slice(4)); } catch (e) { return; } + Promise.resolve().then(() => { + if (obj.type === 'workerMessage' && self.onmessage) { + self.onmessage({ data: obj.message }); + } else if (obj.type === 'result' && obj.exception && self.onerror) { + self.onerror(obj.exception); + } + }); + } + + this.postMessage = function(msg) { + var line = 'MSG:' + JSON.stringify({ type: 'workerMessage', message: msg }); + if (channelWrite) channelWrite(line); + else sendBuffer.push(line); + }; + + this.terminate = function() { + if (terminated) return; + terminated = true; + if (channelTerminate) channelTerminate(); + fireEnd(); + }; + + spawnAndGetHandle(handleLine).then((handle) => { + if (terminated) { handle.terminate(); return; } + channelWrite = handle.write; + channelTerminate = handle.terminate; + for (var i = 0; i < sendBuffer.length; i++) channelWrite(sendBuffer[i]); + sendBuffer = []; + }).catch((e) => { + console.error('PythonMonkeyEvaluator: spawn failed:', e); + terminated = true; + fireEnd(); + }); + }; +} +""") + + +def install(loop: asyncio.AbstractEventLoop): + """Registers globalThis.__pmEvaluatorCtor, backed by real child + processes. Call once, before dcp-client's own init runs (matching the + ordering the original in-process version required -- unverified + whether that ordering constraint still applies here, but preserved + out of caution until tested otherwise).""" + + async def _spawn_and_get_handle(on_line_js_callback): + channel = EvaluatorChannel(loop) + channel.on_line = on_line_js_callback + await loop.run_in_executor(None, channel.spawn_and_connect) + return {"write": channel.write_line, "terminate": channel.terminate} + + pm.globalThis["__pmEvaluatorCtor"] = _ctor_factory(_spawn_and_get_handle) diff --git a/dcp/api/job.py b/dcp/api/job.py index c942127..59cefee 100644 --- a/dcp/api/job.py +++ b/dcp/api/job.py @@ -175,8 +175,204 @@ def exec(self, *args): return results def wait(self): + # NOTE: for a localExec() job, registering these listeners here + # is too late to ever see the 'complete' event -- the real JS + # localExec() Promise (awaited inside localExec() below) does + # not resolve until the WHOLE job (including the 'complete' + # event this function listens for) has already finished, so by + # the time _wait() runs, that event has already fired and been + # missed (EventEmitters don't replay past events to newly-added + # listeners). For localExec(), use its own return value instead + # of calling .wait() afterward (matching the real Node.js usage + # pattern -- no separate .wait() call at all). This method + # remains correct and necessary for real distributed jobs via + # exec()/aio.exec(), where 'complete' genuinely arrives later, + # well after this registration. return dry.aio.blockify(self._wait)() + def _route_arguments_locally(self): + """ + localExec() otherwise unconditionally routes job data through + the real scheduler: jobArguments once a payload exceeds a size + threshold (a scheduler-hosted URL), and slice *values* via a + dedicated, unconditional bulk upload (addSlices(), during the + "uploading" state) that runs regardless of size unless + jobRef.marshaledDataValues is already set. Real Node localExec() + avoids this entirely by writing both to local temp files and + granting the local worker a narrow, path-scoped file:// origin + per file instead -- this mirrors that. + + Must run SYNCHRONOUSLY, immediately after _before_exec() + populates jobArguments/jobInputData and before the real JS + localExec() call (and therefore deployJob()'s upload) ever + runs: the scheduler snapshots jobArguments during deploy, and + the slice-value upload is scheduled immediately once deploy + completes -- both well before localWorker/originManager exist + or any async callback would get a turn to run. + + Returns the list of (path, purpose) grants still needing + origin access once the local worker exists (see + _grant_local_file_origins, called separately, async, after + this returns). + """ + import tempfile + import os as _os + + written_paths = [] + + def _write_temp_file(encoded_str): + fd, path = tempfile.mkstemp(prefix="bifrost2-localExec-arg-", suffix=".kvin") + with _os.fdopen(fd, "w", encoding="utf-8") as f: + f.write(encoded_str) + written_paths.append(path) + return path + + pm.globalThis["__pmWriteLocalArgFile"] = _write_temp_file + rewrite = pm.eval(""" + (jobRef) => { + const KVIN = new (require('kvin').KVIN)(); + const grants = []; + const toLocalURL = (v, purpose) => { + if (v instanceof URL) return v; + const encoded = KVIN.stringify(v); + const filePath = globalThis.__pmWriteLocalArgFile(encoded); + grants.push({ path: filePath, purpose }); + return new URL('file://' + filePath); + }; + if (jobRef.jobArguments) { + jobRef.jobArguments = jobRef.jobArguments.map((v) => toLocalURL(v, 'fetchArguments')); + } + if (Array.isArray(jobRef.jobInputData)) { + jobRef.marshaledDataValues = jobRef.jobInputData.map((v) => toLocalURL(v, 'fetchData')); + } + return grants; + } + """) + grants = [{"path": g["path"], "purpose": g["purpose"]} for g in rewrite(self.js_ref)] + return grants, written_paths + + def _grant_local_file_origins(self, grants, poll_interval=0.05, max_attempts=200): + """ + The local worker doesn't fetch its arguments/slice values until + well after deploy/upload finishes, so (unlike + _route_arguments_locally) this half can safely poll + asynchronously for job.js_ref.localWorker.originManager to + exist, then grant one narrow, path-scoped origin per file, with + the purpose matching what it actually is ('fetchArguments' vs + 'fetchData') -- never a blanket grant. + + NEEDS TESTING once services are back up: this is a direct port + of the proven-working pm_localexec_setup.py version, adapted to + be a real Job method rather than an opt-in hook, but has not + itself been re-tested against a live job yet (blocked on a + planned DCP services outage -- see PYTHONMONKEY_EVALUATOR_PLAN.md + sec5c). In particular: does this timing assumption + (job.js_ref.localWorker.originManager appearing) still hold + with the new separate-process evaluator, where localWorker + construction now involves a real spawned child process instead + of synchronous in-process dispatch? + """ + grant_origin = pm.eval(""" + (originManager, filePath, purpose) => { + originManager.add(new URL('file://' + filePath).pathname, purpose, null); + } + """) + + def _try_grant(attempts_left): + try: + worker = self.js_ref["localWorker"] + if worker is None or isinstance(worker, pm.null.__class__): + if attempts_left > 0: + dry.aio.loop.call_later(poll_interval, _try_grant, attempts_left - 1) + return + origin_manager = worker["originManager"] + if origin_manager is None or isinstance(origin_manager, pm.null.__class__): + if attempts_left > 0: + dry.aio.loop.call_later(poll_interval, _try_grant, attempts_left - 1) + return + for grant in grants: + grant_origin(origin_manager, grant["path"], grant["purpose"]) + except Exception: + if attempts_left > 0: + dry.aio.loop.call_later(poll_interval, _try_grant, attempts_left - 1) + + dry.aio.loop.call_later(poll_interval, _try_grant, max_attempts) + + def localExec(self, *args, **kwargs): + """ + localExec() was otherwise inherited unmodified from the generic + JS-proxy wrapper (dry/class_manager.py __getattr__), which calls + straight through to self.js_ref['localExec'] and skips + _before_exec() entirely. For the pyodide worktime, _before_exec() + is what rewrites workFunctionURI into the real bifrost2-wrapped + script (imports, serializers, and the dcp.set_slice_handler() + registration) -- without it the raw user Python source is sent + as-is, which never calls dcp.set_slice_handler() (-> + ENOSLICEHANDLER in the pyodide worktime). Mirrors _exec()'s + setup, matching the real Node.js usage pattern + (`const results = await job.localExec()`) instead of requiring + a separate `.wait()` call afterward -- see wait()'s own note + for why calling .wait() after .localExec() doesn't work anyway. + + NEEDS TESTING once services are back up (see + PYTHONMONKEY_EVALUATOR_PLAN.md sec6): whether + force_job_completion_when_done()/raise_on_first_work_error()'s + concerns (local Worker never signaling a terminating 'stop'; + work-function errors getting swallowed) still apply under the + new separate-process evaluator, or whether a real child process + now surfaces these correctly on its own via + StandaloneWorker-compatible 'result'/error messages. Not yet + ported here pending that verification -- porting them + unconditionally without checking would risk reintroducing + exactly the kind of unverified, "looks done but isn't" change + this whole rework exists to avoid. + """ + self._before_exec() + self._wrapper_set_attribute("_exec_called", True) + + grants, written_paths = self._route_arguments_locally() + self._grant_local_file_origins(grants) + + try: + try: + ret_val = dry.aio.blockify(self.js_ref.localExec)(*args, **kwargs) + except Exception as e: + # The real error message (a genuine Python traceback + # pointing at the user's own work function) is already + # present on e.jsError.message -- pythonmonkey.SpiderMonkeyError + # exposes the underlying JS Error as .jsError. Without + # this, the caller sees a SpiderMonkeyError whose + # message is buried under dcp-client's own internal JS + # stack frames. Re-raise with just the real message so + # job.localExec() fails pointing at the user's own code. + js_error = getattr(e, 'jsError', None) + message = getattr(js_error, 'message', None) if js_error is not None else None + if message: + raise RuntimeError(message) from None + raise + finally: + # These files hold real job argument/slice-value data + # (potentially sensitive) and are not otherwise cleaned up. + # The local worker has already read them by the time + # localExec()'s own promise settles (succeeded or failed). + import os as _os + for path in written_paths: + try: + _os.unlink(path) + except OSError: + pass + + # ret_val resolves to the job's ResultHandle -- a Proxy whose + # get/has/ownKeys traps make pythonmonkey report + # `typeof ret_val === "function"`, which doesn't support the + # generic subscript-based __getattr__ wrap_obj() normally relies + # on. Convert it to a plain JS array the same way Node's own + # usage does (Array.from(results)) and deserialize each value + # exactly like _wait()'s handle_complete does. + to_array = pm.eval("(rh) => Array.from(rh)") + raw_values = to_array(ret_val) + return [deserialize(v, self.serializers) for v in raw_values] + def on(self, *args): # deserialize job on event parameters before passing them to user defined callback def cb_deserialize_wrapper(callback): diff --git a/dcp/dry/aio.py b/dcp/dry/aio.py index 7e91223..3554152 100644 --- a/dcp/dry/aio.py +++ b/dcp/dry/aio.py @@ -14,9 +14,34 @@ import asyncio import inspect -# TODO: should we always do this? Is this user's responsibility? -import nest_asyncio -nest_asyncio.apply() +# LOCAL PATCH (Python 3.14 compatibility -- credit: Tom Tang's diagnosis): +# nest_asyncio.apply() was previously called unconditionally here. Its +# monkey-patch doesn't propagate asyncio's "current task" context correctly +# under 3.14's changed internals, which breaks any aiohttp call using +# timeout= (including pythonmonkey's XMLHttpRequest-internal.py, which +# backs DCP's socket.io polling transport) with "RuntimeError: Timeout +# should be used inside a task" *before* the request is even sent -- +# dcp-client reports this up the stack as the much more confusing +# "DCPError: no transports defined" (DCPC-1014). Confirmed independent of +# dcp with a minimal pythonmonkey+XMLHttpRequest repro (works fine under +# asyncio.run(), breaks only via this reentrant-loop patch). nest_asyncio's +# last release (1.6.0, Jan 2024) only ever claimed testing through Python +# 3.12. +# +# Fix: only apply the patch when a loop is ALREADY running in the current +# thread at import time -- i.e. only when reentrant run_until_complete() +# support is actually needed (Jupyter, a web server, a GUI app already +# running its own loop). A plain script importing dcp normally has no loop +# running yet at this point, never needed the patch in the first place, and +# now skips the code path that's broken on 3.14. Verified both branches: +# plain scripts (patch skipped, fixes 3.14) and inside a real Jupyter +# kernel via nbclient (patch still applies, notebook usage unaffected). +try: + asyncio.get_running_loop() + import nest_asyncio + nest_asyncio.apply() +except RuntimeError: + pass # no loop already running in this thread; nothing to patch loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) diff --git a/dcp/initialization.py b/dcp/initialization.py index ebe286b..4668b9b 100644 --- a/dcp/initialization.py +++ b/dcp/initialization.py @@ -15,6 +15,7 @@ from .dry import class_manager, aio from . import js from . import api +from ._pm_evaluator import evaluator as _pm_evaluator_module # state INIT_MEMO = None @@ -93,6 +94,17 @@ def init(**kwargs) -> Module: if INIT_MEMO is not None: return INIT_MEMO + # Makes pythonmonkey a real dcp-client platform: registers + # globalThis.__pmEvaluatorCtor, the SandboxConstructor localExec() + # uses (a real separate child process per job, not an in-process + # simulation -- see PYTHONMONKEY_EVALUATOR_PLAN.md for the full + # architecture). Must run BEFORE dcp-client's own init below -- + # preserved out of caution, matching an ordering requirement the + # previous in-process-simulation prototype needed; not yet + # independently reconfirmed as still required for this + # separate-process version specifically. + _pm_evaluator_module.install(aio.loop) + # initialize dcp js.dcp_client['init'](**kwargs) INIT_MEMO = True From 2fa35816254613c660159f20c100184d039196dc Mon Sep 17 00:00:00 2001 From: Daniel Desjardins Date: Sun, 20 Sep 2026 07:00:10 -0400 Subject: [PATCH 2/6] Fix six real bugs in the separate-process localExec() evaluator Real end-to-end job.localExec() now works with zero manual patch wiring (matching dcp_sample_job.py's shape, results = job.localExec() as the only line differing from job.exec()/job.wait()), confirmed against both a simple job and a heavier pycomod stress test. Getting there required finding and fixing six real bugs in dcp/_pm_evaluator/: - url.js's Invalid scheme crash: pythonmonkey ships a native URL (Node's bare vm sandbox doesn't), so url.js takes a branch that reads globalThis.location, which nothing set. Seed a placeholder. - pythonmonkey.null is truthy in Python: pm.eval("null") returns the type itself, not an instance, so "if handler:" can't detect "no handler yet". Check explicitly against pm.null too. - The wire protocol is asymmetric: LOG:/DIE:/MSG: prefixes are child->parent only; parent->child traffic is bare JSON. Fixed in three places that assumed symmetry. - A real race: the parent sends 'describe' before the child's bootstrap even starts, and the actual handler for it doesn't register until file #21 of 22 -- have to buffer every incoming line until the whole bootstrap finishes, not just until the first low-level handler exists. - Replaying that buffer has to happen inside an actively-running Python event loop, since the real handler is async and needs pythonmonkey's Python/JS bridge. - access-lists.js masks our own wire-protocol bridge globals as "unsafe" (Node's sandbox never has extras like these to mask); fixed by making them non-configurable, which its own masking logic already knows to skip. Also fixes the raise_on_first_work_error concern flagged as untested: job.py's localExec() returned normally instead of raising when a work function failed -- the exception landed as a value in the results list, not a promise rejection. Verified against several failure modes (every slice raises, one of several raises, an accidental NameError, an unserializable return value) -- all now raise a clean RuntimeError with the real Python traceback, no dcp-client internal JS stack noise. Co-Authored-By: Claude Sonnet 5 --- PYTHONMONKEY_EVALUATOR_PLAN.md | 248 +++++++++++++++++++--------- dcp/_pm_evaluator/_test_plumbing.py | 20 ++- dcp/_pm_evaluator/channel.py | 19 +-- dcp/_pm_evaluator/child.py | 133 ++++++++++----- dcp/_pm_evaluator/evaluator.py | 11 +- dcp/api/job.py | 49 +++--- 6 files changed, 309 insertions(+), 171 deletions(-) diff --git a/PYTHONMONKEY_EVALUATOR_PLAN.md b/PYTHONMONKEY_EVALUATOR_PLAN.md index 355b00b..39b6f68 100644 --- a/PYTHONMONKEY_EVALUATOR_PLAN.md +++ b/PYTHONMONKEY_EVALUATOR_PLAN.md @@ -349,84 +349,180 @@ With DCP services down (planned outage, sec5c), continued with code-only work th 6. If a job with more than a handful of slices, or a job that legitimately takes a while, hangs after all real results have arrived: that's confirmation `force_job_completion_when_done`'s concern is still needed too (port from lines 832-928). 7. Once real end-to-end success is confirmed (including error propagation): retest `pycomod_localexec_test.py` (heavier stress test) per sec6 step 6, then prepare the real PRs per sec6 step 7 -- monorepo MR (the `pythonmonkeyEvaluatorFactory()` + platform-gate change, sec4 item 1) and the bifrost2 PR (this branch), with a proper handoff doc. +## 5f. RESOLVED: real end-to-end success, six real bugs found and fixed + +Services came back up (planned outage ended). Picked up exactly per sec5e's +sequence. **The target API shape from sec1 now works, end to end, with +zero manual patch wiring**, for both the simple job and the heavy pycomod +one -- confirmed by directly running `dcp_local_job_test_no_timer_hacks.py` +-style and a new `pycomod/pycomod_localexec_test_CLEAN.py` (identical to +the old hacky version, minus every `pm_localexec_setup` import/call). +Both printed correct results (`YELLING!`; pycomod's `values[:5] = [25. 25. +25. 25. 25.]`, `dtype=float32`, matching the old hacky version's output +exactly) with `results = job.localExec()` as the only line differing from +`job.exec()`/`job.wait()`. + +Getting there required finding and fixing **six real, previously-unknown +bugs** -- this section exists so nobody re-derives them from scratch. Each +was confirmed by testing, not guessed: + +1. **`url.js`'s `Invalid scheme` crash.** `url.js` does `if (typeof URL + === 'undefined') { } else if (!('searchParams' + in new URL(globalThis.location))) {...}`. Node's real sandbox + (`evaluator-node.js`) runs inside a bare `vm.createContext({...})` with + no native `URL` at all, so it always takes the first branch and never + touches `location`. pythonmonkey ships its own native/polyfilled `URL` + (visible in the crash's own stack trace, `.../pythonmonkey/node_modules/ + core-js/...`), so this child takes the second branch -- and nothing + anywhere in this pipeline ever sets `globalThis.location` for a + non-browser platform. Fixed in `child.py`: seed a placeholder + (`new URL('file:///')`) before the bootstrap runs. +2. **`pythonmonkey.null` is truthy in Python.** `pm.eval("null")` returns + the `pythonmonkey.null` *class* itself, not an instance -- truthy like + any class. `if handler:` alone doesn't detect "no handler registered + yet"; a message arriving before one exists tried to *call* the class, + raising `TypeError: cannot create 'pythonmonkey.null' instances`. Fixed: + explicitly check `handler is not pm.null` too. +3. **The wire protocol is asymmetric, not symmetric.** Confirmed against + the real source (`lib/standaloneWorker.js`'s `postMessage`/`terminate`, + `sandbox/sa-ww-simulation.js`'s `receiveLine`): `LOG:`/`DIE:`/`MSG:` + prefixes exist *only* in the child-to-parent direction. Parent-to-child + traffic is bare JSON, no prefix. This file's original design assumed + symmetry in three places, each fixed: `evaluator.py`'s `postMessage` + (dropped the stray `'MSG:'` prefix it was prepending), `channel.py`'s + `terminate()` (sends real `{"type":"die"}` JSON, not a literal `"DIE:"` + string), and `child.py`'s `_socket_reader`/`_dispatch_incoming` (every + parent-to-child line is bare JSON for the handler, no prefix routing). +4. **A real race: `describe` arrives before the child can handle it, and + is lost, not delayed.** `Sandbox.start()` calls `describe()` + immediately after constructing the evaluator -- before the child has + even started its 22-file bootstrap. Node's real evaluator never hits + this: it runs its *entire* bootstrap synchronously before ever calling + `inputStream.on('data', ...)`, so anything sent early just sits in the + OS pipe's receive buffer, unread, until every file (including the one + that answers `describe`) has registered its listeners. This child + starts actively reading and dispatching *before* the bootstrap begins. + First fix attempt (replay once the low-level `onreadln` handler + registers, at bootstrap file #2 of 22) was necessary but not + sufficient -- the actual `describe` handler lives in + `calculate-capabilities.js`, file #21 of 22, so the message was still + lost. Real fix: buffer *every* incoming line unconditionally until the + *whole* bootstrap finishes, then replay all of them in order. +5. **Replaying that buffer has to happen inside an actively-running Python + event loop.** `calculate-capabilities.js`'s real `describe` handler is + `async (event) => { ... await protectedStorage.webGPUInitialization(); + ... }` -- it needs pythonmonkey's Python/JS event-loop bridge + (`PyEventLoop::getRunningLoop()`, the exact mechanism PythonMonkey PR + #509 patched). That bridge requires a loop that is *actually running on + this thread right now* -- false in the gap between two + `run_until_complete()` calls. Fixed by moving the buffered-message + replay into `_wait_for_exit()`, which the final `run_until_complete()` + actually drives, instead of calling it in the synchronous gap before + that call starts. +6. **`access-lists.js` masks our own wire-protocol bridge globals as + "unsafe."** `applyAllAccessLists()` (triggered by the real + `applyRequirements` message, after bootstrap) walks every property on + `globalThis` and masks any *configurable* one not on its allowlist, + turning it into an accessor that returns `undefined` until something + sets it. Node's real sandbox never has this problem -- its + `sandboxGlobal` is a bare object with only the specific properties it + explicitly assigns, so there's nothing extra to mask. This child adds + its own bridge functions (`__pmChildWriteLine`, `__pmOnReadlnHandler`, + `__pmChildDie`) directly onto the *same* `globalThis` the sandbox code + enumerates, so access-lists.js masked them right along with everything + else -- confirmed via an exact crash: masking `__pmChildWriteLine` + broke the very next `writeln()` call (the success-acknowledgment for + the `applyRequirements` message itself). Fixed using + `applyAccessLists()`'s own built-in escape hatch: it skips any property + where `Object.getOwnPropertyDescriptor(obj, prop)?.configurable` is + false. Defined all three bridge globals as non-configurable (still + writable, so `__pmOnReadlnHandler` can still be reassigned). + +**Also corrected**: this document's own rationale ("matches how Node.js's +`localExec()` actually works -- spawns a separate OS process") was +factually wrong. `node-localExec.js` shows Node's real `localExec()` uses +a same-process, pipe-connected worker (`standaloneWorker.js`'s +`workerFactory`), not a spawned child process. This doesn't change the +architecture decision here -- a genuinely separate process is *more* +isolated, which is the property that actually matters (no shared global, +period) -- but the doc should say that instead of the incorrect "matches +Node" claim. + +**Category A/B status**: `force_job_completion_when_done` did not need +porting -- confirmed, the real evaluator's `Sandbox`/`DistributiveWorker` +machinery completes jobs correctly on its own. + +`raise_on_first_work_error`'s concern was **real and confirmed present**, +deliberately tested with two work functions designed to fail (all slices +raise; one of several slices raises). In both cases `job.localExec()` +returned *normally* -- the raised exception landed as a `SpiderMonkeyError` +value inside the `results` list (at the correct index; unlike the old +in-process design, indexing was not misaligned) instead of failing the +call. A caller doing `''.join(results)` would hit a confusing `TypeError` +and never see the real bug. + +**Fixed** in `job.py`: `localExec()` now checks each deserialized result +for `isinstance(result, BaseException)` and raises a clean `RuntimeError` +using the same `.jsError.message` extraction the rejection-path handler +already used (factored into a shared `_clean_js_error_message()` helper). +Verified: both failing cases now raise `RuntimeError` with just the real +traceback (`File "", line N, in work_function` / the actual +exception) -- no dcp-client-bundle.js stack noise. Both passing tests +(simple + heavy) re-verified with no regression. Test scripts: +`C:\Users\danie\DCP\test_error_all_slices_raise.py`, +`C:\Users\danie\DCP\test_error_one_slice_raises.py`. + +**Two things this session did NOT do**, still open: +- **The real dcp monorepo platform-gate fix** (sec4 item 1, sec6 step 4). + Both successful runs went through `site-packages`' `dcp-client` bundle, + which already carries an old hand-patch from the original investigation + letting `pythonmonkey` through `job/index.js`'s platform check. The real + monorepo source (confirmed identical on `develop` and on + `pythonmonkey-websocket-transport`) still hard-throws `'localExec is not + supported on this platform'` for `dcpEnv.platform === 'pythonmonkey'`. A + fresh, unpatched `dcp-client` install would still fail. The fix itself + is narrow and already scoped: widen the check at `job/index.js` line + 644-645, add a `pythonmonkey` branch to the `SandboxConstructor` ternary + at 647-649 (a small `pythonmonkeyEvaluatorFactory()` returning + `globalThis.__pmEvaluatorCtor`, exported from `worker/evaluators/ + index.js`), and *deliberately leave* the `if (dcpEnv.platform === + 'nodejs')` local-file-marshaling block (696-800) untouched -- `job.py`'s + `_route_arguments_locally`/`_grant_local_file_origins` already replicate + its effect from the Python side, and that block depends on real Node + `fs`/`tmpfiles` pythonmonkey's `require()` doesn't provide. +- **Committing/pushing this work.** Everything above is still uncommitted + working-tree edits on this branch (`pythonmonkey-platform-support`). + ## 6. Remaining steps, in order -1. **Build the real child process** (`child.py` stage 2). Replace the - stage-1 stub with: import `pythonmonkey`, resolve the dcp-client - bootstrap-file paths *portably* (not the hardcoded - `C:\Users\danie\AppData\Roaming\...` paths `pm_localexec_setup.py` - uses — resolve relative to the installed `dcp` package's own location, - e.g. via `importlib.util.find_spec('dcp')`), run the same - `BOOTSTRAP_FILES` list (copy the list from - `pm_localexec_setup.py` lines 25-48, fix the paths), wire real - `writeln`/`onreadln`/`die` globals to the actual socket (write out - `LOG:`/`MSG:` lines instead of in-process dispatch; feed incoming - socket lines to whatever `onreadln` registered). Test whether this - genuinely-isolated bootstrap needs ANY of Category A's fixes — don't - assume it doesn't just because the global is no longer shared; test it. -2. **Build the JS-side evaluator constructor** (parent, in bifrost2, - injected via `pm.eval()` the same way `pm_localexec_setup.py` did) that - wraps `EvaluatorChannel`: `postMessage()` writes a `MSG:` line via a - Python bridge function exposed on `globalThis`; incoming lines - (delivered via `channel.on_line`, called on the main loop thanks to - `call_soon_threadsafe`) get parsed and dispatched to `onmessage`/ - `onerror`, matching `Sandbox.start()`'s expected contract exactly (same - shape as the existing, proven `__pmEvaluatorCtor` in - `pm_localexec_setup.py` lines 112-193 — reuse that shape, just change - the internals to talk to a real channel instead of fake dispatch). - Since spawning is not instant, buffer any `postMessage()` calls issued - before the channel finishes connecting (same pattern already proven in - `WebSocket.js`'s constructor this session — buffer-then-flush on - connect). -3. **Test the JS↔Python bridge directly** against the stage-2 child, - without going through real `Sandbox.start()`/dcp-client yet — construct - the evaluator directly via `pm.eval()`, call `postMessage`, confirm - `onmessage` fires with real data from the child's actual pythonmonkey - instance. -4. **Monorepo changes** (`C:\Users\danie\DCP\dcp-monorepo`, need a fresh - feature branch off `develop`): - - `src/dcp-client/job/index.js` ~line 645: accept `"pythonmonkey"` - platform (§4 item 1 above). - - New file, modeled on `src/dcp-client/worker/evaluators/node-localExec.js`: - a `pythonmonkeyEvaluatorFactory()` — but note, unlike Node's version, - this doesn't need to do the spawning itself (that's already handled - Python-side via `globalThis.__pmEvaluatorCtor`, exactly like the - existing FIXES_SUMMARY.md §3.1 diff already established) — likely a - much smaller file than `node-localExec.js`, just returning - `globalThis.__pmEvaluatorCtor`. - - Investigate/implement the `"complete"`/`"workError"` signal question - from §4 Category B — determine whether `Sandbox.start()`'s existing - onmessage handling (feeding into `StandaloneWorker`'s already-correct - `'result'`/error dispatch) is now sufent once real socket messages are - flowing, or whether the two bespoke hooks (§4 item 3) are still - needed. -5. **bifrost2 changes**: - - `dcp/initialization.py`: call the new evaluator-registration code - automatically in `init()`, before `js.dcp_client['init'](**kwargs)` - runs (ordering matters — confirmed by the original investigation). - Only port whatever Category A fixes step 1's testing showed are still - genuinely needed. - - `dcp/api/job.py`: rewrite `localExec()` cleanly, folding in whatever - Category B fixes are still needed after step 4's investigation - (route-arguments-locally at minimum; completion/error detection only - if still needed after checking `StandaloneWorker`'s existing handling). -6. **End-to-end test** against a real job. Use - `C:\Users\danie\DCP\dcp_local_job_test.py` as the reference test job - (uppercase-8-letters, Pyodide work function, `demo`/`dcp` compute - group) — but the goal is for the **exact clean script in §1** to work, - not a script with any manual patch wiring. Also retest - `pycomod_localexec_test.py` (heavier stress test — filesystem shipping, - extra Pyodide modules, cloudpickle round-trip) once the basic case - works. -7. **Only once real, end-to-end, tested** — prepare PRs: - - dcp monorepo MR (step 4's changes) - - bifrost2 PR (step 5's changes) — this is the real replacement for the - closed PR #49; make sure it's actually complete and tested this time, - not another partial cut. - - Write a handoff doc (matching the pattern of - `SPIDERMONKEY_VERSION_BUMP.md`/`LOCALEXEC_PYODIDE_FIX_NOTES.md`) - covering the full architecture, what was tested, what's flagged as - needing review. +Steps 1-3 and 6 below are **done** as of sec5f (real child process, real +JS-side evaluator constructor, real JS↔Python bridge, real end-to-end test +against both the simple and heavy jobs, both passing). What's left: + +1. **Monorepo changes** (`C:\Users\danie\DCP\dcp-monorepo`, need a fresh + feature branch off `develop`) -- the one piece that's still simulated + via an old hand-patched `site-packages` bundle rather than the real + source: + - `src/dcp-client/job/index.js` line 644-645: widen the platform check + to also accept `dcpEnv.platform === 'pythonmonkey'`. + - Line 647-649: add a `pythonmonkey` branch to the `SandboxConstructor` + ternary. New file, modeled on `worker/evaluators/node-localExec.js` + but much smaller: a `pythonmonkeyEvaluatorFactory()` that just returns + `globalThis.__pmEvaluatorCtor` (the spawning is already fully handled + Python-side). Export it from `worker/evaluators/index.js`. + - Deliberately leave the `if (dcpEnv.platform === 'nodejs')` block + (696-800) untouched -- see sec5f for why. + - Re-run both end-to-end tests against a fresh, unpatched `dcp-client` + install (not the hand-patched `site-packages` copy) to confirm the + real fix, not the old hand-patch, is what makes this work. +2. ~~Spot-check `raise_on_first_work_error`'s concern~~ -- **done** (sec5f): + found real, fixed in `job.py`, verified. +3. **Commit and push.** Everything in sec5f is still uncommitted + working-tree edits on this branch. Write a handoff doc (matching + `SPIDERMONKEY_VERSION_BUMP.md`'s pattern) covering the architecture, + the six bugs, what was tested, what's flagged as needing review + (`raise_on_first_work_error`, above). Open the bifrost2 PR (the real + replacement for closed PR #49) and, once step 1 lands, the monorepo MR. --- diff --git a/dcp/_pm_evaluator/_test_plumbing.py b/dcp/_pm_evaluator/_test_plumbing.py index 9459279..97797f1 100644 --- a/dcp/_pm_evaluator/_test_plumbing.py +++ b/dcp/_pm_evaluator/_test_plumbing.py @@ -1,8 +1,6 @@ """ -STAGE 1 smoke test: proves spawn+socket-connect+message-exchange+terminate -works at all on this machine, with zero pythonmonkey/bifrost2 involvement -- -isolates subprocess/socket mechanics from everything else before adding -that complexity on top. +Smoke test: spawn+socket-connect+message-exchange+terminate against the +real child.py. For a more thorough bootstrap check, see _test_bootstrap.py. Run directly: python -m dcp._pm_evaluator._test_plumbing """ @@ -27,16 +25,22 @@ async def main(): channel.spawn_and_connect() print(f"[parent] child connected in {time.time()-t0:.2f}s, pid=", channel.proc.pid, flush=True) - channel.write_line('MSG:{"type":"workerMessage","message":"hello from parent"}') + # No onreadln handler exists yet at this point in the bootstrap, so this + # is a no-op -- it only exercises the plumbing, not message dispatch. + # Bare JSON, no "MSG:" prefix (see child.py's _socket_reader). + channel.write_line('{"type":"workerMessage","message":"hello from parent"}') - await asyncio.sleep(2.0) + await asyncio.sleep(5.0) channel.terminate() await asyncio.sleep(0.5) assert any("pythonmonkey ready" in l for l in lines_received), "child pythonmonkey did not start" - assert any("child JS onreadln got" in l for l in lines_received), "child's real pythonmonkey JS did not receive our message via onreadln" - print("STAGE 2 PLUMBING TEST PASSED (real pythonmonkey in child, real socket round-trip)") + assert any("[bootstrap] all 22 files loaded successfully" in l for l in lines_received), "child did not finish loading the sandbox bootstrap" + # False is correct: sa-ww-simulation.js deletes these once captured + # privately (deliberate cleanup, not a clobbering bug). + assert any("writeln/onreadln still intact after bootstrap: False" in l for l in lines_received), "expected writeln/onreadln removed post-bootstrap; got True" + print("PLUMBING TEST PASSED (real pythonmonkey in child, full sandbox bootstrap, real socket round-trip)") if __name__ == "__main__": diff --git a/dcp/_pm_evaluator/channel.py b/dcp/_pm_evaluator/channel.py index 3781dbb..91ef06f 100644 --- a/dcp/_pm_evaluator/channel.py +++ b/dcp/_pm_evaluator/channel.py @@ -2,17 +2,13 @@ Parent-side channel: spawns the child evaluator process, listens for its connection, and bridges line-based traffic to/from it. -Uses a background thread with plain blocking sockets, NOT asyncio streams --- bifrost2's shared loop (dry.aio.loop) has nest_asyncio applied for -pythonmonkey's reentrant event-loop needs, and nest_asyncio's patched loop -was confirmed (empirically, via a real hang with no exception) to break -plain asyncio task/timeout scheduling in ways not worth fighting. A -background thread reading a blocking socket sidesteps that entirely; each -received line is handed back to the main loop/thread via -loop.call_soon_threadsafe(), which is the same safe cross-thread pattern -pythonmonkey's own C++ side uses to reach the event loop (see JobQueue.cc's -dispatchToEventLoop in the pythonmonkey-src rebuild this session). +Uses a background thread with plain blocking sockets, not asyncio streams: +bifrost2's shared loop has nest_asyncio applied (for pythonmonkey's +reentrant event-loop needs), which silently breaks asyncio task/timeout +scheduling. A thread reading a blocking socket sidesteps that; each line +is handed to the main loop via loop.call_soon_threadsafe(). """ +import json import os import socket import subprocess @@ -78,7 +74,8 @@ def write_line(self, line: str): pass def terminate(self): - self.write_line("DIE:") + # Bare JSON, not a "DIE:" line -- see child.py's _socket_reader. + self.write_line(json.dumps({"type": "die"})) self._stop = True if self.proc: try: diff --git a/dcp/_pm_evaluator/child.py b/dcp/_pm_evaluator/child.py index b3d9c1a..28193e3 100644 --- a/dcp/_pm_evaluator/child.py +++ b/dcp/_pm_evaluator/child.py @@ -1,16 +1,7 @@ """ Entry point for the separate-process pythonmonkey evaluator's child side. Invoked as: python --port (direct script path, NOT -`-m dcp._pm_evaluator.child` -- see channel.py's comment on _CHILD_SCRIPT -for why: `-m` forces importing the whole `dcp` package first, which this -process does not need and which was confirmed to slow/complicate startup). - -STAGE 2 (current): adds a real pythonmonkey instance and wires writeln/ -onreadln/die to the real socket, matching the wire protocol StandaloneWorker -(dcp-client's lib/standaloneWorker.js) expects on its read side. Does NOT -yet run the 22-file sandbox bootstrap -- that's the next stage, added only -once this stage is confirmed working (pythonmonkey starts reliably in a -spawned subprocess, and the writeln/onreadln bridge works end to end). +`-m dcp._pm_evaluator.child` -- see channel.py's comment on _CHILD_SCRIPT). """ import argparse import asyncio @@ -22,9 +13,8 @@ def _bootstrap_files(): - """Portable resolution of the 22 sandbox control-code files -- NOT the - hardcoded C:\\Users\\danie\\AppData\\... paths pm_localexec_setup.py - used. Resolved relative to wherever `dcp` is actually installed.""" + """Resolved relative to the installed `dcp` package, not a hardcoded path, + since this runs as a standalone child process that may live anywhere.""" spec = importlib.util.find_spec("dcp") dcp_dir = os.path.dirname(spec.origin) js_root = os.path.join(dcp_dir, "js", "node_modules") @@ -75,24 +65,71 @@ def _write_line(line: str): except OSError: pass - pm.globalThis["__pmChildWriteLine"] = _write_line + # access-lists.js masks every *configurable* global not on its allowlist + # (see below) -- including our own bridge functions, since a real Node + # sandbox never has extras like these to mask. Non-configurable is the + # escape hatch its own masking code already checks for. + def _define_protected_global(name, value): + pm.globalThis[name] = value + pm.eval(f""" + Object.defineProperty(globalThis, {name!r}, {{ + value: globalThis[{name!r}], + writable: true, + configurable: false, + enumerable: false, + }}); + """) + + _define_protected_global("__pmChildWriteLine", _write_line) + + # url.js only reads globalThis.location on engines that already have a + # native URL (pythonmonkey does; Node's bare sandbox doesn't, so it never + # hits this branch there). Nothing else in this pipeline sets `location`. + pm.eval("globalThis.location = new URL('file:///');") + + should_exit = threading.Event() + + def _die(): + # Tell the parent we're dying, then actually stop this process -- + # writing the socket line alone doesn't end the event loop. + _write_line("DIE:") + should_exit.set() + + _define_protected_global("__pmChildDie", _die) + + # The parent sends 'describe' the instant the socket connects, before + # this process has even started its bootstrap -- and the handler that + # answers it (calculate-capabilities.js) isn't registered until file #21 + # of 22. Node's real evaluator avoids this because it runs its whole + # bootstrap before ever reading its input stream, so early messages just + # sit in the OS pipe buffer. We read eagerly instead, so anything that + # arrives before the bootstrap fully finishes must be buffered and + # replayed afterward (_flush_pending, called from _wait_for_exit below). + _pending_lines = [] + _bootstrap_done = threading.Event() - # Same writeln contract as the real sandbox control-code files expect - # (see pm_localexec_setup.py's install(), which this mirrors) -- but - # now genuinely writing to a real socket instead of fake in-process - # dispatch. pm.eval(""" globalThis.writeln = function(line) { globalThis.__pmChildWriteLine(line); }; - globalThis.__pmOnReadlnHandler = null; + // Same non-configurable protection as the bridge globals above -- + // access-lists.js would otherwise mask this too. Stays writable so + // onreadln() can keep reassigning it. + Object.defineProperty(globalThis, '__pmOnReadlnHandler', { + value: null, + writable: true, + configurable: false, + enumerable: false, + }); globalThis.onreadln = function(fn) { globalThis.__pmOnReadlnHandler = fn; }; - globalThis.die = function() { globalThis.__pmChildWriteLine('DIE:'); }; + globalThis.die = function() { globalThis.__pmChildDie(); }; """) - should_exit = threading.Event() - def _socket_reader(): + # The wire protocol is asymmetric: LOG:/DIE:/MSG: prefixes are only + # used in the child->parent direction (sa-ww-simulation.js's send()). + # The parent always writes bare JSON, so every line here goes + # straight to the onreadln handler with no prefix routing. buf = b"" while True: try: @@ -105,32 +142,37 @@ def _socket_reader(): while b"\n" in buf: raw, buf = buf.split(b"\n", 1) text = raw.decode("utf-8", errors="replace") - if text.startswith("DIE:"): - should_exit.set() - return - elif text.startswith("MSG:"): - loop.call_soon_threadsafe(_dispatch_incoming, text[4:]) + if text: + loop.call_soon_threadsafe(_dispatch_incoming, text) should_exit.set() - def _dispatch_incoming(json_text: str): + def _dispatch_incoming(line: str): + if not _bootstrap_done.is_set(): + _pending_lines.append(line) + return + handler = pm.eval("globalThis.__pmOnReadlnHandler") + # pm.eval("null") returns the `pythonmonkey.null` type itself, which + # is truthy like any class -- `if handler:` alone can't tell "no + # handler yet" from a real one, and calling the type raises instead. + if handler and handler is not pm.null: + handler(line) + else: + _pending_lines.append(line) + + def _flush_pending(): + # Runs once the bootstrap is fully done, so every listener + # (including calculate-capabilities.js's) is registered. + _bootstrap_done.set() handler = pm.eval("globalThis.__pmOnReadlnHandler") - if handler: - # onreadln handlers, per the real protocol, receive the RAW - # "MSG:\n" line, not the decoded payload -- matches - # pm_localexec_setup.py's writeln() parsing the raw line itself. - handler("MSG:" + json_text + "\n") + if handler and handler is not pm.null: + while _pending_lines: + handler(_pending_lines.pop(0)) reader_thread = threading.Thread(target=_socket_reader, daemon=True) reader_thread.start() _write_line("LOG:pm evaluator child: pythonmonkey ready, pid=%d" % __import__("os").getpid()) - # STAGE 3: run the real 22-file sandbox bootstrap, in TRUE isolation -- - # this process's JS global is used for NOTHING else (no shared - # "Supervisor" role), so console/require/timer clobbering and - # access-lists masking should not be able to break anything else the - # way they did in the old in-process simulation. Verifying that - # empirically here, not assuming it. async def _run_bootstrap(): for f in _bootstrap_files(): _write_line(f"LOG:[bootstrap] loading {f}") @@ -150,15 +192,18 @@ async def _run_bootstrap(): sock.close() sys.exit(1) - # Confirm writeln/onreadln/die are still OUR functions after the - # bootstrap ran (i.e. nothing in the 22 files redefined them out from - # under us) -- this is exactly the kind of clobbering that broke the - # in-process simulation; check it explicitly rather than assume - # isolation fixed it. + # sa-ww-simulation.js deletes its own writeln/onreadln/die once captured + # privately -- expected to read False here, not a clobbering bug. still_ours = pm.eval("typeof globalThis.writeln === 'function' && typeof globalThis.onreadln === 'function'") _write_line(f"LOG:writeln/onreadln still intact after bootstrap: {still_ours}") async def _wait_for_exit(): + # Must run inside this active loop, not in the synchronous gap + # before it: calculate-capabilities.js's 'describe' handler is + # async and needs pythonmonkey's Python/JS event-loop bridge + # (PyEventLoop::getRunningLoop()), which only finds a loop that is + # actually running on this thread right now. + _flush_pending() while not should_exit.is_set(): await asyncio.sleep(0.05) diff --git a/dcp/_pm_evaluator/evaluator.py b/dcp/_pm_evaluator/evaluator.py index 53ea793..8a81835 100644 --- a/dcp/_pm_evaluator/evaluator.py +++ b/dcp/_pm_evaluator/evaluator.py @@ -59,7 +59,9 @@ } this.postMessage = function(msg) { - var line = 'MSG:' + JSON.stringify({ type: 'workerMessage', message: msg }); + // No LOG:/DIE:/MSG: prefix here -- those only apply child->parent + // (see handleLine above). The parent always writes bare JSON. + var line = JSON.stringify({ type: 'workerMessage', message: msg }); if (channelWrite) channelWrite(line); else sendBuffer.push(line); }; @@ -88,11 +90,8 @@ def install(loop: asyncio.AbstractEventLoop): - """Registers globalThis.__pmEvaluatorCtor, backed by real child - processes. Call once, before dcp-client's own init runs (matching the - ordering the original in-process version required -- unverified - whether that ordering constraint still applies here, but preserved - out of caution until tested otherwise).""" + """Registers globalThis.__pmEvaluatorCtor. Call before dcp-client's own + init runs, so it's available by the time a job needs it.""" async def _spawn_and_get_handle(on_line_js_callback): channel = EvaluatorChannel(loop) diff --git a/dcp/api/job.py b/dcp/api/job.py index 59cefee..b4d5d16 100644 --- a/dcp/api/job.py +++ b/dcp/api/job.py @@ -27,6 +27,14 @@ import urllib from .pyodide_work_function import get_work_function_string +def _clean_js_error_message(e): + """The real traceback pointing at the user's own work function is on + e.jsError.message; str(e) alone is buried under dcp-client's internal + JS stack frames.""" + js_error = getattr(e, 'jsError', None) + message = getattr(js_error, 'message', None) if js_error is not None else None + return message or str(e) + def job_maker(super_class): class Job(super_class): def __init__(self, job_js): @@ -314,18 +322,10 @@ def localExec(self, *args, **kwargs): a separate `.wait()` call afterward -- see wait()'s own note for why calling .wait() after .localExec() doesn't work anyway. - NEEDS TESTING once services are back up (see - PYTHONMONKEY_EVALUATOR_PLAN.md sec6): whether - force_job_completion_when_done()/raise_on_first_work_error()'s - concerns (local Worker never signaling a terminating 'stop'; - work-function errors getting swallowed) still apply under the - new separate-process evaluator, or whether a real child process - now surfaces these correctly on its own via - StandaloneWorker-compatible 'result'/error messages. Not yet - ported here pending that verification -- porting them - unconditionally without checking would risk reintroducing - exactly the kind of unverified, "looks done but isn't" change - this whole rework exists to avoid. + force_job_completion_when_done()'s concern (the local Worker + never signaling a terminating 'stop') did not need porting -- + confirmed by testing, the real evaluator's Sandbox/Worker + machinery completes jobs correctly on its own. """ self._before_exec() self._wrapper_set_attribute("_exec_called", True) @@ -337,19 +337,7 @@ def localExec(self, *args, **kwargs): try: ret_val = dry.aio.blockify(self.js_ref.localExec)(*args, **kwargs) except Exception as e: - # The real error message (a genuine Python traceback - # pointing at the user's own work function) is already - # present on e.jsError.message -- pythonmonkey.SpiderMonkeyError - # exposes the underlying JS Error as .jsError. Without - # this, the caller sees a SpiderMonkeyError whose - # message is buried under dcp-client's own internal JS - # stack frames. Re-raise with just the real message so - # job.localExec() fails pointing at the user's own code. - js_error = getattr(e, 'jsError', None) - message = getattr(js_error, 'message', None) if js_error is not None else None - if message: - raise RuntimeError(message) from None - raise + raise RuntimeError(_clean_js_error_message(e)) from None finally: # These files hold real job argument/slice-value data # (potentially sensitive) and are not otherwise cleaned up. @@ -371,7 +359,16 @@ def localExec(self, *args, **kwargs): # exactly like _wait()'s handle_complete does. to_array = pm.eval("(rh) => Array.from(rh)") raw_values = to_array(ret_val) - return [deserialize(v, self.serializers) for v in raw_values] + results = [deserialize(v, self.serializers) for v in raw_values] + + # A slice whose work function raised resolves normally here -- + # the exception lands as a value in `results`, not a rejection + # of the promise above -- so it has to be checked explicitly. + for result in results: + if isinstance(result, BaseException): + raise RuntimeError(_clean_js_error_message(result)) from None + + return results def on(self, *args): # deserialize job on event parameters before passing them to user defined callback From 206cab42cbb8f84b70c0caf04aa6ce6463fe8564 Mon Sep 17 00:00:00 2001 From: Daniel Desjardins Date: Sun, 20 Sep 2026 07:14:27 -0400 Subject: [PATCH 3/6] Document PythonMonkey PR #509 as a hard prerequisite, not a footnote Neither the plan doc nor the code said clearly that this branch requires PR #509 installed first -- it was one passing mention in a list of "related work already done". Confirmed by testing: swapping the old pythonmonkey.pyd back in shows SharedArrayBuffer/Atomics are simply absent, so Pyodide's threaded WASM build can't link, and separately the JobQueue checkpoint fixes in PR #509's JSFunctionProxy.cc/JSMethodProxy.cc are needed or the evaluator's async JS callbacks hang. Checking out just this branch on an unpatched pythonmonkey doesn't error -- it hangs. Added a loud prerequisites section at the top of the plan doc, and pointers at the actual code (dcp/_pm_evaluator/__init__.py, dcp/initialization.py's install() call site) for anyone who won't read the whole doc first. Also corrected __init__.py's "matches Node.js" claim (Node's real localExec() uses a same-process pipe-connected worker, not a separate process -- already caught and fixed in the plan doc earlier, but this file's own docstring still had the old claim) and its wire-protocol description, which didn't mention the protocol is asymmetric between the two directions. Co-Authored-By: Claude Sonnet 5 --- PYTHONMONKEY_EVALUATOR_PLAN.md | 43 ++++++++++++++++++++++++++++++++++ dcp/_pm_evaluator/__init__.py | 35 +++++++++++++++------------ dcp/initialization.py | 13 ++++++---- 3 files changed, 71 insertions(+), 20 deletions(-) diff --git a/PYTHONMONKEY_EVALUATOR_PLAN.md b/PYTHONMONKEY_EVALUATOR_PLAN.md index 39b6f68..e66d0bf 100644 --- a/PYTHONMONKEY_EVALUATOR_PLAN.md +++ b/PYTHONMONKEY_EVALUATOR_PLAN.md @@ -9,6 +9,49 @@ works. --- +## PREREQUISITE: this branch does not work with the pythonmonkey currently +## on PyPI/npm/site-packages -- it requires PythonMonkey PR #509 installed +## first, not just "eventually." + +https://github.com/Distributive-Network/PythonMonkey/pull/509 + +Not optional, not a future nice-to-have -- confirmed by literally swapping +the old `pythonmonkey.pyd` back in and testing: +`pm.eval('typeof SharedArrayBuffer')` → `undefined` on the old build (was +`'function'` on PR #509's build). Two separate, both-required reasons: + +1. **Pyodide's threaded WASM build cannot link without SharedArrayBuffer/ + Atomics** ("LinkError: shared memory is disabled" otherwise) -- this is + unconditional for *any* Pyodide work function, regardless of anything in + this bifrost2 branch. PR #509 is what adds this. +2. **The JobQueue rewrite that comes with PR #509's SpiderMonkey rebuild + needs its own checkpoint fixes to not hang** -- specifically the + `js::RunJobs(cx)` calls added in `JSFunctionProxy.cc`/`JSMethodProxy.cc` + (PR #509's most recent commit). Without them, the real evaluator's async + JS callbacks (e.g. `calculate-capabilities.js`'s `describe` handler) + hang forever, the exact same symptom class this whole session's work + here fixed on the bifrost2 side. + +**Checking out just this bifrost2 branch, on top of an unpatched +pythonmonkey, will not produce a useful error -- it will hang.** Install +the full current state of PR #509 first (build it, or get a build from +whoever has one) before testing anything in this document. + +--- + +## ALSO REQUIRED, NOT YET DONE: the dcp monorepo platform-gate change + +`src/dcp-client/job/index.js` still hard-throws `'localExec is not +supported on this platform'` for `pythonmonkey` (line 644-645) in the real +monorepo source. Every real end-to-end test in this document passed +against `site-packages`'s `dcp-client` bundle, which carries an *old hand +patch* from an earlier investigation working around this exact throw -- +not the real fix. See sec5f/sec6 below for the exact, narrow change needed +(a few lines in `job/index.js` + one small new file). **Deliberately left +undone in this branch** -- intended to be a separate monorepo MR. + +--- + ## 0. Why this document exists `job.localExec()` under pythonmonkey was gotten working during an earlier diff --git a/dcp/_pm_evaluator/__init__.py b/dcp/_pm_evaluator/__init__.py index 75f234c..299f882 100644 --- a/dcp/_pm_evaluator/__init__.py +++ b/dcp/_pm_evaluator/__init__.py @@ -1,19 +1,24 @@ """ -A real, separate-process evaluator for job.localExec() under pythonmonkey, -matching how localExec() actually works on Node.js -- a genuinely separate -process, not a shared JS global. +A real, separate-process evaluator for job.localExec() under pythonmonkey. +Spawns a real child process for the sandbox rather than sharing one JS +global with the Supervisor -- more isolated than Node's own localExec() +(which uses a same-process, pipe-connected worker; see +src/dcp-client/worker/evaluators/node-localExec.js), but the isolation is +what actually matters here, not matching Node's specific mechanism. -See src/dcp-client/worker/evaluators/node-localExec.js (dcp monorepo) for -the reference architecture this mirrors: a SandboxConstructor that spawns -a worker connected over a socket, wrapped to satisfy the same postMessage/ -onmessage/onerror/terminate/addEventListener contract Sandbox.start() -expects from any platform's evaluator. +REQUIRES PythonMonkey PR #509 installed first (SpiderMonkey rebuilt for +SharedArrayBuffer/Atomics, plus JobQueue checkpoint fixes) -- see +PYTHONMONKEY_EVALUATOR_PLAN.md's top section. Without it, jobs hang, they +don't error cleanly. -Wire protocol (matches lib/standaloneWorker.js's StandaloneWorker exactly, -so a real dcp-worker-shaped child speaks a protocol dcp-client already -understands): newline-delimited, three line-prefixes: - - "LOG:" -- debug/log line, informational only - - "DIE:" -- child is shutting down - - "MSG:" -- {"type": "workerMessage", "message": ...} - (postMessage payload) or {"type": "result", ...} +Wire protocol (matches lib/standaloneWorker.js's StandaloneWorker, so a +real dcp-worker-shaped child speaks a protocol dcp-client already +understands) -- ASYMMETRIC, not the same both directions: + - child -> parent: newline-delimited, prefixed -- + "LOG:" -- debug/log line, informational only + "DIE:" -- child is shutting down + "MSG:" -- {"type": "workerMessage", "message": ...} (a + postMessage payload) or {"type": "result", ...} + - parent -> child: newline-delimited, bare JSON, NO prefix -- e.g. + {"type": "workerMessage", "message": ...} or {"type": "die"} """ diff --git a/dcp/initialization.py b/dcp/initialization.py index 4668b9b..8212f25 100644 --- a/dcp/initialization.py +++ b/dcp/initialization.py @@ -98,11 +98,14 @@ def init(**kwargs) -> Module: # globalThis.__pmEvaluatorCtor, the SandboxConstructor localExec() # uses (a real separate child process per job, not an in-process # simulation -- see PYTHONMONKEY_EVALUATOR_PLAN.md for the full - # architecture). Must run BEFORE dcp-client's own init below -- - # preserved out of caution, matching an ordering requirement the - # previous in-process-simulation prototype needed; not yet - # independently reconfirmed as still required for this - # separate-process version specifically. + # architecture, prerequisites, and known bugs). Must run BEFORE + # dcp-client's own init below -- preserved out of caution, matching + # an ordering requirement the previous in-process-simulation + # prototype needed; not yet independently reconfirmed as still + # required for this separate-process version specifically. + # + # REQUIRES PythonMonkey PR #509 installed first -- jobs hang, not + # error, without it. See PYTHONMONKEY_EVALUATOR_PLAN.md's top section. _pm_evaluator_module.install(aio.loop) # initialize dcp From a5251491195e0f19fb2767eb719af8126ed04a8d Mon Sep 17 00:00:00 2001 From: Daniel Desjardins Date: Sun, 20 Sep 2026 07:22:59 -0400 Subject: [PATCH 4/6] Fail fast in localExec() on old pythonmonkey instead of hanging exec() never loads pyodide locally, so it works fine on any pythonmonkey. localExec() does run pyodide in-process and needs SharedArrayBuffer/Atomics (PR #509+), which an unpatched build has no clean way to signal -- it just hangs deep inside sandbox startup. Check the actual capability at the top of localExec() instead of gating dcp.init() for everyone, since most callers never use localExec() at all. Co-Authored-By: Claude Sonnet 5 --- dcp/api/job.py | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/dcp/api/job.py b/dcp/api/job.py index b4d5d16..6ea8b06 100644 --- a/dcp/api/job.py +++ b/dcp/api/job.py @@ -35,6 +35,21 @@ def _clean_js_error_message(e): message = getattr(js_error, 'message', None) if js_error is not None else None return message or str(e) +def _check_localexec_pm_capability(): + """localExec() runs the work function's pyodide sandbox in-process, + which needs SharedArrayBuffer/Atomics (PythonMonkey PR #509+) to link. + exec() never loads pyodide locally -- it only dispatches -- so this + check is scoped to localExec() rather than gating dcp.init() for + everyone. Fails fast with a clear error instead of the silent hang an + unpatched pythonmonkey produces deep inside sandbox startup.""" + if pm.eval("typeof SharedArrayBuffer") == "undefined": + raise RuntimeError( + "job.localExec() requires a PythonMonkey build with " + "SharedArrayBuffer/Atomics support (PythonMonkey PR #509 or " + "later) -- the pythonmonkey installed here doesn't have it. " + "job.exec() is unaffected and works normally." + ) + def job_maker(super_class): class Job(super_class): def __init__(self, job_js): @@ -327,6 +342,8 @@ def localExec(self, *args, **kwargs): confirmed by testing, the real evaluator's Sandbox/Worker machinery completes jobs correctly on its own. """ + _check_localexec_pm_capability() + self._before_exec() self._wrapper_set_attribute("_exec_called", True) From 258c79692f58443fe56f38929b7c5dd2c4404f19 Mon Sep 17 00:00:00 2001 From: Daniel Desjardins Date: Sun, 20 Sep 2026 07:52:48 -0400 Subject: [PATCH 5/6] Let job.localExec(); results = job.wait() mirror exec()'s pattern localExec() already runs the whole job to completion before it returns, so a listener registered inside wait() afterward always misses the 'complete' event -- EventEmitters don't replay past events. Cache localExec()'s results on the job and have wait() return them directly when present, instead of trying (and failing) to listen for an event that already fired. This means swapping exec() for localExec() never requires moving where `results =` appears: job.exec(); results = job.wait() and job.localExec(); results = job.wait() are now identical in shape. The original results = job.localExec() single-call form still works unchanged. Co-Authored-By: Claude Sonnet 5 --- dcp/api/job.py | 53 ++++++++++++++++++++++++++++++++++---------------- 1 file changed, 36 insertions(+), 17 deletions(-) diff --git a/dcp/api/job.py b/dcp/api/job.py index 6ea8b06..080edf9 100644 --- a/dcp/api/job.py +++ b/dcp/api/job.py @@ -61,6 +61,7 @@ def __init__(self, job_js): job_js.modules = [] #TODO: why is this only done this way for job modules? self._wrapper_set_attribute("fs", JobFS()) self._wrapper_set_attribute("_exec_called", False) + self._wrapper_set_attribute("_local_results", None) self.aio.exec = self._exec; self.aio.wait = self._wait; @@ -180,6 +181,20 @@ def handle_accepted(): def _wait(self): if not self._exec_called: raise Exception("Wait called before exec()") + + # localExec() already ran the whole job to completion, 'complete' + # event included, before it returned (see its own docstring) -- + # by the time wait() could register a listener the event has + # already fired and gone unheard (EventEmitters don't replay). + # Returning the results it already collected -- rather than + # trying to listen for an event that already happened -- is what + # lets `job.localExec(); results = job.wait()` mirror the real + # `job.exec(); results = job.wait()` pattern exactly. + if self._local_results is not None: + complete_future = asyncio.Future() + complete_future.set_result(self._local_results) + return complete_future + complete_future = asyncio.Future() def handle_complete(resultHandle): serialized_results = resultHandle["values"]() @@ -198,19 +213,16 @@ def exec(self, *args): return results def wait(self): - # NOTE: for a localExec() job, registering these listeners here - # is too late to ever see the 'complete' event -- the real JS - # localExec() Promise (awaited inside localExec() below) does - # not resolve until the WHOLE job (including the 'complete' - # event this function listens for) has already finished, so by - # the time _wait() runs, that event has already fired and been - # missed (EventEmitters don't replay past events to newly-added - # listeners). For localExec(), use its own return value instead - # of calling .wait() afterward (matching the real Node.js usage - # pattern -- no separate .wait() call at all). This method - # remains correct and necessary for real distributed jobs via - # exec()/aio.exec(), where 'complete' genuinely arrives later, - # well after this registration. + # For a localExec() job, the real JS localExec() Promise + # (awaited inside localExec() below) doesn't resolve until the + # WHOLE job -- 'complete' event included -- has already + # finished, so a *fresh* listener registered here would always + # miss it (EventEmitters don't replay past events). _wait() + # handles this by returning localExec()'s already-cached + # results instead of listening for an event that already fired. + # For real distributed jobs via exec()/aio.exec(), no such cache + # exists yet at this point, so this registers listeners and + # blocks for 'complete' normally. return dry.aio.blockify(self._wait)() def _route_arguments_locally(self): @@ -332,10 +344,16 @@ def localExec(self, *args, **kwargs): registration) -- without it the raw user Python source is sent as-is, which never calls dcp.set_slice_handler() (-> ENOSLICEHANDLER in the pyodide worktime). Mirrors _exec()'s - setup, matching the real Node.js usage pattern - (`const results = await job.localExec()`) instead of requiring - a separate `.wait()` call afterward -- see wait()'s own note - for why calling .wait() after .localExec() doesn't work anyway. + setup. + + Works both as a single call (`results = job.localExec()`, + matching Node's own usage) and, since the results are cached on + the job once collected, as `job.localExec(); results = + job.wait()` -- matching `job.exec(); results = job.wait()` + exactly, so swapping exec for localExec never requires moving + where `results =` appears. See wait()'s own note for why a + *fresh* listener registered after the fact would otherwise miss + the 'complete' event. force_job_completion_when_done()'s concern (the local Worker never signaling a terminating 'stop') did not need porting -- @@ -385,6 +403,7 @@ def localExec(self, *args, **kwargs): if isinstance(result, BaseException): raise RuntimeError(_clean_js_error_message(result)) from None + self._wrapper_set_attribute("_local_results", results) return results def on(self, *args): From fb3c3c9aae4167ec25a0a459f49f1a945b717405 Mon Sep 17 00:00:00 2001 From: Daniel Desjardins Date: Sun, 20 Sep 2026 08:21:10 -0400 Subject: [PATCH 6/6] Close two more localExec() data-locality gaps: work function source and range-shaped input sets Node's own localExec() rewrites the work function source and a range-shaped input set (e.g. MultiRangeObject) to local files before deploy, exactly like it already does for job arguments and array input sets -- but that rewrite is gated `if (dcpEnv.platform === 'nodejs')`, so pythonmonkey never got it. Two real leaks resulted: the full work function source (embedded in workFunctionURI as a data: URI) and a range object's raw contents both went to the scheduler unmodified even for a "local" job. Extends _route_arguments_locally() to also rewrite both to local temp files with scoped origin grants, using the same real dcp-client classes (globalThis.dcp['range-object'], globalThis.dcp.compute.RemoteDataSet) Node's own code uses for range detection, rather than reimplementing that logic separately. Verified: work function source now lands in a local file (confirmed by reading it back), a synthetic MultiRangeObject-shaped input routes to a local file instead of inlining, and the full existing regression suite (simple job, pycomod with JobFS, error paths, wait() symmetry) still passes unchanged. Co-Authored-By: Claude Sonnet 5 --- dcp/api/job.py | 70 ++++++++++++++++++++++++++++++++++++++------------ 1 file changed, 54 insertions(+), 16 deletions(-) diff --git a/dcp/api/job.py b/dcp/api/job.py index 080edf9..e7446c4 100644 --- a/dcp/api/job.py +++ b/dcp/api/job.py @@ -229,21 +229,27 @@ def _route_arguments_locally(self): """ localExec() otherwise unconditionally routes job data through the real scheduler: jobArguments once a payload exceeds a size - threshold (a scheduler-hosted URL), and slice *values* via a + threshold (a scheduler-hosted URL), slice *values* via a dedicated, unconditional bulk upload (addSlices(), during the "uploading" state) that runs regardless of size unless - jobRef.marshaledDataValues is already set. Real Node localExec() - avoids this entirely by writing both to local temp files and - granting the local worker a narrow, path-scoped file:// origin - per file instead -- this mirrors that. + jobRef.marshaledDataValues is already set, a range-shaped input + set (e.g. MultiRangeObject) inlined as-is in the deploy payload, + and the work function source embedded directly in + workFunctionURI as a data: URI. Real Node localExec() avoids all + of this by writing each to a local temp file and granting the + local worker a narrow, path-scoped file:// origin per file + instead -- this mirrors that for every case Node handles + (job/index.js's own platform-gated rewrite only runs for + 'nodejs', so pythonmonkey needs its own copy). Must run SYNCHRONOUSLY, immediately after _before_exec() - populates jobArguments/jobInputData and before the real JS - localExec() call (and therefore deployJob()'s upload) ever - runs: the scheduler snapshots jobArguments during deploy, and - the slice-value upload is scheduled immediately once deploy - completes -- both well before localWorker/originManager exist - or any async callback would get a turn to run. + populates jobArguments/jobInputData/workFunctionURI and before + the real JS localExec() call (and therefore deployJob()'s + upload) ever runs: the scheduler snapshots all of this during + deploy, and the slice-value upload is scheduled immediately + once deploy completes -- both well before localWorker/ + originManager exist or any async callback would get a turn to + run. Returns the list of (path, purpose) grants still needing origin access once the local worker exists (see @@ -255,8 +261,8 @@ def _route_arguments_locally(self): written_paths = [] - def _write_temp_file(encoded_str): - fd, path = tempfile.mkstemp(prefix="bifrost2-localExec-arg-", suffix=".kvin") + def _write_temp_file(encoded_str, suffix): + fd, path = tempfile.mkstemp(prefix="bifrost2-localExec-arg-", suffix=suffix) with _os.fdopen(fd, "w", encoding="utf-8") as f: f.write(encoded_str) written_paths.append(path) @@ -266,20 +272,52 @@ def _write_temp_file(encoded_str): rewrite = pm.eval(""" (jobRef) => { const KVIN = new (require('kvin').KVIN)(); + const { SuperRangeObject, MultiRangeObject } = globalThis.dcp['range-object']; + const { RemoteDataSet } = globalThis.dcp.compute; const grants = []; const toLocalURL = (v, purpose) => { if (v instanceof URL) return v; const encoded = KVIN.stringify(v); - const filePath = globalThis.__pmWriteLocalArgFile(encoded); + const filePath = globalThis.__pmWriteLocalArgFile(encoded, '.kvin'); grants.push({ path: filePath, purpose }); return new URL('file://' + filePath); }; if (jobRef.jobArguments) { jobRef.jobArguments = jobRef.jobArguments.map((v) => toLocalURL(v, 'fetchArguments')); } - if (Array.isArray(jobRef.jobInputData)) { - jobRef.marshaledDataValues = jobRef.jobInputData.map((v) => toLocalURL(v, 'fetchData')); + + // Same detection marshalInputData() uses (job/index.js) to + // decide dataRange vs dataValues -- a range object encodes + // job design (e.g. combinatorial dimensions), not raw slice + // values, but Node treats it as equally local-only, so this + // matches that rather than treating it as safe to inline. + const inputData = jobRef.jobInputData; + const isRange = inputData instanceof SuperRangeObject + || (inputData && inputData.hasOwnProperty && inputData.hasOwnProperty('ranges') && inputData.ranges instanceof MultiRangeObject) + || (inputData && inputData.hasOwnProperty && inputData.hasOwnProperty('start') && inputData.hasOwnProperty('end')); + + if (isRange) { + const filePath = globalThis.__pmWriteLocalArgFile(JSON.stringify(inputData), '.json'); + grants.push({ path: filePath, purpose: 'fetchData' }); + jobRef.marshaledDataRange = new RemoteDataSet([new URL('file://' + filePath)]); + jobRef.rangeLength = inputData.length; + } else if (Array.isArray(inputData)) { + jobRef.marshaledDataValues = inputData.map((v) => toLocalURL(v, 'fetchData')); + } + + // workFunctionURI is a separate field from jobArguments, set + // directly by bifrost2's _before_exec() as a data: URI + // embedding the full wrapped work function source -- Node's + // own rewrite (job/index.js, gated to 'nodejs') never runs + // for pythonmonkey, so nothing else does this. + if (typeof jobRef.workFunctionURI === 'string' && jobRef.workFunctionURI.startsWith('data:')) { + const commaIndex = jobRef.workFunctionURI.indexOf(','); + const source = decodeURIComponent(jobRef.workFunctionURI.slice(commaIndex + 1)); + const filePath = globalThis.__pmWriteLocalArgFile(source, '.js'); + grants.push({ path: filePath, purpose: 'fetchWorkFunctions' }); + jobRef.workFunctionURI = new URL('file://' + filePath).href; } + return grants; } """)