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 2c2c0e5034 test(workflow-operator): close the AsterixDB source
descriptor's branch arms (#7958)
2c2c0e5034 is described below
commit 2c2c0e5034e1e45fadfcdea5fa2affd122642db9
Author: Xinyuan Lin <[email protected]>
AuthorDate: Tue Aug 25 06:57:24 2026 +0000
test(workflow-operator): close the AsterixDB source descriptor's branch
arms (#7958)
### What changes were proposed in this PR?
Two files in the AsterixDB source family. 52 tests → 68
(`AsterixDBSourceOpDescSpec` 5 → 18, `AsterixDBSourceOpExecSpec` 47 →
50), counts read from the JUnit XML rather than estimated.
| File | Codecov | JaCoCo line-hit | Branch arms |
|---|---|---|---|
| `AsterixDBSourceOpDesc.scala` | 42/50 = 84.0% → **49/50 = 98.0%** |
100% → 100% | 12 missed → **1** |
| `AsterixDBSourceOpExec.scala` | 91/104 = 87.5% → **92/104 = 88.5%** |
99% → 99% | 19 missed → **15** |
**+8 fully-covered lines and 15 branch arms closed.** Note the line-hit
metric does not move at all: every gain is a branch arm flipping a
Codecov "partial" into a "hit". That is the case Codecov penalises and
line-hit hides, and it is the whole content of this PR.
The `OpExec` half was assessed earlier at about 1 line, and that held —
it contributes +1. The `OpDesc` half carries the bundle.
### Verification
30 mutations, **29 killed, 1 survivor.**
The first draft claimed no survivors; **five real, semantic,
non-equivalent mutants survived it**, all five re-run and confirmed
before anything was changed. Its counts were also stale (`testsAfter:
64` against an actual 68) and two mutation rows undercounted their
failures.
Three repairs are worth naming:
- **The stub answered every statement identically**, which made the
query assertions vacuous — an exchanged query passed. It now records the
decoded statements and asserts on them.
- **A mapping assertion was deleted rather than kept**:
`asterixDBVersionMapping.get(host) shouldBe Some("0.9.9")` asserted the
test's own fixture, not production.
- Two new `OpDesc` tests pin ordering that nothing constrained: that
`updatePort` is applied by `sourceSchema` *before* it issues any query,
and that the unset-field error names host, port, database in that order.
**The survivor, and it is a real gap rather than an equivalent mutant:**
adding `cachedTuple = None` to `AsterixDBSourceOpExec.close()` survives
all 100 tests. Today, after `hasNext` then `close()`, a fresh
`produceTuple()` still hands out the previously peeked row; under the
mutant it does not. `close()` clears the iterator and the query string
but not the cached tuple. **Neither side is pinned deliberately** — the
behaviour looks unintended, and asserting either way would cement a
decision nobody has made.
### Deliberately not included, with bytecode evidence
- **`AsterixDBSourceOpExec:331` can never be fully covered.** `javap`
shows offset 225 pushing `iconst_1` unconditionally for the `| _`
alternate, so the `ifeq` has a permanently dead side. It went from 4
missed arms to 1 and still reads as partial, moving Codecov by exactly
zero.
- **Lines 227 and 239** are the null-comparison arms of scalac's `==`
expansion on `attr.getType`; reaching them needs an `Attribute` carrying
a null `AttributeType`.
- **Line 152's `if (values == null) return null` is dead**: `javap`
confirms `CSVParser.parse` returns a `scala.Option`, so `values` is
never a Java null.
- **Line 268** re-checks a condition its only caller already guards, so
the false arm cannot occur.
- `AsterixDBSourceOpDesc` line 177 keeps one arm (mb=1/cb=3) for the
same structural reason.
No production file is touched, and no stray `test_large_binary.txt` was
left behind.
### Any related issues, documentation, discussions?
Closes #7956
### How was this PR tested?
```
sbt "WorkflowOperator/testOnly
org.apache.texera.amber.operator.source.sql.asterixdb.AsterixDBSourceOpDescSpec
org.apache.texera.amber.operator.source.sql.asterixdb.AsterixDBSourceOpExecSpec"
```
```
[info] Total number of tests run: 68
[info] Tests: succeeded 68, failed 0, canceled 0, ignored 0, pending 0
```
The whole `asterixdb` package is green at 100 tests.
`Test/scalafmtCheck` passes.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Opus 5)
---------
Signed-off-by: Xinyuan Lin <[email protected]>
Co-authored-by: Copilot Autofix powered by AI
<[email protected]>
---
.../sql/asterixdb/AsterixDBSourceOpDescSpec.scala | 325 ++++++++++++++++++++-
.../sql/asterixdb/AsterixDBSourceOpExecSpec.scala | 65 ++++-
2 files changed, 383 insertions(+), 7 deletions(-)
diff --git
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpDescSpec.scala
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpDescSpec.scala
index 2f8a935247..2898acf64f 100644
---
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpDescSpec.scala
+++
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpDescSpec.scala
@@ -19,7 +19,9 @@
package org.apache.texera.amber.operator.source.sql.asterixdb
+import com.sun.net.httpserver.{HttpExchange, HttpServer}
import org.apache.texera.amber.core.executor.OpExecWithClassName
+import org.apache.texera.amber.core.tuple.{AttributeType, Schema}
import org.apache.texera.amber.core.workflow.WorkflowContext.{
DEFAULT_EXECUTION_ID,
DEFAULT_WORKFLOW_ID
@@ -27,10 +29,130 @@ import
org.apache.texera.amber.core.workflow.WorkflowContext.{
import org.apache.texera.amber.operator.LogicalOp
import org.apache.texera.amber.operator.metadata.OperatorGroupConstants
import org.apache.texera.amber.util.JSONUtils.objectMapper
+import org.scalatest.BeforeAndAfterAll
import org.scalatest.flatspec.AnyFlatSpec
import org.scalatest.matchers.should.Matchers
-class AsterixDBSourceOpDescSpec extends AnyFlatSpec with Matchers {
+import java.net.{InetSocketAddress, URLDecoder}
+import java.nio.charset.StandardCharsets
+import scala.collection.mutable
+import scala.util.Try
+
+class AsterixDBSourceOpDescSpec extends AnyFlatSpec with Matchers with
BeforeAndAfterAll {
+
+ //
---------------------------------------------------------------------------
+ // In-process AsterixDB stub
+ //
+ // sourceSchema() resolves the dataset's datatype over HTTP through
+ // AsterixDBConnUtil, so the type-mapping tests need a reachable server. The
+ // stub answers the two metadata statements sourceSchema() issues and nothing
+ // else; the same approach is used by AsterixDBConnUtilSpec and
+ // AsterixDBSourceOpExecSpec. Binding port 0 keeps it off any fixed port.
+ //
---------------------------------------------------------------------------
+
+ /** Field name -> AsterixDB type the stub reports for the dataset's
datatype. */
+ @volatile private var datatypeFields: Seq[(String, String)] = Seq.empty
+
+ /** Decoded `statement` form field of every /query/service request, in
order. */
+ private val recordedStatements = mutable.Buffer[String]()
+
+ private val server: HttpServer = HttpServer.create(new InetSocketAddress(0),
0)
+ server.createContext(
+ "/admin/version",
+ (exchange: HttpExchange) => respond(exchange,
"""{"git.build.version":"0.9.9"}""")
+ )
+ server.createContext(
+ "/query/service",
+ (exchange: HttpExchange) => {
+ val is = exchange.getRequestBody
+ val body =
+ try new String(is.readAllBytes(), StandardCharsets.UTF_8)
+ finally is.close()
+ val statement = formField(body, "statement")
+ recordedStatements.synchronized { recordedStatements += statement }
+ respond(exchange, responseFor(statement))
+ }
+ )
+
+ private val host = "localhost"
+ private def port: String = server.getAddress.getPort.toString
+
+ private def responseFor(statement: String): String =
+ if (statement.contains("Metadata.`Datatype`")) {
+ val fields = datatypeFields
+ .map { case (name, tpe) =>
s"""{"FieldName":"$name","FieldType":"$tpe"}""" }
+ .mkString(",")
+ s"""{"results":[{"Fields":[$fields]}]}"""
+ } else if (statement.contains("Metadata.`Dataset`")) {
+ """{"results":[{"DatatypeName":"tweetType"}]}"""
+ } else {
+ // Deliberately not a fall-through onto the dataset answer: a statement
+ // aimed at the wrong metadata table must come back empty, so a query the
+ // descriptor mistargets cannot still yield a plausible schema.
+ """{"results":[]}"""
+ }
+
+ private def respond(exchange: HttpExchange, body: String): Unit = {
+ val bytes = body.getBytes(StandardCharsets.UTF_8)
+ exchange.getResponseHeaders.add("Content-Type", "application/json")
+ exchange.sendResponseHeaders(200, bytes.length.toLong)
+ val os = exchange.getResponseBody
+ try os.write(bytes)
+ finally os.close()
+ }
+
+ private def formField(body: String, name: String): String =
+ body
+ .split("&")
+ .filter(_.contains("="))
+ .map { pair =>
+ val idx = pair.indexOf('=')
+ URLDecoder.decode(pair.substring(0, idx), StandardCharsets.UTF_8) ->
+ URLDecoder.decode(pair.substring(idx + 1), StandardCharsets.UTF_8)
+ }
+ .toMap
+ .getOrElse(name, "")
+
+ override protected def beforeAll(): Unit = {
+ super.beforeAll()
+ server.start()
+ }
+
+ override protected def afterAll(): Unit = {
+ try {
+ server.stop(0)
+ // Clean up AsterixDBConnUtil's host-keyed version cache to avoid leaking
+ // state into other tests that may reuse the same host key.
+ AsterixDBConnUtil.asterixDBVersionMapping -= host
+ } finally super.afterAll()
+ }
+
+ /** The schema sourceSchema() derives when the datatype reports these
fields. */
+ private def schemaOf(fields: (String, String)*): Schema = {
+ datatypeFields = fields
+ recordedStatements.synchronized { recordedStatements.clear() }
+ val d = configured()
+ d.port = port
+ d.sourceSchema()
+ }
+
+ /** The statements the stub saw, in order, since the last schemaOf(). */
+ private def statements: List[String] =
+ recordedStatements.synchronized { recordedStatements.toList }
+
+ /** A descriptor with every connection field filled in, ready to be broken.
*/
+ private def configured(): AsterixDBSourceOpDesc = {
+ val d = new AsterixDBSourceOpDesc
+ d.host = host
+ d.port = "19002"
+ d.database = "test"
+ d.table = "twitter"
+ d
+ }
+
+ //
---------------------------------------------------------------------------
+ // Operator metadata and serialization
+ //
---------------------------------------------------------------------------
"AsterixDBSourceOpDesc.operatorInfo" should
"advertise the AsterixDB source in the Database Connector group with no
input and one output" in {
@@ -56,11 +178,6 @@ class AsterixDBSourceOpDescSpec extends AnyFlatSpec with
Matchers {
d.interval shouldBe 0L
}
- "AsterixDBSourceOpDesc.sourceSchema" should "prompt for connection details
before a connection is configured" in {
- val ex = intercept[IllegalArgumentException]((new
AsterixDBSourceOpDesc).sourceSchema())
- ex.getMessage should include("host")
- }
-
"AsterixDBSourceOpDesc.getPhysicalOp" should
"wire the AsterixDB exec as a source op with no input port and one output
port" in {
val d = new AsterixDBSourceOpDesc
@@ -100,4 +217,200 @@ class AsterixDBSourceOpDescSpec extends AnyFlatSpec with
Matchers {
r.username shouldBe null
r.password shouldBe null
}
+
+ //
---------------------------------------------------------------------------
+ // updatePort
+ //
---------------------------------------------------------------------------
+
+ "AsterixDBSourceOpDesc.updatePort" should
+ "resolve the sentinel `default` to AsterixDB's HTTP API port and leave any
other port alone" in {
+ val d = new AsterixDBSourceOpDesc
+ d.port = "default"
+ d.updatePort()
+ d.port shouldBe "19002"
+
+ // The sentinel is recognised through surrounding whitespace.
+ d.port = " default "
+ d.updatePort()
+ d.port shouldBe "19002"
+
+ // An explicit port is left exactly as configured.
+ d.port = "19004"
+ d.updatePort()
+ d.port shouldBe "19004"
+ }
+
+ it should "be applied by sourceSchema before it issues any query" in {
+ // sourceSchema resolves `default` to 19002; stand up a minimal stub there
so
+ // the test doesn't depend on an external AsterixDB process.
+ val defaultPortServer = Try(HttpServer.create(new
InetSocketAddress(19002), 0)).getOrElse {
+ cancel("port 19002 is unavailable; cannot verify updatePort integration
without a local stub")
+ }
+ try {
+ defaultPortServer.createContext(
+ "/admin/version",
+ (exchange: HttpExchange) => respond(exchange,
"""{"git.build.version":"0.9.9"}""")
+ )
+ defaultPortServer.createContext(
+ "/query/service",
+ (exchange: HttpExchange) => {
+ val is = exchange.getRequestBody
+ val body =
+ try new String(is.readAllBytes(), StandardCharsets.UTF_8)
+ finally is.close()
+ val statement = formField(body, "statement")
+ respond(exchange, responseFor(statement))
+ }
+ )
+ defaultPortServer.start()
+
+ val d = configured()
+ d.port = "default"
+ d.sourceSchema()
+ d.port shouldBe "19002"
+ } finally defaultPortServer.stop(0)
+ }
+
+ //
---------------------------------------------------------------------------
+ // sourceSchema - connection validation
+ //
+ // Every one of these fails before any HTTP call is made, so the stub is not
+ // involved. Each assertion pins the whole message, because the only job of
+ // these guards is to tell the user which field is wrong.
+ //
---------------------------------------------------------------------------
+
+ "AsterixDBSourceOpDesc.sourceSchema" should
+ "prompt for connection details before a connection is configured" in {
+ val ex = intercept[IllegalArgumentException]((new
AsterixDBSourceOpDesc).sourceSchema())
+ ex.getMessage shouldBe "requirement failed: Please enter a valid host name
for AsterixDB."
+ }
+
+ it should "reject a host that is only whitespace" in {
+ val d = configured()
+ d.host = " "
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid host name for AsterixDB."
+ }
+
+ it should "name the port when the port is missing" in {
+ val d = configured()
+ d.port = null
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid port for AsterixDB."
+ }
+
+ it should "name the port when the port is only whitespace" in {
+ val d = configured()
+ d.port = " "
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid port for AsterixDB."
+ }
+
+ it should "name the database when the database is missing" in {
+ val d = configured()
+ d.database = null
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid database name for AsterixDB."
+ }
+
+ it should "name the database when the database is only whitespace" in {
+ val d = configured()
+ d.database = " "
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid database name for AsterixDB."
+ }
+
+ it should "name the table when the table is missing" in {
+ val d = configured()
+ d.table = null
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid table name for AsterixDB."
+ }
+
+ it should "name the table when the table is only whitespace" in {
+ val d = configured()
+ d.table = " "
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid table name for AsterixDB."
+ }
+
+ it should "name the first unset field, checking host, port, database, then
table" in {
+ // Each test above breaks exactly one field, which pins the message
attached
+ // to each guard but not the order the guards run in. `require` reports
only
+ // the FIRST failure, so on a half-configured operator that order is the
+ // whole user-visible behaviour: fixing one field must surface the next.
+ val d = configured()
+ d.port = " "
+ d.database = null
+ d.table = ""
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid port for AsterixDB."
+ d.port = "19002"
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid database name for AsterixDB."
+ d.database = "test"
+ intercept[IllegalArgumentException](d.sourceSchema()).getMessage shouldBe
+ "requirement failed: Please enter a valid table name for AsterixDB."
+ }
+
+ //
---------------------------------------------------------------------------
+ // sourceSchema - metadata resolution
+ //
---------------------------------------------------------------------------
+
+ it should "look the dataset's datatype up first, then that datatype's
fields" in {
+ // Two round trips, and the second has to key off the DatatypeName the
first
+ // one returned. Without this the whole first round trip could be discarded
+ // and the schema would still come out right against a stub that answers
any
+ // datatype name - but against a real instance the datatype lookup would
miss
+ // and sourceSchema would silently return an EMPTY schema.
+ schemaOf("id" -> "int64")
+ statements should have length 2
+ statements.head shouldBe
+ "SELECT DatatypeName FROM Metadata.`Dataset` ds where
ds.`DatasetName`='twitter';"
+ statements(1) should include("dt.DatatypeName = 'tweetType'")
+ }
+
+ //
---------------------------------------------------------------------------
+ // sourceSchema - AsterixDB datatype -> Texera attribute type
+ //
---------------------------------------------------------------------------
+
+ it should "map every AsterixDB scalar type it knows onto the matching
attribute type" in {
+ val schema = schemaOf(
+ "a_boolean" -> "boolean",
+ "b_int32" -> "int32",
+ "c_int64" -> "int64",
+ "d_float" -> "float",
+ "e_double" -> "double",
+ "f_datetime" -> "datetime",
+ "g_date" -> "date",
+ "h_string" -> "string"
+ )
+ // Attributes arrive sorted by field name.
+ schema.getAttributeNames shouldBe List(
+ "a_boolean",
+ "b_int32",
+ "c_int64",
+ "d_float",
+ "e_double",
+ "f_datetime",
+ "g_date",
+ "h_string"
+ )
+ schema.getAttribute("a_boolean").getType shouldBe AttributeType.BOOLEAN
+ schema.getAttribute("b_int32").getType shouldBe AttributeType.INTEGER
+ schema.getAttribute("c_int64").getType shouldBe AttributeType.LONG
+ // float and double both widen to DOUBLE.
+ schema.getAttribute("d_float").getType shouldBe AttributeType.DOUBLE
+ schema.getAttribute("e_double").getType shouldBe AttributeType.DOUBLE
+ // date carries no time of day, but is still surfaced as a TIMESTAMP.
+ schema.getAttribute("f_datetime").getType shouldBe AttributeType.TIMESTAMP
+ schema.getAttribute("g_date").getType shouldBe AttributeType.TIMESTAMP
+ schema.getAttribute("h_string").getType shouldBe AttributeType.STRING
+ }
+
+ it should "fall back to STRING for an AsterixDB type it does not recognise"
in {
+ // `uuid` is a real AsterixDB scalar type that this mapping does not name.
+ // The catch-all keeps the column readable instead of failing the schema.
+ schemaOf("id" -> "uuid").getAttribute("id").getType shouldBe
AttributeType.STRING
+ }
}
diff --git
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpExecSpec.scala
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpExecSpec.scala
index 87b12ed8eb..b1023437de 100644
---
a/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpExecSpec.scala
+++
b/common/workflow-operator/src/test/scala/org/apache/texera/amber/operator/source/sql/asterixdb/AsterixDBSourceOpExecSpec.scala
@@ -517,6 +517,21 @@ class AsterixDBSourceOpExecSpec
}.getMessage should startWith("Unexpected type:")
}
+ it should "reject a batch column whose type is untyped or opaque" in {
+ // ANY and BINARY cannot come out of an AsterixDB datatype today, but the
+ // catch-all has to keep rejecting them rather than rendering them into the
+ // sliding-window predicate as a bare literal.
+ val exec = newExec()
+ exec.batchByAttribute = Some(new Attribute("anything", AttributeType.ANY))
+ intercept[IllegalArgumentException] {
+ exec.batchAttributeToString(java.lang.Long.valueOf(1L))
+ }.getMessage should startWith("Unexpected type:")
+ exec.batchByAttribute = Some(new Attribute("blob", AttributeType.BINARY))
+ intercept[IllegalArgumentException] {
+ exec.batchAttributeToString(java.lang.Long.valueOf(1L))
+ }.getMessage shouldBe "Unexpected type: binary"
+ }
+
it should "reject the call when no batch column is resolved" in {
val exec = newExec(_.batchByColumn = Some("created_at"))
exec.batchByAttribute shouldBe None
@@ -573,7 +588,28 @@ class AsterixDBSourceOpExecSpec
tableNameRows = Seq("\"other\"\n", "\"twitter\"\n")
exec.open()
exec.tableNames.toList shouldBe List("other", "twitter")
- AsterixDBConnUtil.asterixDBVersionMapping.get(host) shouldBe Some("0.9.9")
+ }
+
+ it should "refresh a stale cached API version when it opens" in {
+ val exec = newExec()
+ // Asserting the mapping straight after open() would be vacuous: the
+ // constructor's own `schema = desc.sourceSchema()` already went through
+ // queryAsterixDB, which fills a MISSING host entry, so `0.9.9` is cached
+ // before open() is ever entered. Overwriting the entry first makes
open()'s
+ // explicit refresh the only thing that can restore it - queryAsterixDB
+ // never rewrites an entry that is already present, so a stale version
(which
+ // selects a different `format` field) would otherwise survive forever in
+ // this singleton.
+ val prev = AsterixDBConnUtil.asterixDBVersionMapping.get(host)
+ try {
+ AsterixDBConnUtil.asterixDBVersionMapping += (host -> "0.0.0")
+ exec.open()
+ AsterixDBConnUtil.asterixDBVersionMapping.get(host) shouldBe
Some("0.9.9")
+ } finally {
+ prev.fold(AsterixDBConnUtil.asterixDBVersionMapping -= host)(v =>
+ AsterixDBConnUtil.asterixDBVersionMapping += (host -> v)
+ )
+ }
}
it should "reject a table that the AsterixDB instance does not expose" in {
@@ -615,6 +651,24 @@ class AsterixDBSourceOpExecSpec
dataStatements.head should endWith(";")
}
+ it should "answer a repeated hasNext from the cached row instead of
consuming another" in {
+ val exec = newExec()
+ dataRows = Seq(
+ "1,2023-11-13T10:15:30,1,1.0,first,true",
+ "2,2023-11-13T10:15:30,2,2.0,second,true"
+ )
+ val iterator = exec.produceTuple()
+ // The first peek pulls a row into the cache; every further peek has to
read
+ // that cache, or the peeked row is dropped on the floor.
+ iterator.hasNext shouldBe true
+ iterator.hasNext shouldBe true
+ iterator.hasNext shouldBe true
+ iterator.next().asInstanceOf[Tuple].getField[Any]("text") shouldBe "first"
+ iterator.next().asInstanceOf[Tuple].getField[Any]("text") shouldBe "second"
+ iterator.hasNext shouldBe false
+ dataStatements should have length 1
+ }
+
it should "map the literal token null, in any column, to a null field" in {
val exec = newExec()
dataRows = Seq("null,null,null,null,null,null")
@@ -718,6 +772,15 @@ class AsterixDBSourceOpExecSpec
// so the second row is never delivered.
exec.produceTuple().hasNext shouldBe false
dataStatements should have length 1
+ // NOT covered here, on purpose: this drives next() directly, so
`cachedTuple`
+ // is empty by the time close() runs. The engine's own loop is
+ // `while (it.hasNext) it.next()`, i.e. hasNext always peeks first, and
+ // close() clears curResultIterator/curQueryString but not the inherited
+ // cachedTuple (neither does SQLSourceOpExec.close()) - so a peeked but
+ // unconsumed row is still handed out after close(). That is left unpinned
on
+ // purpose: asserting it would cement what looks like a leak, and asserting
+ // the opposite would fail, so the contract has to be decided before it is
+ // tested.
}
it should "leave a partially consumed scan resumable when it is not closed"
in {