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-8327-fd09f20ec9cdf17f1c9fe18af735a22d95774d98 in repository https://gitbox.apache.org/repos/asf/texera.git
commit 4d7fd493a8f47b28cc2155bd79c946699cec89a1 Author: Kary Zheng <[email protected]> AuthorDate: Fri Sep 11 04:18:26 2026 +0000 feat(workflow-compiling-service): export a workflow as a standalone Python script (#8327) ### What changes were proposed in this PR? A workflow can be built and read in the editor, but there is no form of it that runs anywhere else. This adds one: given a plan, the compiling service returns a single Python file that reads the same sources, applies the same operators in the same order, and prints its results. An operator says how it reads outside the engine by implementing `StandaloneCodeGenerator`, returning a block of pandas that names its inputs and outputs by position — `in1df`, `in2df`, `out1df`. The translator walks the plan in topological order, gives every output port a variable, substitutes those placeholders for the variables its upstreams were given, and prints the leaves. Union takes the whole list of upstreams rather than a fixed count, since any count an operator states would be wrong for some workflow. An operator that has no generator yet leaves a commented placeholder rather than a line that looks like it works, so the export is useful before every operator implements the trait. Five operators implement it here to show the shape and to give the translator something real to walk: Distinct, Filter, Limit, Projection and Union. The rest of the operator set follows a family at a time. ### Any related issues, documentation, discussions? Part of #8325, 1 of 26; that issue lists the set in order. Closes #8407, the task this change is the whole of. ### How was this PR tested? `WorkflowToPythonTranslatorSpec` covers what the translator does with a plan: the topological order, the variable each port is given, the placeholder substitution, the variadic port, and the operator that has no generator. Each of the five operators asserts the block it emits in its own spec. ### 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]> --- .../amber/operator/StandaloneCodeGenerator.scala | 112 +++++++ workflow-compiling-service/build.sbt | 20 +- .../translator/WorkflowToPythonTranslator.scala | 344 +++++++++++++++++++++ .../WorkflowToPythonTranslatorSpec.scala | 175 +++++++++++ 4 files changed, 650 insertions(+), 1 deletion(-) diff --git a/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala new file mode 100644 index 0000000000..3853d1cdd6 --- /dev/null +++ b/common/workflow-operator/src/main/scala/org/apache/texera/amber/operator/StandaloneCodeGenerator.scala @@ -0,0 +1,112 @@ +/* + * 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.operator + +import org.apache.texera.amber.core.tuple.{AttributeType, Schema} +import org.apache.texera.amber.core.workflow.PortIdentity + +import java.net.URLDecoder +import java.nio.charset.StandardCharsets + +trait StandaloneCodeGenerator { + + /** + * The Python this operator contributes to an exported script. + * + * Frames are named by placeholder rather than outright: `in1df`, `in2df` and + * so on for what the operator reads, `out1df` and so on for what it writes. + * A file the operator writes is named the same way, `outputHtml` or + * `outputJson`. The translator puts the names it assigned in their place. + */ + def generateStandaloneCode(): String + + /** + * The same Python, for an operator that has to know what the columns it reads + * are DECLARED as. See [[renderedAsText]] for what the file cannot carry. + * Defaults to the schema-free form. + */ + def generateStandaloneCode(inputSchemas: Map[PortIdentity, Schema]): String = + generateStandaloneCode() + + /** + * A column as the text the engine's `toString` would have produced. + * + * A hole costs the column its type: pandas reads a holed integer column as a + * float and a holed boolean one as 1.0 and 0.0, so 6 renders as "6.0" and true + * as "1.0". Only the declared type can say which was meant, a real DOUBLE + * holding 6.0 looking the same. `None` reads the column as it arrives. + * + * The caller declares [[StandaloneHelpers.AttributeCasts]] itself; doing it + * here would emit the helper into every script. + */ + protected def renderedAsText(column: String, declared: Option[AttributeType]): String = { + val narrowed = declared match { + case Some(AttributeType.INTEGER) | Some(AttributeType.LONG) => + s"""$column.astype("Int64")""" + case Some(AttributeType.BOOLEAN) => s"""$column.astype("boolean")""" + case _ => column + } + s"""_texera_cast_string($narrowed)""" + } + + /** + * The file's own name, for a script that reads it from its own directory + * rather than through Texera's resolved URI. + * + * Taken from the last path segment instead of by parsing the whole string as a + * URI: the resolver percent-encodes the file-relative segments but leaves the + * repository and version names as the user typed them, so a dataset version + * called `v3 - with long text` makes `new URI` throw on the space and no code + * is generated at all. + */ + protected def sourceBasename(rawPath: String): String = { + val segment = rawPath.split("/").lastOption.getOrElse("") + // Percent-decoding only, matching what `URI.getPath` used to return here: form + // decoding would also turn a literal `+` in a file name into a space. + URLDecoder.decode(segment.replace("+", "%2B"), StandardCharsets.UTF_8) + } + + def producesDataFrame(): Boolean = true + + /** + * Definitions this operator's standalone code depends on, emitted once near + * the top of the script rather than inline. + * + * The translator concatenates operator bodies into a single module, so an + * operator needing a helper class has nowhere to put it that another operator + * would not duplicate. Helpers returned here are collected across the whole + * plan and deduplicated by their text, so two sampling operators in one + * workflow yield one copy of the generator they share. + */ + def standaloneHelpers(): Seq[String] = Seq.empty + + /** + * Modules this operator's standalone code needs, written as the import + * statements themselves, collected across the plan and emitted once at the + * top of the script. + * + * pandas is not named here: the translator emits it for every script, since + * an operator body reads and writes frames whatever else it does. What an + * operator states here is what it needs beyond that, so a script built from + * operators that only reshape a table does not require a plotting library + * to start. + */ + def standaloneImports(): Seq[String] = Seq.empty +} diff --git a/workflow-compiling-service/build.sbt b/workflow-compiling-service/build.sbt index 2af92efc4d..1c440d14b0 100644 --- a/workflow-compiling-service/build.sbt +++ b/workflow-compiling-service/build.sbt @@ -41,9 +41,27 @@ ThisBuild / semanticdbVersion := scalafixSemanticdb.revision // Manage dependency conflicts by always using the latest revision ThisBuild / conflictManager := ConflictManager.latestRevision -// Restrict parallel execution of tests to avoid conflicts +// Restrict parallel execution of tests to avoid conflicts. This caps how many +// test *suites* run concurrently; ParallelTestExecution still parallelizes the +// tests *within* a suite (e.g. OperatorBehaviorSpec) via ScalaTest's own pool. Global / concurrentRestrictions += Tags.limit(Tags.Test, 1) +// -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 +// fixed 4 stays deterministic across machines (incl. CI runners) while still +// running ~3x faster than serial, and it matches PythonWorkerPool's own default +// worker cap so the two bounds agree rather than multiply. Unconditional, so a +// local run reproduces the concurrency CI runs at instead of a faster one that +// flakes differently; WCS_TEST_FILTER selects which tests run, which is a +// separate question from how many run at once. The fast-unit job is unaffected +// either way, since OperatorBehaviorSpec is the only spec here that +// parallelizes and that job excludes it. It lives here rather than in the +// shared helper so that helper stays identical for every module. sbt +// concatenates the ScalaTest arguments of every testOptions entry, so this +// lands in the same argument list as the -n above. +Test / testOptions += Tests.Argument(TestFrameworks.ScalaTest, "-P4") + ///////////////////////////////////////////////////////////////////////////// // Compiler Options ///////////////////////////////////////////////////////////////////////////// diff --git a/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala b/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala new file mode 100644 index 0000000000..8e87340378 --- /dev/null +++ b/workflow-compiling-service/src/main/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslator.scala @@ -0,0 +1,344 @@ +/* + * 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 + +import com.typesafe.scalalogging.LazyLogging +import org.apache.texera.amber.core.tuple.Schema +import org.apache.texera.amber.core.virtualidentity.OperatorIdentity +import org.apache.texera.amber.core.workflow.PortIdentity +import org.apache.texera.common.compiler.model.LogicalPlan +import org.apache.texera.amber.operator.StandaloneCodeGenerator + +import scala.collection.mutable +import scala.collection.mutable.ArrayBuffer +import scala.jdk.CollectionConverters._ + +class WorkflowToPythonTranslator extends LazyLogging { + + // Output-port-level key. An operator with N output ports gets N entries + // (e.g. Split has port 0 and port 1, each with its own assigned dfN var). + private type PortKey = (String, Int) // (opId, portIdx) + + /** + * @param outputSchemas what each operator's output ports carry, as + * [[org.apache.texera.common.compiler.WorkflowCompiler]] + * reports it. Folded into input-port schemas for the + * generators that need a column's declared type, which the + * file the script reads cannot carry. Empty is allowed. + */ + def translate( + logicalPlan: LogicalPlan, + outputSchemas: Map[OperatorIdentity, Map[PortIdentity, Option[Schema]]] = Map.empty + ): String = { + // Track downstream connections per (opId, fromPortIdx). A port is a leaf + // if it has no outgoing edges — operator-level "no outgoing links" is too + // coarse for multi-output ops (Split's port 0 may have downstream while + // port 1 doesn't, or vice versa). + val outgoingFromPort = mutable.Map[PortKey, Int]().withDefaultValue(0) + logicalPlan.links.foreach { link => + outgoingFromPort((link.fromOpId.id, link.fromPortId.id)) += 1 + } + + val outputVar = mutable.Map[PortKey, String]() + var varCounter = 1 + // How many operators of each name have already been given output files. The + // whole plan runs as one program in one directory, so a chart naming its own + // file would be overwritten by the next chart naming the same one. + val fileBaseCounts = mutable.Map[String, Int]() + val script = ArrayBuffer[String]() + + // getTopologicalOpIds() uses jgrapht internally — no need for a custom topo sort + val topoOrder = logicalPlan.getTopologicalOpIds.asScala.toList + + // pandas is the one module every generator uses: an operator body reads and + // writes frames whatever else it does. Everything beyond that is asked of the + // operators in the plan, so a script that draws nothing does not require a + // plotting library to start. + script += "import pandas as pd" + topoOrder + .map(logicalPlan.getOperator) + .collect { case gen: StandaloneCodeGenerator => gen.standaloneImports() } + .flatten + .distinct + .foreach(script += _) + script += "" + + // Helper definitions the operator bodies below refer to. Collected across the + // whole plan and deduplicated by text, so a workflow holding two operators + // that share one helper still emits it once. Order follows the topological + // order, which keeps the script stable for a given plan. + val helpers = topoOrder + .map(logicalPlan.getOperator) + .collect { case gen: StandaloneCodeGenerator => gen.standaloneHelpers() } + .flatten + .distinct + if (helpers.nonEmpty) { + helpers.foreach { helper => script += helper; script += "" } + } + + for (opIdentity <- topoOrder) { + val opId = opIdentity.id + val op = logicalPlan.getOperator(opIdentity) + val displayName = op.operatorInfo.userFriendlyName + + // Resolve upstream inputs in the consuming operator's input-port order + // (link.toPortId), NOT the order links happen to appear in the plan's + // link list. This makes in1df/in2df/... deterministic and correct for + // multi-input operators (joins, set ops) where port 0 vs port 1 carries + // semantics (e.g. build vs probe side). Ties on the same toPortId keep + // link order — relevant for variadic single-port operators like Union. + // Each upstream link is resolved via (fromOpId, fromPortId) so that a + // multi-output upstream (Split) hands each downstream the correct DF. + val inVars = logicalPlan + .getUpstreamLinks(opIdentity) + .sortBy(link => (link.toPortId.id, link.toPortId.internal)) + .map(link => outputVar((link.fromOpId.id, link.fromPortId.id))) + + // Allocate one dfN per declared output port. Existing single-output + // operators have outputPorts.size == 1, so they get exactly one var and + // their behavior is identical to the previous flat scheme. + val outVars = op.operatorInfo.outputPorts.map { port => + val v = s"df$varCounter" + varCounter += 1 + outputVar((opId, port.id.id)) = v + v + } + + script += s"# [$displayName]" + + // Jackson deserializes each operator into its concrete subclass via @JsonSubTypes on LogicalOp, + // so the pattern match below will resolve to the correct descriptor (e.g. BarChartOpDesc). + op match { + case gen: StandaloneCodeGenerator => + // Each upstream link carries its source port's schema to the port it + // arrives at. An unresolved source is left out rather than guessed at. + val inputSchemas = logicalPlan + .getUpstreamLinks(opIdentity) + .flatMap { link => + outputSchemas + .get(link.fromOpId) + .flatMap(_.get(link.fromPortId)) + .flatten + .map(link.toPortId -> _) + } + .toMap + + // generateStandaloneCode() returns a code block using in{N}df / out{N}df + // placeholders; substituteVars() replaces them with the assigned vars. + script += substituteVars( + gen.generateStandaloneCode(inputSchemas), + inVars, + outVars, + fileBase(displayName, fileBaseCounts), + displayName + ) + + case _ => + logger.warn( + s"Operator '$displayName' does not implement StandaloneCodeGenerator. Skipping." + ) + script += s"# TODO: '$displayName' is not yet supported by the translator." + outVars.zipWithIndex.foreach { + case (v, i) => script += s"# $v = <output port $i of $displayName>" + } + } + + script += "" + } + + // Leaf detection runs at the port level: a (opId, port) pair is a leaf + // if no link consumes it. For Split with one downstream port and one + // dangling port, only the dangling port is treated as a leaf to print. + val leafPorts = outputVar.keys.toList + .sortBy { case (_, portIdx) => portIdx } + .filter(key => outgoingFromPort(key) == 0) + val dataFrameLeafPorts = leafPorts.filter { + case (opId, _) => + logicalPlan.getOperator(OperatorIdentity(opId)) match { + case gen: StandaloneCodeGenerator => gen.producesDataFrame() + case _ => false + } + } + + if (dataFrameLeafPorts.nonEmpty) { + script += "# --- Output ---" + // Print in topological order of the producing operator so multi-port + // operators print contiguously and the order matches the script flow. + val topoIndex = topoOrder.map(_.id).zipWithIndex.toMap + dataFrameLeafPorts + .sortBy { case (opId, portIdx) => (topoIndex.getOrElse(opId, Int.MaxValue), portIdx) } + .foreach { + case (opId, portIdx) => + val varName = outputVar((opId, portIdx)) + val displayName = + logicalPlan.getOperator(OperatorIdentity(opId)).operatorInfo.userFriendlyName + val portSuffix = if (outputVar.keys.count(_._1 == opId) > 1) s" port $portIdx" else "" + script += s"""print("\\n[$displayName$portSuffix] $varName:")""" + // The frame itself rather than head(): pandas already elides the + // middle of a long one, and it states the row and column count, + // which head() hides. + script += s"print($varName)" + script += "" + } + } + + script.mkString("\n") + } + + // The stem of the files one operator writes, taken from its name so a reader + // can tell whose picture is whose. Numbered from one even when the plan holds + // a single chart, so the name a script writes does not depend on what else + // the plan happens to contain. + private def fileBase(displayName: String, counts: mutable.Map[String, Int]): String = { + val slug = displayName.toLowerCase.replaceAll("[^a-z0-9]+", "_").replaceAll("^_|_$", "") + val stem = if (slug.isEmpty) "output" else slug + val n = counts.getOrElse(stem, 0) + 1 + counts(stem) = n + s"${stem}_$n" + } + + // Replaces in{N}df / out{N}df placeholders with concrete variable names. + // Substitutes in reverse index order to prevent partial matches (e.g. in1df + // inside in10df). Only the code parts are rewritten: a generator writes a + // column name as a string literal, and a column may be named `in1df`. + // After substitution, scans the code for any leftover placeholders and logs + // a warning — that signals a mismatch between an operator's declared port + // count and what its generateStandaloneCode actually emits. + private def substituteVars( + code: String, + inVars: List[String], + outVars: List[String], + fileBase: String, + displayName: String + ): String = { + def substitute(fragment: String): String = { + var result = fragment + + // An operator that writes a file names it outputHtml or outputJson and + // gets back a name of its own, for the reason given where fileBaseCounts + // is declared. The stem holds letters, digits and underscores only, so it + // carries nothing replaceAll would read as a group reference. + result = result.replaceAll("""\boutputHtml\b""", "\"" + fileBase + ".html\"") + result = result.replaceAll("""\boutputJson\b""", "\"" + fileBase + ".json\"") + + // A variadic port takes as many upstream links as the user draws, and an + // operator reading one cannot name them: `in1df`/`in2df` state a count, and + // whichever count it states is wrong for every other workflow. This one + // placeholder becomes the whole list, so the operator writes the same line + // whether it is fed one table or five. + result = result.replaceAll("""\binAlldf\b""", inVars.mkString("[", ", ", "]")) + inVars.zipWithIndex.reverse.foreach { + case (v, idx) => result = result.replaceAll(s"\\bin${idx + 1}df\\b", v) + } + outVars.zipWithIndex.reverse.foreach { + case (v, idx) => result = result.replaceAll(s"\\bout${idx + 1}df\\b", v) + } + result + } + + val substituted = splitOffLiterals(code).map { + case (isCode, text) => (isCode, if (isCode) substitute(text) else text) + } + + // The code segments are joined by a newline for the scan, so the text + // either side of a skipped literal cannot spell a placeholder nobody wrote. + val scanned = substituted.collect { case (true, text) => text }.mkString("\n") + val leftoverIn = """\bin\d+df\b""".r.findAllIn(scanned).toSet + val leftoverOut = """\bout\d+df\b""".r.findAllIn(scanned).toSet + if (leftoverIn.nonEmpty || leftoverOut.nonEmpty) { + logger.warn( + s"Operator '$displayName' emitted placeholders that don't match its port " + + s"count: leftover inputs=$leftoverIn, leftover outputs=$leftoverOut. " + + s"Generated script will reference unbound variables." + ) + } + + substituted.map(_._2).mkString + } + + // Splits a block into (isCode, text) segments, where a plain string literal + // and a comment are not code and everything else is. An f-string counts as + // code: its braces hold expressions, and a generator that renders a column + // name renders it as a plain literal. + private def splitOffLiterals(code: String): List[(Boolean, String)] = { + val segments = ArrayBuffer[(Boolean, String)]() + val pending = new StringBuilder + var i = 0 + + def flushCode(): Unit = { + if (pending.nonEmpty) { + segments += ((true, pending.toString)) + pending.setLength(0) + } + } + + while (i < code.length) { + val c = code.charAt(i) + if (c == '#') { + flushCode() + val newline = code.indexOf('\n', i) + val end = if (newline < 0) code.length else newline + segments += ((false, code.substring(i, end))) + i = end + } else if (c == '\'' || c == '"') { + val triple = c.toString * 3 + val delim = if (code.startsWith(triple, i)) triple else c.toString + val end = endOfLiteral(code, i + delim.length, delim) + if (isFormatted(code, i)) pending ++= code.substring(i, end) + else { + flushCode() + segments += ((false, code.substring(i, end))) + } + i = end + } else { + pending += c + i += 1 + } + } + flushCode() + segments.toList + } + + // Whether the quote at `quoteIdx` opens an f-string. Its prefix is whatever + // letters run up to it; nothing else may sit against a quote in Python. + private def isFormatted(code: String, quoteIdx: Int): Boolean = { + var j = quoteIdx + var formatted = false + while (j > 0 && code.charAt(j - 1).isLetter) { + j -= 1 + if (code.charAt(j) == 'f' || code.charAt(j) == 'F') formatted = true + } + formatted + } + + // The index just past the closing delimiter, or the end of the block if the + // literal is never closed. A backslash escapes the next character even in a + // raw string, where it still keeps the quote from closing the literal. + private def endOfLiteral(code: String, from: Int, delim: String): Int = { + var j = from + var end = -1 + while (end < 0 && j < code.length) { + if (code.charAt(j) == '\\') j += 2 + else if (code.startsWith(delim, j)) end = j + delim.length + else j += 1 + } + if (end < 0) code.length else end + } +} diff --git a/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala new file mode 100644 index 0000000000..ba0c3f3ae7 --- /dev/null +++ b/workflow-compiling-service/src/test/scala/org/apache/texera/amber/translator/WorkflowToPythonTranslatorSpec.scala @@ -0,0 +1,175 @@ +/* + * 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 + +import org.apache.texera.amber.core.virtualidentity.{ExecutionIdentity, WorkflowIdentity} +import org.apache.texera.amber.core.workflow.{InputPort, OutputPort, PhysicalOp, PortIdentity} +import org.apache.texera.amber.operator.metadata.{OperatorGroupConstants, OperatorInfo} +import org.apache.texera.amber.operator.{LogicalOp, StandaloneCodeGenerator} +import org.apache.texera.common.compiler.model.{LogicalLink, LogicalPlan} +import org.scalatest.flatspec.AnyFlatSpec +import org.scalatest.matchers.should.Matchers + +/** The placeholder substitution, which is where an operator's generated code + * meets the variables the script actually binds. A variadic port is the case + * the numbered placeholders cannot state, so it is the case worth pinning. + * + * Every operator here is a stub. What the translator does is place a block and + * bind the variables around it, so a stub that says exactly which block to + * place keeps these assertions off any real operator's emitted text, which + * would otherwise turn a change to that operator into a failure here. + */ +class WorkflowToPythonTranslatorSpec extends AnyFlatSpec with Matchers { + + private class StubOp(block: String) extends LogicalOp with StandaloneCodeGenerator { + override def getPhysicalOp( + workflowId: WorkflowIdentity, + executionId: ExecutionIdentity + ): PhysicalOp = + throw new UnsupportedOperationException("the translator never builds a physical op") + + override def operatorInfo: OperatorInfo = + OperatorInfo( + "Stub", + "Stands in for an operator that implements the trait", + OperatorGroupConstants.UTILITY_GROUP, + inputPorts = List(InputPort()), + outputPorts = List(OutputPort()) + ) + + override def generateStandaloneCode(): String = block + } + + private def stub(id: String, block: String): LogicalOp = { + val op = new StubOp(block) + op.setOperatorId(id) + op + } + + private def upstream(id: String): LogicalOp = stub(id, "out1df = in1df.copy()") + + private def link(from: LogicalOp, to: LogicalOp): LogicalLink = + LogicalLink(from.operatorIdentifier, PortIdentity(0), to.operatorIdentifier, PortIdentity(0)) + + /** `n` upstreams, all drawn into one port, which is what a variadic port looks + * like in a plan. + */ + private def variadicOf(n: Int): String = { + val sink = stub("sink", "out1df = pd.concat(inAlldf, ignore_index=True)") + val ups = (1 to n).map(i => upstream(s"up$i")) + new WorkflowToPythonTranslator().translate( + LogicalPlan(ups.toList :+ sink, ups.map(link(_, sink)).toList) + ) + } + + "WorkflowToPythonTranslator" should "hand a variadic port every upstream it was drawn" in { + variadicOf(3) should include("pd.concat([df1, df2, df3], ignore_index=True)") + } + + it should "hand a variadic port a one-element list when only one link is drawn" in { + // The case the old fixed `[in1df, in2df]` got wrong in the other direction: + // it named a second frame the script never bound. + variadicOf(1) should include("pd.concat([df1], ignore_index=True)") + } + + it should "leave no placeholder behind for a variadic port" in { + variadicOf(2) should not include "inAlldf" + } + + // head() shows five rows and does not say how many there were, so a script whose + // leaf holds more reads as if that were the whole answer. + it should "print the leaf frame rather than its first rows" in { + val script = variadicOf(2) + script should include("print(df3)") + script should not include ".head())" + } + + // A script that only reshapes a table should run wherever pandas is installed, + // so an import no operator in the plan asked for must not be in the header. + it should "import pandas alone for a plan that asks for nothing else" in { + val script = variadicOf(2) + script should include("import pandas as pd") + script should not include "import plotly" + } + + // Two operators naming the same module yield one import, the way two operators + // sharing one helper yield one copy of it. + it should "emit an operator's declared import once per plan" in { + val ops = List("a", "b").map { id => + val op = new StubOp("out1df = in1df.copy()") { + override def standaloneImports(): Seq[String] = Seq("import numpy as np") + } + op.setOperatorId(id) + op + } + val script = new WorkflowToPythonTranslator().translate(LogicalPlan(ops, List.empty)) + script.linesIterator.count(_ == "import numpy as np") shouldBe 1 + } + + it should "still resolve a numbered placeholder against its own upstream" in { + // The variadic form is an addition, not a replacement: a chain of ordinary + // single-input operators has to keep reading `in1df` as its predecessor. + val first = upstream("first") + val second = upstream("second") + val script = new WorkflowToPythonTranslator().translate( + LogicalPlan(List(first, second), List(link(first, second))) + ) + script should include("df2 = df1.copy()") + } + + /** Nothing stops a column from being named after a placeholder. The + * substitution rewrites the variable a block reads with, never the column + * name it asks that variable for. + */ + it should "leave a column named after a placeholder alone" in { + val source = upstream("source") + val reader = stub("reader", """out1df = in1df[["in1df"]].copy()""") + val script = new WorkflowToPythonTranslator().translate( + LogicalPlan(List(source, reader), List(link(source, reader))) + ) + script should include("""df2 = df1[["in1df"]].copy()""") + } + + /** Two operators that write a file write two of them. The plan runs as one + * program in one directory, so a name either of them had chosen for itself + * would leave one picture where the workflow drew two. + */ + it should "give each operator writing a file a name of its own" in { + val ops = List("a", "b").map { id => + val op = new StubOp("fig.write_json(outputJson)\nfig.write_html(outputHtml)") + op.setOperatorId(id) + op + } + val script = new WorkflowToPythonTranslator().translate(LogicalPlan(ops, List.empty)) + script should include("""fig.write_json("stub_1.json")""") + script should include("""fig.write_html("stub_1.html")""") + script should include("""fig.write_html("stub_2.html")""") + } + + /** The translator's own contract when it meets an operator it cannot render: + * a comment rather than a silently wrong line. + */ + it should "leave a TODO for an operator with no standalone code generator" in { + val op = new org.apache.texera.amber.operator.udf.python.PythonUDFOpDescV2 + op.setOperatorId("udf") + val script = new WorkflowToPythonTranslator().translate(LogicalPlan(List(op), List.empty)) + script should include("# TODO:") + } +}
