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 @@ -77,8 +77,10 @@ public Optional<QueryEstimateSnapshot> estimateSubmittedQuery(UUID queryRequestI
var request = buildRequest(snapshot);
var dryRun = queryExecutor.dryRun(request);
Long affectedRows = countAffectedRows(request, snapshot.queryType());
return persistAndPublish(queryRequestId,
toCommand(snapshot, dryRun, affectedRows, durationMs(start)), true);
boolean redactPredicates = !request.rowSecurityPredicates().isEmpty()
|| !dryRun.appliedRowSecurityPolicyIds().isEmpty();
return persistAndPublish(queryRequestId, toCommand(snapshot, dryRun, affectedRows,
redactPredicates, durationMs(start)), true);
} catch (RuntimeException ex) {
log.warn("Cost estimate failed for query {}: {}", queryRequestId, ex.getMessage());
var command = new PersistQueryEstimateCommand(null, snapshot.queryType(), false, null,
Expand Down Expand Up @@ -119,7 +121,7 @@ private Long countAffectedRows(QueryExecutionRequest request, QueryType queryTyp

private PersistQueryEstimateCommand toCommand(QueryRequestSnapshot snapshot,
QueryDryRunResult dryRun, Long affectedRows,
int durationMs) {
boolean redactPredicates, int durationMs) {
if (!dryRun.supported()) {
var reason = dryRun.unsupportedReason() != null
? dryRun.unsupportedReason()
Expand All @@ -142,8 +144,8 @@ private PersistQueryEstimateCommand toCommand(QueryRequestSnapshot snapshot,
affectedRows,
access != null ? truncateTo(access.operation(), 128) : null,
root != null ? root.estimatedCost() : null,
planJson(root),
dryRun.rawPlan(), null, false, null, durationMs);
planJson(root, redactPredicates),
redactPredicates ? null : dryRun.rawPlan(), null, false, null, durationMs);
}

/**
Expand Down Expand Up @@ -176,15 +178,19 @@ private static PersistQueryEstimateCommand withAffectedRows(PersistQueryEstimate
command.unsupportedReason(), command.failed(), command.errorMessage(), durationMs);
}

/** Serializes the plan tree with explicit snake_case keys — the frontend's PlanTree shape. */
private String planJson(QueryPlanNode root) {
/**
* Serializes the plan tree with explicit snake_case keys — the frontend's PlanTree shape. With
* {@code redactPredicates} every node's {@code detail} is dropped: the dry-run bound the
* submitter's row-security values, and engines inline them into predicate text (#1092).
*/
private String planJson(QueryPlanNode root, boolean redactPredicates) {
if (root == null) {
return null;
}
return objectMapper.writeValueAsString(planNode(root));
return objectMapper.writeValueAsString(planNode(root, redactPredicates));
}

private ObjectNode planNode(QueryPlanNode node) {
private ObjectNode planNode(QueryPlanNode node, boolean redactPredicates) {
var out = objectMapper.createObjectNode();
out.put("operation", node.operation());
out.put("target", node.target());
Expand All @@ -198,10 +204,10 @@ private ObjectNode planNode(QueryPlanNode node) {
} else {
out.putNull("estimated_cost");
}
out.put("detail", node.detail());
out.put("detail", redactPredicates ? null : node.detail());
var children = out.putArray("children");
for (var child : node.children()) {
children.add(planNode(child));
children.add(planNode(child, redactPredicates));
}
return out;
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
-- #1092: the cost-estimate dry-run binds the submitter's row-security values, and engines inline
-- them into plan predicate text (PostgreSQL Index Cond / Filter, MySQL attached_condition,
-- MongoDB stage filters). New estimates drop that text whenever row security applied; which
-- existing rows had it was never recorded, so strip every stored plan's per-node `detail` and
-- raw plan. Operation, target, row and cost figures are kept.

CREATE FUNCTION pg_temp.af_strip_plan_detail(node JSONB) RETURNS JSONB
LANGUAGE plpgsql IMMUTABLE AS $$
BEGIN
IF node IS NULL OR jsonb_typeof(node) <> 'object' THEN
RETURN node;
END IF;
RETURN node || jsonb_build_object(
'detail', NULL::JSONB,
'children', COALESCE(
(SELECT jsonb_agg(pg_temp.af_strip_plan_detail(child) ORDER BY ord)
FROM jsonb_array_elements(
CASE WHEN jsonb_typeof(node -> 'children') = 'array'
THEN node -> 'children' ELSE '[]'::JSONB END)
WITH ORDINALITY AS c(child, ord)),
'[]'::JSONB));
END;
$$;

UPDATE query_estimates
SET plan = pg_temp.af_strip_plan_detail(plan),
raw_plan = NULL
WHERE plan IS NOT NULL OR raw_plan IS NOT NULL;

DROP FUNCTION pg_temp.af_strip_plan_detail(JSONB);
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
package com.bablsoft.accessflow.core.internal.persistence;

import org.flywaydb.core.Flyway;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.springframework.jdbc.core.ConnectionCallback;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.datasource.DriverManagerDataSource;
import org.testcontainers.postgresql.PostgreSQLContainer;
import tools.jackson.databind.json.JsonMapper;

import java.util.Map;
import java.util.UUID;

import static org.assertj.core.api.Assertions.assertThat;

/**
* Upgrade test for V188 (#1092): drives Flyway to V187 on a private container, seeds estimates
* whose plans carry inlined row-security values, then applies V188.
*/
class QueryEstimatePlanRedactionMigrationIntegrationTest {

private static final Map<String, String> PLACEHOLDERS = Map.of(
"app_role", "accessflow",
"audit_role", "accessflow_audit",
"rag_pgvector_dimensions", "1536");

@SuppressWarnings("resource")
static PostgreSQLContainer postgres = new PostgreSQLContainer("pgvector/pgvector:pg18")
.withInitScript("db/test-init-audit-roles.sql");

static JdbcTemplate jdbc;

@BeforeAll
static void migrateToTheVersionBeforeV188() {
postgres.start();
jdbc = new JdbcTemplate(new DriverManagerDataSource(
postgres.getJdbcUrl(), postgres.getUsername(), postgres.getPassword()));
flyway("187").migrate();
}

@AfterAll
static void stopContainer() {
postgres.stop();
}

private static Flyway flyway(String target) {
return Flyway.configure()
.dataSource(postgres.getJdbcUrl(), postgres.getUsername(), postgres.getPassword())
.locations("classpath:db/migration")
.placeholders(PLACEHOLDERS)
.target(target)
.load();
}

@Test
void v188StripsEveryNodeDetailAndRawPlanButKeepsPlanFigures() {
var nested = estimate("""
{"operation": "Nested Loop", "target": null, "estimated_rows": 10.0,
"estimated_cost": 8.5, "detail": "(o.user_id = u.id)",
"children": [
{"operation": "Index Scan", "target": "users", "estimated_rows": 1.0,
"estimated_cost": 2.0,
"detail": "((email)::text = 'dana@acme.example'::text)", "children": []},
{"operation": "Seq Scan", "target": "orders", "estimated_rows": 9.0,
"estimated_cost": 4.0, "detail": null, "children": []}]}
""", "[{\"Plan\": {\"Filter\": \"'dana@acme.example'\"}}]");
var rawOnly = estimate(null, "EXPLAIN text 'dana@acme.example'");
var unsupported = estimate(null, null);

flyway("188").migrate();

String plan = jdbc.queryForObject(
"SELECT plan::text FROM query_estimates WHERE id = ?", String.class, nested);
assertThat(plan).doesNotContain("dana@acme.example").doesNotContain("o.user_id");
var root = JsonMapper.builder().build().readTree(plan);
assertThat(root.get("detail").isNull()).isTrue();
assertThat(root.get("operation").asString()).isEqualTo("Nested Loop");
assertThat(root.get("estimated_cost").asDouble()).isEqualTo(8.5);
assertThat(root.get("children")).hasSize(2);
assertThat(root.get("children").get(0).get("operation").asString())
.isEqualTo("Index Scan");
assertThat(root.get("children").get(0).get("target").asString()).isEqualTo("users");
assertThat(root.get("children").get(0).get("detail").isNull()).isTrue();
assertThat(root.get("children").get(1).get("operation").asString())
.isEqualTo("Seq Scan");
assertThat(rawPlan(nested)).isNull();
assertThat(rawPlan(rawOnly)).isNull();
assertThat(jdbc.queryForObject("SELECT count(*) FROM query_estimates WHERE id = ?",
Long.class, unsupported)).isEqualTo(1L);
}

private static String rawPlan(UUID id) {
return jdbc.queryForObject("SELECT raw_plan FROM query_estimates WHERE id = ?",
String.class, id);
}

/** Seeds an estimate without its FK chain — the migration only rewrites query_estimates. */
private static UUID estimate(String planJson, String rawPlan) {
var id = UUID.randomUUID();
jdbc.execute((ConnectionCallback<Void>) connection -> {
try (var statement = connection.createStatement()) {
statement.execute("SET session_replication_role = replica");
}
try (var insert = connection.prepareStatement("""
INSERT INTO query_estimates (id, query_request_id, supported, plan, raw_plan)
VALUES (?, ?, true, CAST(? AS JSONB), ?)
""")) {
insert.setObject(1, id);
insert.setObject(2, UUID.randomUUID());
insert.setString(3, planJson);
insert.setString(4, rawPlan);
insert.executeUpdate();
}
try (var statement = connection.createStatement()) {
statement.execute("SET session_replication_role = DEFAULT");
}
return null;
});
return id;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@
import com.bablsoft.accessflow.core.api.QueryRequestSnapshot;
import com.bablsoft.accessflow.core.api.QueryStatus;
import com.bablsoft.accessflow.core.api.QueryType;
import com.bablsoft.accessflow.core.api.ResolvedRowSecurityPredicate;
import com.bablsoft.accessflow.core.api.RowSecurityOperator;
import com.bablsoft.accessflow.core.api.RowSecurityResolutionService;
import com.bablsoft.accessflow.core.events.QueryEstimateCompletedEvent;
import com.bablsoft.accessflow.core.events.QueryEstimateFailedEvent;
Expand Down Expand Up @@ -189,6 +191,72 @@ void writePlanDescendsIntoAccessNodeForScanTypeAndEstimate() {
assertThat(captor.getValue().estimatedRows()).isEqualTo(2_400_000L);
}

@Test
void keepsPredicateDetailAndRawPlanWithoutRowSecurity() {
stubSelectDryRun(List.of(), Set.of());

service.estimateSubmittedQuery(queryRequestId);

var command = capturePersisted();
assertThat(command.planJson()).contains("(active = false)").contains("(id = 7)");
assertThat(command.rawPlan()).isEqualTo("[raw (active = false)]");
}

@Test
void dropsPredicateDetailAndRawPlanWhenRowSecurityResolved() {
var policyId = UUID.randomUUID();
stubSelectDryRun(List.of(new ResolvedRowSecurityPredicate(policyId, "public.users",
"email", RowSecurityOperator.EQUALS, List.of("dana@acme.example"))),
Set.of(policyId));

service.estimateSubmittedQuery(queryRequestId);

var command = capturePersisted();
assertThat(command.planJson()).doesNotContain("active = false")
.doesNotContain("id = 7")
.contains("\"operation\":\"Nested Loop\"")
.contains("\"operation\":\"Index Scan\"")
.contains("\"detail\":null");
assertThat(command.rawPlan()).isNull();
assertThat(command.scanType()).isEqualTo("Nested Loop");
assertThat(command.estimatedRows()).isEqualTo(10L);
}

@Test
void dropsPredicateDetailWhenEngineReportsAppliedPolicyWithoutResolvedDirective() {
stubSelectDryRun(List.of(), Set.of(UUID.randomUUID()));

service.estimateSubmittedQuery(queryRequestId);

var command = capturePersisted();
assertThat(command.planJson()).doesNotContain("active = false");
assertThat(command.rawPlan()).isNull();
}

private void stubSelectDryRun(List<ResolvedRowSecurityPredicate> predicates,
Set<UUID> appliedPolicyIds) {
when(lookupService.findByQueryRequestId(queryRequestId))
.thenReturn(Optional.empty())
.thenReturn(Optional.of(persistedSnapshot()));
when(queryRequestLookupService.findById(queryRequestId))
.thenReturn(Optional.of(snapshot(QueryType.SELECT, false)));
when(rowSecurityResolutionService.resolveApplicable(organizationId, datasourceId, userId))
.thenReturn(predicates);
var index = new QueryPlanNode("Index Scan", "orders", 1.0, 2.0, "(id = 7)");
var root = new QueryPlanNode("Nested Loop", null, 10.0, 8.0, "(active = false)",
List.of(index));
when(queryExecutor.dryRun(any())).thenReturn(QueryDryRunResult.of("postgresql",
QueryType.SELECT, 10L, root, "[raw (active = false)]", appliedPolicyIds,
Duration.ZERO));
when(persistenceService.persist(eq(queryRequestId), any())).thenReturn(estimateId);
}

private PersistQueryEstimateCommand capturePersisted() {
var captor = ArgumentCaptor.forClass(PersistQueryEstimateCommand.class);
verify(persistenceService).persist(eq(queryRequestId), captor.capture());
return captor.getValue();
}

@Test
void selectSkipsAffectedRowCount() {
when(lookupService.findByQueryRequestId(queryRequestId))
Expand Down
4 changes: 2 additions & 2 deletions docs/03-data-model.md
Original file line number Diff line number Diff line change
Expand Up @@ -1158,8 +1158,8 @@ Persisted pre-flight cost / blast-radius estimate for a submitted query (AF-624,
| `affected_row_count` | BIGINT nullable — exact governed count for UPDATE/DELETE; null when the shape can't be provably counted (joins, `USING`, MERGE, …) or the engine doesn't support counting |
| `scan_type` | VARCHAR(128) nullable — the plan's root operation (e.g. `Seq Scan`, `COLLSCAN`) |
| `estimated_cost` | DOUBLE PRECISION nullable — the plan's root cost when the engine exposes one |
| `plan` | JSONB nullable — the snake_case plan-node tree (same shape as the dry-run endpoint's `plan`) |
| `raw_plan` | TEXT nullable — the engine's raw plan output |
| `plan` | JSONB nullable — the snake_case plan-node tree (same shape as the dry-run endpoint's `plan`). Every node's `detail` is null when row security applied to the submitter, because engines inline the bound values into predicate text (#1092; V188 stripped all rows stored earlier) |
| `raw_plan` | TEXT nullable — the engine's raw plan output; null when row security applied (#1092) |
| `unsupported_reason` | VARCHAR(500) nullable — localized reason when `supported=false` |
| `failed` | BOOLEAN NOT NULL DEFAULT false — true when the computation hit an unexpected error (sentinel row) |
| `error_message` | VARCHAR(500) nullable — failure reason when `failed=true` |
Expand Down
Loading
Loading