Class QueryExecutor
SemanticModel.
The engine is streaming-first. The executeStreaming methods
hand each query statement's output to a RowStreamConsumer as a
lazy Stream of rows, so a query's result is never fully buffered by the
executor itself — a consumer that processes rows incrementally (writing them
out, counting them, etc.) runs in constant memory regardless of result size.
The convenience execute methods are thin wrappers that collect each
stream into a QueryResult for callers that genuinely want the whole
result in memory.
For each query statement in SemanticModel.rootQueries() a
result is produced in declaration order. Queries are executed one at a time:
a query's row stream is fully consumed (and closed) before the next query runs.
Precondition
The model (or semantic result) passed to any execute/
executeStreaming method must be fully valid — i.e.
SemanticResult.isFullyValid() must return true. If a
SemanticResult with errors is supplied, an
IllegalArgumentException is thrown immediately before any execution
begins. Passing a raw SemanticModel that was produced from an invalid
analysis may cause EvaluationExceptions that would have been caught
during semantic analysis.
Usage
SemanticResult result = analyzer.analyze(inputStream);
// Streaming: rows are written out incrementally, never all buffered.
executor.executeStreaming(result, (label, schema, rows) ->
rows.forEach(row -> writer.write(format(row))));
// Convenience: collect the whole result into memory.
List<QueryResult> results = executor.execute(result);
Thread safety
This class is stateless and therefore thread-safe.
-
Constructor Summary
ConstructorsConstructorDescriptionCreates an executor reading the system UTC clock.QueryExecutor(Clock clock) Creates an executor whose current-time built-ins readclock. -
Method Summary
Modifier and TypeMethodDescriptionexecute(SemanticResult result) Executes allquerystatements inresultusing an inline-only connector and collects each result into aQueryResult.execute(SemanticResult result, DataSourceConnector connector) Executes allquerystatements inresult, using the supplied connector, and collects each result into aQueryResult.execute(SemanticModel model) Executes allquerystatements inmodelusing an inline-only connector and collects each result into aQueryResult.execute(SemanticModel model, DataSourceConnector connector) Executes allquerystatements inmodel, using the supplied connector, and collects each result into aQueryResult.voidexecuteLineage(SemanticModel model, DataSourceConnector connector, ProvenanceConsumer<Polynomial> consumer) Evaluates eachquerystatement inmodelas a full why-provenance K-relation over the polynomial-lineage semiringℕ[X], handing each result toconsumer.executeOptimized(SemanticModel model, List<RelNode> optimizedRoots) Executes optimizer-rewritten query trees using an inline-only connector and collects each result into aQueryResult.executeOptimized(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector) Executes optimizer-rewritten query trees, using the supplied connector, and collects each result into aQueryResult.voidexecuteOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer) Streams the results of optimizer-rewritten query trees, one per root query inSemanticModel.rootQueries(), using the supplied connector.voidexecuteOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds) LikeexecuteOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer)but aborts anyFIXfixpoint that exceedsmaxFixpointRoundsiterations.voidexecuteOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds, QueryEventListener listener) LikeexecuteOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer, int)but reports the run'sPLANandEXECUTEdecisions tolistener, each tagged with the owning query's label — the observing counterpart for a caller that has already optimized.voidexecuteOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, RowStreamConsumer consumer) Streams the results of optimizer-rewritten query trees using an inline-only connector.<K> voidexecuteProvenance(SemanticModel model, DataSourceConnector connector, Semiring<K> semiring, ProvenanceConsumer<K> consumer) Evaluates eachquerystatement inmodelas an annotatedK-relationoversemiring, handing each result toconsumerin declaration order.<K> voidexecuteProvenance(SemanticModel model, DataSourceConnector connector, Semiring<K> semiring, String weightColumn, int maxFixpointRounds, ProvenanceConsumer<K> consumer) Evaluates eachqueryas a K-relation oversemiring, supplying per-edge weights fromweightColumnfor any transitive closure in the query — the semiring-weighted-closure path.voidexecuteStreaming(SemanticResult result, DataSourceConnector connector, RowStreamConsumer consumer) Streams the results of allquerystatements inresulttoconsumer, using the supplied connector for external relations.voidexecuteStreaming(SemanticResult result, RowStreamConsumer consumer) Streams the results of allquerystatements inresulttoconsumerusing an inline-only connector.voidexecuteStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer) Streams the results of allquerystatements inmodeltoconsumer, using the supplied connector for external relations.voidexecuteStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds) LikeexecuteStreaming(SemanticModel, DataSourceConnector, RowStreamConsumer)but aborts anyFIXfixpoint that exceedsmaxFixpointRoundsiterations with anEvaluationException.voidexecuteStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds, QueryEventListener listener) LikeexecuteStreaming(SemanticModel, DataSourceConnector, RowStreamConsumer, int)but reports the run'sPLANandEXECUTEdecisions tolistener, each tagged with the owning query's label.voidexecuteStreaming(SemanticModel model, RowStreamConsumer consumer) Streams the results of allquerystatements inmodeltoconsumerusing an inline-only connector.voidexplain(SemanticModel model, PlanConsumer consumer) Renders the physical plan of eachquerystatement inmodeland hands it toconsumer, without executing anything.voidexplain(SemanticModel model, List<RelNode> roots, PlanConsumer consumer) Renders the physical plan of each (optimizer-rewritten) root tree and hands it toconsumer, without executing anything.voidplan(SemanticModel model, List<RelNode> roots, QueryEventListener listener, PhysicalPlanConsumer consumer) Plans each (optimizer-rewritten) root tree — without executing it — and hands the resultingPhysicalNodetoconsumer, while emitting the planner's physical-decisionQueryEvents tolistener(each tagged with the owning query's label).static voidrequireBounded(List<RelNode> roots, BoundednessSource boundedness) Refuses to collect a relation that provably never ends.voidtrace(SemanticModel model, List<RelNode> roots, QueryEventListener listener) Plans each (optimizer-rewritten) root tree — without executing it — so the planner emits its physical-decisionQueryEvents tolistener.voidtraceExecute(SemanticModel model, List<RelNode> roots, DataSourceConnector connector, QueryEventListener listener) Executes each (optimizer-rewritten) root tree purely for itsQueryEvent.Stage.EXECUTEside-effect events — discarding the rows — emitting each tolistenertagged with the owning query's label.
-
Constructor Details
-
QueryExecutor
public QueryExecutor()Creates an executor reading the system UTC clock. -
QueryExecutor
Creates an executor whose current-time built-ins readclock.Pinning the clock makes a run reproducible: a script with a relative window (
NOW() - DURATION 'PT1H') selects the same rows however much later it is replayed. The CLI exposes this as--now/RELIX_NOW.- Parameters:
clock- the clock to read; must not benull
-
-
Method Details
-
executeStreaming
public void executeStreaming(SemanticResult result, DataSourceConnector connector, RowStreamConsumer consumer) Streams the results of allquerystatements inresulttoconsumer, using the supplied connector for external relations.- Parameters:
result- a fully-valid semantic analysis result; must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row stream- Throws:
IllegalArgumentException- ifresultcontains semantic errorsEvaluationException- if a runtime data-level error occurs while a stream is read
-
executeStreaming
Streams the results of allquerystatements inresulttoconsumerusing an inline-only connector.- Parameters:
result- a fully-valid semantic analysis result; must not be nullconsumer- receives each query's label, schema, and lazy row stream- Throws:
IllegalArgumentException- ifresultcontains semantic errorsEvaluationException- if a runtime error occurs, including external-source access
-
executeStreaming
public void executeStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer) Streams the results of allquerystatements inmodeltoconsumer, using the supplied connector for external relations.Each query's row stream is opened, passed to
consumer, and closed before the next query begins. The consumer must consume the stream within the callback (seeRowStreamConsumer).- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row stream- Throws:
EvaluationException- if a runtime data-level error occurs while a stream is read
-
executeStreaming
public void executeStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds) LikeexecuteStreaming(SemanticModel, DataSourceConnector, RowStreamConsumer)but aborts anyFIXfixpoint that exceedsmaxFixpointRoundsiterations with anEvaluationException.- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row streammaxFixpointRounds- maximum semi-naïve fixpoint iterations; useExecutionContext.UNLIMITED_FIXPOINT_ROUNDSfor no cap- Throws:
EvaluationException- if a runtime data-level error occurs or a fixpoint cap is hit
-
executeStreaming
public void executeStreaming(SemanticModel model, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds, QueryEventListener listener) LikeexecuteStreaming(SemanticModel, DataSourceConnector, RowStreamConsumer, int)but reports the run'sPLANandEXECUTEdecisions tolistener, each tagged with the owning query's label.This is the plain execution path's observability hook: unlike
trace(com.darkcollective.relix.semantic.SemanticModel, java.util.List<com.darkcollective.relix.ast.RelNode>, com.darkcollective.relix.events.QueryEventListener), it neither plans a second time nor drains rows on the caller's behalf, so a host can watch an ordinary run without changing what that run does. NoOPTIMIZEevents reach the feed here — this method runs whatever trees it is given, and rewriting them is the optimizer's own pass, which reports to its own listener.- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row streammaxFixpointRounds- maximum semi-naïve fixpoint iterations; useExecutionContext.UNLIMITED_FIXPOINT_ROUNDSfor no caplistener- notified on each planning/execution decision; must not be null- Throws:
EvaluationException- if a runtime data-level error occurs or a fixpoint cap is hit
-
executeStreaming
Streams the results of allquerystatements inmodeltoconsumerusing an inline-only connector.- Parameters:
model- the semantic model; must not be nullconsumer- receives each query's label, schema, and lazy row stream- Throws:
EvaluationException- if a runtime error occurs, including external-source access
-
execute
Executes allquerystatements inresult, using the supplied connector, and collects each result into aQueryResult.- Parameters:
result- a fully-valid semantic analysis result; must not be nullconnector- provides rows for external data sources; must not be null- Returns:
- an unmodifiable list of results, one per
querystatement, in declaration order - Throws:
IllegalArgumentException- ifresultcontains semantic errorsEvaluationException- if a runtime data-level error occurs
-
execute
Executes allquerystatements inresultusing an inline-only connector and collects each result into aQueryResult.- Parameters:
result- a fully-valid semantic analysis result; must not be null- Returns:
- an unmodifiable list of results, one per
querystatement - Throws:
IllegalArgumentException- ifresultcontains semantic errorsEvaluationException- if a runtime error occurs, including external-source access
-
execute
Executes allquerystatements inmodel, using the supplied connector, and collects each result into aQueryResult.- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be null- Returns:
- an unmodifiable list of results in declaration order
- Throws:
EvaluationException- if a runtime data-level error occurs
-
execute
Executes allquerystatements inmodelusing an inline-only connector and collects each result into aQueryResult.- Parameters:
model- the semantic model; must not be null- Returns:
- an unmodifiable list of results in declaration order
- Throws:
EvaluationException- if a runtime error occurs
-
executeProvenance
public <K> void executeProvenance(SemanticModel model, DataSourceConnector connector, Semiring<K> semiring, ProvenanceConsumer<K> consumer) Evaluates eachquerystatement inmodelas an annotatedK-relationoversemiring, handing each result toconsumerin declaration order.Provenance is an explicit, opt-in mode: the chosen semiring threads through the positive-algebra operators (σ/π/×/⋈/∪) while every non-positive operator is read as an opaque base relation — see
ProvenanceEvaluator. The query is evaluated from its raw logical tree (no optimizer, no pushdown): annotation tracking is in-engine, above the federation boundary. Because a K-relation is canonical, each result is fully materialised rather than streamed.- Type Parameters:
K- the semiring annotation type- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be nullsemiring- the annotation semiring; must not be nullconsumer- receives each query's label and annotated relation; must not be null- Throws:
EvaluationException- if a runtime data-level error occurs during evaluation
-
executeProvenance
public <K> void executeProvenance(SemanticModel model, DataSourceConnector connector, Semiring<K> semiring, String weightColumn, int maxFixpointRounds, ProvenanceConsumer<K> consumer) Evaluates eachqueryas a K-relation oversemiring, supplying per-edge weights fromweightColumnfor any transitive closure in the query — the semiring-weighted-closure path.When
weightColumnis non-null, each base edge is annotated with the value of that column coerced into the semiring (adoublefor tropical shortest path, a multiplicity for ℕ; the boolean/security semirings ignore the weight). A null/absent/non-numeric weight, and every non-closure base tuple, lift toone. WhenweightColumnis null this is the plainone()lift — boolean reachability and ℕ path-counting need no weight column.maxFixpointRoundsbounds the weighted-closure iteration (seeExecutionContext.maxFixpointRounds()); it matters only for a non-idempotent semiring (ℕ) over a cyclic graph, which otherwise never converges. UseExecutionContext.UNLIMITED_FIXPOINT_ROUNDSfor no cap.- Type Parameters:
K- the semiring annotation type- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be nullsemiring- the annotation semiring; must not be nullweightColumn- the per-edge weight column for weighted closures, or nullmaxFixpointRounds- the weighted-closure iteration cap (≥ 1)consumer- receives each query's label and annotated relation; must not be null- Throws:
EvaluationException- if a runtime data-level error occurs during evaluation
-
executeLineage
public void executeLineage(SemanticModel model, DataSourceConnector connector, ProvenanceConsumer<Polynomial> consumer) Evaluates eachquerystatement inmodelas a full why-provenance K-relation over the polynomial-lineage semiringℕ[X], handing each result toconsumer.Unlike the cheap semirings, lineage mints a distinct provenance variable per base-tuple occurrence (
<relation>#<ordinal>), so each output tuple's annotation is the polynomial recording exactly which input tuples produced it and how they were combined. The representation is bounded: a polynomial that exceedsPolynomialSemiring.MAX_MONOMIALSterms is truncated and flagged. Like all provenance evaluation this runs on the raw logical tree, in-engine, never pushed down.- Parameters:
model- the semantic model; must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label and lineage-annotated relation; must not be null- Throws:
EvaluationException- if a runtime data-level error occurs during evaluation
-
executeOptimizedStreaming
public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer) Streams the results of optimizer-rewritten query trees, one per root query inSemanticModel.rootQueries(), using the supplied connector.optimizedRootsmust be parallel tomodel.rootQueries(): elementiis the rewritten tree to execute in place of root queryi. Because the optimizer produces freshRelNodeinstances absent fromSemanticModel.nodeSchemas()(which is keyed by object identity), this method re-infers schema annotations for each rewritten tree viaSchemaInference.annotate(com.darkcollective.relix.symbol.table.SymbolTable, com.darkcollective.relix.ast.RelNode, com.darkcollective.relix.semantic.SchemaAnnotations, com.darkcollective.relix.function.FunctionCatalog), retaining the original annotations so that nested view bodies referenced by — but not rewritten within — the optimized tree remain executable.- Parameters:
model- the semantic model the trees were optimized from; must not be nulloptimizedRoots- rewritten trees, parallel tomodel.rootQueries(); must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row stream- Throws:
IllegalArgumentException- ifoptimizedRootssize differs from the root-query countEvaluationException- if a runtime data-level error occurs while a stream is read
-
executeOptimizedStreaming
public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds) LikeexecuteOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer)but aborts anyFIXfixpoint that exceedsmaxFixpointRoundsiterations.- Parameters:
model- the semantic model; must not be nulloptimizedRoots- rewritten trees, parallel tomodel.rootQueries(); must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row streammaxFixpointRounds- maximum semi-naïve fixpoint iterations; useExecutionContext.UNLIMITED_FIXPOINT_ROUNDSfor no cap- Throws:
IllegalArgumentException- ifoptimizedRootssize differs from the root-query countEvaluationException- if a runtime data-level error occurs or a fixpoint cap is hit
-
executeOptimizedStreaming
public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector, RowStreamConsumer consumer, int maxFixpointRounds, QueryEventListener listener) LikeexecuteOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer, int)but reports the run'sPLANandEXECUTEdecisions tolistener, each tagged with the owning query's label — the observing counterpart for a caller that has already optimized. TheOPTIMIZEevents belong to that earlier pass and are collected there.- Parameters:
model- the semantic model the trees were optimized from; must not be nulloptimizedRoots- rewritten trees, parallel tomodel.rootQueries(); must not be nullconnector- provides rows for external data sources; must not be nullconsumer- receives each query's label, schema, and lazy row streammaxFixpointRounds- maximum semi-naïve fixpoint iterations; useExecutionContext.UNLIMITED_FIXPOINT_ROUNDSfor no caplistener- notified on each planning/execution decision; must not be null- Throws:
IllegalArgumentException- ifoptimizedRootssize differs from the root-query countEvaluationException- if a runtime data-level error occurs or a fixpoint cap is hit
-
executeOptimizedStreaming
public void executeOptimizedStreaming(SemanticModel model, List<RelNode> optimizedRoots, RowStreamConsumer consumer) Streams the results of optimizer-rewritten query trees using an inline-only connector.- Parameters:
model- the semantic model the trees were optimized from; must not be nulloptimizedRoots- rewritten trees, parallel tomodel.rootQueries(); must not be nullconsumer- receives each query's label, schema, and lazy row stream- Throws:
IllegalArgumentException- ifoptimizedRootssize differs from the root-query countEvaluationException- if a runtime error occurs, including external-source access
-
executeOptimized
public List<QueryResult> executeOptimized(SemanticModel model, List<RelNode> optimizedRoots, DataSourceConnector connector) Executes optimizer-rewritten query trees, using the supplied connector, and collects each result into aQueryResult. SeeexecuteOptimizedStreaming(SemanticModel, List, DataSourceConnector, RowStreamConsumer).- Parameters:
model- the semantic model the trees were optimized from; must not be nulloptimizedRoots- rewritten trees, parallel tomodel.rootQueries(); must not be nullconnector- provides rows for external data sources; must not be null- Returns:
- an unmodifiable list of results in declaration order
- Throws:
IllegalArgumentException- ifoptimizedRootssize differs from the root-query countEvaluationException- if a runtime data-level error occurs
-
executeOptimized
Executes optimizer-rewritten query trees using an inline-only connector and collects each result into aQueryResult.- Parameters:
model- the semantic model the trees were optimized from; must not be nulloptimizedRoots- rewritten trees, parallel tomodel.rootQueries(); must not be null- Returns:
- an unmodifiable list of results in declaration order
- Throws:
IllegalArgumentException- ifoptimizedRootssize differs from the root-query countEvaluationException- if a runtime error occurs, including external-source access
-
explain
Renders the physical plan of eachquerystatement inmodeland hands it toconsumer, without executing anything. This shows the planner's physical decisions — notably which sub-trees were pushed to the database as aPhysicalNode.PushedScanand each join's algorithm and build side.- Parameters:
model- the semantic model; must not be nullconsumer- receives each query's label and rendered plan; must not be null
-
explain
Renders the physical plan of each (optimizer-rewritten) root tree and hands it toconsumer, without executing anything. Use this overload after optimization so the explained plan matches whatexecuteOptimized*would run; passmodel.rootQueries()'s resolved nodes (viaexplain(SemanticModel, PlanConsumer)) for the un-optimized plan.- Parameters:
model- the semantic model the trees belong to; must not be nullroots- the trees to plan, parallel tomodel.rootQueries(); must not be nullconsumer- receives each query's label and rendered plan; must not be null- Throws:
IllegalArgumentException- ifrootssize differs from the root-query count
-
trace
Plans each (optimizer-rewritten) root tree — without executing it — so the planner emits its physical-decisionQueryEvents tolistener. Each event is tagged with the owning query's label.Pair this with the optimizer's listener-aware
optimizeoverload on the same listener to get one feed of both the optimizer's rule firings and the planner's decisions.- Parameters:
model- the semantic model the trees belong to; must not be nullroots- the trees to plan, parallel tomodel.rootQueries(); must not be nulllistener- notified on each physical decision; must not be null- Throws:
IllegalArgumentException- ifrootssize differs from the root-query count
-
traceExecute
public void traceExecute(SemanticModel model, List<RelNode> roots, DataSourceConnector connector, QueryEventListener listener) Executes each (optimizer-rewritten) root tree purely for itsQueryEvent.Stage.EXECUTEside-effect events — discarding the rows — emitting each tolistenertagged with the owning query's label.This is the execution-stage companion to
trace(com.darkcollective.relix.semantic.SemanticModel, java.util.List<com.darkcollective.relix.ast.RelNode>, com.darkcollective.relix.events.QueryEventListener): wheretracesurfaces the planner's decisions without running anything, this surfaces decisions made during execution (e.g. a declarative- optimisation group skipped as infeasible). Planning inside execution usesQueryEventListener.NONE, so noPLANevents are double-emitted; only theEXECUTEfeed flows here.- Parameters:
model- the semantic model the trees belong to; must not be nullroots- the trees to run, parallel tomodel.rootQueries(); must not be nullconnector- provides rows for external data sources; must not be nulllistener- notified on each execution-stage decision; must not be null- Throws:
IllegalArgumentException- ifrootssize differs from the root-query countEvaluationException- if a runtime data-level error occurs while a stream is read
-
plan
public void plan(SemanticModel model, List<RelNode> roots, QueryEventListener listener, PhysicalPlanConsumer consumer) Plans each (optimizer-rewritten) root tree — without executing it — and hands the resultingPhysicalNodetoconsumer, while emitting the planner's physical-decisionQueryEvents tolistener(each tagged with the owning query's label).This is the structured counterpart to
explain(SemanticModel, java.util.List, PlanConsumer)(which renders text): use it when the caller needs the plan object — for example to serialize it to JSON for the query bundle. PassQueryEventListener.NONEwhen the events are not wanted.- Parameters:
model- the semantic model the trees belong to; must not be nullroots- the trees to plan, parallel tomodel.rootQueries(); must not be nulllistener- notified on each physical decision; must not be nullconsumer- receives each query's label and planned physical tree; must not be null- Throws:
IllegalArgumentException- ifrootssize differs from the root-query count
-
requireBounded
Refuses to collect a relation that provably never ends.BoundednessCheckerruns in the planner and catches a blocking operator over an unbounded input. It cannot catch this one, because here the blocking is done by the consumer:σ (Naturals)contains no blocking node at all, so the plan is legal and it iscollect(java.util.function.Consumer<com.darkcollective.relix.processor.internal.RowStreamConsumer>, int)that never returns. Streaming callers are right to be allowed it — emitting an endless relation is what a generator is for — so the guard belongs on the collecting entry points alone, which is why it lives here rather than in the planner.The remedy named is the planner's own, so a user who has met either message has met both.
Package-private rather than private so it can be tested against a real
GeneratorBoundednessSourcewithout standing up a whole analysed script; the streaming entry points deliberately never call it.- Parameters:
roots- the trees about to be materialisedboundedness- the per-leaf boundedness for this run- Throws:
BoundednessException- if any root is provably unbounded
-