This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-8356-5042d96ec85d98ed18bde841b2e74922c43f5c3b
in repository https://gitbox.apache.org/repos/asf/texera.git

commit 0449b467b0fbf0314dd307dac1a5b4f61d5b6441
Author: Kary Zheng <[email protected]>
AuthorDate: Mon Sep 14 19:26:45 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 c0107997b3..0614d54ae0 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
@@ -1053,6 +1099,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()
+  }
+}

Reply via email to