Skip to content

localExec(): keep slice results off the scheduler - #50

Closed
dan-distributive wants to merge 1 commit into
pythonmonkey-platform-supportfrom
localexec-local-result-storage
Closed

dan-distributive wants to merge 1 commit into
pythonmonkey-platform-supportfrom
localexec-local-result-storage

Conversation

@dan-distributive

Copy link
Copy Markdown

Summary

localExec() already kept job arguments, input sets, and the work function local (see the parent branch). Slice results were the one remaining gap: the local worker POSTs each computed value to the real resultSubmitter service, and the job receives it back via the scheduler's pubsub relay — the actual value left the machine, not just a completion signal.

This closes that gap using dcp-client's own existing, public job.setResultStorage(url, postParams) API — normally used to redirect result storage to a self-hosted location (S3, Dropbox, etc; see "DCP Job Architecture for data routing.pptx", slides 5-8). Pointing it at a small local HTTP server keeps the real value on loopback; the server responds with a small opaque token, and that token — not the real value — is all that reaches the scheduler and comes back through the completely unmodified resultSubmitter/pubsub path. This is exactly how the feature already behaves for any other off-prem storage target, not a special case.

Results are captured via KVIN (application/x-kvin), not plain JSON — JSON would silently mangle binary/pickled results and strip a failed slice's error object down to an inert dict, which would have broken the existing raise_on_first_work_error behavior.

No public API changes — job.localExec() behaves identically from the caller's side; this is entirely internal plumbing, wired into localExec() itself (no new methods for callers to remember).

Base branch note: this is stacked on pythonmonkey-platform-support since it depends on localExec()'s existing local-routing machinery.

Verification

  • A sentinel value proven, via a live test against the real scheduler, to never appear in the scheduler-relayed ResultHandle contents, while job.localExec()'s returned value is still correct.
  • Full existing regression suite passes unchanged: simple job, pycomod (cloudpickled/numpy results via JobFS), all four error scenarios (all-raise, one-raises, NameError, unserializable return), and job.wait() symmetry.

Test plan

  • Review _setup_local_result_storage()/_grant_send_results_origin() in dcp/api/job.py
  • Run a real localExec() job and confirm results are correct
  • Confirm a failing slice still raises a clean RuntimeError with the real traceback

🤖 Generated with Claude Code

localExec() already kept job arguments, input sets, and the work
function local, but slice *results* still round-tripped through the
real scheduler: the local worker POSTs each computed value to the
resultSubmitter service, and the job receives it back via the
scheduler's pubsub relay -- the actual value left the machine, not
just a completion signal.

dcp-client already ships a public API for exactly this class of
problem: job.setResultStorage(url, postParams) redirects a slice's
result upload to a self-hosted location instead of the scheduler's
own storage (normally used for S3, Dropbox, etc. -- see "DCP Job
Architecture for data routing.pptx", slides 5-8). Pointing it at a
small local HTTP server keeps the real value on loopback; the server
responds with a small opaque token, and that token -- not the real
value -- is all that travels to the scheduler and back through the
completely unmodified resultSubmitter/pubsub path, exactly how the
feature already behaves for any other off-prem storage target.

Results are captured via KVIN (application/x-kvin), not plain JSON:
JSON would silently mangle binary/pickled results and strip a failed
slice's error object down to an inert dict, breaking the existing
raise_on_first_work_error behavior. Confirmed both cases explicitly:
a sentinel value proven to never appear in the scheduler-relayed
ResultHandle, and the full existing regression suite (simple job,
pycomod with cloudpickled/numpy results via JobFS, all four error
scenarios, job.wait() symmetry) passing unchanged.

No public API changes -- job.localExec() behaves identically from
the caller's side; this is entirely internal plumbing.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
@dan-distributive

Copy link
Copy Markdown
Author

Closing per a decision made with Ryan: the default scheduler round-trip for slice results is fine as-is. If a job needs results kept off the scheduler, that's already achievable today via the existing, public job.setResultStorage() API — no bifrost2-specific plumbing needed. This also brings localExec() closer to matching Node's own behavior, which doesn't do anything special for results either.

The other locality fixes (job arguments, input sets, work function source — PR #51) are unaffected by this: those mirror what Node's own localExec() already does for itself, not something invented here, so they stay.

@dan-distributive
dan-distributive deleted the localexec-local-result-storage branch September 22, 2026 16:36
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant