diff --git a/teaql-runtime/src/main/java/io/teaql/runtime/DefaultUserContext.java b/teaql-runtime/src/main/java/io/teaql/runtime/DefaultUserContext.java index e0f2f968..3f8f68b0 100644 --- a/teaql-runtime/src/main/java/io/teaql/runtime/DefaultUserContext.java +++ b/teaql-runtime/src/main/java/io/teaql/runtime/DefaultUserContext.java @@ -56,19 +56,7 @@ public Stream internalExecuteForStream(SearchRequest searc @Override @SuppressWarnings("unchecked") public Stream internalExecuteForStream(SearchRequest searchRequest, int batchSize) { - // Check whether the backing executor declares streaming capability. - // If it does, delegate to a lazy batch-pull iterator. - // If not, fall back to a fully-materialized list (safe but not lazy). - EntityDescriptor descriptor = runtime.getMetadata().resolveEntityDescriptor(searchRequest.getTypeName()); - String route = descriptor != null ? descriptor.getDataService() : null; - if (route == null || route.isEmpty()) { - route = "default"; - } - DataServiceExecutor executor = runtime.getRegistry().resolve(route); - if (executor instanceof StreamingQueryExecutor) { - return ((StreamingQueryExecutor) executor).queryForStream(this, searchRequest); - } - throw new TeaQLRuntimeException("Streaming query is not supported for route: " + route); + return runtime.internalExecuteForStream(this, searchRequest); } @Override @@ -95,12 +83,12 @@ public SmartList executeForPage( @Override public Stream executeForStream(ExecutableRequest request) { - return internalExecuteForStream(request.request()); + return runtime.executeForStream(this, request.request()); } @Override public Stream executeForStream(ExecutableRequest request, int enhanceBatchSize) { - return internalExecuteForStream(request.request(), enhanceBatchSize); + return runtime.executeForStream(this, request.request()); } @Override diff --git a/teaql-runtime/src/main/java/io/teaql/runtime/TeaQLRuntime.java b/teaql-runtime/src/main/java/io/teaql/runtime/TeaQLRuntime.java index 91a16c0f..91b43bd9 100644 --- a/teaql-runtime/src/main/java/io/teaql/runtime/TeaQLRuntime.java +++ b/teaql-runtime/src/main/java/io/teaql/runtime/TeaQLRuntime.java @@ -15,6 +15,7 @@ import io.teaql.core.reference.RoundTripReferenceCodec; import io.teaql.core.reference.RoundTripReferenceProvider; import java.util.*; +import java.util.stream.Stream; public class TeaQLRuntime { private final EntityMetaFactory metadata; @@ -216,6 +217,50 @@ public SmartList executeForPage( return rows; } + /** Executes a business-facing streaming query after the same policy gate as list queries. */ + public Stream executeForStream( + UserContext context, SearchRequest request) { + if (request.purpose() == null || request.purpose().trim().isEmpty()) { + throw new TeaQLRuntimeException( + "[PURPOSE REQUIRED] Missing .purpose() on streaming query execution."); + } + if (requestPolicy != null) { + requestPolicy.enforceSelect(context, request); + } + return executeForStreamResolved(context, request); + } + + /** + * Executes a framework-owned nested stream under an already-authorized root query. The + * nested request still passes through RequestPolicy before reaching the provider. + */ + public Stream internalExecuteForStream( + UserContext context, SearchRequest request) { + if (context.getTraceChain() == null || context.getTraceChain().isEmpty()) { + throw new TeaQLRuntimeException( + "[INTERNAL QUERY CONTEXT REQUIRED] Nested streaming query has no authorized root trace."); + } + if (requestPolicy != null) { + requestPolicy.enforceSelect(context, request); + } + return executeForStreamResolved(context, request); + } + + @SuppressWarnings("unchecked") + private Stream executeForStreamResolved( + UserContext context, SearchRequest request) { + EntityDescriptor descriptor = metadata.resolveEntityDescriptor(request.getTypeName()); + String route = descriptor != null ? descriptor.getDataService() : null; + if (route == null || route.isEmpty()) { + route = "default"; + } + DataServiceExecutor executor = registry.resolve(route); + if (executor instanceof StreamingQueryExecutor streamingQueryExecutor) { + return streamingQueryExecutor.queryForStream(context, request); + } + throw new TeaQLRuntimeException("Streaming query is not supported for route: " + route); + } + /** * Executes a framework-owned nested query under the trace established by its * already-authorized root request. Nested relation requests are generated as diff --git a/teaql-runtime/src/test/java/io/teaql/runtime/TeaQLRuntimeTest.java b/teaql-runtime/src/test/java/io/teaql/runtime/TeaQLRuntimeTest.java index 1233071f..cbb7f7ca 100644 --- a/teaql-runtime/src/test/java/io/teaql/runtime/TeaQLRuntimeTest.java +++ b/teaql-runtime/src/test/java/io/teaql/runtime/TeaQLRuntimeTest.java @@ -15,6 +15,7 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicInteger; import java.util.Map; +import java.util.stream.Stream; public class TeaQLRuntimeTest { @@ -124,6 +125,21 @@ public DataServiceCapabilities capabilities() { } } + public static class RecordingStreamingQueryExecutor implements StreamingQueryExecutor { + private SearchRequest request; + + @Override + public Stream queryForStream( + UserContext context, SearchRequest request) { + this.request = request; + return Stream.empty(); + } + + @Override public String name() { return "dummy"; } + + @Override public DataServiceCapabilities capabilities() { return null; } + } + public static class RecordingMutationExecutor implements MutationExecutor { public final List requests = new ArrayList<>(); @@ -392,6 +408,74 @@ public void pagedExecutionReturnsRowsAndExactPolicyFilteredTotal() { executor.requests.get(1).getSearchCriteria()); } + @Test + public void streamingExecutionAppliesRequestPolicyBeforeProvider() { + RecordingStreamingQueryExecutor executor = new RecordingStreamingQueryExecutor(); + AtomicInteger policyCalls = new AtomicInteger(); + TeaQLRuntime runtime = TeaQLRuntime.builder() + .metadata(new DummyMetaFactory()) + .dataService("dummy", executor) + .requestPolicy(new RequestPolicy() { + @Override public void enforceSelect( + UserContext context, SearchRequest query) { + policyCalls.incrementAndGet(); + BaseRequest request = (BaseRequest) query; + request.appendSearchCriteria(request.createBasicSearchCriteria( + "status", Operator.EQUAL, "ACTIVE")); + } + }) + .build(); + BaseRequest request = new BaseRequest<>(DummyEntity.class) { + { internalComment("stream active entities"); } + @Override public String getTypeName() { return "Dummy"; } + }; + + try (Stream ignored = request + .purpose("prove policy cannot be bypassed by streaming") + .executeForStream(new DefaultUserContext(runtime))) { + Assert.assertEquals(0, ignored.count()); + } + + Assert.assertEquals(1, policyCalls.get()); + Assert.assertSame(request, executor.request); + Assert.assertNotNull(executor.request.getSearchCriteria()); + } + + @Test + public void internalStreamingRequiresAuthorizedRootAndAppliesRequestPolicy() { + RecordingStreamingQueryExecutor executor = new RecordingStreamingQueryExecutor(); + AtomicInteger policyCalls = new AtomicInteger(); + TeaQLRuntime runtime = TeaQLRuntime.builder() + .metadata(new DummyMetaFactory()) + .dataService("dummy", executor) + .requestPolicy(new RequestPolicy() { + @Override public void enforceSelect( + UserContext context, SearchRequest query) { + policyCalls.incrementAndGet(); + } + }) + .build(); + DefaultUserContext context = new DefaultUserContext(runtime); + SearchRequest request = bareDummyRequest(); + + try { + context.internalExecuteForStream(request); + Assert.fail("internal stream without a trusted root trace must be rejected"); + } catch (TeaQLRuntimeException expected) { + Assert.assertTrue(expected.getMessage().contains("INTERNAL QUERY CONTEXT REQUIRED")); + } + Assert.assertEquals(0, policyCalls.get()); + + context.pushTrace("authorized root query"); + try (Stream ignored = context.internalExecuteForStream(request)) { + Assert.assertEquals(0, ignored.count()); + } finally { + context.popTrace(); + } + Assert.assertEquals(1, policyCalls.get()); + Assert.assertSame(request, executor.request); + } + @Test public void materializedListHardLimitRejectsClientOverride() { TeaQLRuntime runtime = TeaQLRuntime.builder().metadata(new DummyMetaFactory())