Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
597 changes: 597 additions & 0 deletions PYTHONMONKEY_EVALUATOR_PLAN.md

Large diffs are not rendered by default.

24 changes: 24 additions & 0 deletions dcp/_pm_evaluator/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
"""
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.

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, 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:" <text> -- debug/log line, informational only
"DIE:" -- child is shutting down
"MSG:" <json> -- {"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"}
"""
62 changes: 62 additions & 0 deletions dcp/_pm_evaluator/_test_bootstrap.py
Original file line number Diff line number Diff line change
@@ -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())
48 changes: 48 additions & 0 deletions dcp/_pm_evaluator/_test_exec_control.py
Original file line number Diff line number Diff line change
@@ -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")
47 changes: 47 additions & 0 deletions dcp/_pm_evaluator/_test_plumbing.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
"""
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
"""
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)

# 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(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("[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__":
asyncio.run(main())
57 changes: 57 additions & 0 deletions dcp/_pm_evaluator/_test_real_job.py
Original file line number Diff line number Diff line change
@@ -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")
89 changes: 89 additions & 0 deletions dcp/_pm_evaluator/channel.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
"""
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 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
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):
# 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:
self.proc.wait(timeout=5)
except subprocess.TimeoutExpired:
self.proc.kill()
if self.sock:
try:
self.sock.close()
except OSError:
pass
Loading
Loading