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 91a701b7c1 fix(amber): carry the loop StateFrame envelope through JVM
operators (#6661)
91a701b7c1 is described below
commit 91a701b7c1dc16f3a896971ddef4c1859a50e51f
Author: Xinyuan Lin <[email protected]>
AuthorDate: Tue Jul 21 21:04:45 2026 -0700
fix(amber): carry the loop StateFrame envelope through JVM operators (#6661)
### What changes were proposed in this PR?
A loop with a Java/Scala built-in operator in its body could never
iterate — the Loop End worker crashed on the first iteration:
```
RuntimeError: no loop-back state URI configured for LoopStart ''
```
**Root cause.** Per-iteration loop state rides a `StateFrame` envelope —
`loop_counter` (nesting-depth countdown) and `loop_start_id` (the
back-jump target) — materialized as their own columns next to `content`
(#5900). The Python worker carries the envelope on every hop; the Scala
side stripped it at every seam, so any JVM hop re-emitted states with
the "no loop" defaults `(0, "")`:
```
LoopStart ──(0, LS-id)──▶ Limit(Scala) ──(0, "")──▶ LoopEnd
▲ envelope stripped │
▼
back-jump lookup by "" fails
```
In a nested loop the zeroed `loop_counter` additionally makes the inner
Loop End mis-consume the enclosing loop's state instead of passing it
through.
**Fix.** Loop operators are Python-only, so a JVM operator only ever
needs to *carry the envelope through unchanged* (the +1/−1 bookkeeping
lives in the Python runtime). The columns already exist in the
materialized State schema and on the Arrow flight wire — no format
change, purely Scala-side plumbing:
| Seam | Change |
|---|---|
| `StateFrame` (`DataPayload.scala`) | gains `loopCounter` /
`loopStartId` fields (defaults `0` / `""`, mirroring Python's
`StateFrame`) |
| `State` (`State.scala`) | `loopCounterFrom` / `loopStartIdFrom`
extractors read the envelope columns back off a row |
| `InputPortMaterializationReaderThread` | rebuilds the envelope when
replaying materialized states |
| `DataProcessor` → `OutputManager` | the incoming envelope is threaded
through `processInputState` to `emitState` /
`saveStateToStorageIfNeeded` / `sendState` |
| `PythonProxyClient` / `PythonProxyServer` | envelope preserved across
the Arrow flight bridge in both directions |
Scala-originated states (start/end-channel handlers) keep the defaults —
correct, since a JVM operator can never be a Loop Start.
**Size:** the fix itself is **+38 / −17 lines of code** (excluding
comments/blanks; +78 / −20 gross) across the 8 source files above; the
remaining **+255 lines** are test cases (4 unit specs + 4 e2e
workflows).
### Any related issues, documentation, discussions?
Closes #6660
### How was this PR tested?
- **Unit** — each seam is pinned locally: `StateSpec` (envelope column
round-trip + no-loop defaults), `DataProcessorSpec` (a state
pass-through emits the exact incoming envelope, not just *a*
`StateFrame`), `OutputManagerSpec` (`emitState` stamps the envelope onto
every buffer), `NetworkOutputBufferSpec` (`sendState` stamps the frame).
- **E2E** — four new `LoopIntegrationSpec` cases run real workflows with
JVM operators in the loop body:
1. single loop with `Limit` (`TextInput → LoopStart → Limit → LoopEnd`,
asserts 3 accumulated iterations — reproduces the reported failure
verbatim, fails against the pre-fix code);
2. nested 3×3 loop with `Limit` in the inner body (asserts 9 outer / 3
inner rows — pins that the counter *magnitude* survives the JVM hop, not
just the id);
3. single loop with a **chain of two JVM operators** (`Limit →
Sleep(0)`) — pins the JVM→JVM state handoff (one Scala worker writes the
envelope columns, another reads them back), which the single-op cases
never exercise;
4. nested 3×3 loop with the **chain inside the inner body** — the
strongest combination: nested (counter magnitude) × multi-hop × both
operator kinds.
Manually tested with the this workflow:
[loop.json](https://github.com/user-attachments/files/30251886/loop.json)
<img width="570" height="180" alt="Screenshot 2026-07-21 at 20 56 23"
src="https://github.com/user-attachments/assets/0c002e9e-694e-445b-9ca7-b95a2cb405ae"
/>
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Fable 5)
---
.../messaginglayer/OutputManager.scala | 24 +++-
.../pythonworker/PythonProxyClient.scala | 7 +-
.../pythonworker/PythonProxyServer.scala | 9 +-
.../sendsemantics/partitioners/Partitioner.scala | 4 +-
.../engine/architecture/worker/DataProcessor.scala | 17 ++-
.../InputPortMaterializationReaderThread.scala | 12 +-
.../engine/common/ambermessage/DataPayload.scala | 13 +-
.../amber/engine/e2e/LoopIntegrationSpec.scala | 156 +++++++++++++++++++++
.../messaginglayer/OutputManagerSpec.scala | 23 +++
.../partitioners/NetworkOutputBufferSpec.scala | 12 ++
.../architecture/worker/DataProcessorSpec.scala | 47 +++++++
.../org/apache/texera/amber/core/state/State.scala | 12 +-
.../apache/texera/amber/core/state/StateSpec.scala | 17 +++
13 files changed, 333 insertions(+), 20 deletions(-)
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManager.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManager.scala
index b594cc8ab0..b15e899254 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManager.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManager.scala
@@ -192,9 +192,16 @@ class OutputManager(
buffersToFlush.foreach(_.flush())
}
- def emitState(state: State): Unit = {
- networkOutputBuffers.foreach(kv => kv._2.sendState(state))
- saveStateToStorageIfNeeded(state)
+ /**
+ * Emit a State to every network buffer and (if configured) the state
+ * storage. `loopCounter` / `loopStartId` are the loop envelope riding
+ * alongside the State (see `StateFrame`); a JVM hop inside a loop body
+ * passes the incoming envelope through unchanged, while a Scala-originated
+ * state (start/end-channel handlers) uses the "no loop" defaults.
+ */
+ def emitState(state: State, loopCounter: Long = 0L, loopStartId: String =
""): Unit = {
+ networkOutputBuffers.foreach(kv => kv._2.sendState(state, loopCounter,
loopStartId))
+ saveStateToStorageIfNeeded(state, loopCounter, loopStartId)
}
def addPort(portId: PortIdentity, schema: Schema, storageURIBaseOption:
Option[URI]): Unit = {
@@ -236,13 +243,18 @@ class OutputManager(
})
}
- private def saveStateToStorageIfNeeded(state: State): Unit = {
+ private def saveStateToStorageIfNeeded(
+ state: State,
+ loopCounter: Long,
+ loopStartId: String
+ ): Unit = {
// The same state row is fanned out to every output port's state
// table. This mirrors the broadcast-to-all-workers behavior on the
// emit side: state is shared context, not per-key data, so every
// downstream operator (and every worker reading the materialization)
- // needs the full set.
- stateWriterThreads.values.foreach(_.queue.put(Left(state.toTuple())))
+ // needs the full set. The loop envelope is materialized as its own
+ // columns so the downstream reader can rebuild it.
+
stateWriterThreads.values.foreach(_.queue.put(Left(state.toTuple(loopCounter,
loopStartId))))
}
/**
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyClient.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyClient.scala
index 144c3ac57a..d3ea86cdd4 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyClient.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyClient.scala
@@ -124,8 +124,11 @@ class PythonProxyClient(portNumberPromise: Promise[Int],
val actorId: ActorVirtu
dataPayload match {
case DataFrame(frame) =>
writeArrowStream(mutable.Queue(ArraySeq.unsafeWrapArray(frame): _*),
from, "Data")
- case StateFrame(state) =>
- writeArrowStream(mutable.Queue(state.toTuple()), from, "State")
+ case StateFrame(state, loopCounter, loopStartId) =>
+ // The Arrow wire format for states IS the State row (content +
+ // loop_counter + loop_start_id), so the envelope rides its own
+ // columns to the Python worker -- see network_receiver.py.
+ writeArrowStream(mutable.Queue(state.toTuple(loopCounter,
loopStartId)), from, "State")
}
}
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyServer.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyServer.scala
index 2ff866365b..82b5cea85e 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyServer.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/pythonworker/PythonProxyServer.scala
@@ -128,7 +128,14 @@ private class AmberProducer(
dataHeader.payloadType match {
case "State" =>
assert(root.getRowCount == 1)
- outputPort.sendTo(to,
StateFrame(State.fromTuple(ArrowUtils.getTexeraTuple(0, root))))
+ // The row is a full State row (content + loop_counter + loop_start_id;
+ // see network_sender.py): rebuild the loop envelope alongside the
+ // content so it survives the hop through this JVM proxy.
+ val row = ArrowUtils.getTexeraTuple(0, root)
+ outputPort.sendTo(
+ to,
+ StateFrame(State.fromTuple(row), State.loopCounterFrom(row),
State.loopStartIdFrom(row))
+ )
case "ECM" =>
assert(root.getRowCount == 1)
outputPort.sendTo(
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/Partitioner.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/Partitioner.scala
index 39065ca693..620eb93651 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/Partitioner.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/Partitioner.scala
@@ -49,9 +49,9 @@ class NetworkOutputBuffer(
}
}
- def sendState(state: State): Unit = {
+ def sendState(state: State, loopCounter: Long = 0L, loopStartId: String =
""): Unit = {
flush()
- dataOutputPort.sendTo(to, StateFrame(state))
+ dataOutputPort.sendTo(to, StateFrame(state, loopCounter, loopStartId))
flush()
}
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessor.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessor.scala
index e88b001610..a86d16c926 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessor.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessor.scala
@@ -121,11 +121,20 @@ class DataProcessor(
}
}
- private[this] def processInputState(state: State, port: Int): Unit = {
+ private[this] def processInputState(
+ state: State,
+ port: Int,
+ loopCounter: Long,
+ loopStartId: String
+ ): Unit = {
try {
val outputState = executor.processState(state, port)
if (outputState.isDefined) {
- outputManager.emitState(outputState.get)
+ // Carry the incoming loop envelope through unchanged: loop operators
+ // are Python-only, so a JVM operator inside a loop body only ever
+ // forwards the envelope (the +1/-1 bookkeeping lives in the Python
+ // runtime's _process_state_frame).
+ outputManager.emitState(outputState.get, loopCounter, loopStartId)
}
} catch safely {
case e =>
@@ -220,8 +229,8 @@ class DataProcessor(
)
inputManager.initBatch(channelId, tuples)
processInputTuple(inputManager.getNextTuple)
- case StateFrame(state) =>
- processInputState(state, portId.id)
+ case StateFrame(state, loopCounter, loopStartId) =>
+ processInputState(state, portId.id, loopCounter, loopStartId)
}
statisticsManager.increaseDataProcessingTime(System.nanoTime() -
dataProcessingStartTime)
}
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/managers/InputPortMaterializationReaderThread.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/managers/InputPortMaterializationReaderThread.scala
index 2f4386c1d8..f1ece04b5a 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/managers/InputPortMaterializationReaderThread.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/worker/managers/InputPortMaterializationReaderThread.scala
@@ -104,9 +104,17 @@ class InputPortMaterializationReaderThread(
.asInstanceOf[VirtualDocument[Tuple]]
val stateReadIterator = stateDocument.get()
while (stateReadIterator.hasNext) {
- val state = State.fromTuple(stateReadIterator.next())
+ val row = stateReadIterator.next()
+ // Rebuild the loop envelope (loop_counter / loop_start_id) alongside
+ // the content: a JVM operator inside a loop body must carry it through
+ // unchanged, or the matching LoopEnd loses its back-jump target.
+ val stateFrame = StateFrame(
+ State.fromTuple(row),
+ State.loopCounterFrom(row),
+ State.loopStartIdFrom(row)
+ )
inputMessageQueue.put(
- FIFOMessageElement(WorkflowFIFOMessage(channelId, getSequenceNumber,
StateFrame(state)))
+ FIFOMessageElement(WorkflowFIFOMessage(channelId, getSequenceNumber,
stateFrame))
)
}
diff --git
a/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/DataPayload.scala
b/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/DataPayload.scala
index 54f577a0be..b8ac640ad5 100644
---
a/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/DataPayload.scala
+++
b/amber/src/main/scala/org/apache/texera/amber/engine/common/ambermessage/DataPayload.scala
@@ -24,7 +24,18 @@ import org.apache.texera.amber.core.tuple.Tuple
sealed trait DataPayload extends WorkflowFIFOMessagePayload {}
-final case class StateFrame(frame: State) extends DataPayload
+/**
+ * A single State travelling between operators. `loopCounter` / `loopStartId`
+ * are the loop envelope owned by the (Python) worker runtime -- carried
+ * alongside the State payload (not inside it) so it never collides with user
+ * state, and materialized/transported as their own columns parallel to the
+ * content (see `State.toTuple`). Loop operators are Python-only, so a JVM
+ * operator inside a loop body only ever carries the envelope through
+ * unchanged; the defaults are the "no loop" values for all non-loop state.
+ * Mirrors the Python `StateFrame` (core/models/payload.py).
+ */
+final case class StateFrame(frame: State, loopCounter: Long = 0L, loopStartId:
String = "")
+ extends DataPayload
final case class DataFrame(frame: Array[Tuple]) extends DataPayload {
val inMemSize: Long = {
diff --git
a/amber/src/test/integration/org/apache/texera/amber/engine/e2e/LoopIntegrationSpec.scala
b/amber/src/test/integration/org/apache/texera/amber/engine/e2e/LoopIntegrationSpec.scala
index aa17d6c7a8..2c81f78045 100644
---
a/amber/src/test/integration/org/apache/texera/amber/engine/e2e/LoopIntegrationSpec.scala
+++
b/amber/src/test/integration/org/apache/texera/amber/engine/e2e/LoopIntegrationSpec.scala
@@ -41,7 +41,9 @@ import org.apache.texera.amber.engine.e2e.TestUtils.{
workflowContext
}
import org.apache.texera.amber.operator.LogicalOp
+import org.apache.texera.amber.operator.limit.LimitOpDesc
import org.apache.texera.amber.operator.loop.{LoopEndOpDesc, LoopStartOpDesc}
+import org.apache.texera.amber.operator.sleep.SleepOpDesc
import org.apache.texera.amber.operator.source.scan.text.TextInputSourceOpDesc
import org.apache.texera.amber.tags.IntegrationTest
import org.apache.texera.workflow.LogicalLink
@@ -162,6 +164,26 @@ class LoopIntegrationSpec
op
}
+ private def limit(n: Int): LimitOpDesc = {
+ val op = new LimitOpDesc()
+ op.limit = n
+ op
+ }
+
+ private def sleep(seconds: Int): SleepOpDesc = {
+ val op = new SleepOpDesc()
+ op.sleepTime = seconds
+ op
+ }
+
+ // NOTE: a JavaUDF (runtime-compiled Java) case cannot run in this suite.
+ // The integration tests run forkless under sbt (fork := false), where
+ // `javax.tools` sees only the sbt launcher on `java.class.path` -- the
+ // Texera classes live in sbt's layered classloaders, so
+ // JavaRuntimeCompilation fails with "package ... does not exist" before
+ // any loop logic runs. Production workers launch with a real classpath,
+ // so this is a harness limitation, not an engine one.
+
private def link(from: LogicalOp, to: LogicalOp): LogicalLink =
LogicalLink(from.operatorIdentifier, PortIdentity(),
to.operatorIdentifier, PortIdentity())
@@ -231,4 +253,138 @@ class LoopIntegrationSpec
s"expected 3, got $innerRows (all: $materialized)"
)
}
+
+ it should "run a loop whose body contains a JVM (Scala) operator" in {
+ // TextInput -> LoopStart -> Limit (a Scala built-in) -> LoopEnd.
+ //
+ // The loop envelope (loop_counter / loop_start_id) must survive the hop
+ // through a JVM operator: the Scala worker replays the materialized state,
+ // forwards it through the default processState, and re-materializes it for
+ // the LoopEnd. Before the envelope was carried through on the Scala side
+ // (StateFrame fields + reader/emit/save plumbing), this hop zeroed the
+ // counter and blanked the id, so the LoopEnd captured loop_start_id = ""
+ // and the back-jump failed with "no loop-back state URI configured for
+ // LoopStart ''" -- the loop never iterated. Limit(10) passes every row
+ // through, so the loop semantics are identical to the single-loop test.
+ val src = textInput("1\n2\n3")
+ val start = loopStart("i = 0", "table.iloc[i]")
+ val mid = limit(10)
+ val end = loopEnd("i += 1", "i < len(table)")
+ val materialized = runAndGetMaterializedRowCounts(
+ List(src, start, mid, end),
+ List(link(src, start), link(start, mid), link(mid, end))
+ )
+ val endRows = materialized.getOrElse(end.operatorIdentifier, -1L)
+ assert(
+ endRows == 3,
+ s"LoopEnd must accumulate all 3 iterations with a Scala operator in the
" +
+ s"loop body: expected 3, got $endRows (all: $materialized)"
+ )
+ }
+
+ it should "run a nested loop whose inner body contains a JVM (Scala)
operator" in {
+ // TextInput -> OuterStart -> InnerStart -> Limit -> InnerEnd -> OuterEnd.
+ //
+ // The nested variant of the JVM-hop test: the OUTER loop's state crosses
+ // the Scala operator stamped (loop_counter = 1, loop_start_id = outer), so
+ // this pins that the JVM hop preserves the counter MAGNITUDE too, not just
+ // the id. A zeroed counter would make the inner LoopEnd mis-consume the
+ // outer state (KeyError on the missing 'table' key / clobbered jump id)
+ // instead of passing it through with -1.
+ val src = textInput("1\n2\n3")
+ val outerStart = loopStart("i = 0", "table")
+ val innerStart = loopStart("j = 0", "table.iloc[j]")
+ val mid = limit(10)
+ val innerEnd = loopEnd("j += 1", "j < len(table)")
+ val outerEnd = loopEnd("i += 1", "i < len(table)")
+ val materialized = runAndGetMaterializedRowCounts(
+ List(src, outerStart, innerStart, mid, innerEnd, outerEnd),
+ List(
+ link(src, outerStart),
+ link(outerStart, innerStart),
+ link(innerStart, mid),
+ link(mid, innerEnd),
+ link(innerEnd, outerEnd)
+ )
+ )
+ val outerRows = materialized.getOrElse(outerEnd.operatorIdentifier, -1L)
+ val innerRows = materialized.getOrElse(innerEnd.operatorIdentifier, -1L)
+ assert(
+ outerRows == 9,
+ s"outer LoopEnd must accumulate all 9 inner-iteration rows with a Scala
" +
+ s"operator in the inner body: expected 9, got $outerRows (all:
$materialized)"
+ )
+ assert(
+ innerRows == 3,
+ s"inner LoopEnd must reset per outer iteration (3 rows, not 9): " +
+ s"expected 3, got $innerRows (all: $materialized)"
+ )
+ }
+
+ it should "run a nested loop whose inner body chains multiple JVM operators"
in {
+ // TextInput -> OuterStart -> InnerStart -> Limit -> Sleep(0) -> InnerEnd
-> OuterEnd.
+ //
+ // The strongest combination: nested (the OUTER loop's state crosses the
+ // JVM hops stamped with loop_counter = 1, so the counter MAGNITUDE must
+ // survive both hops) x multi-hop (the JVM-to-JVM state handoff between
+ // the two Scala workers) x both operator kinds. A zeroed counter or a
+ // blanked id at either hop would mis-route the outer state at the inner
+ // LoopEnd.
+ val src = textInput("1\n2\n3")
+ val outerStart = loopStart("i = 0", "table")
+ val innerStart = loopStart("j = 0", "table.iloc[j]")
+ val first = limit(10)
+ val second = sleep(0)
+ val innerEnd = loopEnd("j += 1", "j < len(table)")
+ val outerEnd = loopEnd("i += 1", "i < len(table)")
+ val materialized = runAndGetMaterializedRowCounts(
+ List(src, outerStart, innerStart, first, second, innerEnd, outerEnd),
+ List(
+ link(src, outerStart),
+ link(outerStart, innerStart),
+ link(innerStart, first),
+ link(first, second),
+ link(second, innerEnd),
+ link(innerEnd, outerEnd)
+ )
+ )
+ val outerRows = materialized.getOrElse(outerEnd.operatorIdentifier, -1L)
+ val innerRows = materialized.getOrElse(innerEnd.operatorIdentifier, -1L)
+ assert(
+ outerRows == 9,
+ s"outer LoopEnd must accumulate all 9 inner-iteration rows with two " +
+ s"chained JVM operators in the inner body: expected 9, got $outerRows
" +
+ s"(all: $materialized)"
+ )
+ assert(
+ innerRows == 3,
+ s"inner LoopEnd must reset per outer iteration (3 rows, not 9): " +
+ s"expected 3, got $innerRows (all: $materialized)"
+ )
+ }
+
+ it should "run a loop whose body chains multiple JVM operators" in {
+ // TextInput -> LoopStart -> Limit -> Sleep(0) -> LoopEnd.
+ //
+ // Two consecutive JVM hops pin the JVM-to-JVM state handoff: the first
+ // Scala worker WRITES the envelope columns to its output state table and
+ // the second Scala worker READS them back. The single-JVM-op tests never
+ // exercise that pairing (there, each hop faces a Python worker on at
+ // least one side). Sleep(0) is a pure identity pass-through.
+ val src = textInput("1\n2\n3")
+ val start = loopStart("i = 0", "table.iloc[i]")
+ val first = limit(10)
+ val second = sleep(0)
+ val end = loopEnd("i += 1", "i < len(table)")
+ val materialized = runAndGetMaterializedRowCounts(
+ List(src, start, first, second, end),
+ List(link(src, start), link(start, first), link(first, second),
link(second, end))
+ )
+ val endRows = materialized.getOrElse(end.operatorIdentifier, -1L)
+ assert(
+ endRows == 3,
+ s"LoopEnd must accumulate all 3 iterations with two chained JVM " +
+ s"operators in the loop body: expected 3, got $endRows (all:
$materialized)"
+ )
+ }
}
diff --git
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManagerSpec.scala
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManagerSpec.scala
index fa9a52c430..943dd98ab6 100644
---
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManagerSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/OutputManagerSpec.scala
@@ -307,6 +307,29 @@ class OutputManagerSpec extends AnyFlatSpec with
MockFactory {
outputManager.emitState(state)
}
+ it should "carry the loop envelope onto every emitted StateFrame" in {
+ // The loop envelope (loop_counter / loop_start_id) is owned by the Python
+ // runtime; a JVM hop only forwards it. emitState must stamp the incoming
+ // envelope onto the outgoing frames (and the storage write -- covered by
+ // State.toTuple in StateSpec), not reset it to the no-loop defaults.
+ val outputManager: OutputManager = wire[OutputManager]
+ val portA = PortIdentity(7)
+ outputManager.addPort(portA, schema, None)
+ // Fresh receiver: the spec-level gateway tracks sequence numbers per
+ // channel, so reusing another test's receiver would start at seq 1.
+ val rec = ActorVirtualIdentity("state-recEnvelope")
+ outputManager.addPartitionerWithPartitioning(
+ PhysicalLink(physicalOpId(), portA, physicalOpId(), portA),
+ OneToOnePartitioning(10, Seq(channelTo(rec)))
+ )
+
+ val state = State(Map("k" -> 1))
+ (mockHandler.apply _).expects(
+ WorkflowFIFOMessage(channelTo(rec), 0, StateFrame(state, 2L,
"outer-loop"))
+ )
+ outputManager.emitState(state, loopCounter = 2L, loopStartId =
"outer-loop")
+ }
+
// -- addPort / storage no-ops / getSingleOutputPortIdentity --------------
"addPort" should "be idempotent for a duplicate port id" in {
diff --git
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/NetworkOutputBufferSpec.scala
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/NetworkOutputBufferSpec.scala
index 33f89a3348..b3f8817bf4 100644
---
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/NetworkOutputBufferSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/sendsemantics/partitioners/NetworkOutputBufferSpec.scala
@@ -177,6 +177,18 @@ class NetworkOutputBufferSpec extends AnyFlatSpec {
assert(cap.messages.head.payload == StateFrame(state))
}
+ it should "stamp the loop envelope onto the sent StateFrame" in {
+ // A JVM hop inside a loop body forwards the loop envelope unchanged
+ // (loop operators are Python-only); sendState must put the caller's
+ // loop_counter / loop_start_id on the frame instead of the no-loop
+ // defaults, or the matching LoopEnd's back-jump loses its target.
+ val (buf, cap) = newBuffer()
+ val state = State(Map("k" -> "v"))
+ buf.sendState(state, loopCounter = 2L, loopStartId = "outer-loop")
+ assert(cap.messages.size == 1)
+ assert(cap.messages.head.payload == StateFrame(state, 2L, "outer-loop"))
+ }
+
it should "leave the tuple buffer empty after sendState (trailing flush
no-op)" in {
// sendState calls flush() AFTER sending the state too. Pin that the
// trailing flush doesn't double-send and that subsequent addTuple
diff --git
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessorSpec.scala
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessorSpec.scala
index 414b0a46a1..53f167c4ea 100644
---
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessorSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/worker/DataProcessorSpec.scala
@@ -286,6 +286,53 @@ class DataProcessorSpec extends AnyFlatSpec with
MockFactory with Matchers with
assert(emitted.exists(_.payload.isInstanceOf[StateFrame]))
}
+ "data processor" should "carry the loop envelope through a state
pass-through unchanged" in {
+ // Loop operators are Python-only, so a JVM operator inside a loop body
+ // only ever FORWARDS the StateFrame loop envelope (loop_counter /
+ // loop_start_id); the +1/-1 bookkeeping lives in the Python runtime.
+ // A dropped envelope zeroes the counter and blanks the id, which breaks
+ // the matching LoopEnd's back-jump ("no loop-back state URI configured
+ // for LoopStart ''") -- so pin exact envelope equality on the emitted
+ // frame, not just that a StateFrame came out.
+ val dp = mkDataProcessor
+ dp.executor = executor
+ dp.stateManager.transitTo(READY)
+ val emitted = scala.collection.mutable.ArrayBuffer[WorkflowFIFOMessage]()
+ (outputHandler.apply _)
+ .expects(*)
+ .onCall { (m: Either[MainThreadDelegateMessage, WorkflowFIFOMessage]) =>
+ m.foreach(emitted += _); ()
+ }
+ .anyNumberOfTimes()
+ val inputState = State(Map("i" -> 1))
+ (
+ (
+ state: State,
+ port: Int
+ ) => executor.processState(state, port)
+ )
+ .expects(inputState, 0)
+ .returning(Some(inputState))
+ dp.inputManager.addPort(inputPortId, schema, List.empty, List.empty)
+ dp.inputGateway
+ .getChannel(ChannelIdentity(senderWorkerId, testWorkerId, isControl =
false))
+ .setPortId(inputPortId)
+ dp.outputManager.addPort(outputPortId, schema, None)
+ dp.outputManager.addPartitionerWithPartitioning(
+ PhysicalLink(testOpId, outputPortId, upstreamOpId, inputPortId),
+ OneToOnePartitioning(1, Seq(ChannelIdentity(testWorkerId,
senderWorkerId, isControl = false)))
+ )
+ dp.processDataPayload(
+ ChannelIdentity(senderWorkerId, testWorkerId, isControl = false),
+ StateFrame(inputState, loopCounter = 2L, loopStartId = "outer-loop")
+ )
+ val stateFrames = emitted.map(_.payload).collect { case sf: StateFrame =>
sf }
+ assert(
+ stateFrames.toList == List(StateFrame(inputState, 2L, "outer-loop")),
+ s"envelope must ride through unchanged, got: $stateFrames"
+ )
+ }
+
"data processor" should "not emit when processState yields None" in {
val dp = mkDataProcessor
dp.executor = executor
diff --git
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/state/State.scala
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/state/State.scala
index f2896c8e19..aa88ce1660 100644
---
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/state/State.scala
+++
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/state/State.scala
@@ -45,10 +45,18 @@ object State {
private val Content = "content"
// loop-control bookkeeping owned by the (Python) worker runtime; not user
// state and never in the content JSON. Materialized as its own columns,
- // parallel to content. Scala never originates loop state (loop operators are
- // Python-only), so toTuple defaults these to the "no loop" values.
+ // parallel to content. Scala never ORIGINATES loop state (loop operators are
+ // Python-only), so toTuple defaults these to the "no loop" values -- but a
+ // JVM operator inside a loop body must CARRY them through unchanged, so the
+ // extractors below read them back off a materialized/transported row.
private val LoopCounter = "loop_counter"
private val LoopStartId = "loop_start_id"
+
+ /** Read the loop-envelope counter off a State row (see `toTuple`). */
+ def loopCounterFrom(row: Tuple): Long =
row.getField[java.lang.Long](LoopCounter).longValue()
+
+ /** Read the loop-envelope LoopStart id off a State row (see `toTuple`). */
+ def loopStartIdFrom(row: Tuple): String = row.getField[String](LoopStartId)
private val BytesTypeMarker = "__texera_type__"
private val BytesValue = "bytes"
private val PayloadMarker = "payload"
diff --git
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/state/StateSpec.scala
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/state/StateSpec.scala
index f0af0b30c8..8856d48ea9 100644
---
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/state/StateSpec.scala
+++
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/state/StateSpec.scala
@@ -125,6 +125,23 @@ class StateSpec extends AnyFlatSpec {
assert(tuple.getField[String]("content") == """{"x":1}""")
}
+ it should "read the loop envelope back off a materialized tuple" in {
+ // A JVM operator inside a loop body replays materialized states and must
+ // carry loop_counter / loop_start_id through unchanged -- the columns
+ // exist precisely so the envelope survives storage. These extractors are
+ // what InputPortMaterializationReaderThread / PythonProxyServer read; a
+ // rename of the columns on either side must break this.
+ val tuple = State(Map("i" -> 1L)).toTuple(2L, "outer-loop")
+ assert(State.loopCounterFrom(tuple) == 2L)
+ assert(State.loopStartIdFrom(tuple) == "outer-loop")
+ }
+
+ it should "default the loop envelope to the no-loop values" in {
+ val tuple = State(Map("i" -> 1L)).toTuple()
+ assert(State.loopCounterFrom(tuple) == 0L)
+ assert(State.loopStartIdFrom(tuple) == "")
+ }
+
it should "decode a payload encoded by the Python serializer" in {
// Wire-format compatibility check: the bytes-marker keys and the
// single-row "content" column must match what core/models/state.py