From 4562b6306c8661454addb63a93bee8734317774f Mon Sep 17 00:00:00 2001 From: kary zheng Date: Wed, 2 Sep 2026 16:30:44 -0700 Subject: [PATCH 01/21] test(verify): run the script the export produces The other side of the comparison. `StandaloneRunner` writes the script the operator's generator emits, binds its inputs to the files the fixture wrote, runs it, and reads the frames it leaves behind. The script is kept where it ran, so an operator whose two answers differ can be opened as generated rather than described second-hand. `HarnessSpec` covers all three pieces on one operator whose answer is short enough to state in full. Co-Authored-By: Claude Opus 5 (1M context) --- .../resources/python/standalone_worker.py | 122 ++++++ .../amber/translator/verify/HarnessSpec.scala | 115 ++++++ .../translator/verify/StandaloneRunner.scala | 367 ++++++++++++++++++ 3 files changed, 604 insertions(+) create mode 100644 workflow-compiling-service/src/test/resources/python/standalone_worker.py create mode 100644 workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala create mode 100644 workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala diff --git a/workflow-compiling-service/src/test/resources/python/standalone_worker.py b/workflow-compiling-service/src/test/resources/python/standalone_worker.py new file mode 100644 index 00000000000..946fec23206 --- /dev/null +++ b/workflow-compiling-service/src/test/resources/python/standalone_worker.py @@ -0,0 +1,122 @@ +#!/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 `: 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": "", "workDir": ""}\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. ---------------------- +# These mirror the imports StandaloneRunner injects at the top of every +# rendered script. Pre-importing them populates sys.modules, so each executed +# script's own `import pandas as pd` / `import plotly...` is a cache hit. +# numpy is intentionally NOT imported (see StandaloneRunner.renderScript: the +# production translator only provides pandas + plotly, so an operator needing +# numpy must import it itself — we must not mask that). +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 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() diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala new file mode 100644 index 00000000000..685b6f056db --- /dev/null +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -0,0 +1,115 @@ +/* + * 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. + */ + +package org.apache.texera.amber.translator.verify + +import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema, Tuple} +import org.apache.texera.amber.core.workflow.PortIdentity +import org.apache.texera.amber.operator.distinct.DistinctOpDesc +import org.scalatest.Tag +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers + +import java.nio.file.{Files, Path} + +/** The two ways of running one operator, and the file format they meet in. + * + * `Distinct` is the operator under test throughout, because what is being + * tested is the harness rather than the operator: it takes one input, needs no + * configuration, and its answer is short enough to state in full. + */ +class HarnessSpec extends AnyFlatSpec with Matchers { + + /** Only the standalone run needs an interpreter, so only it is held back from + * the job that provisions none. The other two are JVM-side and run there. + */ + private val NeedsPython = + Tag("org.apache.texera.amber.translator.verify.tags.IntegrationTest") + + private val schema = new Schema( + new Attribute("id", AttributeType.INTEGER), + new Attribute("name", AttributeType.STRING) + ) + + private def tuple(id: Int, name: String): Tuple = { + val b = Tuple.builder(schema) + b.add(schema.getAttribute("id"), Int.box(id)) + b.add(schema.getAttribute("name"), name) + b.build() + } + + /** Four rows, the last a repeat of the second. */ + private val rows = Seq(tuple(1, "a"), tuple(2, "b"), tuple(3, "c"), tuple(2, "b")) + + private def withInput(test: (Path, Path) => Unit): Unit = { + val dir = Files.createTempDirectory("harness-spec-") + val input = dir.resolve("input_port_0.jsonl") + TupleIO.writeTuples(input, rows.iterator, schema) + test(dir, input) + } + + "TupleIO" should "read back the rows and the schema it wrote" in { + withInput { (_, input) => + // The schema travels in a sidecar rather than in the JSONL, which carries + // values alone and so cannot say a column is INTEGER rather than a number. + TupleIO.readSchemaSidecar(input) shouldBe schema + val read = TupleIO.readTuples(input, schema).toSeq + read should have length 4 + read.map(_.getField[Integer]("id").intValue) shouldBe Seq(1, 2, 3, 2) + } + } + + "OpExecHarness" should "run an operator and write one file per output port" in { + withInput { (dir, input) => + val out = dir.resolve("actual") + val result = + OpExecHarness.execute(new DistinctOpDesc, Map(PortIdentity(0) -> input), out) + + result.outputs should have size 1 + val produced = result.outputs(PortIdentity(0)) + Files.exists(produced) shouldBe true + + val written = TupleIO.readTuples(produced, result.outputSchemas(PortIdentity(0))).toSeq + written.map(_.getField[Integer]("id").intValue) shouldBe Seq(1, 2, 3) + } + } + + "StandaloneRunner" should "run the generated script and reach the same answer" taggedAs NeedsPython in { + withInput { (dir, input) => + val work = dir.resolve("standalone") + Files.createDirectories(work) + val result = StandaloneRunner.run( + opDesc = new DistinctOpDesc, + inputs = Map(1 -> input), + outputPortCount = 1, + workDir = work + ) + + // The script is kept where it ran, so a failing operator can be opened as + // generated rather than described second-hand. + Files.exists(work.resolve("script.py")) shouldBe true + + val produced = result.outputs(1) + val lines = Files.readAllLines(produced) + lines should have size 3 + lines.get(0) should include("\"id\":1") + lines.get(2) should include("\"id\":3") + } + } +} diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala new file mode 100644 index 00000000000..359bad51e5a --- /dev/null +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -0,0 +1,367 @@ +/* + * 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. + */ + +package org.apache.texera.amber.translator.verify + +import com.typesafe.scalalogging.LazyLogging +import org.apache.texera.amber.core.tuple.AttributeType +import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} +import org.apache.texera.amber.util.python.PythonWorkerPool + +import java.nio.charset.StandardCharsets +import java.nio.file.{Files, Path} +import scala.collection.mutable.ArrayBuffer +import scala.sys.process._ + +/** + * Executes the Python code an OpDesc's [[StandaloneCodeGenerator]] emits and + * captures its DataFrame outputs as JSONL files (compatible with + * [[TupleIO]]'s sidecar-schema format on the comparison side). + * + * Wraps the operator's raw generated code with: + * + * ── prologue ────────────────────────────────────────────── + * in1df = pd.read_json("input_port_0.jsonl", lines=True) + * in2df = pd.read_json("input_port_1.jsonl", lines=True) + * ... + * inAlldf = [in1df, in2df] + * ── operator body (verbatim from generateStandaloneCode) ── + * out1df = in1df[in1df["age"] > 18] + * ── epilogue ───────────────────────────────────────────── + * out1df.to_json("output_port_0.jsonl", orient='records', lines=True) + * ... + * + * Port indexing matches the placeholder convention used by the translator: + * `inNdf`/`outNdf` is 1-based and corresponds to the operator's N-th external + * input/output port in declaration order. The harness key (a 1-based Int) is + * what the placeholder uses; the caller is responsible for ordering inputs + * the same way the operator's `generateStandaloneCode()` expects. + * + * The subprocess inherits the caller's environment so the Python interpreter + * picks up whatever pandas/plotly the test fixture installed. + */ +object StandaloneRunner extends LazyLogging { + + /** The value both paths seed numpy's global RNG with. Any fixed number does; + * what matters is that the two agree, so it is declared once here and + * referenced by name from py_op_driver's comment. + */ + private[verify] val VerifySeed: Int = 20260811 + + /** + * @param outputs paths to the per-port output JSONL files. Empty map iff + * the operator's `producesDataFrame()` returned false + * (visualizations, etc.) — caller handles those separately. + * @param stdout raw subprocess stdout (useful for failure diagnostics) + * @param stderr raw subprocess stderr + */ + final case class Result(outputs: Map[Int, Path], stdout: String, stderr: String) + + /** + * Generate, write, and execute the standalone Python script for `opDesc`. + * + * @param opDesc must mix in [[StandaloneCodeGenerator]]; otherwise we throw + * since there's nothing to test. + * @param inputs map from 1-based port index → JSONL fixture path. The + * script reads each into `inNdf`. + * @param outputPortCount how many `outNdf` variables the operator declares. + * Caller derives this from the OpDesc's output ports. + * @param workDir directory used for the generated `script.py` and output + * JSONL files. Created if missing. + * @param pythonExe path to the Python 3.12 interpreter. Defaults to + * the env var `UDF_PYTHON_PATH`, then `python3.12`, then + * `python3`. The same fallback chain used by the rest of + * the Texera test suite for Python-backed operators. + */ + def run( + opDesc: LogicalOp, + inputs: Map[Int, Path], + outputPortCount: Int, + workDir: Path, + pythonExe: String = resolvePython() + ): Result = { + val gen = opDesc match { + case g: StandaloneCodeGenerator => g + case other => + throw new IllegalArgumentException( + s"OpDesc ${other.getClass.getSimpleName} does not implement " + + s"StandaloneCodeGenerator; nothing to verify" + ) + } + + Files.createDirectories(workDir) + val scriptPath = workDir.resolve("script.py") + val outputPaths: Map[Int, Path] = + if (gen.producesDataFrame()) + (1 to outputPortCount).map(i => i -> workDir.resolve(s"output_port_${i - 1}.jsonl")).toMap + else Map.empty + + val source = + renderScript(gen.generateStandaloneCode(), inputs, outputPaths, gen.standaloneHelpers()) + Files.write(scriptPath, source.getBytes(StandardCharsets.UTF_8)) + + val (exit, stdout, stderr) = execute(scriptPath, workDir, pythonExe) + if (exit != 0) { + throw new StandaloneExecutionException(exit, scriptPath, source, stdout, stderr) + } + Result(outputPaths, stdout, stderr) + } + + private val WorkerResourcePath = "/python/standalone_worker.py" + + // Run the rendered script and return (exitCode, stdout, stderr). Prefers a + // pooled persistent worker (imports pandas/plotly once, ~18x faster per op — + // see PythonWorkerPool); a rare hard worker crash falls back to a one-shot + // subprocess so behavior is never worse than the original path. Both paths + // run with cwd = workDir and read results from files, so they are + // interchangeable — the executed script is byte-identical. + private def execute(scriptPath: Path, workDir: Path, pythonExe: String): (Int, String, String) = { + if (PythonWorkerPool.enabled) { + try { + val req = org.apache.texera.amber.util.JSONUtils.objectMapper.createObjectNode() + req.put("scriptPath", scriptPath.toString) + req.put("workDir", workDir.toString) + val o = PythonWorkerPool.run(WorkerResourcePath, Seq.empty, pythonExe, req) + return (o.exit, o.stdout, o.stderr) + } catch { + case e: PythonWorkerPool.WorkerDiedException => + logger.warn( + s"Standalone worker unavailable; falling back to one-shot subprocess " + + s"for $scriptPath: ${e.getMessage}" + ) + } + } + runSubprocess(scriptPath, workDir, pythonExe) + } + + // Original one-process-per-operator path. Retained as the fallback and as the + // behavior selected by TEXERA_TEST_PYTHON_WORKER=0. + private def runSubprocess( + scriptPath: Path, + workDir: Path, + pythonExe: String + ): (Int, String, String) = { + // Capture stdout/stderr separately. ProcessLogger's append is called from + // the subprocess's I/O thread, so we collect into ArrayBuffer (thread-safe + // append is fine for this serial use) and join at the end. + val outBuf = ArrayBuffer.empty[String] + val errBuf = ArrayBuffer.empty[String] + val logger = ProcessLogger(line => outBuf += line, line => errBuf += line) + // cwd = workDir so generated code using *relative* paths (e.g. CSVScan's + // basename-stripped `pd.read_csv("sample.csv")`) resolves against workDir. + // Absolute paths written by the prologue/epilogue are unaffected. + val exit = Process(Seq(pythonExe, scriptPath.toString), Some(workDir.toFile)).!(logger) + (exit, outBuf.mkString("\n"), errBuf.mkString("\n")) + } + + // Builds the full Python source: imports + prologue + verbatim operator body + // + epilogue. We intentionally do NOT substitute the inNdf/outNdf placeholders + // — the body keeps them so the var-bindings the prologue/epilogue introduce + // (also named inNdf/outNdf) reference the same names. + private def renderScript( + body: String, + inputs: Map[Int, Path], + outputs: Map[Int, Path], + helpers: Seq[String] + ): String = { + val sb = new StringBuilder + + sb.append("# Auto-generated by StandaloneRunner. Do not commit.\n") + sb.append("import json\n") + sb.append("import sys\n") + sb.append("import base64\n") + sb.append("import pickle\n") + // NOTE: numpy is intentionally NOT injected here. The production translator + // (WorkflowToPythonTranslator) only provides pandas + plotly to standalone + // scripts, so any operator whose standalone code needs numpy must import it + // itself. Injecting numpy here would mask that class of bug in verify tests. + sb.append("import pandas as pd\n") + sb.append("import plotly.express as px\n") + sb.append("import plotly.graph_objects as go\n") + sb.append("import plotly.io\n") + // Same seed as py_op_driver's run_config, for the reason given there. Bound + // under a private name and deleted so the note above still holds: a script + // that wants numpy has to import it, and this does not hand it one. + sb.append(s"import numpy as _texera_np; _texera_np.random.seed($VerifySeed); del _texera_np\n") + sb.append("\n") + + // Object columns holding non-primitive values (e.g. a trained sklearn model + // in a BINARY output column) can't go through to_json. Pickle+base64 them so + // the JSONL matches py_op_driver's BINARY write path exactly. Primitives + // (str/int/float/bool/None) pass through unchanged, so ordinary DataFrame + // outputs are unaffected. + sb.append("def _texera_encode_obj_cols(df):\n") + sb.append(" for _c in df.columns:\n") + sb.append(" if df[_c].dtype == object:\n") + sb.append( + " df[_c] = df[_c].map(lambda _v: base64.b64encode(pickle.dumps(_v)).decode('ascii') " + + "if not isinstance(_v, (str, int, float, bool, type(None))) else _v)\n" + ) + sb.append(" return df\n") + sb.append("\n") + + // TIMESTAMP columns are handed to the operator as datetime64 (see the + // prologue below) to match the schema-typed runtime path, but the runtime + // path serializes a TIMESTAMP back out with java.sql.Timestamp.toString — + // "yyyy-mm-dd hh:mm:ss.f", trailing zeros trimmed to at least one digit — + // whereas pandas' to_json would emit epoch millis. Convert datetime columns + // back to that exact form before writing so both paths' JSONL agree. + sb.append("def _texera_ts_str(_v):\n") + sb.append(" if pd.isna(_v):\n") + sb.append(" return None\n") + sb.append(" _s = _v.strftime('%Y-%m-%d %H:%M:%S.%f').rstrip('0')\n") + sb.append(" return _s + '0' if _s.endswith('.') else _s\n") + sb.append("\n") + sb.append("def _texera_encode_ts_cols(df):\n") + sb.append(" for _c in df.columns:\n") + sb.append(" if pd.api.types.is_datetime64_any_dtype(df[_c]):\n") + sb.append(" df[_c] = df[_c].map(_texera_ts_str)\n") + sb.append(" return df\n") + sb.append("\n") + + // Prologue: load each external input into in{N}df. Note: pd.read_json with + // lines=True correctly handles empty files (returns empty DataFrame). + // convert_dates=False: pd.read_json otherwise auto-coerces ISO-ish strings + // and columns named like dates ("date", "*_at", …) to datetime64, which the + // schema-typed runtime path (STRING) does not do — that divergence would + // make a plain date string column serialize as "...T00:00:00" on only one + // side. Operators that genuinely need datetimes convert explicitly, so both + // paths stay in sync. + // precise_float=True: pd.read_json's default (ujson) fast double parser is + // lossy in the last few ULPs, so a DOUBLE column would load slightly + // different values than the schema-typed runtime path (which parses doubles + // exactly). Operators that stringify raw cell values (e.g. Radar hover text) + // then diverge; precise_float=True keeps both paths bit-identical. + // The blanket convert_dates=False also leaves genuine TIMESTAMP columns as + // strings, which the runtime path delivers as datetime64 — a divergence for + // any operator that renders or computes on them. The fixture's schema + // sidecar says which columns those are, so cast exactly those back. + inputs.toSeq.sortBy(_._1).foreach { + case (n, path) => + sb.append( + s"in${n}df = pd.read_json(${py(path.toString)}, lines=True, convert_dates=False, precise_float=True)\n" + ) + timestampColumns(path).foreach { col => + sb.append(s"if ${py(col)} in in${n}df.columns:\n") + sb.append(s" in${n}df[${py(col)}] = pd.to_datetime(in${n}df[${py(col)}])\n") + } + doubleColumns(path).foreach { col => + sb.append(s"if ${py(col)} in in${n}df.columns:\n") + sb.append(s" in${n}df[${py(col)}] = in${n}df[${py(col)}].astype('float64')\n") + } + } + // The variadic placeholder, bound here for the same reason the numbered ones + // are: this script leaves the body's placeholders alone and defines names to + // match them, so an operator reading a variadic port finds its list here the + // way the translator would have written one out. + if (inputs.nonEmpty) { + sb.append( + inputs.keys.toSeq.sorted.map(n => s"in${n}df").mkString("inAlldf = [", ", ", "]\n") + ) + } + sb.append("\n") + + // Body verbatim — placeholders left in place. + // Emitted ahead of the body the way the translator does, so an operator that + // declares a helper is exercised here exactly as it runs in a real script. + helpers.foreach { helper => + sb.append(helper) + if (!helper.endsWith("\n")) sb.append('\n') + sb.append('\n') + } + + sb.append("# ── operator body ──\n") + sb.append(body) + if (!body.endsWith("\n")) sb.append('\n') + sb.append("\n") + + // Epilogue: dump each out{N}df to JSONL. When producesDataFrame() is false + // (visualization ops), `outputs` is empty and this block is a no-op — the + // caller is expected to verify viz outputs by other means. + outputs.toSeq.sortBy(_._1).foreach { + case (n, path) => + sb.append( + s"_texera_encode_obj_cols(_texera_encode_ts_cols(out${n}df))" + + s".to_json(${py(path.toString)}, orient='records', lines=True)\n" + ) + } + + sb.toString + } + + // TIMESTAMP-typed column names from a fixture's `.jsonl.schema.json` sidecar. + // A missing or unreadable sidecar means no casts — the prologue then behaves + // exactly as before. + private def timestampColumns(input: Path): Seq[String] = + columnsOfType(input, AttributeType.TIMESTAMP) + + // DOUBLE-typed column names. pd.read_json narrows a float column whose values + // are all integral to int64, while the runtime path keeps the schema's DOUBLE, + // so a column like 7.0 stringifies as "7" on one side and "7.0" on the other — + // invisible to numeric comparison, visible the moment an operator uses the + // column as a label (a trace name, a legend entry, hover text). + private def doubleColumns(input: Path): Seq[String] = + columnsOfType(input, AttributeType.DOUBLE) + + private def columnsOfType(input: Path, attributeType: AttributeType): Seq[String] = + scala.util + .Try(TupleIO.readSchemaSidecar(input)) + .toOption + .toSeq + .flatMap( + _.getAttributes.filter(_.getType == attributeType).map(_.getName) + ) + + // Python string literal, single-quoted with backslashes escaped. We + // deliberately don't use repr() in Scala (no such thing) — JSON.toString + // would also work but introduces double-quote escaping when the path has + // spaces. + private def py(s: String): String = + "'" + s.replace("\\", "\\\\").replace("'", "\\'") + "'" + + // Resolution chain mirrors the rest of the Texera test infra: env var first + // (set by CI / the shared-venv setup), then conventional names. + private def resolvePython(): String = { + val fromEnv = sys.env.get("UDF_PYTHON_PATH").filter(_.nonEmpty) + fromEnv.getOrElse { + // We don't try to probe `which` here — if neither env var nor a literal + // `python3.12` is on PATH, the subprocess invocation will fail and the + // error path below surfaces it. + "python3.12" + } + } +} + +final class StandaloneExecutionException( + val exitCode: Int, + val scriptPath: Path, + val source: String, + val stdout: String, + val stderr: String +) extends RuntimeException( + // The script path goes first in the message so a failing CI log makes it + // immediately obvious which file to open. stderr ends the message because + // the Python traceback (if any) is the most actionable signal. + s"""Standalone Python script exited with code $exitCode. + |Script: $scriptPath + |--- stdout --- + |$stdout + |--- stderr --- + |$stderr""".stripMargin + ) From 60fbc77992d04701842853a6f970347851708ce9 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 3 Sep 2026 15:45:59 -0700 Subject: [PATCH 02/21] test(verify): say it once, in the shape the code does not already give Three kinds of comment came out. A drawing of the string the code below assembles. A restatement of a branch the reader can see. And the word MVP, which dated the scope to a moment rather than stating it. What replaces them says the same thing shorter, or says what the code cannot: which cases the harness does not drive and why none of them has an operator asking for it. Co-Authored-By: Claude Opus 5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 15 +++------------ 1 file changed, 3 insertions(+), 12 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 359bad51e5a..24a7c8fe405 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -34,18 +34,9 @@ import scala.sys.process._ * captures its DataFrame outputs as JSONL files (compatible with * [[TupleIO]]'s sidecar-schema format on the comparison side). * - * Wraps the operator's raw generated code with: - * - * ── prologue ────────────────────────────────────────────── - * in1df = pd.read_json("input_port_0.jsonl", lines=True) - * in2df = pd.read_json("input_port_1.jsonl", lines=True) - * ... - * inAlldf = [in1df, in2df] - * ── operator body (verbatim from generateStandaloneCode) ── - * out1df = in1df[in1df["age"] > 18] - * ── epilogue ───────────────────────────────────────────── - * out1df.to_json("output_port_0.jsonl", orient='records', lines=True) - * ... + * The operator's code is wrapped in a prologue that reads each input file into + * an `inNdf` and an epilogue that writes each `outNdf` back out, with the + * generated body verbatim between them. * * Port indexing matches the placeholder convention used by the translator: * `inNdf`/`outNdf` is 1-based and corresponds to the operator's N-th external From 4ab2c9ff5df685a61923edf55dccb3778ef1924a Mon Sep 17 00:00:00 2001 From: kary zheng Date: Fri, 4 Sep 2026 12:39:28 -0700 Subject: [PATCH 03/21] test(verify): build the script's header the way the export builds it The harness wrote plotly's three modules into every rendered script, so an operator that draws with plotly and never says so still ran here and would have failed on export. The header now carries pandas plus what the operator declares, which is what the translator emits. The pooled worker keeps pre-importing plotly, since that is a cache warmer rather than a name the script can reach: each script runs in a fresh namespace, so the missing declaration still surfaces as a NameError. Co-Authored-By: Claude Opus 5 (1M context) --- .../resources/python/standalone_worker.py | 12 +++++----- .../translator/verify/StandaloneRunner.scala | 24 ++++++++++++------- 2 files changed, 21 insertions(+), 15 deletions(-) diff --git a/workflow-compiling-service/src/test/resources/python/standalone_worker.py b/workflow-compiling-service/src/test/resources/python/standalone_worker.py index 946fec23206..ff4d2aa23f5 100644 --- a/workflow-compiling-service/src/test/resources/python/standalone_worker.py +++ b/workflow-compiling-service/src/test/resources/python/standalone_worker.py @@ -59,12 +59,12 @@ from contextlib import redirect_stderr, redirect_stdout # --- Pay the heavy import cost ONCE, here, at startup. ---------------------- -# These mirror the imports StandaloneRunner injects at the top of every -# rendered script. Pre-importing them populates sys.modules, so each executed -# script's own `import pandas as pd` / `import plotly...` is a cache hit. -# numpy is intentionally NOT imported (see StandaloneRunner.renderScript: the -# production translator only provides pandas + plotly, so an operator needing -# numpy must import it itself — we must not mask that). +# 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 diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 24a7c8fe405..e39990c0768 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -104,7 +104,13 @@ object StandaloneRunner extends LazyLogging { else Map.empty val source = - renderScript(gen.generateStandaloneCode(), inputs, outputPaths, gen.standaloneHelpers()) + renderScript( + gen.generateStandaloneCode(), + inputs, + outputPaths, + gen.standaloneHelpers(), + gen.standaloneImports() + ) Files.write(scriptPath, source.getBytes(StandardCharsets.UTF_8)) val (exit, stdout, stderr) = execute(scriptPath, workDir, pythonExe) @@ -169,7 +175,8 @@ object StandaloneRunner extends LazyLogging { body: String, inputs: Map[Int, Path], outputs: Map[Int, Path], - helpers: Seq[String] + helpers: Seq[String], + imports: Seq[String] ): String = { val sb = new StringBuilder @@ -178,14 +185,13 @@ object StandaloneRunner extends LazyLogging { sb.append("import sys\n") sb.append("import base64\n") sb.append("import pickle\n") - // NOTE: numpy is intentionally NOT injected here. The production translator - // (WorkflowToPythonTranslator) only provides pandas + plotly to standalone - // scripts, so any operator whose standalone code needs numpy must import it - // itself. Injecting numpy here would mask that class of bug in verify tests. + // NOTE: nothing beyond pandas is injected here. The production translator + // (WorkflowToPythonTranslator) emits pandas for every script and then only + // what the operators in the plan ask for, so an operator whose standalone + // code needs numpy, or plotly, must say so. Injecting either here would mask + // that class of bug: the script would run in verify and fail on export. sb.append("import pandas as pd\n") - sb.append("import plotly.express as px\n") - sb.append("import plotly.graph_objects as go\n") - sb.append("import plotly.io\n") + imports.foreach(line => sb.append(line).append("\n")) // Same seed as py_op_driver's run_config, for the reason given there. Bound // under a private name and deleted so the note above still holds: a script // that wants numpy has to import it, and this does not hand it one. From dd36970ad74fc7351d909c41fb4b2d6a7c3ca34a Mon Sep 17 00:00:00 2001 From: kary zheng Date: Tue, 8 Sep 2026 23:49:58 -0700 Subject: [PATCH 04/21] test(verify): read a script's own exit as the exit it asked for An operator that stops early on an input it cannot draw ends itself with SystemExit. Under `python script.py` nothing catches that and the interpreter exits with the code; here the catch-all below read it as a crash, so a run that succeeded was reported as a failure. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/test/resources/python/standalone_worker.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/workflow-compiling-service/src/test/resources/python/standalone_worker.py b/workflow-compiling-service/src/test/resources/python/standalone_worker.py index ff4d2aa23f5..b49d1733a40 100644 --- a/workflow-compiling-service/src/test/resources/python/standalone_worker.py +++ b/workflow-compiling-service/src/test/resources/python/standalone_worker.py @@ -93,6 +93,12 @@ def _run_one(script_path: str, work_dir: str) -> "dict[str, object]": 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() From 3e87f6bcf77fee2768346a8506743e4fa21f7740 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Wed, 9 Sep 2026 16:17:38 -0700 Subject: [PATCH 05/21] test(verify): keep the input schema's string columns as strings pd.read_json infers a column of numeric-looking strings as a number, so a STRING column holding 001 loaded as 1 and one holding only nulls loaded as NaN. The engine path calls str() on the cell and leaves a null alone, so the two disagreed about the operator's input rather than its output. The fixture's schema sidecar already says which columns those are. Co-Authored-By: Claude Opus 5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 29 +++++++++++++++++-- 1 file changed, 26 insertions(+), 3 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index e39990c0768..0f3ffef541a 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -249,10 +249,18 @@ object StandaloneRunner extends LazyLogging { // strings, which the runtime path delivers as datetime64 — a divergence for // any operator that renders or computes on them. The fixture's schema // sidecar says which columns those are, so cast exactly those back. + // read_json also infers a column of numeric-looking strings as a number, so + // a STRING column holding "001" arrives as 1, and one holding only nulls as + // NaN rather than None. The sidecar's STRING columns are pinned to object. inputs.toSeq.sortBy(_._1).foreach { case (n, path) => + val dtype = stringColumns(path) match { + case Seq() => "" + case cols => cols.map(c => s"${py(c)}: 'object'").mkString(", dtype={", ", ", "}") + } sb.append( - s"in${n}df = pd.read_json(${py(path.toString)}, lines=True, convert_dates=False, precise_float=True)\n" + s"in${n}df = pd.read_json(${py(path.toString)}, lines=True, " + + s"convert_dates=False, precise_float=True$dtype)\n" ) timestampColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") @@ -272,6 +280,11 @@ object StandaloneRunner extends LazyLogging { inputs.keys.toSeq.sorted.map(n => s"in${n}df").mkString("inAlldf = [", ", ", "]\n") ) } + // The file placeholders, bound the same way. A script this runner builds + // holds one operator, so the plain names are unambiguous and the comparison + // knows where to look. + sb.append("outputHtml = \"output.html\"\n") + sb.append("outputJson = \"output.json\"\n") sb.append("\n") // Body verbatim — placeholders left in place. @@ -316,6 +329,10 @@ object StandaloneRunner extends LazyLogging { private def doubleColumns(input: Path): Seq[String] = columnsOfType(input, AttributeType.DOUBLE) + // STRING-typed column names, for the read_json dtype map above. + private def stringColumns(input: Path): Seq[String] = + columnsOfType(input, AttributeType.STRING) + private def columnsOfType(input: Path, attributeType: AttributeType): Seq[String] = scala.util .Try(TupleIO.readSchemaSidecar(input)) @@ -328,9 +345,15 @@ object StandaloneRunner extends LazyLogging { // Python string literal, single-quoted with backslashes escaped. We // deliberately don't use repr() in Scala (no such thing) — JSON.toString // would also work but introduces double-quote escaping when the path has - // spaces. + // spaces. Control characters are escaped too: a fixture column name can hold + // a literal newline, which raw would end the literal and break the script. private def py(s: String): String = - "'" + s.replace("\\", "\\\\").replace("'", "\\'") + "'" + s.map { + case '\\' => "\\\\" + case '\'' => "\\'" + case c if c.toInt < 0x20 => f"\\x${c.toInt}%02x" + case c => c.toString + }.mkString("'", "", "'") // Resolution chain mirrors the rest of the Texera test infra: env var first // (set by CI / the shared-venv setup), then conventional names. From 6e2e71d7750e2f16cd05d4ee3275168bb630b2dd Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 10 Sep 2026 13:08:58 -0700 Subject: [PATCH 06/21] test(verify): hand the script the schema, and a JVM path's integers exactly The sidecar already says what each input column was declared as, and the generator now has an overload that takes it, so the script is told what a JSONL file cannot carry. `read_json` parses a column holding a null through float64, so a LONG of 9007199254740993 reached the operator as 9007199254740992 while Path A still held the tuple. The engine's own Python side refuses a float outside that window rather than accept a corrupted rendition; the harness was quietly accepting one. Only where Path A is the JVM harness: a Python operator's table is built by pandas from the same tuples, so the float is what BOTH sides see there. Writing needs two more unwrappings for the dtype that follows: an operator rebuilding rows out of an integer column hands back numpy scalars, and a nullable column carries pd.NA. Neither is a Python scalar, so both were being written into the JSONL as base64 pickles. Co-Authored-By: Claude Opus 5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 78 +++++++++++++++++-- 1 file changed, 73 insertions(+), 5 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 0f3ffef541a..b27a329d0a0 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -20,7 +20,8 @@ package org.apache.texera.amber.translator.verify import com.typesafe.scalalogging.LazyLogging -import org.apache.texera.amber.core.tuple.AttributeType +import org.apache.texera.amber.core.tuple.{AttributeType, Schema} +import org.apache.texera.amber.core.workflow.PortIdentity import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} import org.apache.texera.amber.util.python.PythonWorkerPool @@ -85,7 +86,8 @@ object StandaloneRunner extends LazyLogging { inputs: Map[Int, Path], outputPortCount: Int, workDir: Path, - pythonExe: String = resolvePython() + pythonExe: String = resolvePython(), + exactIntegers: Boolean = false ): Result = { val gen = opDesc match { case g: StandaloneCodeGenerator => g @@ -105,11 +107,14 @@ object StandaloneRunner extends LazyLogging { val source = renderScript( - gen.generateStandaloneCode(), + // The sidecar says what each column was DECLARED as, which the JSONL + // cannot carry. + gen.generateStandaloneCode(schemasOf(inputs)), inputs, outputPaths, gen.standaloneHelpers(), - gen.standaloneImports() + gen.standaloneImports(), + exactIntegers ) Files.write(scriptPath, source.getBytes(StandardCharsets.UTF_8)) @@ -176,7 +181,8 @@ object StandaloneRunner extends LazyLogging { inputs: Map[Int, Path], outputs: Map[Int, Path], helpers: Seq[String], - imports: Seq[String] + imports: Seq[String], + exactIntegers: Boolean ): String = { val sb = new StringBuilder @@ -204,8 +210,19 @@ object StandaloneRunner extends LazyLogging { // (str/int/float/bool/None) pass through unchanged, so ordinary DataFrame // outputs are unaffected. sb.append("def _texera_encode_obj_cols(df):\n") + // A numpy scalar is unwrapped rather than pickled, and pd.NA is written as + // null. An operator that rebuilds rows out of an integer column hands back + // neither a Python scalar nor None, so the branch below would put a base64 + // pickle into a column the schema declares a number. + sb.append(" def _texera_plain(_v):\n") + sb.append(" if _v is pd.NA:\n") + sb.append(" return None\n") + sb.append(" if hasattr(_v, 'item') and getattr(_v, 'ndim', None) == 0:\n") + sb.append(" return _v.item()\n") + sb.append(" return _v\n") sb.append(" for _c in df.columns:\n") sb.append(" if df[_c].dtype == object:\n") + sb.append(" df[_c] = df[_c].map(_texera_plain)\n") sb.append( " df[_c] = df[_c].map(lambda _v: base64.b64encode(pickle.dumps(_v)).decode('ascii') " + "if not isinstance(_v, (str, int, float, bool, type(None))) else _v)\n" @@ -213,6 +230,27 @@ object StandaloneRunner extends LazyLogging { sb.append(" return df\n") sb.append("\n") + // Both paths have to be handed the same numbers. `read_json` parses a column + // holding a null through float64, so a LONG of 9007199254740993 arrives as + // 9007199254740992 while the engine still has the tuple. Python's json reads + // it exactly. Only for a JVM Path A: a Python operator's own table goes + // through pandas too, so there the float is what BOTH sides see. + if (exactIntegers) { + sb.append("def _texera_exact_ints(_path, _columns):\n") + sb.append(" _values = {_c: [] for _c in _columns}\n") + sb.append(" with open(_path, 'r', encoding='utf-8') as _f:\n") + sb.append(" for _line in _f:\n") + sb.append(" _line = _line.strip()\n") + sb.append(" if not _line:\n") + sb.append(" continue\n") + sb.append(" _row = json.loads(_line)\n") + sb.append(" for _c in _columns:\n") + sb.append(" _v = _row.get(_c)\n") + sb.append(" _values[_c].append(pd.NA if _v is None else int(_v))\n") + sb.append(" return {_c: pd.array(_v, dtype='Int64') for _c, _v in _values.items()}\n") + sb.append("\n") + } + // TIMESTAMP columns are handed to the operator as datetime64 (see the // prologue below) to match the schema-typed runtime path, but the runtime // path serializes a TIMESTAMP back out with java.sql.Timestamp.toString — @@ -270,6 +308,18 @@ object StandaloneRunner extends LazyLogging { sb.append(s"if ${py(col)} in in${n}df.columns:\n") sb.append(s" in${n}df[${py(col)}] = in${n}df[${py(col)}].astype('float64')\n") } + // Only where the reader lost the value: a column with no holes already + // came back exact, and replacing it would hand the operator a nullable + // dtype the run never had. + if (exactIntegers) integerColumns(path) match { + case Seq() => () + case cols => + val names = cols.map(py).mkString(", ") + sb.append(s"for _c, _v in _texera_exact_ints(${py(path.toString)}, [$names]).items():\n") + sb.append(s" if _c in in${n}df.columns and not pd.api.types.is_integer_dtype(") + sb.append(s"in${n}df[_c]):\n") + sb.append(s" in${n}df[_c] = _v\n") + } } // The variadic placeholder, bound here for the same reason the numbered ones // are: this script leaves the body's placeholders alone and defines names to @@ -330,6 +380,24 @@ object StandaloneRunner extends LazyLogging { columnsOfType(input, AttributeType.DOUBLE) // STRING-typed column names, for the read_json dtype map above. + /** The declared schema behind each input port, read from the sidecars. A port the + * sidecar does not cover is left out rather than guessed at. + */ + private def schemasOf(inputs: Map[Int, Path]): Map[PortIdentity, Schema] = + inputs.flatMap { + case (port, path) => + scala.util + .Try(TupleIO.readSchemaSidecar(path)) + .toOption + .map(schema => PortIdentity(port - 1) -> schema) + } + + /** The columns the sidecar declares integral, both widths: the loss is the same + * for either. + */ + private def integerColumns(input: Path): Seq[String] = + columnsOfType(input, AttributeType.INTEGER) ++ columnsOfType(input, AttributeType.LONG) + private def stringColumns(input: Path): Seq[String] = columnsOfType(input, AttributeType.STRING) From d7385cee42475c14a7c4a13bcd07ead46797fff9 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 10 Sep 2026 13:54:35 -0700 Subject: [PATCH 07/21] test(verify): give an empty input its declared columns `read_json` has no rows to read column names off a file with none in it, so it produced a frame of no columns while the engine hands the operator the port's declared ones. A Projection naming a column then raised KeyError on a table the workflow handled. The frame is rebuilt from the sidecar, dtypes included, so an empty table is the same table on both paths. Co-Authored-By: Claude Opus 5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 34 +++++++++++++++++++ 1 file changed, 34 insertions(+) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index b27a329d0a0..406517822d0 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -300,6 +300,18 @@ object StandaloneRunner extends LazyLogging { s"in${n}df = pd.read_json(${py(path.toString)}, lines=True, " + s"convert_dates=False, precise_float=True$dtype)\n" ) + // read_json has no rows to read column names off a file with none in it, + // so it produces a frame of no columns, while the engine hands the + // operator the port's declared ones (see Table.empty_of). Rebuild it + // from the sidecar, dtypes included, so an empty table is the same table + // on both paths. + emptyFrameColumns(path) match { + case Seq() => () + case cols => + val fields = cols.map { case (c, d) => s"${py(c)}: pd.Series(dtype=${py(d)})" } + sb.append(s"if in${n}df.empty:\n") + sb.append(s" in${n}df = pd.DataFrame({${fields.mkString(", ")}})\n") + } timestampColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") sb.append(s" in${n}df[${py(col)}] = pd.to_datetime(in${n}df[${py(col)}])\n") @@ -392,6 +404,28 @@ object StandaloneRunner extends LazyLogging { .map(schema => PortIdentity(port - 1) -> schema) } + /** Each declared column with the pandas dtype an Arrow round-trip gives it, which + * is what the engine's own empty table carries. Only the types a fixture can hold + * are named; anything else falls back to object, the dtype an inferred column of + * unknown content would have had anyway. + */ + private def emptyFrameColumns(input: Path): Seq[(String, String)] = + scala.util + .Try(TupleIO.readSchemaSidecar(input)) + .toOption + .toSeq + .flatMap(_.getAttributes.map { attr => + val dtype = attr.getType match { + case AttributeType.INTEGER => "int32" + case AttributeType.LONG => "int64" + case AttributeType.DOUBLE => "float64" + case AttributeType.BOOLEAN => "bool" + case AttributeType.TIMESTAMP => "datetime64[us]" + case _ => "object" + } + attr.getName -> dtype + }) + /** The columns the sidecar declares integral, both widths: the loss is the same * for either. */ From 06ffa1893d946e9665b3c2f953d70f85e9ceac63 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 10 Sep 2026 16:39:39 -0700 Subject: [PATCH 08/21] test(verify): say the reader's four options once each Twenty lines in one paragraph read as a wall. Same reasons, one short paragraph apiece. Co-Authored-By: Claude Opus 5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 33 ++++++++----------- 1 file changed, 13 insertions(+), 20 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 406517822d0..c378a7619a1 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -270,26 +270,19 @@ object StandaloneRunner extends LazyLogging { sb.append(" return df\n") sb.append("\n") - // Prologue: load each external input into in{N}df. Note: pd.read_json with - // lines=True correctly handles empty files (returns empty DataFrame). - // convert_dates=False: pd.read_json otherwise auto-coerces ISO-ish strings - // and columns named like dates ("date", "*_at", …) to datetime64, which the - // schema-typed runtime path (STRING) does not do — that divergence would - // make a plain date string column serialize as "...T00:00:00" on only one - // side. Operators that genuinely need datetimes convert explicitly, so both - // paths stay in sync. - // precise_float=True: pd.read_json's default (ujson) fast double parser is - // lossy in the last few ULPs, so a DOUBLE column would load slightly - // different values than the schema-typed runtime path (which parses doubles - // exactly). Operators that stringify raw cell values (e.g. Radar hover text) - // then diverge; precise_float=True keeps both paths bit-identical. - // The blanket convert_dates=False also leaves genuine TIMESTAMP columns as - // strings, which the runtime path delivers as datetime64 — a divergence for - // any operator that renders or computes on them. The fixture's schema - // sidecar says which columns those are, so cast exactly those back. - // read_json also infers a column of numeric-looking strings as a number, so - // a STRING column holding "001" arrives as 1, and one holding only nulls as - // NaN rather than None. The sidecar's STRING columns are pinned to object. + // Prologue: load each external input into in{N}df. Every option below undoes + // an inference that would otherwise hand the two paths different data. + // + // convert_dates=False: read_json reads ISO-ish strings, and columns merely + // NAMED like dates, as datetime64, so a plain date string would serialize + // with a "T00:00:00" on one side only. It leaves real TIMESTAMP columns as + // strings too, which the sidecar names and the casts below restore. + // + // precise_float=True: the default ujson parser is lossy in the last few + // ULPs, and an operator that renders a DOUBLE as text prints the difference. + // + // dtype=object on the sidecar's STRING columns: "001" would arrive as 1, and + // a column holding only nulls as NaN rather than None. inputs.toSeq.sortBy(_._1).foreach { case (n, path) => val dtype = stringColumns(path) match { From f0616ccc2ecaad1a228e1a405d6452d98050cb68 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 10 Sep 2026 16:40:20 -0700 Subject: [PATCH 09/21] style: wrap a line scalafmt rejects Co-Authored-By: Claude Opus 5 (1M context) --- .../texera/amber/translator/verify/StandaloneRunner.scala | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index c378a7619a1..39a87b1797b 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -320,7 +320,9 @@ object StandaloneRunner extends LazyLogging { case Seq() => () case cols => val names = cols.map(py).mkString(", ") - sb.append(s"for _c, _v in _texera_exact_ints(${py(path.toString)}, [$names]).items():\n") + sb.append( + s"for _c, _v in _texera_exact_ints(${py(path.toString)}, [$names]).items():\n" + ) sb.append(s" if _c in in${n}df.columns and not pd.api.types.is_integer_dtype(") sb.append(s"in${n}df[_c]):\n") sb.append(s" in${n}df[_c] = _v\n") From f99d1ef8056c838e3aefc7cee5831fa7bcdcf167 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 10 Sep 2026 16:46:27 -0700 Subject: [PATCH 10/21] test(verify): cut the doc to what the code cannot say Co-Authored-By: Claude Opus 5 (1M context) --- .../amber/translator/verify/StandaloneRunner.scala | 13 +++---------- 1 file changed, 3 insertions(+), 10 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 39a87b1797b..c60ea71e0e9 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -37,16 +37,9 @@ import scala.sys.process._ * * The operator's code is wrapped in a prologue that reads each input file into * an `inNdf` and an epilogue that writes each `outNdf` back out, with the - * generated body verbatim between them. - * - * Port indexing matches the placeholder convention used by the translator: - * `inNdf`/`outNdf` is 1-based and corresponds to the operator's N-th external - * input/output port in declaration order. The harness key (a 1-based Int) is - * what the placeholder uses; the caller is responsible for ordering inputs - * the same way the operator's `generateStandaloneCode()` expects. - * - * The subprocess inherits the caller's environment so the Python interpreter - * picks up whatever pandas/plotly the test fixture installed. + * generated body verbatim between them. `N` is 1-based and counts the + * operator's external ports in declaration order, the translator's own + * convention, so the caller has to hand inputs over in that order. */ object StandaloneRunner extends LazyLogging { From 04f55557af6be70c8b1f80cf54b2904e48bbcd55 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Fri, 18 Sep 2026 16:45:40 -0700 Subject: [PATCH 11/21] test(verify): hand the script a boolean column with a hole as booleans pandas has no plain boolean column that carries a null, so read_json reads one with a hole as float64 and the script was handed 1.0 and 0.0 where the run had true and false. Seventeen operators compared unequal on that column alone. The prologue takes such a column to the nullable boolean dtype, which carries the two values and the hole, and leaves a column without a hole alone: that one arrives as bool already, and the nullable dtype would be one the run never had. The same rule the exact-integer restore follows. Co-Authored-By: Claude Opus 5 (1M context) --- .../amber/translator/verify/HarnessSpec.scala | 36 +++++++++++++++++++ .../translator/verify/StandaloneRunner.scala | 18 ++++++++++ 2 files changed, 54 insertions(+) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala index 685b6f056db..69c198639a9 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -112,4 +112,40 @@ class HarnessSpec extends AnyFlatSpec with Matchers { lines.get(2) should include("\"id\":3") } } + + // pandas has no plain boolean column that carries a null, so read_json reads + // one with a hole as float64 and the operator is handed 1.0 and 0.0 where the + // run had true and false. The prologue takes the column back to the nullable + // boolean dtype, which carries the two values and the hole. + it should "hand a boolean column with a hole to the script as booleans" taggedAs NeedsPython in { + val holed = new Schema( + new Attribute("id", AttributeType.INTEGER), + new Attribute("flag", AttributeType.BOOLEAN) + ) + def row(id: Int, flag: java.lang.Boolean): Tuple = { + val b = Tuple.builder(holed) + b.add(holed.getAttribute("id"), Int.box(id)) + b.add(holed.getAttribute("flag"), flag) + b.build() + } + + val dir = Files.createTempDirectory("harness-spec-boolean-") + val input = dir.resolve("input_port_0.jsonl") + TupleIO.writeTuples(input, Iterator(row(1, true), row(2, null), row(3, false)), holed) + val work = dir.resolve("standalone") + Files.createDirectories(work) + + val result = StandaloneRunner.run( + opDesc = new DistinctOpDesc, + inputs = Map(1 -> input), + outputPortCount = 1, + workDir = work + ) + + val lines = Files.readAllLines(result.outputs(1)) + lines should have size 3 + lines.get(0) should include("\"flag\":true") + lines.get(1) should include("\"flag\":null") + lines.get(2) should include("\"flag\":false") + } } diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index c60ea71e0e9..cff9ef6059d 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -306,6 +306,18 @@ object StandaloneRunner extends LazyLogging { sb.append(s"if ${py(col)} in in${n}df.columns:\n") sb.append(s" in${n}df[${py(col)}] = in${n}df[${py(col)}].astype('float64')\n") } + // A boolean column the fixture punched a hole in comes back as float64, + // so the operator is handed 1.0 and 0.0 where the run had true and + // false. The nullable dtype carries both the values and the hole. One + // without a hole already arrives as bool and is left alone, the way an + // exact integer is: the nullable dtype would be one the run never had. + booleanColumns(path).foreach { col => + sb.append( + s"if ${py(col)} in in${n}df.columns and not pd.api.types.is_bool_dtype(" + + s"in${n}df[${py(col)}]):\n" + ) + sb.append(s" in${n}df[${py(col)}] = in${n}df[${py(col)}].astype('boolean')\n") + } // Only where the reader lost the value: a column with no holes already // came back exact, and replacing it would hand the operator a nullable // dtype the run never had. @@ -423,6 +435,12 @@ object StandaloneRunner extends LazyLogging { private def stringColumns(input: Path): Seq[String] = columnsOfType(input, AttributeType.STRING) + /** BOOLEAN-typed column names. A hole makes read_json read the whole column + * as float64, so the run's true and false reach the operator as 1.0 and 0.0. + */ + private def booleanColumns(input: Path): Seq[String] = + columnsOfType(input, AttributeType.BOOLEAN) + private def columnsOfType(input: Path, attributeType: AttributeType): Seq[String] = scala.util .Try(TupleIO.readSchemaSidecar(input)) From 2c96035e342cb1bcc933f359b4e8a34fe4581356 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Fri, 18 Sep 2026 17:04:45 -0700 Subject: [PATCH 12/21] test(verify): write an integral column the way the engine's writer does The engine writes a tuple through the schema, so a column it declares INTEGER leaves as an integer however the operator held it. pandas has no plain integer that carries a null, so the same column left the script as 6.0 as soon as a row was missing, and the two sides then disagreed on every value in it. Seven operators were red on that alone. The epilogue writes the declared columns as integers, hole kept, and only where the values are whole: a column the script filled with 6.5 where the run had 6 is a real disagreement, and rounding it here would report the two as equal. An operator that rebuilds its rows hands back objects rather than floats, so the values are read one at a time. The schemas come from the caller, which has the ones the run itself produced. Asking the operator to propagate a second time builds it a second physical plan, and Aggregate's two-stage one does not survive that. Co-Authored-By: Claude Opus 5 (1M context) --- .../amber/translator/verify/HarnessSpec.scala | 38 ++++++++++ .../translator/verify/StandaloneRunner.scala | 76 +++++++++++++++++-- 2 files changed, 109 insertions(+), 5 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala index 69c198639a9..9836c1e8cc3 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -148,4 +148,42 @@ class HarnessSpec extends AnyFlatSpec with Matchers { lines.get(1) should include("\"flag\":null") lines.get(2) should include("\"flag\":false") } + + // The engine writes a tuple through the schema, so a column it declares + // INTEGER leaves as an integer. pandas has no plain integer that carries a + // null, so the same column left the script as 6.0 as soon as a row was + // missing, and the two sides disagreed on every value in it. + it should "write an integral column with a hole the way the engine writes it" taggedAs NeedsPython in { + val holed = new Schema( + new Attribute("id", AttributeType.INTEGER), + new Attribute("n", AttributeType.INTEGER) + ) + def row(id: Int, n: Integer): Tuple = { + val b = Tuple.builder(holed) + b.add(holed.getAttribute("id"), Int.box(id)) + b.add(holed.getAttribute("n"), n) + b.build() + } + + val dir = Files.createTempDirectory("harness-spec-integral-") + val input = dir.resolve("input_port_0.jsonl") + TupleIO.writeTuples(input, Iterator(row(1, 6), row(2, null), row(3, 7)), holed) + val work = dir.resolve("standalone") + Files.createDirectories(work) + + val result = StandaloneRunner.run( + opDesc = new DistinctOpDesc, + inputs = Map(1 -> input), + outputPortCount = 1, + workDir = work, + outputSchemas = Map(PortIdentity(0) -> holed) + ) + + val lines = Files.readAllLines(result.outputs(1)) + lines should have size 3 + lines.get(0) should include("\"n\":6") + lines.get(0) should not include "\"n\":6.0" + lines.get(1) should include("\"n\":null") + lines.get(2) should include("\"n\":7") + } } diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index cff9ef6059d..61ea12bac22 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -80,7 +80,8 @@ object StandaloneRunner extends LazyLogging { outputPortCount: Int, workDir: Path, pythonExe: String = resolvePython(), - exactIntegers: Boolean = false + exactIntegers: Boolean = false, + outputSchemas: Map[PortIdentity, Schema] = Map.empty ): Result = { val gen = opDesc match { case g: StandaloneCodeGenerator => g @@ -98,16 +99,18 @@ object StandaloneRunner extends LazyLogging { (1 to outputPortCount).map(i => i -> workDir.resolve(s"output_port_${i - 1}.jsonl")).toMap else Map.empty + val inputSchemas = schemasOf(inputs) val source = renderScript( // The sidecar says what each column was DECLARED as, which the JSONL // cannot carry. - gen.generateStandaloneCode(schemasOf(inputs)), + gen.generateStandaloneCode(inputSchemas), inputs, outputPaths, gen.standaloneHelpers(), gen.standaloneImports(), - exactIntegers + exactIntegers, + integralOutputColumns(outputSchemas, outputPaths.keys.toSeq) ) Files.write(scriptPath, source.getBytes(StandardCharsets.UTF_8)) @@ -175,7 +178,8 @@ object StandaloneRunner extends LazyLogging { outputs: Map[Int, Path], helpers: Seq[String], imports: Seq[String], - exactIntegers: Boolean + exactIntegers: Boolean, + integralOutputs: Map[Int, Seq[String]] ): String = { val sb = new StringBuilder @@ -263,6 +267,41 @@ object StandaloneRunner extends LazyLogging { sb.append(" return df\n") sb.append("\n") + // The engine writes a tuple through the schema, so a column it declares + // INTEGER leaves as an integer however the operator held it. pandas has no + // plain integer that carries a null, so the same column leaves the script + // as 6.0 whenever a row is missing, and the two sides then disagree on + // every value in it. The declared columns are written the way the engine + // writes them: as integers, with the hole kept. + // + // Only where the values ARE whole. A column the script filled with 6.5 + // where the run had 6 is a real disagreement, and rounding it here would + // report the two sides as equal. + if (integralOutputs.values.exists(_.nonEmpty)) { + sb.append("def _texera_int_cols(df, _columns):\n") + // An operator that rebuilds its rows hands back a column of objects + // rather than a float column, so the values are read one at a time: a + // numpy scalar unwrapped, anything that is not a plain number left for + // the branch below to refuse. + sb.append(" def _texera_num(_v):\n") + sb.append(" if hasattr(_v, 'item') and getattr(_v, 'ndim', None) == 0:\n") + sb.append(" _v = _v.item()\n") + sb.append(" if isinstance(_v, bool) or not isinstance(_v, (int, float)):\n") + sb.append(" return None\n") + sb.append(" return _v\n") + sb.append(" for _c in _columns:\n") + sb.append(" if _c not in df.columns or pd.api.types.is_integer_dtype(df[_c]):\n") + sb.append(" continue\n") + sb.append(" _known = [_texera_num(_v) for _v in df[_c] if not pd.isna(_v)]\n") + sb.append(" if any(_v is None or float(_v) % 1 != 0 for _v in _known):\n") + sb.append(" continue\n") + sb.append(" df[_c] = pd.array(\n") + sb.append(" [pd.NA if pd.isna(_v) else int(_v) for _v in df[_c]], dtype='Int64'\n") + sb.append(" )\n") + sb.append(" return df\n") + sb.append("\n") + } + // Prologue: load each external input into in{N}df. Every option below undoes // an inference that would otherwise hand the two paths different data. // @@ -368,8 +407,12 @@ object StandaloneRunner extends LazyLogging { // caller is expected to verify viz outputs by other means. outputs.toSeq.sortBy(_._1).foreach { case (n, path) => + val frame = integralOutputs.getOrElse(n, Seq.empty) match { + case Seq() => s"out${n}df" + case cols => s"_texera_int_cols(out${n}df, [${cols.map(py).mkString(", ")}])" + } sb.append( - s"_texera_encode_obj_cols(_texera_encode_ts_cols(out${n}df))" + + s"_texera_encode_obj_cols(_texera_encode_ts_cols($frame))" + s".to_json(${py(path.toString)}, orient='records', lines=True)\n" ) } @@ -441,6 +484,29 @@ object StandaloneRunner extends LazyLogging { private def booleanColumns(input: Path): Seq[String] = columnsOfType(input, AttributeType.BOOLEAN) + /** Per output port, the columns the schema declares integral, which is what + * the engine's writer goes by. + * + * The schemas are the ones the run itself produced rather than a propagation + * of this operator's own: asking an operator to propagate a second time + * builds it a second physical plan, and Aggregate's two-stage one does not + * survive that. A caller with no schemas to hand over names no column, and + * the epilogue then writes each frame as pandas holds it. + */ + private def integralOutputColumns( + outputSchemas: Map[PortIdentity, Schema], + ports: Seq[Int] + ): Map[Int, Seq[String]] = + ports.map { port => + val declared = outputSchemas + .get(PortIdentity(port - 1)) + .toSeq + .flatMap(_.getAttributes) + .filter(a => a.getType == AttributeType.INTEGER || a.getType == AttributeType.LONG) + .map(_.getName) + port -> declared + }.toMap + private def columnsOfType(input: Path, attributeType: AttributeType): Seq[String] = scala.util .Try(TupleIO.readSchemaSidecar(input)) From f4b9496ebd7c1c4d0a423e62d3003b74ada3d1a7 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Fri, 18 Sep 2026 20:54:34 -0700 Subject: [PATCH 13/21] fix(verify): hand a binary column to the script as the bytes the run gets JSON carries bytes as base64 text, which is how the fixture carries a BINARY column, and the engine decodes it before the operator sees the field. The script was handed the text: its operator held a str where the run's held bytes, and anything it did with them, a length or a decode or a digest, answered for the base64 rather than for the value. Distinct answers the same either way, which is why the column looked right on the way out. Decoding it on the way in is half of it. The writer pickles whatever is not a plain scalar, so the bytes would have left as a base64 pickle of themselves where the engine's writer writes the base64 of the value. Bytes are now written as that base64. Co-Authored-By: Claude Opus 5 (1M context) --- .../amber/translator/verify/HarnessSpec.scala | 50 +++++++++++++++++++ .../translator/verify/StandaloneRunner.scala | 28 +++++++++++ 2 files changed, 78 insertions(+) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala index 9836c1e8cc3..2189964fccc 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -26,6 +26,7 @@ import org.scalatest.Tag import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers +import java.nio.charset.StandardCharsets import java.nio.file.{Files, Path} /** The two ways of running one operator, and the file format they meet in. @@ -186,4 +187,53 @@ class HarnessSpec extends AnyFlatSpec with Matchers { lines.get(1) should include("\"n\":null") lines.get(2) should include("\"n\":7") } + + // JSON carries bytes as base64 text, and the engine decodes it before the + // operator sees the field (TupleIO.readTuples). Left as text, the script's + // operator is handed a str where the run's was handed bytes, and anything it + // does with them answers for the base64 rather than for the value. + it should "hand a binary column to the script as bytes" taggedAs NeedsPython in { + val withBlob = new Schema( + new Attribute("id", AttributeType.INTEGER), + new Attribute("blob", AttributeType.BINARY) + ) + def row(id: Int, blob: Array[Byte]): Tuple = { + val b = Tuple.builder(withBlob) + b.add(withBlob.getAttribute("id"), Int.box(id)) + b.add(withBlob.getAttribute("blob"), blob) + b.build() + } + + val dir = Files.createTempDirectory("harness-spec-binary-") + val input = dir.resolve("input_port_0.jsonl") + TupleIO.writeTuples( + input, + Iterator(row(1, "hi".getBytes(StandardCharsets.UTF_8)), row(2, null)), + withBlob + ) + val work = dir.resolve("standalone") + Files.createDirectories(work) + + val result = StandaloneRunner.run( + opDesc = new DistinctOpDesc, + inputs = Map(1 -> input), + outputPortCount = 1, + workDir = work + ) + + // The prologue decodes that column and no other. Distinct answers the same + // whether it is handed the bytes or the base64 spelling of them, so what the + // body is given has to be read off the script rather than off the output. + val script = Files.readString(work.resolve("script.py")) + script should include("in1df['blob'] = in1df['blob'].map(") + script should include("base64.b64decode") + script should not include "in1df['id'] = in1df['id'].map(" + + // And the column still leaves as the base64 the engine's writer wrote: + // decoding it on the way in must not turn it into a pickle on the way out. + val lines = Files.readAllLines(result.outputs(1)) + lines should have size 2 + lines.get(0) should include("\"blob\":\"aGk=\"") + lines.get(1) should include("\"blob\":null") + } } diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 61ea12bac22..1fac3680eaa 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -220,6 +220,14 @@ object StandaloneRunner extends LazyLogging { sb.append(" for _c in df.columns:\n") sb.append(" if df[_c].dtype == object:\n") sb.append(" df[_c] = df[_c].map(_texera_plain)\n") + // Bytes are written as the base64 the engine's writer writes, not as a + // base64 pickle of them: the prologue hands a BINARY column over as bytes, + // the way the engine hands one to an operator, and the branch below would + // otherwise answer a column the run wrote as `aGk=` with a pickle of it. + sb.append( + " df[_c] = df[_c].map(lambda _v: base64.b64encode(_v).decode('ascii') " + + "if isinstance(_v, (bytes, bytearray)) else _v)\n" + ) sb.append( " df[_c] = df[_c].map(lambda _v: base64.b64encode(pickle.dumps(_v)).decode('ascii') " + "if not isinstance(_v, (str, int, float, bool, type(None))) else _v)\n" @@ -357,6 +365,20 @@ object StandaloneRunner extends LazyLogging { ) sb.append(s" in${n}df[${py(col)}] = in${n}df[${py(col)}].astype('boolean')\n") } + // The fixture carries a BINARY column as base64 text, which is how JSON + // carries bytes at all, and the engine decodes it before the operator + // sees it (TupleIO.readTuples). Left as text, the script's operator is + // handed a str where the run's was handed bytes, and anything it does + // with them -- a length, a decode, a digest -- answers for the base64 + // rather than for the value. + binaryColumns(path).foreach { col => + sb.append(s"if ${py(col)} in in${n}df.columns:\n") + sb.append( + s" in${n}df[${py(col)}] = in${n}df[${py(col)}].map(\n" + + s" lambda _v: None if pd.isna(_v) else base64.b64decode(_v)\n" + + s" )\n" + ) + } // Only where the reader lost the value: a column with no holes already // came back exact, and replacing it would hand the operator a nullable // dtype the run never had. @@ -484,6 +506,12 @@ object StandaloneRunner extends LazyLogging { private def booleanColumns(input: Path): Seq[String] = columnsOfType(input, AttributeType.BOOLEAN) + /** BINARY-typed column names. JSON carries bytes as base64 text, and the + * engine decodes it before the operator sees the field. + */ + private def binaryColumns(input: Path): Seq[String] = + columnsOfType(input, AttributeType.BINARY) + /** Per output port, the columns the schema declares integral, which is what * the engine's writer goes by. * From 36f79d8f988edbde4f12b655a133e1cd91c75d15 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Sat, 19 Sep 2026 00:49:04 -0700 Subject: [PATCH 14/21] test(verify): name the file a source reads, which its body leaves to the caller A source writes a placeholder where its file should be named and leaves the naming to whoever assembles the script, so that two sources reading different files whose paths end alike do not both ask for the same one. The translator does that across a whole plan. This runner assembles a single operator and bound nothing, so every source's script stopped on a `sourceFile` that was never defined. It binds the name the source offers, there being nobody to collide with in a one-operator script. No operator spec could have caught this: they assert the text the operator emits, and the text is right. Co-Authored-By: Claude Opus 5 (1M context) --- .../amber/translator/verify/HarnessSpec.scala | 59 ++++++++++++++++++- .../translator/verify/StandaloneRunner.scala | 16 ++++- 2 files changed, 72 insertions(+), 3 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala index 2189964fccc..5dc3e896e85 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -20,8 +20,11 @@ package org.apache.texera.amber.translator.verify import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema, Tuple} -import org.apache.texera.amber.core.workflow.PortIdentity +import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} +import org.apache.texera.amber.core.workflow.{OutputPort, PhysicalOp, PortIdentity} import org.apache.texera.amber.operator.distinct.DistinctOpDesc +import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} +import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} import org.scalatest.Tag import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers @@ -65,6 +68,32 @@ class HarnessSpec extends AnyFlatSpec with Matchers { test(dir, input) } + /** A source, in the one respect this spec is about: it reads no input port and + * names the file it reads by placeholder, leaving the naming to whoever + * assembles the script. + */ + private class StubSource(path: Path) extends LogicalOp with StandaloneCodeGenerator { + override def getPhysicalOp( + workflowId: WorkflowIdentity, + executionId: ExecutionIdentity + ): PhysicalOp = + throw new UnsupportedOperationException("the harness never builds a physical op") + + override def operatorInfo: OperatorInfo = + OperatorInfo( + "StubSource", + "Stands in for a source that names its file by placeholder", + OperatorGroupConstants.INPUT_GROUP, + inputPorts = List.empty, + outputPorts = List(OutputPort()) + ) + + override def standaloneSourcePath(): Option[String] = Some(path.toUri.toString) + + override def generateStandaloneCode(): String = + s"out1df = pd.read_json(${StandaloneCodeGenerator.SourceFilePlaceholder}, lines=True)" + } + "TupleIO" should "read back the rows and the schema it wrote" in { withInput { (_, input) => // The schema travels in a sidecar rather than in the JSONL, which carries @@ -114,6 +143,34 @@ class HarnessSpec extends AnyFlatSpec with Matchers { } } + // A source writes a placeholder where its file should be named, since only + // whoever assembles the whole script can settle on a name. Nothing bound it + // here, so every source's script stopped on a `sourceFile` that was never + // defined, which no operator spec could see: they assert the text the operator + // emits, and the text is right. + it should "name the file a source reads, which the body leaves to it" taggedAs NeedsPython in { + val dir = Files.createTempDirectory("harness-source-") + // Beside the script, under the name the source offers, which is how an + // exported script is meant to find its data. + val data = dir.resolve("rows.jsonl") + TupleIO.writeTuples(data, rows.iterator, schema) + + val result = StandaloneRunner.run( + opDesc = new StubSource(data), + inputs = Map.empty, + outputPortCount = 1, + workDir = dir + ) + + val script = Files.readString(dir.resolve("script.py")) + script should include("sourceFile = 'rows.jsonl'") + + val lines = Files.readAllLines(result.outputs(1)) + lines should have size 4 + lines.get(0) should include("\"id\":1") + lines.get(3) should include("\"id\":2") + } + // pandas has no plain boolean column that carries a null, so read_json reads // one with a hole as float64 and the operator is handed 1.0 and 0.0 where the // run had true and false. The prologue takes the column back to the nullable diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 1fac3680eaa..8d2513f1784 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -110,7 +110,8 @@ object StandaloneRunner extends LazyLogging { gen.standaloneHelpers(), gen.standaloneImports(), exactIntegers, - integralOutputColumns(outputSchemas, outputPaths.keys.toSeq) + integralOutputColumns(outputSchemas, outputPaths.keys.toSeq), + gen.standaloneSourceName() ) Files.write(scriptPath, source.getBytes(StandardCharsets.UTF_8)) @@ -179,7 +180,8 @@ object StandaloneRunner extends LazyLogging { helpers: Seq[String], imports: Seq[String], exactIntegers: Boolean, - integralOutputs: Map[Int, Seq[String]] + integralOutputs: Map[Int, Seq[String]], + sourceFileName: Option[String] ): String = { val sb = new StringBuilder @@ -419,6 +421,16 @@ object StandaloneRunner extends LazyLogging { sb.append('\n') } + // A source names the file it reads by placeholder and leaves the naming to + // whoever assembles the script, the translator doing it across a whole plan + // so that two sources reading different files whose paths end alike do not + // both ask for the same one. This runner assembles a single operator, so + // there is nobody to collide with and the name the source offers stands. + // Without this the body reads a bare `sourceFile` that was never bound. + sourceFileName.foreach { name => + sb.append(s"${StandaloneCodeGenerator.SourceFilePlaceholder} = ${py(name)}\n\n") + } + sb.append("# ── operator body ──\n") sb.append(body) if (!body.endsWith("\n")) sb.append('\n') From 67f361f64efa8c8bed3b82d110e6604607b942f2 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Tue, 22 Sep 2026 15:53:18 -0700 Subject: [PATCH 15/21] test(verify): let an empty input frame keep the columns pandas gives it apache/texera#8488 is closed: the engine hands an operator a table with no columns when a port carried no rows, so rebuilding them from the sidecar would make the generated script differ from the run it is compared against. Co-Authored-By: Claude Opus 5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 34 ------------------- 1 file changed, 34 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 8d2513f1784..288904f10c8 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -335,18 +335,6 @@ object StandaloneRunner extends LazyLogging { s"in${n}df = pd.read_json(${py(path.toString)}, lines=True, " + s"convert_dates=False, precise_float=True$dtype)\n" ) - // read_json has no rows to read column names off a file with none in it, - // so it produces a frame of no columns, while the engine hands the - // operator the port's declared ones (see Table.empty_of). Rebuild it - // from the sidecar, dtypes included, so an empty table is the same table - // on both paths. - emptyFrameColumns(path) match { - case Seq() => () - case cols => - val fields = cols.map { case (c, d) => s"${py(c)}: pd.Series(dtype=${py(d)})" } - sb.append(s"if in${n}df.empty:\n") - sb.append(s" in${n}df = pd.DataFrame({${fields.mkString(", ")}})\n") - } timestampColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") sb.append(s" in${n}df[${py(col)}] = pd.to_datetime(in${n}df[${py(col)}])\n") @@ -481,28 +469,6 @@ object StandaloneRunner extends LazyLogging { .map(schema => PortIdentity(port - 1) -> schema) } - /** Each declared column with the pandas dtype an Arrow round-trip gives it, which - * is what the engine's own empty table carries. Only the types a fixture can hold - * are named; anything else falls back to object, the dtype an inferred column of - * unknown content would have had anyway. - */ - private def emptyFrameColumns(input: Path): Seq[(String, String)] = - scala.util - .Try(TupleIO.readSchemaSidecar(input)) - .toOption - .toSeq - .flatMap(_.getAttributes.map { attr => - val dtype = attr.getType match { - case AttributeType.INTEGER => "int32" - case AttributeType.LONG => "int64" - case AttributeType.DOUBLE => "float64" - case AttributeType.BOOLEAN => "bool" - case AttributeType.TIMESTAMP => "datetime64[us]" - case _ => "object" - } - attr.getName -> dtype - }) - /** The columns the sidecar declares integral, both widths: the loss is the same * for either. */ From 84fedadae3e298f66728052f99591e9964f4b7ca Mon Sep 17 00:00:00 2001 From: kary zheng Date: Wed, 23 Sep 2026 15:52:33 -0700 Subject: [PATCH 16/21] test(verify): hand the script a timestamp the year 2500 and nine digits both reach The harness read every TIMESTAMP column through pd.to_datetime, whose nanoseconds reach only 1677 to 2262, so a year-2500 input failed before the operator ran. Java's Timestamp holds it. A column outside that window is now read cell by cell at microseconds. On the way out, %f wrote six digits of the fraction where Timestamp.toString writes nine, so an unchanged value could disagree with the run because of the harness alone. The fraction is now built from the nanoseconds. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 27 ++++++++++++++++--- 1 file changed, 24 insertions(+), 3 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 288904f10c8..35cc1ac37bc 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -264,11 +264,14 @@ object StandaloneRunner extends LazyLogging { // "yyyy-mm-dd hh:mm:ss.f", trailing zeros trimmed to at least one digit — // whereas pandas' to_json would emit epoch millis. Convert datetime columns // back to that exact form before writing so both paths' JSONL agree. + // + // The fraction is built from the nanoseconds and not from %f, which stops at + // six digits, while Timestamp.toString writes all nine. sb.append("def _texera_ts_str(_v):\n") sb.append(" if pd.isna(_v):\n") sb.append(" return None\n") - sb.append(" _s = _v.strftime('%Y-%m-%d %H:%M:%S.%f').rstrip('0')\n") - sb.append(" return _s + '0' if _s.endswith('.') else _s\n") + sb.append(" _f = f'{_v.microsecond * 1000 + _v.nanosecond:09d}'.rstrip('0') or '0'\n") + sb.append(" return _v.strftime('%Y-%m-%d %H:%M:%S.') + _f\n") sb.append("\n") sb.append("def _texera_encode_ts_cols(df):\n") sb.append(" for _c in df.columns:\n") @@ -277,6 +280,24 @@ object StandaloneRunner extends LazyLogging { sb.append(" return df\n") sb.append("\n") + // A TIMESTAMP is a java.sql.Timestamp on the engine side, which holds the + // year 2500 as readily as 2024. pandas' default nanoseconds reach only 1677 + // to 2262, so a column outside that window is read at microseconds instead + // of failing before the operator runs. A column inside it keeps the + // nanoseconds, which Timestamp counts too. + // + // Cell by cell through dateutil and not astype: astype takes its unit from + // the text, so a fraction of nine digits is parsed as nanoseconds again and + // fails the same way. A datetime holds microseconds, so the rest is dropped. + sb.append("def _texera_read_ts(_s):\n") + sb.append(" try:\n") + sb.append(" return pd.to_datetime(_s)\n") + sb.append(" except pd.errors.OutOfBoundsDatetime:\n") + sb.append(" from dateutil.parser import parse as _parse_date\n") + sb.append(" _read = _s.map(lambda _v: None if pd.isna(_v) else _parse_date(_v))\n") + sb.append(" return _read.astype('datetime64[us]')\n") + sb.append("\n") + // The engine writes a tuple through the schema, so a column it declares // INTEGER leaves as an integer however the operator held it. pandas has no // plain integer that carries a null, so the same column leaves the script @@ -337,7 +358,7 @@ object StandaloneRunner extends LazyLogging { ) timestampColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") - sb.append(s" in${n}df[${py(col)}] = pd.to_datetime(in${n}df[${py(col)}])\n") + sb.append(s" in${n}df[${py(col)}] = _texera_read_ts(in${n}df[${py(col)}])\n") } doubleColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") From 658b1a6875a1bc8f7cc3cc04b41837a5e937afb4 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Wed, 23 Sep 2026 23:14:53 -0700 Subject: [PATCH 17/21] test(verify): hand the script an empty input with its columns read_json finds no column names in a file with no rows, so the script was handed a frame of no columns. An exported script never sees one: a filter that matches nothing and a header-only CSV both leave the columns in place, and each operator's code is handed that frame. Without them 21 operators that handle an empty table raised KeyError on the script side only, a failure the test made rather than one the export has. The columns are rebuilt from the sidecar again. The earlier removal took the engine's columnless empty table as the reference for this side too, but the engine is what Path A stands for, not the script. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../translator/verify/StandaloneRunner.scala | 33 +++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 35cc1ac37bc..42786c4d4dc 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -356,6 +356,18 @@ object StandaloneRunner extends LazyLogging { s"in${n}df = pd.read_json(${py(path.toString)}, lines=True, " + s"convert_dates=False, precise_float=True$dtype)\n" ) + // read_json finds no column names in a file with no rows, so it produces a + // frame of no columns. In an exported script an empty frame keeps its + // columns, since a filter that matches nothing and a header-only CSV both + // leave them in place. Rebuild them from the sidecar so the script is + // handed the table it would get in a real run. + emptyFrameColumns(path) match { + case Seq() => () + case cols => + val fields = cols.map { case (c, d) => s"${py(c)}: pd.Series(dtype=${py(d)})" } + sb.append(s"if in${n}df.empty:\n") + sb.append(s" in${n}df = pd.DataFrame({${fields.mkString(", ")}})\n") + } timestampColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") sb.append(s" in${n}df[${py(col)}] = _texera_read_ts(in${n}df[${py(col)}])\n") @@ -490,6 +502,27 @@ object StandaloneRunner extends LazyLogging { .map(schema => PortIdentity(port - 1) -> schema) } + /** Each declared column with the pandas dtype it holds when it has values in it. + * Only the types a fixture can hold are named; anything else falls back to + * object, the dtype an inferred column of unknown content would have had anyway. + */ + private def emptyFrameColumns(input: Path): Seq[(String, String)] = + scala.util + .Try(TupleIO.readSchemaSidecar(input)) + .toOption + .toSeq + .flatMap(_.getAttributes.map { attr => + val dtype = attr.getType match { + case AttributeType.INTEGER => "int32" + case AttributeType.LONG => "int64" + case AttributeType.DOUBLE => "float64" + case AttributeType.BOOLEAN => "bool" + case AttributeType.TIMESTAMP => "datetime64[us]" + case _ => "object" + } + attr.getName -> dtype + }) + /** The columns the sidecar declares integral, both widths: the loss is the same * for either. */ From 7b3c37a9df1252ad5775b96d7a3fc3e11bfea9ae Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 24 Sep 2026 12:55:21 -0700 Subject: [PATCH 18/21] test(verify): hand a pickled binary cell to the script as the object The worker's ArrowTableTupleProvider unpickles a BINARY cell that starts with the cast's pickle marker, so a model column reaches the next operator as a model. The script's prologue decoded the base64 and stopped there, so an operator reading a model port was handed bytes. The prologue now makes the same test and slice, so the script gets the object the run gets and any other bytes stay bytes. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../amber/translator/verify/HarnessSpec.scala | 55 ++++++++++++++++++- .../translator/verify/StandaloneRunner.scala | 15 ++++- 2 files changed, 66 insertions(+), 4 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala index 5dc3e896e85..cb2348899a5 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -21,7 +21,7 @@ package org.apache.texera.amber.translator.verify import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema, Tuple} import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} -import org.apache.texera.amber.core.workflow.{OutputPort, PhysicalOp, PortIdentity} +import org.apache.texera.amber.core.workflow.{InputPort, OutputPort, PhysicalOp, PortIdentity} import org.apache.texera.amber.operator.distinct.DistinctOpDesc import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} @@ -94,6 +94,47 @@ class HarnessSpec extends AnyFlatSpec with Matchers { s"out1df = pd.read_json(${StandaloneCodeGenerator.SourceFilePlaceholder}, lines=True)" } + /** Reports the Python type of each cell of its `blob` column. */ + private class BlobTypeOp extends LogicalOp with StandaloneCodeGenerator { + override def getPhysicalOp( + workflowId: WorkflowIdentity, + executionId: ExecutionIdentity + ): PhysicalOp = + throw new UnsupportedOperationException("the harness never builds a physical op") + + override def operatorInfo: OperatorInfo = + OperatorInfo( + "BlobType", + "Reports the type of each binary cell", + OperatorGroupConstants.UTILITY_GROUP, + inputPorts = List(InputPort()), + outputPorts = List(OutputPort()) + ) + + override def generateStandaloneCode(): String = + """out1df = pd.DataFrame({"kind": [type(_v).__name__ for _v in in1df["blob"]]})""" + } + + private def blobKind(cell: Array[Byte]): String = { + val blobOnly = new Schema(new Attribute("blob", AttributeType.BINARY)) + val dir = Files.createTempDirectory("harness-spec-blob-kind-") + val input = dir.resolve("input_port_0.jsonl") + TupleIO.writeTuples( + input, + Iterator(Tuple.builder(blobOnly).add(blobOnly.getAttribute("blob"), cell).build()), + blobOnly + ) + val result = StandaloneRunner.run( + opDesc = new BlobTypeOp, + inputs = Map(1 -> input), + outputPortCount = 1, + workDir = dir + ) + val lines = Files.readAllLines(result.outputs(1)) + lines should have size 1 + lines.get(0) + } + "TupleIO" should "read back the rows and the schema it wrote" in { withInput { (_, input) => // The schema travels in a sidecar rather than in the JSONL, which carries @@ -293,4 +334,16 @@ class HarnessSpec extends AnyFlatSpec with Matchers { lines.get(0) should include("\"blob\":\"aGk=\"") lines.get(1) should include("\"blob\":null") } + + // A model column arrives as the marker followed by the pickle, and the worker + // unpickles it before the operator sees it. This is pickle.dumps(["a"], + // protocol=0), which stands in for a fitted estimator. + it should "hand a pickled binary cell to the script as the object" taggedAs NeedsPython in { + val pickled = "pickle ".getBytes("US-ASCII") ++ "(lp0\nVa\np1\na.".getBytes("US-ASCII") + blobKind(pickled) should include("\"kind\":\"list\"") + } + + it should "hand any other binary cell to the script as bytes" taggedAs NeedsPython in { + blobKind(Array[Byte](0, 1, 2)) should include("\"kind\":\"bytes\"") + } } diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 42786c4d4dc..74df136d2c4 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -237,6 +237,17 @@ object StandaloneRunner extends LazyLogging { sb.append(" return df\n") sb.append("\n") + // A model column arrives as the cast's pickle marker followed by the pickle, + // and the worker's ArrowTableTupleProvider unpickles such a cell before the + // operator sees it. Same test and slice, so the script is handed the model + // the run was handed and any other bytes stay bytes. + sb.append("def _texera_read_binary(_v):\n") + sb.append(" if pd.isna(_v):\n") + sb.append(" return None\n") + sb.append(" _b = base64.b64decode(_v)\n") + sb.append(" return pickle.loads(_b[10:]) if _b[:6] == b'pickle' else _b\n") + sb.append("\n") + // Both paths have to be handed the same numbers. `read_json` parses a column // holding a null through float64, so a LONG of 9007199254740993 arrives as // 9007199254740992 while the engine still has the tuple. Python's json reads @@ -397,9 +408,7 @@ object StandaloneRunner extends LazyLogging { binaryColumns(path).foreach { col => sb.append(s"if ${py(col)} in in${n}df.columns:\n") sb.append( - s" in${n}df[${py(col)}] = in${n}df[${py(col)}].map(\n" + - s" lambda _v: None if pd.isna(_v) else base64.b64decode(_v)\n" + - s" )\n" + s" in${n}df[${py(col)}] = in${n}df[${py(col)}].map(_texera_read_binary)\n" ) } // Only where the reader lost the value: a column with no holes already From d8440e2f8849e7726d13c585f0cafe063fcc0f80 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Thu, 24 Sep 2026 20:59:03 -0700 Subject: [PATCH 19/21] test(verify): hand a real fitted model to the script in the harness spec The pickled-cell test used a pickled list in place of a model. It now pickles a fitted DecisionTreeClassifier with the interpreter the script runs and has the script call predict on it, which is the call that raised AttributeError when the prologue handed over bytes. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../amber/translator/verify/HarnessSpec.scala | 33 +++++++++++++++---- 1 file changed, 26 insertions(+), 7 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala index cb2348899a5..163309ed556 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/HarnessSpec.scala @@ -31,6 +31,8 @@ import org.scalatest.matchers.should.Matchers import java.nio.charset.StandardCharsets import java.nio.file.{Files, Path} +import java.util.Base64 +import scala.sys.process._ /** The two ways of running one operator, and the file format they meet in. * @@ -115,7 +117,26 @@ class HarnessSpec extends AnyFlatSpec with Matchers { """out1df = pd.DataFrame({"kind": [type(_v).__name__ for _v in in1df["blob"]]})""" } - private def blobKind(cell: Array[Byte]): String = { + /** Calls predict on each cell of its `blob` column, as a model's consumer does. */ + private class ModelPredictOp extends BlobTypeOp { + override def generateStandaloneCode(): String = + """out1df = pd.DataFrame({"kind": [str(_v.predict([[1.0]])[0]) for _v in in1df["blob"]]})""" + } + + // A fitted DecisionTreeClassifier behind the marker the cast writes, pickled + // by the interpreter the script runs so the two share a sklearn. + private def fittedTreeCell(): Array[Byte] = { + val script = + """import base64, pickle, sys + |from sklearn.tree import DecisionTreeClassifier + |model = DecisionTreeClassifier(random_state=0).fit([[0.0], [1.0]], [0, 1]) + |sys.stdout.write(base64.b64encode(b"pickle " + pickle.dumps(model)).decode("ascii")) + |""".stripMargin + val encoded = Process(Seq(PyOpExecHarness.resolvePython(), "-c", script)).!! + Base64.getDecoder.decode(encoded.trim) + } + + private def blobKind(cell: Array[Byte], opDesc: BlobTypeOp = new BlobTypeOp): String = { val blobOnly = new Schema(new Attribute("blob", AttributeType.BINARY)) val dir = Files.createTempDirectory("harness-spec-blob-kind-") val input = dir.resolve("input_port_0.jsonl") @@ -125,7 +146,7 @@ class HarnessSpec extends AnyFlatSpec with Matchers { blobOnly ) val result = StandaloneRunner.run( - opDesc = new BlobTypeOp, + opDesc = opDesc, inputs = Map(1 -> input), outputPortCount = 1, workDir = dir @@ -336,11 +357,9 @@ class HarnessSpec extends AnyFlatSpec with Matchers { } // A model column arrives as the marker followed by the pickle, and the worker - // unpickles it before the operator sees it. This is pickle.dumps(["a"], - // protocol=0), which stands in for a fitted estimator. - it should "hand a pickled binary cell to the script as the object" taggedAs NeedsPython in { - val pickled = "pickle ".getBytes("US-ASCII") ++ "(lp0\nVa\np1\na.".getBytes("US-ASCII") - blobKind(pickled) should include("\"kind\":\"list\"") + // unpickles it before the operator sees it, so predict works on it. + it should "hand a pickled model cell to the script as the model" taggedAs NeedsPython in { + blobKind(fittedTreeCell(), new ModelPredictOp) should include("\"kind\":\"1\"") } it should "hand any other binary cell to the script as bytes" taggedAs NeedsPython in { From 83a6c1e53c03531559858eec5bee5891de7bab62 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Fri, 25 Sep 2026 01:11:11 -0700 Subject: [PATCH 20/21] test(verify): write down the dtype each output column was left in The JSONL an output is written to carries values and not their dtypes, so a float32 and a float64 write the same text and a zoned timestamp loses its zone on the way. Both still reach the next operator as what they are. The script now writes each column's dtype beside the output, taken before the encoding rewrites the frame, for a runner to check. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../texera/amber/translator/verify/StandaloneRunner.scala | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 74df136d2c4..3c7d526b5f5 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -475,6 +475,14 @@ object StandaloneRunner extends LazyLogging { case Seq() => s"out${n}df" case cols => s"_texera_int_cols(out${n}df, [${cols.map(py).mkString(", ")}])" } + // The dtype each column is left in, which the JSONL cannot carry: a + // float32 and a float64 write the same text, and a zoned timestamp loses + // its zone on the way to one. Taken before the encoding below rewrites + // the frame. + sb.append( + s"json.dump({str(_c): str(_t) for _c, _t in out${n}df.dtypes.items()}, " + + s"open(${py(path.toString + ".dtypes.json")}, 'w'))\n" + ) sb.append( s"_texera_encode_obj_cols(_texera_encode_ts_cols($frame))" + s".to_json(${py(path.toString)}, orient='records', lines=True)\n" From ccb2046b51463cd089d1aa38b0503a0c41f8ef67 Mon Sep 17 00:00:00 2001 From: kary zheng Date: Sat, 26 Sep 2026 19:25:16 -0700 Subject: [PATCH 21/21] test(verify): say how the script side reads a timestamp and which interpreter it falls back to Co-Authored-By: Claude Opus 5.5 (1M context) --- .../texera/amber/translator/verify/OpExecHarness.scala | 5 ++--- .../texera/amber/translator/verify/StandaloneRunner.scala | 4 ++-- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala index 01e48e40da2..d584d35e04c 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala @@ -386,9 +386,8 @@ object TupleIO { // Timestamps round-trip through the JDBC string form // ("yyyy-mm-dd hh:mm:ss[.f]"), the exact inverse of Timestamp.toString // below — timezone-free, so no shift across write/read. The Python - // side reads this column with convert_dates=False (see - // StandaloneRunner) and treats it as an opaque string, so both paths - // agree on pass-through. + // side parses it back into a datetime and writes it out in this same + // form (see StandaloneRunner), so both paths agree on pass-through. case AttributeType.TIMESTAMP => Timestamp.valueOf(fieldNode.asText()) case other => diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala index 3c7d526b5f5..f59f29aa671 100644 --- a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/StandaloneRunner.scala @@ -70,8 +70,8 @@ object StandaloneRunner extends LazyLogging { * @param workDir directory used for the generated `script.py` and output * JSONL files. Created if missing. * @param pythonExe path to the Python 3.12 interpreter. Defaults to - * the env var `UDF_PYTHON_PATH`, then `python3.12`, then - * `python3`. The same fallback chain used by the rest of + * the env var `UDF_PYTHON_PATH`, then `python3.12`. + * The same fallback chain used by the rest of * the Texera test suite for Python-backed operators. */ def run(