This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new bf356f5a0c test(verify): run an operator through the engine and keep
what it wrote (#8356)
bf356f5a0c is described below
commit bf356f5a0cde39bee4b1cc98cd1ad0fb55910de5
Author: Kary Zheng <[email protected]>
AuthorDate: Tue Sep 15 05:31:35 2026 +0000
test(verify): run an operator through the engine and keep what it wrote
(#8356)
### What changes were proposed in this PR?
An operator's answer cannot be compared against anything until there is
a way
to get one. `OpExecHarness` builds the physical operator a descriptor
describes, feeds it the rows of a JSONL file per input port, and writes
what
each output port produced back out, schema in a sidecar because JSONL
carries
values alone and cannot say a column is an integer rather than a number.
The tests that need a Python interpreter are tagged and split into a job
that
provisions one, so the job that does not stays as fast as it was.
### Any related issues, documentation, discussions?
Part of #8325, 3 of 27; that issue lists the set in order.
Closes #8408, the task this change is the whole of.
### How was this PR tested?
The tests in this change cover it. The whole set is exercised together
once the last piece lands: every operator run through the engine and
through its generated script, and the two answers compared.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 5)
---------
Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
.github/workflows/build.yml | 65 +++
workflow-compiling-service/build.sbt | 9 +
.../translator/verify/tags/IntegrationTest.java | 56 +++
.../amber/translator/verify/OpExecHarness.scala | 441 +++++++++++++++++++++
4 files changed, 571 insertions(+)
diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml
index f4c041a4fa..c3cfd13be3 100644
--- a/.github/workflows/build.yml
+++ b/.github/workflows/build.yml
@@ -783,6 +783,13 @@ jobs:
env:
JAVA_OPTS: -Xms2048M -Xmx2048M -Xss6M -XX:ReservedCodeCacheSize=256M
-Dfile.encoding=UTF-8
JVM_OPTS: -Xms2048M -Xmx2048M -Xss6M -XX:ReservedCodeCacheSize=256M
-Dfile.encoding=UTF-8
+ # Exclude @IntegrationTest-tagged specs (workflow-compiling-service's
+ # OperatorBehaviorSpec forks Python, which this job does not provision).
+ # Those run in the platform-integration job's workflow-compiling-service
+ # matrix entry, which provisions Python. Read only by
+ # workflow-compiling-service/build.sbt; a no-op for the other services
+ # in this matrix.
+ WCS_TEST_FILTER: skip-integration
services:
# Each platform service transitively depends on DAO, which runs JOOQ
# code generation at compile time and needs the live texera schema.
@@ -972,6 +979,45 @@ jobs:
- uses: coursier/cache-action@95e5b1029b6b86e7bac033ee44a0697d8a527d2d #
v8.1.1
with:
extraSbtFiles: '["*.sbt", "project/**.{scala,sbt}",
"project/build.properties" ]'
+ # --- workflow-compiling-service only: provision Python so the verify
+ # spec (OperatorBehaviorSpec) can fork real interpreters. These four
+ # steps mirror the retired standalone
workflow-compiling-service-integration
+ # job; they are no-ops for every other service in the matrix.
+ - name: Setup Python for Scala-Python verification tests
+ if: ${{ matrix.service == 'workflow-compiling-service' }}
+ uses: actions/setup-python@v7
+ with:
+ python-version: "3.12"
+ - name: Install Python dependencies
+ # OperatorBehaviorSpec forks Python subprocesses that import pandas /
+ # pyarrow / plotly and pytexera (the harness adds amber/src/main/python
+ # to PYTHONPATH). Install amber's runtime deps, same as
amber-integration.
+ # dev-requirements.txt provides the betterproto plugin used by
+ # bin/python-proto-gen.sh in the proto-generation step below.
+ if: ${{ matrix.service == 'workflow-compiling-service' }}
+ run: |
+ python -m pip install uv
+ if [ -f amber/requirements.txt ]; then uv pip install --system
--index-strategy unsafe-best-match -r amber/requirements.txt; fi
+ if [ -f amber/operator-requirements.txt ]; then uv pip install
--system --index-strategy unsafe-best-match -r amber/operator-requirements.txt;
fi
+ if [ -f amber/dev-requirements.txt ]; then uv pip install --system
--index-strategy unsafe-best-match -r amber/dev-requirements.txt; fi
+ - name: Install protoc
+ # Path A forks py_op_driver, which imports pyamber and hence the
+ # generated betterproto bindings in amber/src/main/python/proto
(gitignored,
+ # not checked in). Pin protoc to bin/protoc-version.txt via the
upstream
+ # release zip. Linux-only: this job runs on ubuntu-latest.
+ if: ${{ matrix.service == 'workflow-compiling-service' }}
+ run: |
+ PROTOC_VERSION=$(cat bin/protoc-version.txt)
+ curl -fsSL -o /tmp/protoc.zip
"https://github.com/protocolbuffers/protobuf/releases/download/v${PROTOC_VERSION}/protoc-${PROTOC_VERSION}-linux-x86_64.zip"
+ sudo unzip -o /tmp/protoc.zip -d /usr/local
+ sudo chmod +x /usr/local/bin/protoc
+ sudo chmod -R a+rX /usr/local/include/google
+ - name: Generate Python proto bindings
+ # Regenerate amber/src/main/python/proto so the forked py_op_driver can
+ # import pyamber; without this Path A fails with ImportError on proto
+ # symbols (e.g. ChannelIdentity).
+ if: ${{ matrix.service == 'workflow-compiling-service' }}
+ run: bash bin/python-proto-gen.sh
- name: Create Databases
run: |
psql -h localhost -U postgres -f sql/texera_ddl.sql
@@ -1054,6 +1100,25 @@ jobs:
# smoke-boot's verdict is LISTEN-based, never log-scraping (#6332).
TEXERA_SERVICE_LOG_LEVEL: ${{ runner.debug == '1' && 'DEBUG' ||
'WARN' }}
run: .github/scripts/smoke-boot.sh "/tmp/dists/${{ matrix.service
}}-*/bin/${{ matrix.service }}" "${{ matrix.port }}"
+ - name: Run workflow-compiling-service Python e2e (integration) tests
+ # workflow-compiling-service only: the @IntegrationTest-tagged verify
+ # spec (OperatorBehaviorSpec) forks real Python processes and compares
+ # translator-generated standalone code against platform output,
operator
+ # by operator. Reuses this job's compiled WCS (the dist step above) and
+ # postgres instead of a separate integration job.
+ # WCS_TEST_FILTER=integration-only keeps only @IntegrationTest specs
and
+ # bounds ScalaTest's parallel pool to 4 threads (build.sbt). Every
operator
+ # carrying standalone code runs: the spec's narrowing knobs
+ # (VERIFY_ONLY / VERIFY_SKIP) are unset here, as they are in any run
that
+ # does not ask for less, and what stays withheld is a single variant
or an
+ # operator that cannot run, recorded against an issue in the runner
itself.
+ # UDF_PYTHON_PATH points the harness at the provisioned interpreter
(bare
+ # name resolves via PATH); it auto-locates amber/src/main/python.
+ if: ${{ matrix.service == 'workflow-compiling-service' }}
+ env:
+ WCS_TEST_FILTER: integration-only
+ UDF_PYTHON_PATH: python
+ run: sbt "WorkflowCompilingService/test"
pyamber:
if: ${{ inputs.run_pyamber }}
diff --git a/workflow-compiling-service/build.sbt
b/workflow-compiling-service/build.sbt
index 1c440d14b0..2427f479a1 100644
--- a/workflow-compiling-service/build.sbt
+++ b/workflow-compiling-service/build.sbt
@@ -46,6 +46,15 @@ ThisBuild / conflictManager := ConflictManager.latestRevision
// tests *within* a suite (e.g. OperatorBehaviorSpec) via ScalaTest's own pool.
Global / concurrentRestrictions += Tags.limit(Tags.Test, 1)
+// The fast-unit / integration test split; the selection logic itself is shared
+// in project/TestFilters.scala. The tag it names is the one this change adds,
+// so the two arrive together and the filter never selects on a tag nothing
+// carries.
+Test / testOptions ++= TestFilters.integrationSplit(
+ envVar = "WCS_TEST_FILTER",
+ tag = "org.apache.texera.amber.translator.verify.tags.IntegrationTest"
+)
+
// -P4 bounds ScalaTest's ParallelTestExecution pool, and only this module
wants
// it: OperatorBehaviorSpec forks a Python subprocess per operator, and at
// core-count concurrency (e.g. 12) resource contention caused rare flakes. A
diff --git
a/workflow-compiling-service/src/test/java/org/apache/texera/amber/translator/verify/tags/IntegrationTest.java
b/workflow-compiling-service/src/test/java/org/apache/texera/amber/translator/verify/tags/IntegrationTest.java
new file mode 100644
index 0000000000..4da3aa9cd2
--- /dev/null
+++
b/workflow-compiling-service/src/test/java/org/apache/texera/amber/translator/verify/tags/IntegrationTest.java
@@ -0,0 +1,56 @@
+/*
+ * 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.tags;
+
+import java.lang.annotation.ElementType;
+import java.lang.annotation.Retention;
+import java.lang.annotation.RetentionPolicy;
+import java.lang.annotation.Target;
+
+import org.scalatest.TagAnnotation;
+
+/**
+ * Class-level marker tag for workflow-compiling-service ScalaTest specs that
+ * exercise both Scala and Python end-to-end (they fork a real Python process
+ * to run and compare the translator-generated code). Routing to the
+ * {@code workflow-compiling-service-integration} CI job is by ScalaTest tag
+ * filtering, controlled by the {@code WCS_TEST_FILTER} env var in
+ * {@code workflow-compiling-service/build.sbt}: the lighter
+ * {@code workflow-compiling-service} test run uses {@code skip-integration}
+ * (which passes {@code -l
org.apache.texera.amber.translator.verify.tags.IntegrationTest}
+ * to ScalaTest), and the integration job uses {@code integration-only} (which
+ * passes {@code -n} for the same tag).
+ *
+ * <p>Mirrors amber's {@code org.apache.texera.amber.tags.IntegrationTest}.
Only
+ * {@code OperatorBehaviorSpec} carries this tag today — it is the sole spec
that
+ * spawns Python; the other verify specs only exercise pure-JVM classification
+ * and comparison logic.
+ *
+ * <p>Written in Java rather than Scala because ScalaTest detects tag
+ * annotations via {@code java.lang.annotation} reflection. A Scala
+ * {@code class extends StaticAnnotation} does not produce a JVM annotation
+ * interface that {@code @TagAnnotation} can attach to, so the tag would be
+ * invisible to ScalaTest at runtime.
+ */
+@TagAnnotation
+@Retention(RetentionPolicy.RUNTIME)
+@Target({ElementType.METHOD, ElementType.TYPE})
+public @interface IntegrationTest {
+}
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
new file mode 100644
index 0000000000..01e48e40da
--- /dev/null
+++
b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/verify/OpExecHarness.scala
@@ -0,0 +1,441 @@
+/*
+ * 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.fasterxml.jackson.databind.node.ObjectNode
+import com.typesafe.scalalogging.LazyLogging
+import org.apache.texera.amber.core.executor.{ExecFactory,
OpExecWithClassName, OperatorExecutor}
+import org.apache.texera.amber.core.tuple.{AttributeType, Schema, Tuple,
TupleLike}
+import org.apache.texera.amber.core.virtualidentity.{
+ ExecutionIdentity,
+ PhysicalOpIdentity,
+ WorkflowIdentity
+}
+import org.apache.texera.amber.core.workflow.{PhysicalOp, PhysicalPlan,
PortIdentity}
+import org.apache.texera.amber.operator.LogicalOp
+import org.apache.texera.amber.util.JSONUtils.objectMapper
+
+import java.nio.file.{Files, Path}
+import java.sql.Timestamp
+import java.util.Base64
+import scala.collection.mutable
+import scala.jdk.CollectionConverters._
+
+/**
+ * Generic harness that drives an OpDesc's OpExec(s) directly, bypassing the
+ * Pekko/actor runtime.
+ *
+ * The wiring, meaning how many OpExecs there are, the internal links and the
+ * input-port dependency order, comes from `opDesc.getPhysicalPlan(...)`, so a
+ * join's build and probe, a split's two outputs and a source's empty input
map
+ * all work without a harness change. I/O is JSON Lines with a
+ * `*.jsonl.schema.json` sidecar per file.
+ *
+ * It runs one worker, idx=0 of 1, so nothing here coordinates partitioners
+ * across executors. Python UDFs go to [[PyOpExecHarness]], which has a real
+ * worker to give them.
+ */
+object OpExecHarness extends LazyLogging {
+
+ // Test-only workflow / execution IDs. The values don't matter — the harness
+ // never persists state under them — but the PhysicalOp factory needs *some*
+ // IDs to embed in PhysicalOpIdentity.
+ private val TestWorkflowId = WorkflowIdentity(0L)
+ private val TestExecutionId = ExecutionIdentity(0L)
+
+ /**
+ * @param outputs external output port → JSONL file path
+ * @param outputSchemas same keys as outputs, gives each port's [[Schema]]
+ */
+ final case class Result(
+ outputs: Map[PortIdentity, Path],
+ outputSchemas: Map[PortIdentity, Schema]
+ )
+
+ /**
+ * Run `opDesc` against the given inputs and write the outputs to
`outputDir`.
+ *
+ * @param inputs map keyed by the *external* input port identifier (the one
+ * the user-visible LogicalOp exposes). Each path points to a
+ * `.jsonl` file with a sibling `.schema.json`.
+ * @param outputDir destination directory; created if missing. Output files
+ * are named `output_port_<id>.jsonl` per external output.
+ */
+ def execute(
+ opDesc: LogicalOp,
+ inputs: Map[PortIdentity, Path],
+ outputDir: Path
+ ): Result = {
+ Files.createDirectories(outputDir)
+
+ // 1. Compile OpDesc → PhysicalPlan. For most ops this is a
single-PhysicalOp
+ // plan; HashJoin and other multi-stage ops return multiple PhysicalOps
+ // plus internal PhysicalLinks (e.g. build.out → probe.in0).
+ val plan = opDesc.getPhysicalPlan(TestWorkflowId, TestExecutionId)
+
+ // 2. Identify external input ports. A PhysicalOp port is "external" iff no
+ // PhysicalLink in this plan terminates at it. The user's `inputs` map
+ // must cover exactly these (matched by PortIdentity).
+ val externalInputs: Set[(PhysicalOpIdentity, PortIdentity)] =
+ plan.operators.flatMap { phOp =>
+ phOp.inputPorts.keys.collect {
+ case portId if !plan.links.exists(l => l.toOpId == phOp.id &&
l.toPortId == portId) =>
+ (phOp.id, portId)
+ }
+ }
+ validateInputCoverage(externalInputs, inputs.keySet)
+
+ // 3. Load each input file as (Schema, Iterator[Tuple]). We read schemas
+ // eagerly but tuples lazily (saves memory on large fixtures).
+ val inputSchemas: Map[PortIdentity, Schema] =
+ inputs.map { case (portId, path) => portId ->
TupleIO.readSchemaSidecar(path) }
+ val inputTuples: Map[PortIdentity, () => Iterator[Tuple]] =
+ inputs.map {
+ case (portId, path) =>
+ val schema = inputSchemas(portId)
+ portId -> (() => TupleIO.readTuples(path, schema))
+ }
+
+ // 4. Propagate schemas. We CAN'T use `plan.propagateSchema(inputSchemas)`
+ // directly because it indexes by PortIdentity globally — for HashJoin
+ // probe's internal port 0 would collide with build's external port 0.
+ // Instead, set schemas only on the external (phOpId, portId) pairs and
+ // let `addLink` propagate the internal links' schemas naturally.
+ val planWithSchemas =
+ propagateExternalSchemas(plan, externalInputs, inputSchemas)
+
+ // 5. Identify external output ports symmetrically: no outgoing
PhysicalLink.
+ val externalOutputs: Set[(PhysicalOpIdentity, PortIdentity)] =
+ planWithSchemas.operators.flatMap { phOp =>
+ phOp.outputPorts.keys.collect {
+ case portId
+ if !planWithSchemas.links
+ .exists(l => l.fromOpId == phOp.id && l.fromPortId == portId)
=>
+ (phOp.id, portId)
+ }
+ }
+
+ // 6. Instantiate OpExec per PhysicalOp. We only support
OpExecWithClassName;
+ // fail loudly otherwise so test authors know to mock Python UDFs out.
+ val opExecs: Map[PhysicalOpIdentity, OperatorExecutor] =
+ planWithSchemas.operators.map { phOp =>
+ phOp.opExecInitInfo match {
+ case OpExecWithClassName(className, descString) =>
+ phOp.id -> ExecFactory.newExecFromJavaClassName(
+ className,
+ descString,
+ idx = 0,
+ workerCount = 1
+ )
+ case other =>
+ throw new UnsupportedOperationException(
+ s"OpExecHarness only supports OpExecWithClassName, got: $other"
+ )
+ }
+ }.toMap
+
+ // 7. Drive each PhysicalOp in topological order. Buffer outputs in memory
+ // keyed by (producer phOpId, output port). Downstream PhysicalOps then
+ // consume from these buffers via the plan's internal links.
+ val producedBuffer =
+ mutable.Map.empty[(PhysicalOpIdentity, PortIdentity),
mutable.ArrayBuffer[Tuple]]
+
+ planWithSchemas.topologicalIterator().foreach { phOpId =>
+ val phOp = planWithSchemas.getOperator(phOpId)
+ val opExec = opExecs(phOpId)
+ runOneOp(
+ phOp,
+ opExec,
+ externalInputProvider = portId =>
inputTuples.get(portId).map(_.apply()),
+ upstreamBuffer = producedBuffer,
+ plan = planWithSchemas,
+ produced = producedBuffer
+ )
+ }
+
+ // 8. Materialize external outputs to JSONL with their propagated schemas.
+ val outputPaths = mutable.Map.empty[PortIdentity, Path]
+ val outputSchemas = mutable.Map.empty[PortIdentity, Schema]
+ externalOutputs.foreach {
+ case (phOpId, portId) =>
+ val schema = planWithSchemas
+ .getOperator(phOpId)
+ .outputPorts(portId)
+ ._3
+ .toOption
+ .getOrElse(
+ throw new IllegalStateException(
+ s"Output schema for ($phOpId, $portId) was not propagated"
+ )
+ )
+ val tuples =
+ producedBuffer.getOrElse((phOpId, portId),
mutable.ArrayBuffer.empty[Tuple])
+ val file = outputDir.resolve(s"output_port_${portId.id}.jsonl")
+ TupleIO.writeTuples(file, tuples.iterator, schema)
+ outputPaths(portId) = file
+ outputSchemas(portId) = schema
+ }
+
+ Result(outputPaths.toMap, outputSchemas.toMap)
+ }
+
+ /**
+ * Drives one PhysicalOp's lifecycle: open → input ports in dependency order
+ * (processTupleMultiPort + onFinishMultiPort per port) → close. Source ops
+ * (no input ports) get a single onFinishMultiPort(0) call which gives their
+ * `produceTuple()`-backed implementation a chance to emit.
+ *
+ * Outputs are bucketed by output PortIdentity. `processTupleMultiPort`'s
+ * `Option[PortIdentity]` return: None means port 0 (the default single-
+ * output convention used by the trait's fallback). Multi-output ops like
+ * Split set it explicitly.
+ */
+ private def runOneOp(
+ phOp: PhysicalOp,
+ opExec: OperatorExecutor,
+ externalInputProvider: PortIdentity => Option[Iterator[Tuple]],
+ upstreamBuffer: mutable.Map[
+ (PhysicalOpIdentity, PortIdentity),
+ mutable.ArrayBuffer[Tuple]
+ ],
+ plan: PhysicalPlan,
+ produced: mutable.Map[
+ (PhysicalOpIdentity, PortIdentity),
+ mutable.ArrayBuffer[Tuple]
+ ]
+ ): Unit = {
+ opExec.open()
+ try {
+ def bucket(emitted: Iterator[(TupleLike, Option[PortIdentity])]): Unit =
{
+ emitted.foreach {
+ case (tupleLike, portOpt) =>
+ // Default: the op's single output port. Most operators have one
+ // output and use the trait's default port-0 wrapping, but
+ // multi-stage plans (e.g. HashJoin build) put their internal
+ // output on PortIdentity(0, internal = true) — a hardcoded
+ // PortIdentity(0, false) would NoSuchElementException here.
+ val outPortId = portOpt.getOrElse {
+ if (phOp.outputPorts.size == 1) phOp.outputPorts.keys.head
+ else PortIdentity(0)
+ }
+ val outSchema = phOp
+ .outputPorts(outPortId)
+ ._3
+ .toOption
+ .getOrElse(
+ throw new IllegalStateException(
+ s"Op ${phOp.id} emitted to port $outPortId before its output
schema was propagated"
+ )
+ )
+ val tuple = tupleLike
+ .asInstanceOf[org.apache.texera.amber.core.tuple.SeqTupleLike]
+ .enforceSchema(outSchema)
+ produced
+ .getOrElseUpdate((phOp.id, outPortId),
mutable.ArrayBuffer.empty[Tuple]) += tuple
+ }
+ }
+
+ // Process each input port in declared dependency order (e.g. HashJoin
+ // probe's build-side port must finish before the data-side port starts).
+ val portOrder =
+ if (phOp.getInputPortDependencyPairs.nonEmpty)
+ phOp.getInputPortDependencyPairs
+ else phOp.inputPorts.keys.toList.sortBy(_.id)
+
+ portOrder.foreach { portId =>
+ val tuples: Iterator[Tuple] =
+ if (externalInputProvider(portId).isDefined) {
+ externalInputProvider(portId).get
+ } else {
+ // Internal port: stitch upstream PhysicalLinks' buffers together
+ val upstream = plan.links
+ .filter(l => l.toOpId == phOp.id && l.toPortId == portId)
+ .toList
+ .sortBy(l => (l.fromOpId.toString, l.fromPortId.id))
+ upstream.iterator
+ .flatMap(l =>
+ upstreamBuffer
+ .getOrElse(
+ (l.fromOpId, l.fromPortId),
+ mutable.ArrayBuffer.empty[Tuple]
+ )
+ .iterator
+ )
+ }
+
+ tuples.foreach { t =>
+ bucket(opExec.processTupleMultiPort(t, portId.id))
+ }
+ bucket(opExec.onFinishMultiPort(portId.id))
+ }
+
+ // Source operator: no input ports. Trigger production via
onFinishMultiPort
+ // on a synthetic port 0 — SourceOperatorExecutor.onFinish ignores the
port
+ // and emits everything from produceTuple().
+ if (phOp.inputPorts.isEmpty) {
+ bucket(opExec.onFinishMultiPort(0))
+ }
+ } finally {
+ opExec.close()
+ }
+ }
+
+ // Walks the plan in topo order, propagating schemas only at the truly
external
+ // input ports. Internal ports get their schema via `addLink` (PhysicalPlan
+ // re-applies the source's output schema to the destination port). This
avoids
+ // a collision when multiple PhysicalOps share a PortIdentity (e.g. HashJoin
+ // probe.in0 internal vs build.in0 external both have PortIdentity(0)).
+ // `private[verify]` rather than private: the Python harness runs the same
plan
+ // through a different executor, and how a plan's external ports get their
+ // schemas does not change with the executor behind them.
+ private[verify] def propagateExternalSchemas(
+ plan: PhysicalPlan,
+ externalPorts: Set[(PhysicalOpIdentity, PortIdentity)],
+ schemas: Map[PortIdentity, Schema]
+ ): PhysicalPlan = {
+ var acc = PhysicalPlan(operators = Set.empty, links = Set.empty)
+ plan.topologicalIterator().map(plan.getOperator).foreach { phOp =>
+ val updated = phOp.inputPorts.keys.foldLeft(phOp) { (op, portId) =>
+ if (externalPorts.contains((phOp.id, portId)) &&
schemas.contains(portId)) {
+ op.propagateSchema(Some((portId, schemas(portId))))
+ } else op
+ }
+ // .propagateSchema() with no arg re-fires output derivation if all
inputs
+ // are now resolved (source ops trigger immediately since inputPorts
empty).
+ acc = acc.addOperator(updated.propagateSchema())
+ plan.getUpstreamPhysicalLinks(phOp.id).foreach { link =>
+ acc = acc.addLink(link)
+ }
+ }
+ acc
+ }
+
+ private[verify] def validateInputCoverage(
+ external: Set[(PhysicalOpIdentity, PortIdentity)],
+ provided: Set[PortIdentity]
+ ): Unit = {
+ val expected = external.map(_._2)
+ val missing = expected -- provided
+ val extra = provided -- expected
+ require(
+ missing.isEmpty,
+ s"Missing input fixtures for external ports: $missing (expected
$expected)"
+ )
+ if (extra.nonEmpty) {
+ logger.warn(s"Input fixtures provided for non-external ports (ignored):
$extra")
+ }
+ }
+}
+
+/**
+ * JSON Lines I/O for Tuples. Each `.jsonl` file holds one record per line and
+ * is paired with a `.jsonl.schema.json` sidecar listing [[Attribute]]s in
+ * column order, which is what carries the types a JSON line cannot.
+ *
+ * pandas symmetry: `pd.read_json(path, lines=True)` and
+ * `df.to_json(path, orient='records', lines=True)` round-trip cleanly for the
+ * supported types (STRING / INTEGER / LONG / DOUBLE / BOOLEAN).
+ */
+object TupleIO {
+
+ private def sidecar(path: Path): Path =
+ path.resolveSibling(path.getFileName.toString + ".schema.json")
+
+ def readSchemaSidecar(path: Path): Schema = {
+ val text = new String(Files.readAllBytes(sidecar(path)))
+ objectMapper.readValue(text, classOf[Schema])
+ }
+
+ def readTuples(path: Path, schema: Schema): Iterator[Tuple] = {
+ // readAllLines closes the underlying handle; safer than Files.lines for
+ // test-scale fixtures where memory cost is negligible.
+ val lines = Files.readAllLines(path).asScala
+ lines.iterator.filter(_.trim.nonEmpty).map { line =>
+ val node = objectMapper.readTree(line)
+ val builder = Tuple.builder(schema)
+ schema.getAttributes.foreach { attr =>
+ val fieldNode = node.get(attr.getName)
+ val v: Any =
+ if (fieldNode == null || fieldNode.isNull) null
+ else
+ attr.getType match {
+ case AttributeType.STRING => fieldNode.asText()
+ case AttributeType.INTEGER => Int.box(fieldNode.asInt())
+ case AttributeType.LONG => Long.box(fieldNode.asLong())
+ case AttributeType.DOUBLE => Double.box(fieldNode.asDouble())
+ case AttributeType.BOOLEAN => Boolean.box(fieldNode.asBoolean())
+ case AttributeType.BINARY =>
+ Base64.getDecoder.decode(fieldNode.asText())
+ // 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.
+ case AttributeType.TIMESTAMP =>
+ Timestamp.valueOf(fieldNode.asText())
+ case other =>
+ throw new UnsupportedOperationException(
+ s"TupleIO MVP doesn't support $other yet"
+ )
+ }
+ builder.add(attr, v)
+ }
+ builder.build()
+ }
+ }
+
+ def writeTuples(path: Path, tuples: Iterator[Tuple], schema: Schema): Unit =
{
+ // Sidecar first so a partial main-file write still has a recoverable
schema.
+ Files.write(sidecar(path), objectMapper.writeValueAsBytes(schema))
+ val writer = Files.newBufferedWriter(path)
+ try {
+ tuples.foreach { t =>
+ val node: ObjectNode = objectMapper.createObjectNode()
+ schema.getAttributes.zipWithIndex.foreach {
+ case (attr, idx) =>
+ val v = t.getField[Any](idx)
+ if (v == null) node.putNull(attr.getName)
+ else
+ attr.getType match {
+ case AttributeType.STRING => node.put(attr.getName,
v.toString)
+ case AttributeType.INTEGER => node.put(attr.getName,
v.asInstanceOf[Int])
+ case AttributeType.LONG => node.put(attr.getName,
v.asInstanceOf[Long])
+ case AttributeType.DOUBLE => node.put(attr.getName,
v.asInstanceOf[Double])
+ case AttributeType.BOOLEAN => node.put(attr.getName,
v.asInstanceOf[Boolean])
+ case AttributeType.BINARY =>
+ node.put(
+ attr.getName,
+
Base64.getEncoder.encodeToString(v.asInstanceOf[Array[Byte]])
+ )
+ case AttributeType.TIMESTAMP =>
+ node.put(attr.getName, v.asInstanceOf[Timestamp].toString)
+ case other =>
+ throw new UnsupportedOperationException(
+ s"TupleIO MVP doesn't support $other yet"
+ )
+ }
+ }
+ writer.write(objectMapper.writeValueAsString(node))
+ writer.newLine()
+ }
+ } finally writer.close()
+ }
+}