Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -56,19 +56,7 @@ public <T extends Entity> Stream<T> internalExecuteForStream(SearchRequest searc
@Override
@SuppressWarnings("unchecked")
public <T extends Entity> Stream<T> 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
Expand All @@ -95,12 +83,12 @@ public <T extends Entity> SmartList<T> executeForPage(

@Override
public <T extends Entity> Stream<T> executeForStream(ExecutableRequest<T> request) {
return internalExecuteForStream(request.request());
return runtime.executeForStream(this, request.request());
}

@Override
public <T extends Entity> Stream<T> executeForStream(ExecutableRequest<T> request, int enhanceBatchSize) {
return internalExecuteForStream(request.request(), enhanceBatchSize);
return runtime.executeForStream(this, request.request());
}

@Override
Expand Down
45 changes: 45 additions & 0 deletions teaql-runtime/src/main/java/io/teaql/runtime/TeaQLRuntime.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -216,6 +217,50 @@ public <T extends Entity> SmartList<T> executeForPage(
return rows;
}

/** Executes a business-facing streaming query after the same policy gate as list queries. */
public <T extends Entity> Stream<T> executeForStream(
UserContext context, SearchRequest<T> 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 <T extends Entity> Stream<T> internalExecuteForStream(
UserContext context, SearchRequest<T> 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 <T extends Entity> Stream<T> executeForStreamResolved(
UserContext context, SearchRequest<T> 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
Expand Down
84 changes: 84 additions & 0 deletions teaql-runtime/src/test/java/io/teaql/runtime/TeaQLRuntimeTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand Down Expand Up @@ -124,6 +125,21 @@ public DataServiceCapabilities capabilities() {
}
}

public static class RecordingStreamingQueryExecutor implements StreamingQueryExecutor {
private SearchRequest<?> request;

@Override
public <T extends Entity> Stream<T> queryForStream(
UserContext context, SearchRequest<T> 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<DefaultMutationRequest> requests = new ArrayList<>();

Expand Down Expand Up @@ -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<DummyEntity> request = new BaseRequest<>(DummyEntity.class) {
{ internalComment("stream active entities"); }
@Override public String getTypeName() { return "Dummy"; }
};

try (Stream<DummyEntity> 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<DummyEntity> 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<DummyEntity> 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())
Expand Down
Loading