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 @@ -8,7 +8,9 @@

/**
* One executed DDL/DELETE operation with its approvers, for the
* {@link ComplianceReportType#REGULATORY_AUDIT_TRAIL} report (#459).
* {@link ComplianceReportType#REGULATORY_AUDIT_TRAIL} report (#459). {@code effectiveSql} is the
* statement as actually executed, bound values redacted as {@code ?}, or {@code null} when no
* row-security / soft-delete rewrite occurred (#937).
*/
public record RegulatoryAuditTrailRow(
UUID queryRequestId,
Expand All @@ -19,7 +21,8 @@ public record RegulatoryAuditTrailRow(
QueryType queryType,
String sqlText,
List<Approver> approvers,
Instant executedAt) {
Instant executedAt,
String effectiveSql) {

public RegulatoryAuditTrailRow {
approvers = approvers == null ? List.of() : List.copyOf(approvers);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,8 @@ private void writeClassifiedAccess(StringBuilder sb, List<ClassifiedAccessReport

private void writeAuditTrail(StringBuilder sb, List<RegulatoryAuditTrailRow> rows) {
writeRow(sb, List.of("query_request_id", "datasource_id", "datasource_name",
"submitter_email", "query_type", "sql_text", "approvers", "executed_at"));
"submitter_email", "query_type", "sql_text", "effective_sql", "approvers",
"executed_at"));
for (var row : rows) {
writeRow(sb, List.of(
str(row.queryRequestId()),
Expand All @@ -62,6 +63,7 @@ private void writeAuditTrail(StringBuilder sb, List<RegulatoryAuditTrailRow> row
nullToEmpty(row.submitterEmail()),
str(row.queryType()),
nullToEmpty(row.sqlText()),
nullToEmpty(row.effectiveSql()),
formatApprovers(row.approvers()),
str(row.executedAt())));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -95,14 +95,16 @@ private void drawClassifiedAccess(Canvas c, List<ClassifiedAccessReportRow> rows
}

private void drawAuditTrail(Canvas c, List<RegulatoryAuditTrailRow> rows) throws IOException {
var headers = List.of("Executed At", "Datasource", "Submitter", "Type", "SQL", "Approvers");
float[] weights = {2f, 2f, 2f, 1f, 4f, 3f};
var headers = List.of("Executed At", "Datasource", "Submitter", "Type", "SQL",
"Effective SQL", "Approvers");
float[] weights = {2f, 2f, 2f, 1f, 3f, 3f, 2f};
var data = new ArrayList<List<String>>();
for (var row : rows) {
data.add(List.of(
str(row.executedAt()), nullToEmpty(row.datasourceName()),
nullToEmpty(row.submitterEmail()), str(row.queryType()),
nullToEmpty(row.sqlText()), formatApprovers(row.approvers())));
nullToEmpty(row.sqlText()), nullToEmpty(row.effectiveSql()),
formatApprovers(row.approvers())));
}
drawTable(c, headers, weights, data);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,8 @@ private ComplianceReport regulatoryTrail(UUID organizationId, ComplianceReportRe
snapshot.queryType(),
snapshot.sqlText(),
reviewDecisionsParser.approvers(snapshot.reviewDecisionsJson()),
snapshot.executedAt()));
snapshot.executedAt(),
snapshot.effectiveSql()));
}

return new ComplianceReport(request.type(), organizationId, request.from(), request.to(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,8 @@ public record RegulatoryAuditTrailRow(
QueryType queryType,
String sqlText,
List<Approver> approvers,
Instant executedAt) {
Instant executedAt,
String effectiveSql) {
}

public record RetentionAdherenceRow(
Expand Down Expand Up @@ -93,7 +94,7 @@ public static ComplianceReportResponse from(ComplianceReport report) {
.map(a -> new Approver(a.email(), a.displayName(), a.decision(),
a.decidedAt()))
.toList(),
r.executedAt()))
r.executedAt(), r.effectiveSql()))
.toList();
var retention = report.retentionAdherence().stream()
.map(r -> new RetentionAdherenceRow(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,8 @@ public record SelectExecutionResult(
Duration duration,
Set<UUID> appliedMaskingPolicyIds,
Set<UUID> appliedRowSecurityPolicyIds,
String truncatedReason) implements QueryExecutionResult {
String truncatedReason,
String effectiveSql) implements QueryExecutionResult {

/** {@link #truncatedReason()} value when the configured row cap cut the result short. */
public static final String TRUNCATED_ROW_LIMIT = "ROW_LIMIT";
Expand All @@ -31,26 +32,46 @@ public record SelectExecutionResult(

public SelectExecutionResult(List<ResultColumn> columns, List<List<Object>> rows, long rowCount,
boolean truncated, Duration duration) {
this(columns, rows, rowCount, truncated, duration, Set.of(), Set.of(), null);
this(columns, rows, rowCount, truncated, duration, Set.of(), Set.of(), null, null);
}

public SelectExecutionResult(List<ResultColumn> columns, List<List<Object>> rows, long rowCount,
boolean truncated, Duration duration,
Set<UUID> appliedMaskingPolicyIds) {
this(columns, rows, rowCount, truncated, duration, appliedMaskingPolicyIds, Set.of(), null);
this(columns, rows, rowCount, truncated, duration, appliedMaskingPolicyIds, Set.of(), null,
null);
}

public SelectExecutionResult(List<ResultColumn> columns, List<List<Object>> rows, long rowCount,
boolean truncated, Duration duration,
Set<UUID> appliedMaskingPolicyIds,
Set<UUID> appliedRowSecurityPolicyIds) {
this(columns, rows, rowCount, truncated, duration, appliedMaskingPolicyIds,
appliedRowSecurityPolicyIds, null);
appliedRowSecurityPolicyIds, null, null);
}

/** Pre-#937 canonical shape — kept so published engine plugins stay binary-compatible. */
public SelectExecutionResult(List<ResultColumn> columns, List<List<Object>> rows, long rowCount,
boolean truncated, Duration duration,
Set<UUID> appliedMaskingPolicyIds,
Set<UUID> appliedRowSecurityPolicyIds, String truncatedReason) {
this(columns, rows, rowCount, truncated, duration, appliedMaskingPolicyIds,
appliedRowSecurityPolicyIds, truncatedReason, null);
}

/** Returns a copy of this result with the given row-security policy ids attached. */
public SelectExecutionResult withRowSecurityPolicyIds(Set<UUID> ids) {
return new SelectExecutionResult(columns, rows, rowCount, truncated, duration,
appliedMaskingPolicyIds, ids, truncatedReason);
appliedMaskingPolicyIds, ids, truncatedReason, effectiveSql);
}

/**
* Returns a copy carrying the statement as actually executed (#937) — the row-security /
* soft-delete rewrite with bound values left as {@code ?}; {@code null} when nothing was
* rewritten.
*/
public SelectExecutionResult withEffectiveSql(String sql) {
return new SelectExecutionResult(columns, rows, rowCount, truncated, duration,
appliedMaskingPolicyIds, appliedRowSecurityPolicyIds, truncatedReason, sql);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@
import java.util.UUID;

public record UpdateExecutionResult(long rowsAffected, Duration duration,
Set<UUID> appliedRowSecurityPolicyIds)
Set<UUID> appliedRowSecurityPolicyIds,
String effectiveSql)
implements QueryExecutionResult {

public UpdateExecutionResult {
Expand All @@ -14,6 +15,12 @@ public record UpdateExecutionResult(long rowsAffected, Duration duration,
}

public UpdateExecutionResult(long rowsAffected, Duration duration) {
this(rowsAffected, duration, Set.of());
this(rowsAffected, duration, Set.of(), null);
}

/** Pre-#937 canonical shape — kept so published engine plugins stay binary-compatible. */
public UpdateExecutionResult(long rowsAffected, Duration duration,
Set<UUID> appliedRowSecurityPolicyIds) {
this(rowsAffected, duration, appliedRowSecurityPolicyIds, null);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ private QueryExecutionResult executeInternal(QueryExecutionRequest request,
}
var rewrite = rowSecurityRewriter.rewrite(request.sql(), request.rowSecurityPredicates(),
request.softDeleteDirectives());
String effectiveSql = effectiveSql(rewrite, request.sql());
// SELECT result cache (AF-457): keyed over the RLS-rewritten SQL + binds + mask/restriction
// directives + row cap, so security scope is part of the key. SELECTs whose referenced
// tables are unknown are never cached (no write-invalidation coverage).
Expand All @@ -127,7 +128,7 @@ private QueryExecutionResult executeInternal(QueryExecutionRequest request,
var hit = resultCache.get(request.datasourceId(), cacheKey, durationSince(start));
if (hit.isPresent()) {
observation.lowCardinalityKeyValue("cache", "hit");
return hit.get();
return hit.get().withEffectiveSql(effectiveSql);
}
}
observation.lowCardinalityKeyValue("cache", cacheable ? "miss" : "off");
Expand All @@ -143,13 +144,13 @@ private QueryExecutionResult executeInternal(QueryExecutionRequest request,
execProps.maxResultBytes(), descriptor.dbType(), start,
request.restrictedColumns(), request.columnMasks(),
rewrite.appliedPolicyIds());
if (cacheable && result instanceof SelectExecutionResult select) {
if (cacheable) {
resultCache.put(request.datasourceId(), cacheKey,
request.referencedTables(), resultCache.ttlFor(descriptor), select);
request.referencedTables(), resultCache.ttlFor(descriptor), result);
}
return result;
return result.withEffectiveSql(effectiveSql);
}
var result = runUpdate(statement, start, rewrite.appliedPolicyIds());
var result = runUpdate(statement, start, rewrite.appliedPolicyIds(), effectiveSql);
// Any successful write drops cached SELECTs over the touched tables (unknown
// tables ⇒ full-datasource purge, fail-safe for DDL).
resultCache.invalidateTables(request.datasourceId(), request.referencedTables());
Expand Down Expand Up @@ -318,7 +319,8 @@ private QueryExecutionResult executeTransactional(QueryExecutionRequest request,
}
throw ex;
}
return new UpdateExecutionResult(totalAffected, durationSince(start), appliedPolicyIds);
return new UpdateExecutionResult(totalAffected, durationSince(start), appliedPolicyIds,
effectiveBatchSql(statements, rewrites));
} catch (SQLException ex) {
log.debug("Transactional SQL execution failed for datasource {}: {}",
request.datasourceId(), ex.getMessage());
Expand Down Expand Up @@ -390,7 +392,28 @@ private static long sumBatchCounts(long[] counts) {
return sum;
}

private QueryExecutionResult runSelect(PreparedStatement statement, int effectiveMaxRows,
/**
* The statement as actually executed (#937): the rewriter already deparses bound predicate
* values as {@code ?} placeholders, so this is the redacted form. {@code null} when nothing was
* rewritten, so an unrewritten query never stores a copy of its own SQL.
*/
private static String effectiveSql(RowSecurityRewriter.RewriteResult rewrite, String original) {
return rewrite.sql().equals(original) ? null : rewrite.sql();
}

/** Whole-batch effective form — every statement, joined — or {@code null} when none changed. */
private static String effectiveBatchSql(List<String> statements,
RowSecurityRewriter.RewriteResult[] rewrites) {
boolean anyRewritten = false;
var joined = new java.util.StringJoiner(";\n");
for (int i = 0; i < rewrites.length; i++) {
anyRewritten |= !rewrites[i].sql().equals(statements.get(i));
joined.add(rewrites[i].sql());
}
return anyRewritten ? joined.toString() : null;
}

private SelectExecutionResult runSelect(PreparedStatement statement, int effectiveMaxRows,
long maxResultBytes, DbType dbType, Instant start,
List<String> restrictedColumns,
List<ColumnMaskDirective> columnMasks,
Expand All @@ -407,10 +430,12 @@ private QueryExecutionResult runSelect(PreparedStatement statement, int effectiv
}

private UpdateExecutionResult runUpdate(PreparedStatement statement, Instant start,
java.util.Set<java.util.UUID> appliedRowSecurityPolicyIds)
java.util.Set<java.util.UUID> appliedRowSecurityPolicyIds,
String effectiveSql)
throws SQLException {
long affected = statement.executeLargeUpdate();
return new UpdateExecutionResult(affected, durationSince(start), appliedRowSecurityPolicyIds);
return new UpdateExecutionResult(affected, durationSince(start), appliedRowSecurityPolicyIds,
effectiveSql);
}

private static void bind(PreparedStatement statement, List<Object> binds) throws SQLException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,11 @@ public interface QuerySnapshotService {
/**
* Records the immutable snapshot for an executed query. Idempotent: a no-op when a snapshot already
* exists for the query (safe under event redelivery). Never throws to its caller — failures are
* logged and swallowed so snapshot capture cannot disrupt query execution.
* logged and swallowed so snapshot capture cannot disrupt query execution. {@code effectiveSql} is the
* statement as actually executed (row-security / soft-delete rewrite, bound values redacted as
* {@code ?}), or {@code null} when nothing was rewritten (#937).
*/
void recordOnExecution(UUID queryRequestId);
void recordOnExecution(UUID queryRequestId, String effectiveSql);

/** Loads the snapshot for a query, scoped to the organization. */
Optional<QuerySnapshotView> find(UUID queryRequestId, UUID organizationId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,9 @@
* Read-side view of an immutable {@code query_snapshots} row (AF-449). The AI analysis and review
* decisions are carried as raw JSON strings (as {@code QueryDetailView} does for {@code issuesJson}) so
* this api type stays free of any third-party dependency. {@code referencedTables} is the set of tables
* the snapshotted query touched, used by the replay schema-compatibility gate.
* the snapshotted query touched, used by the replay schema-compatibility gate. {@code effectiveSql} is the
* statement as actually executed with bound values redacted as {@code ?}, or {@code null} when no
* rewrite occurred (#937).
*/
public record QuerySnapshotView(
UUID id,
Expand All @@ -30,7 +32,8 @@ public record QuerySnapshotView(
Long rowsAffected,
Integer executionDurationMs,
Instant executedAt,
Instant createdAt) {
Instant createdAt,
String effectiveSql) {

public QuerySnapshotView {
referencedTables = referencedTables == null ? List.of() : List.copyOf(referencedTables);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@
* with a runtime failure ({@code finalStatus = FAILED}). Drives the realtime
* {@code query.executed} push to the submitter. {@code recurringParentId} is non-null only for
* occurrence rows of a recurring series (#627) — it gates the result-delivery notification.
* {@code effectiveSql} is the statement as actually executed — row-security / soft-delete
* rewrite, bound values redacted as {@code ?} — or {@code null} when nothing was rewritten (#937).
*
* <p>Published outside any transaction, so consumers must be plain {@code @EventListener}s —
* an {@code @ApplicationModuleListener} (AFTER_COMMIT) would silently never fire.
Expand All @@ -18,11 +20,18 @@ public record QueryExecutedEvent(
Long rowsAffected,
long durationMs,
QueryStatus finalStatus,
UUID recurringParentId) {
UUID recurringParentId,
String effectiveSql) {

/** Backward-compatible constructor without the #627 series back-pointer. */
public QueryExecutedEvent(UUID queryRequestId, Long rowsAffected, long durationMs,
QueryStatus finalStatus) {
this(queryRequestId, rowsAffected, durationMs, finalStatus, null);
this(queryRequestId, rowsAffected, durationMs, finalStatus, null, null);
}

/** Backward-compatible constructor without the #937 effective statement. */
public QueryExecutedEvent(UUID queryRequestId, Long rowsAffected, long durationMs,
QueryStatus finalStatus, UUID recurringParentId) {
this(queryRequestId, rowsAffected, durationMs, finalStatus, recurringParentId, null);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -348,17 +348,20 @@ private ExecutionOutcome doExecute(QueryRequestSnapshot query, UUID actorUserId,
Set<UUID> appliedMaskingPolicyIds = Set.of();
Set<UUID> appliedRowSecurityPolicyIds;
Set<UUID> appliedRowLimitPolicyIds = Set.of();
String effectiveSql;
switch (result) {
case SelectExecutionResult select -> {
rowsAffected = select.rowCount();
appliedRowLimitPolicyIds = bindingRowLimitPolicyIds;
appliedMaskingPolicyIds = select.appliedMaskingPolicyIds();
appliedRowSecurityPolicyIds = select.appliedRowSecurityPolicyIds();
effectiveSql = select.effectiveSql();
persistSelectResult(query.id(), select, durationMs);
}
case UpdateExecutionResult update -> {
rowsAffected = update.rowsAffected();
appliedRowSecurityPolicyIds = update.appliedRowSecurityPolicyIds();
effectiveSql = update.effectiveSql();
}
}
var canonicalSql = sqlCanonicalizer.canonicalize(query.sqlText());
Expand Down Expand Up @@ -404,7 +407,7 @@ private ExecutionOutcome doExecute(QueryRequestSnapshot query, UUID actorUserId,
query.organizationId(), successMetadata);
eventPublisher.publishEvent(new QueryExecutedEvent(
query.id(), rowsAffected, durationMs, QueryStatus.EXECUTED,
query.recurringParentId()));
query.recurringParentId(), effectiveSql));
return new ExecutionOutcome(query.id(), QueryStatus.EXECUTED, rowsAffected, durationMs);
} catch (UnrewritableRowSecurityException | InvalidSqlException ex) {
// A structurally unfilterable (or unparseable) query is a client error. For an
Expand Down
Loading
Loading