Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
4562b63
test(verify): run the script the export produces
kz930 Sep 2, 2026
60fbc77
test(verify): say it once, in the shape the code does not already give
kz930 Sep 3, 2026
4ab2c9f
test(verify): build the script's header the way the export builds it
kz930 Sep 4, 2026
dd36970
test(verify): read a script's own exit as the exit it asked for
kz930 Sep 9, 2026
3e87f6b
test(verify): keep the input schema's string columns as strings
kz930 Sep 9, 2026
6e2e71d
test(verify): hand the script the schema, and a JVM path's integers e…
kz930 Sep 10, 2026
d7385ce
test(verify): give an empty input its declared columns
kz930 Sep 10, 2026
06ffa18
test(verify): say the reader's four options once each
kz930 Sep 10, 2026
f0616cc
style: wrap a line scalafmt rejects
kz930 Sep 10, 2026
f99d1ef
test(verify): cut the doc to what the code cannot say
kz930 Sep 10, 2026
0f78bb1
Merge branch 'main' into feat/verify-run-generated-script
kz930 Sep 15, 2026
04f5555
test(verify): hand the script a boolean column with a hole as booleans
kz930 Sep 18, 2026
2c96035
test(verify): write an integral column the way the engine's writer does
kz930 Sep 19, 2026
f4b9496
fix(verify): hand a binary column to the script as the bytes the run …
kz930 Sep 19, 2026
36f79d8
test(verify): name the file a source reads, which its body leaves to …
kz930 Sep 19, 2026
67f361f
test(verify): let an empty input frame keep the columns pandas gives it
kz930 Sep 22, 2026
84fedad
test(verify): hand the script a timestamp the year 2500 and nine digi…
kz930 Sep 23, 2026
658b1a6
test(verify): hand the script an empty input with its columns
kz930 Sep 24, 2026
7b3c37a
test(verify): hand a pickled binary cell to the script as the object
kz930 Sep 24, 2026
d8440e2
test(verify): hand a real fitted model to the script in the harness spec
kz930 Sep 25, 2026
83a6c1e
test(verify): write down the dtype each output column was left in
kz930 Sep 25, 2026
ccb2046
test(verify): say how the script side reads a timestamp and which int…
kz930 Sep 27, 2026
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
#!/usr/bin/env python3
#
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
"""
Persistent worker for the Path B (standalone) verify path.

Motivation: forking a fresh interpreter per operator pays the pandas/plotly
import cost (~260-310 ms) on every spawn, while the operator's actual compute
on the tiny canonical fixtures is ~4 ms. Imports dominate ~96% of the per-spawn
cost. This worker imports those heavy libraries ONCE at startup, then executes
many operators' generated scripts over its lifetime — so the import cost is
paid once, not once per operator.

It is a drop-in replacement for `python <script.py>`: it runs the exact same
rendered script `StandaloneRunner` already produces (imports + prologue + body
+ epilogue). The script's own top-of-file `import pandas` becomes a ~0 ms
`sys.modules` cache hit.

Protocol (line-delimited JSON, both directions):

startup worker -> parent: {"ready": true}
request parent -> worker: {"scriptPath": "<abs>", "workDir": "<abs>"}\n
response worker -> parent: {"exit": 0, "stdout": "...", "stderr": "..."}\n

`exit` is 0 on success or 1 if the script raised; on 1, `stderr` carries the
traceback — mirroring a nonzero subprocess exit so the Scala side's
StandaloneExecutionException path is unchanged. The worker keeps running after
a script error (only a hard interpreter crash ends it); parent closes stdin
(EOF) to shut it down.

Isolation trade-off (accepted, per design discussion): all jobs share one
interpreter, so module-level state (e.g. pandas display options) can leak
between operators. Each job is exec'd in a FRESH namespace and chdir'd to its
own workDir to contain the common cases; this is weaker than the old
process-per-operator isolation.
"""
from __future__ import annotations

import io
import json
import os
import sys
import traceback
from contextlib import redirect_stderr, redirect_stdout

# --- Pay the heavy import cost ONCE, here, at startup. ----------------------
# Pre-importing populates sys.modules, so a script's own `import pandas as pd`
# or `import plotly...` is a cache hit. It does not hand the script the name:
# each one runs in a fresh namespace, so an operator that draws with plotly and
# forgot to declare it still fails here with a NameError, which is the point.
# numpy is left out for the same reason plotly is only a cache warmer: the
# script has to ask for what it uses (see StandaloneRunner.renderScript).
import pandas as pd # noqa: F401
import plotly.express as px # noqa: F401
import plotly.graph_objects as go # noqa: F401
import plotly.io # noqa: F401


def _run_one(script_path: str, work_dir: str) -> "dict[str, object]":
"""Execute one rendered standalone script and capture its output.

Runs in a fresh namespace with cwd = work_dir (generated code may use
relative paths, e.g. CSVScan's `pd.read_csv("sample.csv")`; absolute paths
written by the prologue/epilogue are unaffected). The script's stdout /
stderr are redirected into buffers so they never corrupt the protocol
channel on real stdout.
"""
out_buf, err_buf = io.StringIO(), io.StringIO()
# __name__ = "__main__" so scripts with a `if __name__ == "__main__"` guard
# still run their body (the translator does not emit one, but it is free
# insurance and matches `python script.py` semantics).
namespace = {"__name__": "__main__", "__file__": script_path}
try:
with open(script_path, "r", encoding="utf-8") as f:
source = f.read()
os.chdir(work_dir)
code = compile(source, script_path, "exec")
with redirect_stdout(out_buf), redirect_stderr(err_buf):
exec(code, namespace) # noqa: S102 (running generated verify code by design)
return {"exit": 0, "stdout": out_buf.getvalue(), "stderr": err_buf.getvalue()}
except SystemExit as e:
# The catch-all below would read this as a crash, so it is answered with
# the code the script asked for: an operator that stops early on an input
# it cannot draw succeeded.
code = e.code if isinstance(e.code, int) else (0 if e.code is None else 1)
return {"exit": code, "stdout": out_buf.getvalue(), "stderr": err_buf.getvalue()}
except BaseException: # noqa: BLE001 — a script error must NOT kill the worker
# Match a nonzero subprocess exit: traceback goes to stderr, exit = 1.
err = err_buf.getvalue() + traceback.format_exc()
return {"exit": 1, "stdout": out_buf.getvalue(), "stderr": err}


def main() -> None:
# Signal readiness only after the heavy imports above have completed, so the
# parent can warm a pool and attribute startup cost deterministically.
sys.stdout.write(json.dumps({"ready": True}) + "\n")
sys.stdout.flush()

for line in sys.stdin:
line = line.strip()
if not line:
continue
try:
req = json.loads(line)
result = _run_one(req["scriptPath"], req["workDir"])
except Exception: # malformed request — report, keep serving
result = {"exit": 1, "stdout": "", "stderr": traceback.format_exc()}
sys.stdout.write(json.dumps(result) + "\n")
sys.stdout.flush()


if __name__ == "__main__":
main()
Loading
Loading