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 09b4e74560 test(amber): make FlowControlSpec's size-cap test actually
assert (#8406)
09b4e74560 is described below
commit 09b4e74560a0cb16c557816711639493564f9d3c
Author: Xinyuan Lin <[email protected]>
AuthorDate: Sat Sep 5 06:56:08 2026 +0000
test(amber): make FlowControlSpec's size-cap test actually assert (#8406)
### What changes were proposed in this PR?
`FlowControlSpec` had a test that asserted nothing:
```scala
"FlowControl" should "trip the size-cap assertion for a message that
exceeds maxByteAllowed" in {
// ... comment conceding it cannot synthesize an oversized payload ...
val fc = new FlowControl()
(1L to 1000L).foreach(i => fc.getMessagesToSend(msg(i)))
succeed
}
```
It sent 1000 messages and then checked no property of the result. The
name claimed a guarantee the file did not pin.
**Measured proof the name was empty.** Deleting the guard the test is
named after — the `assert(creditNeeded <= maxByteAllowed, ...)` block at
the top of `FlowControl.getMessagesToSend` — changed nothing:
| production `FlowControl.scala` | `FlowControlSpec` |
|---|---|
| unmodified | 14 passed, 0 failed |
| `assert(creditNeeded <= maxByteAllowed, ...)` block deleted | 14
passed, 0 failed |
The fixture could never reach the guard: `FixedSizePayload` reports 200
bytes and `flow-control.max-credit-allowed-in-bytes-per-channel` is
1,600,000,000, so `200 <= 1600000000` held on all 1000 iterations. 1000
x 200 = 200,000 bytes does not exhaust credit either, so the messages
never took the stashing path.
**The change: two tests, both with real assertions.**
1. `"FlowControl.getMessagesToSend" should "reject a message larger than
the whole credit cap"` — this one actually trips the guard. An oversized
payload turns out to be cheap to build: `DataFrame.inMemSize` is
`frame.map(_.inMemSize).sum`, which does **not** deduplicate, so an
`Array[Tuple]` holding the same tuple reference N times reports N x its
size. One tuple with a 100,000-char string, repeated `maxBytes /
tupleSize + 1` times, reports over the cap for a few hundred KB of real
memory. The test intercepts the `AssertionError`, checks its message,
and checks the rejection left the channel untouched (`getCredit`
unchanged, not marked overloaded — which is what the out-of-credit
branch below the guard would have done instead).
2. `it should "forward every under-cap message and charge its size
against the credit"` — replaces the old body, pinning the fast path the
guard sits on:
```
per message i in 1..1000:
getMessagesToSend(msg(i)).toList == List(msg(i)) -- forwarded, not
stashed
getCredit == maxBytes - i * msgSize -- charged exactly its
own size
isOverloaded == false
precondition: batch * msgSize < maxBytes -- else these
expectations
describe the
stashing path
```
Expected values are derived from the fixture (`msgSize` from
`WorkflowMessage.getInMemSize`, `maxBytes` from `ApplicationConfig`),
not hard-coded.
**A claim from my own earlier draft, corrected.** An earlier revision of
this PR stated that covering the size-cap guard "needs an injection seam
for `maxByteAllowed` … that is a production change" and listed the guard
as a disclosed, unavoidable gap. **That was wrong**, and review pressure
is what sent me back to check it. `DataFrame`'s non-deduplicating size
sum is the seam, it is test-only, and the guard is now covered — see the
control row in the mutation table below, which went from PASS to FAIL.
No production change was needed.
**What this PR does NOT do.** It does not change `FlowControl` or any
other `src/main` file (`git diff` on `amber/src/main` is empty). It does
not reformat the file or touch the other 13 tests. Two pre-existing
weaknesses in neighbouring tests are left alone as out of scope:
`"eventually drain the stash across many ack cycles"` ends in
`assert(seen == stashed.size)` where `seen` is incremented once per
element of `stashed` in the same loop, so that line is true by
construction; and the suite-constructor `assert(msgSize == 200L)`
hard-codes `WorkflowMessage`'s default, which would abort the whole
suite rather than fail one test if that default ever changed.
### Any related issues, documentation, discussions?
Closes #8402
### How was this PR tested?
All runs: `sbt "WorkflowExecutionService/testOnly ..."` on Java 17,
based on `1cbe857007`.
**Non-vacuity: the new tests can fail, and they catch things the suite
did not already catch.** Four mutations to
`FlowControl.getMessagesToSend`, each run against both the new spec and
a verbatim copy of the old `succeed` test (kept in a throwaway probe
suite in the *same* `testOnly` invocation, then deleted). Failing tests
were read from `amber/target/test-reports/TEST-*.xml` by identity, not
from console counts:
| mutation | old `succeed` test | new spec | which identities failed |
|---|---|---|---|
| M1: drop `inflightCredit += creditNeeded` on the fast path | PASS |
**FAIL** (2) | new fast-path test **+ pre-existing**
`decreaseInflightCredit should free credit equal to the acked amount` |
| M2: fast path returns `Iterable.empty` instead of `Iterable(msg)` |
PASS | **FAIL** (2) | new fast-path test **+ pre-existing**
`getMessagesToSend should forward an incoming message when credit is
available` |
| M7: `if (inflightCredit == 0) inflightCredit += creditNeeded` — charge
only the first message | PASS | **FAIL** (1) | **new fast-path test
only** |
| Control: delete the `assert(creditNeeded <= maxByteAllowed, ...)`
block | PASS | **FAIL** (1) | **new size-cap test only** |
Being explicit about what each row proves, since two of them prove less
than they look:
- M1 and M2 are each **also** caught by one pre-existing neighbour. On
those two rows alone you could not tell whether the rewritten test adds
coverage or merely duplicates it.
- M7 and the control row are the ones that settle it. M7 is a real
behavioural break — flow control stops accounting after the first
message — that **no other test in the file detects**; it fails only
because the new test walks a whole batch instead of one message. The
control row is the original defect: the guard is now pinned, where
before nothing in the file noticed its removal.
Failure messages, for the record:
```
M7 : 1599999800 did not equal 1599999600 after 2 forwarded messages
the credit must be down by 2 * 200
Control : Expected exception java.lang.AssertionError to be thrown, but no
exception was thrown
```
Every mutation was applied from, and reverted to, a pristine copy of
`FlowControl.scala` kept outside the repo (never `git checkout` / `git
restore`), and `git diff -- amber/src/main` was verified empty after
each revert and at the end. The probe suite is deleted; `git status
--porcelain` shows only the one intended test file.
**Regression: baseline first, then compared by failing-test identity,
not by counts.** Scope: every spec in `...architecture.messaginglayer.*`
plus `PekkoMessageTransferServiceSpec`, the only other spec reading the
same credit config, in one invocation.
| run | suites | tests | failed |
|---|---|---|---|
| baseline (file restored to `1cbe857007` content) | 12 | 120 | 0 |
| after this change | 12 | 121 | 0 |
The identity diff is exactly one removal and two additions, with no
other test's name or status changed:
```
- FlowControlSpec :: FlowControl should trip the size-cap assertion for a
message that exceeds maxByteAllowed :: PASS
+ FlowControlSpec :: FlowControl.getMessagesToSend should forward every
under-cap message and charge its size against the credit :: PASS
+ FlowControlSpec :: FlowControl.getMessagesToSend should reject a message
larger than the whole credit cap :: PASS
```
That the multi-run drain test is absent from this diff is the point of
one line in the change: it used `it should`, which bound to the subject
of the test being replaced. It now declares `"FlowControl" should`
explicitly, which is why its identity is byte-identical across the two
runs.
Not run: the full `amber` module. Its `@IntegrationTest` suites spawn
Python workers and cannot run on this Windows host, and several
Iceberg-backed specs fail here regardless of the change, so a
full-module identity comparison would have been noise. The change adds
no globals or shared state, so the within-JVM leakage that makes amber's
serial execution matter does not apply.
Lint, all clean in the same invocation as the final test run:
`WorkflowExecutionService/scalafmtCheck`,
`WorkflowExecutionService/Test/scalafmtCheck`,
`WorkflowExecutionService/scalafixAll --check`.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 5)
---
.../messaginglayer/FlowControlSpec.scala | 78 +++++++++++++++++-----
1 file changed, 62 insertions(+), 16 deletions(-)
diff --git
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/FlowControlSpec.scala
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/FlowControlSpec.scala
index 418a02426c..d17dd3f22b 100644
---
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/FlowControlSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/messaginglayer/FlowControlSpec.scala
@@ -20,9 +20,11 @@
package org.apache.texera.amber.engine.architecture.messaginglayer
import org.apache.texera.common.config.ApplicationConfig
+import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema,
Tuple}
import org.apache.texera.amber.core.virtualidentity.{ActorVirtualIdentity,
ChannelIdentity}
import
org.apache.texera.amber.engine.architecture.common.WorkflowActor.NetworkMessage
import org.apache.texera.amber.engine.common.ambermessage.{
+ DataFrame,
WorkflowFIFOMessage,
WorkflowFIFOMessagePayload,
WorkflowMessage
@@ -49,6 +51,11 @@ class FlowControlSpec extends AnyFlatSpec {
private val maxBytes = ApplicationConfig.maxCreditAllowedInBytesPerChannel
+ // One-field schema, used only to build the oversized DataFrame payload that
+ // trips the size-cap guard.
+ private val payloadAttr = new Attribute("payload", AttributeType.STRING)
+ private val payloadSchema: Schema = Schema().add(payloadAttr)
+
"FlowControl" should "report full credit and not be overloaded initially" in
{
val fc = new FlowControl()
assert(fc.getCredit == maxBytes)
@@ -116,6 +123,33 @@ class FlowControlSpec extends AnyFlatSpec {
assert(fc.isOverloaded, "stash still has msg(2L), so overloaded must
remain true")
}
+ // The fast path below the size-cap guard: while every message stays under
the
+ // cap, each one is forwarded immediately rather than stashed, is charged
exactly
+ // its own size against the credit, and never flips the overloaded flag.
Running
+ // a whole batch rather than a single message is what catches accounting
that is
+ // right for the first message and then drifts. (The guard itself is covered
by
+ // "reject a message larger than the whole credit cap" further down.)
+ it should "forward every under-cap message and charge its size against the
credit" in {
+ val batch = 1000L
+ // Fixture precondition: the whole batch must fit under the cap, otherwise
the
+ // expectations below would be describing the stashing path instead.
+ assert(
+ batch * msgSize < maxBytes,
+ s"fixture no longer exercises the fast path: $batch * $msgSize >=
$maxBytes"
+ )
+
+ val fc = new FlowControl()
+ (1L to batch).foreach { i =>
+ val out = fc.getMessagesToSend(msg(i)).toList
+ assert(out == List(msg(i)), s"message $i must be forwarded, not stashed")
+ assert(
+ fc.getCredit == maxBytes - i * msgSize,
+ s"after $i forwarded messages the credit must be down by $i * $msgSize"
+ )
+ assert(!fc.isOverloaded, s"must not be overloaded after $i under-cap
messages")
+ }
+ }
+
"FlowControl.updateQueuedCredit" should "shrink the available credit" in {
val fc = new FlowControl()
fc.updateQueuedCredit(100L)
@@ -146,25 +180,37 @@ class FlowControlSpec extends AnyFlatSpec {
// Edge / invalid-input cases — credit math under abnormal conditions
//
---------------------------------------------------------------------------
- "FlowControl" should "trip the size-cap assertion for a message that exceeds
maxByteAllowed" in {
- // Build a payload whose getInMemSize returns a value larger than the
- // configured per-channel cap. We do this by ratcheting up the Pekko-side
- // size accounting via an oversized DataFrame stand-in: emulate by
- // exhausting credit to <= 0 and then sending a payload that's already
- // larger than 0 — but the assertion in source compares creditNeeded
- // against `maxByteAllowed`, not credit. Since FixedSizePayload is 200L
- // and maxByteAllowed is multi-GB, we cannot synthesize a too-big payload
- // in unit-test scope without producing terabytes. Instead, lock down
- // the *guard* shape: a message at exactly maxByteAllowed is allowed by
- // the assertion (not strictly greater), so any 200L payload always
- // passes — confirm that 1000 sequential 200L messages all pass the
- // assertion regardless of credit accounting.
+ // The size-cap guard at the top of getMessagesToSend. A payload larger than
the
+ // entire per-channel cap could never be sent, so FlowControl rejects it
outright
+ // rather than stashing it forever. Reaching a multi-GB reported size costs
+ // almost no memory: DataFrame.inMemSize sums Tuple.inMemSize across its
array
+ // without deduplicating, so an array holding the same tuple reference many
+ // times reports an oversized payload while allocating one tuple and N
pointers.
+ "FlowControl.getMessagesToSend" should "reject a message larger than the
whole credit cap" in {
+ val bigTuple = Tuple.builder(payloadSchema).add(payloadAttr, "x" *
100000).build()
+ val copies = (maxBytes / bigTuple.inMemSize + 1).toInt
+ val oversized = NetworkMessage(
+ 1L,
+ WorkflowFIFOMessage(channelId, 1L,
DataFrame(Array.fill(copies)(bigTuple)))
+ )
+ // Fixture precondition: the payload really is over the cap, so the guard
is
+ // the thing under test rather than the out-of-credit branch below it.
+ val reportedSize = WorkflowMessage.getInMemSize(oversized.internalMessage)
+ assert(
+ reportedSize > maxBytes,
+ s"fixture payload is not oversized: $reportedSize <= $maxBytes"
+ )
+
val fc = new FlowControl()
- (1L to 1000L).foreach(i => fc.getMessagesToSend(msg(i)))
- succeed
+ val thrown = intercept[AssertionError](fc.getMessagesToSend(oversized))
+ assert(thrown.getMessage.contains("too big to send through flow control"))
+ // The rejection must leave the channel untouched — not half-charged, and
not
+ // marked overloaded as the out-of-credit branch would have done.
+ assert(fc.getCredit == maxBytes)
+ assert(!fc.isOverloaded)
}
- it should "eventually drain the stash across many ack cycles (multi-run)" in
{
+ "FlowControl" should "eventually drain the stash across many ack cycles
(multi-run)" in {
val fc = new FlowControl()
// Saturate credit and stash a batch of messages.
fc.updateQueuedCredit(maxBytes)