github-actions[bot] commented on code in PR #68499: URL: https://github.com/apache/doris/pull/68499#discussion_r4177410292
########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditPublicationHorizon.java: ########## @@ -0,0 +1,263 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.plugin.audit; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.InternalSchema; +import org.apache.doris.common.FeConstants; +import org.apache.doris.common.util.TimeUtils; +import org.apache.doris.qe.AuditEventProcessor; +import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + +/** + * The cluster-wide AUDIT PUBLICATION HORIZON: the start time (epoch millis, the + * {@code time} column of {@code audit_log}) of the oldest audit event that any FE has + * accepted but not yet PUBLISHED. The SPM capture scans the shared audit table from the + * leader, so it uses this value as a progress FENCE: its next scan window must still + * start at or before it, otherwise a row an FE still owes falls behind the advanced + * watermark and is never captured. + * + * <p>Three layers make the fence complete: + * <ul> + * <li>{@link #localHorizon()} folds THIS FE's whole audit pipeline: completed queries + * still held by the {@link WorkloadRuntimeStatusMgr} (they enter the pipeline + * before any loader sees them), the {@link AuditEventProcessor} queue and its + * in-flight event (a plugin can stall while an event is dequeued), and the + * {@link AuditLoader} queue / assembled batch / not-yet-visible batch (a stream + * load can report Publish Timeout after commit).</li> + * <li>each FE REPORTS its local horizon into the shared + * {@link InternalSchema#SPM_AUDIT_HORIZON_TBL_NAME} table, so a follower's + * backlog is visible to the leader that runs the capture.</li> + * <li>{@link #clusterHorizon()} is the MINIMUM over the local value and the FRESH + * rows of that table; a row the reporter stopped refreshing (its FE died or + * stopped reporting - the events are gone with it) is ignored.</li> + * </ul> + */ +public final class AuditPublicationHorizon { + + private static final Logger LOG = LogManager.getLogger(AuditPublicationHorizon.class); + + /** + * A reported row older than this is IGNORED: its FE stopped refreshing the fence + * (crashed / killed / its reporter thread is gone), so the events it still owed are + * lost with it and fencing progress forever would freeze the capture instead of + * protecting anything. Must be comfortably larger than the reporter's keepalive + * interval ({@link AuditLoader#HORIZON_KEEPALIVE_MILLIS}). + */ + public static final long ROW_STALE_MILLIS = 5 * 60 * 1000L; + + private static final String SELECT_ROWS_SQL = + "SELECT `horizon_ms`, `update_time` FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`"; + private static final String DELETE_OWN_ROW_SQL = "DELETE FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "` WHERE `fe_name` = '${feName}'"; + private static final String INSERT_OWN_ROW_SQL = "INSERT INTO `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`" + + " (`fe_name`, `horizon_ms`, `update_time`) VALUES ('${feName}', ${horizonMs}, '${updateTime}')"; + private static final int IO_TIMEOUT_SECONDS = 10; + + /** + * Test seam: the shared-table read (one row per FE). Null in production. + */ + @VisibleForTesting + static volatile Supplier<List<Object[]>> horizonRowsReaderForTest; + + /** + * Test seam: the shared-table write of this FE's row (delete + optional insert). + * Null in production. + */ + @VisibleForTesting + static volatile Consumer<Long> localHorizonWriterForTest; + + private AuditPublicationHorizon() { + } + + /** + * The oldest audit event THIS FE has accepted but not published, 0 when nothing is + * outstanding: the MINIMUM over every stage of the local pipeline (see the class + * javadoc). Cheap - no I/O - so callers may poll it. + */ + public static long localHorizon() { + long oldest = 0; + oldest = minPositive(oldest, AuditLoader.oldestUnpublishedEventTime()); Review Comment: [P2] Take a consistent snapshot across audit pipeline stages. This method reads the loader before the processor, and `preLoaderHorizon` reads the processor before the runtime manager. An event transferred between either pair after the downstream read but before the upstream read can be missed even when each stage reports its own events correctly. A zero follower fence can let capture skip the unpublished event; preserve ownership across handoffs or use a coordinated snapshot. ########## fe/fe-core/src/main/java/org/apache/doris/resource/workloadschedpolicy/WorkloadRuntimeStatusMgr.java: ########## @@ -163,6 +163,35 @@ public void submitFinishQueryToAudit(AuditEvent event, Set<Long> expectedBackend } } + /** + * Start time (epoch millis, the {@code time} column of {@code audit_log}) of the + * OLDEST completed query this FE still HOLDS for auditing, 0 when it holds none. + * + * <p>A completed query enters {@code queryAuditEventList} BEFORE the audit event + * processor - and therefore before any audit loader or the shared table - sees it, + * and it stays here until {@code query_audit_log_timeout_ms} expires (or the expected + * backends reported). The SPM capture's publication fence must include it: otherwise + * the row's release lands behind the capture's advanced scan watermark and the query + * is never captured (round-36 #3). + */ + public long oldestHeldAuditEventTime() { Review Comment: [P2] Keep ready audit events in the publication horizon through the processor handoff. `getQueryNeedAudit` removes them from `queryAuditEventList` before statistics are rebuilt and `handleAuditEvent` enqueues them; this new reader sees only the list. A follower reporter can send zero in that interval, letting the leader advance past an older-than-overlap event that publishes later. Retain an in-flight fence until the processor accepts each event. ########## fe/fe-core/src/main/java/org/apache/doris/qe/AuditEventProcessor.java: ########## @@ -136,13 +169,19 @@ public void run() { } try { + // the event is OUT of the queue while the plugins run: publish it as the + // in-flight fence so a concurrent horizon read still sees it (see + // oldestQueuedOrInFlightEventTime) + processingEvent = auditEvent; Review Comment: [P2] Make the processor queue-to-in-flight transfer atomic with its horizon read. The worker polls before setting `processingEvent`, while `oldestQueuedOrInFlightEventTime` reads the flag before the queue. A follower can briefly report zero for an unpublished old event and the leader can advance its capture window past it. Keep the event visible across the transfer and test that interleaving. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/capture/PlanCaptureManager.java: ########## @@ -0,0 +1,2275 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.nereids.spm.capture; + +import org.apache.doris.catalog.Env; +import org.apache.doris.common.Config; +import org.apache.doris.common.FeConstants; +import org.apache.doris.common.UserException; +import org.apache.doris.common.util.MasterDaemon; +import org.apache.doris.nereids.spm.BaselinePlan; +import org.apache.doris.nereids.spm.BaselineSource; +import org.apache.doris.nereids.spm.SPMPlanner; +import org.apache.doris.nereids.spm.SPMUtils; +import org.apache.doris.nereids.spm.manager.BaselineManager; +import org.apache.doris.plugin.audit.AuditLoader; +import org.apache.doris.plugin.audit.AuditPublicationHorizon; +import org.apache.doris.qe.AutoCloseConnectContext; +import org.apache.doris.qe.ConnectContext; +import org.apache.doris.qe.GlobalVariable; +import org.apache.doris.qe.SessionVariable; +import org.apache.doris.qe.SqlModeHelper; +import org.apache.doris.qe.VariableMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; +import java.util.function.Supplier; + +/** + * PlanCaptureManager - SPM auto capture scheduler (Phase 2, design doc 7.2.1 / 7.2.4). + * + * A Leader-FE daemon that periodically scans the audit_log internal table and + * automatically creates baselines for high-value queries: + * + * - only queries executed by the Nereids planner are captured; + * - the capture filter (PlanCaptureFilter) enforces the multi-table / table-exists / + * regex / performance-threshold rules; + * - the baseline is built through the Phase 1 flow (SPMPlanner.buildBaselineFromSql: + * SPM-mode optimize + decompile + parameterize) with source = CAPTURE and the actual + * query_time filled for candidate ordering; + * - duplicate (digest, planSql) baselines are skipped (BaselineManager dedup). + * + * The whole cycle is guarded by the global session variable enable_plan_capture + * (default false, tunable via `SET GLOBAL enable_plan_capture = true`), and any + * failure is logged and skipped so auto capture never breaks the cluster. + */ +public class PlanCaptureManager extends MasterDaemon { + + /** + * Test seam replacing the live leadership probe of {@link #persistCheckpoint} (null in + * production). A capture cycle runs on the master, but an in-flight cycle can reach + * its checkpoint write AFTER a handoff (the daemon checks isMaster only at the cycle + * start), which a unit test cannot interleave otherwise. + */ + @VisibleForTesting + public static volatile java.util.function.BooleanSupplier checkpointLeadershipProbeForTest; + + /** + * Statement timeout (seconds) of the checkpoint read / write. The default + * StatisticsUtil overloads assign the ANALYZE timeout (43,200 seconds), so a stalled + * internal-table read or write could hold the single capture cycle for hours and + * delay every later capture / retry. Both operations are latency-sensitive: fail + * fast, keep the cycle consistent, retry next cycle. + */ + static final int CHECKPOINT_IO_TIMEOUT_SECONDS = 10; + + /** Bounded read-back attempts confirming the first reservation is VISIBLE. */ + private static final int CHECKPOINT_VISIBILITY_ATTEMPTS = 5; + + /** Delay between the reservation visibility reads (millis). */ + private static final long CHECKPOINT_VISIBILITY_RETRY_MILLIS = 200L; + + private static final Logger LOG = LogManager.getLogger(PlanCaptureManager.class); + + private static final PlanCaptureManager INSTANCE = new PlanCaptureManager(); + + /** + * Re-scan overlap (millis) applied to the watermark: AuditLoader buffers events + * asynchronously and writes their original event timestamp, so a row can become + * visible AFTER its window has passed (it would otherwise be excluded from every + * future window forever). Re-scanning a lagged/overlapping window plus query-id + * deduplication makes late arrivals capturable without processing an execution + * twice. + */ + private static final long SCAN_WINDOW_OVERLAP_MS = 300_000L; + + /** Upper bound for the processed-query-id dedup map. */ + private static final int MAX_TRACKED_QUERY_IDS = 10000; + + /** + * Audit pages ONE wakeup may consume. The daemon interval (default 3h) bounds how + * often the backlog is drained, so consuming a single page per wakeup left a window + * truncated at the page limit needing one extra interval per page - a window holding + * more than `plan_capture_max_batch_size` eligible rows per interval could never catch + * up. The drain stays BOUNDED so one cycle cannot run unboundedly long (each page is + * one bounded query plus its checkpoint write). + */ + private static final int MAX_PAGES_PER_CYCLE = 50; + + /** + * Wakeup delay used while a window is still pending after a cycle: the backlog drains + * promptly instead of one page per `plan_capture_interval_seconds`. + */ + private static final long PENDING_WINDOW_RESUME_INTERVAL_MS = 5_000L; + + /** + * In-memory budget of the queued retry candidates, in statement characters (the + * dominant part of a queued entry; the statement text is what the leader FE holds). + * One drain of {@link #MAX_PAGES_PER_CYCLE} pages can enqueue up to a page budget of + * failures per page, so a transient external-metadata outage (every capture fails + * with "table metadata unavailable") would otherwise retain tens of thousands of full + * statements - hundreds of megabytes - before the first cohort reaches its third + * attempt. When the budget is reached the DRAIN pauses: the pending window keeps its + * bounds and its cursor (the unconsumed rows stay reachable by the keyset scan) while + * the replay burns the queue down, and the daemon resumes promptly (see + * pendingWindowNeedsPromptResume). + */ + private static final long MAX_QUEUED_FAILURE_CHARS = 64L * 1024 * 1024; + + /** + * Test seam overriding {@link #MAX_QUEUED_FAILURE_CHARS} (null = the production + * budget): a unit test cannot queue tens of megabytes of statements just to reach it. + * Written through {@link #setQueuedFailureBudgetForTest}. + */ + private static volatile Long queuedFailureBudgetForTest; + + /** + * Queued failures replayed in ONE cycle (see {@link #replayQueuedFailures}): replanning + * a whole outage-sized queue every wakeup would keep the FE busy for the length of the + * outage itself. + */ + private static final int MAX_RETRY_REPLAY_PER_CYCLE = 1000; + + /** + * Bounded retries for a FAILED capture: the query id stays retryable for later + * overlapping scans until it either succeeds or reaches this attempt count. Marking + * the id before processing would make a transient failure permanent - the + * overlapping scans would skip the row and the watermark passes it long before the + * dedup map evicts the entry. + */ + private static final int MAX_CAPTURE_ATTEMPTS = 3; + + /** Durable checkpoint key: the internal table holds exactly one row. */ + private static final long CHECKPOINT_ID = 1L; + + /** Upper bound for the retry entries written into the checkpoint row (row size). */ + private static final int MAX_PERSISTED_RETRIES = 64; + + /** Table of the durable capture checkpoint (see InternalSchema). */ + private static final String CHECKPOINT_TABLE = + "`__internal_schema`.`spm_capture_checkpoint`"; + + private static final String CHECKPOINT_SELECT_SQL = + "SELECT `last_scan_timestamp`, `pending_window_start`, `pending_window_end`," + + " `cursor_query_time`, `cursor_time`, `cursor_query_id`," + + " `failed_attempts`, `retry_queue`, `cursor_tail`," + + " `min_query_time_ms`, `min_scan_rows`, `include_pattern`," + + " `exclude_pattern`, `scan_zone` FROM " + CHECKPOINT_TABLE + + " WHERE `id` = " + CHECKPOINT_ID + " ORDER BY `update_time` DESC LIMIT 1"; + + /** + * One UPSERT statement: the table is UNIQUE-key(id) with merge-on-write, so inserting + * the row again REPLACES it atomically. The previous delete-then-insert pair was two + * separately committed statements: a crash / leadership loss / timeout / failed + * INSERT after the DELETE left NO row for the next leader, which then derived a fresh + * window and permanently skipped the deleted pending window's unconsumed tail. + * + * The target columns are listed EXPLICITLY. The VALUES order below follows + * {@link org.apache.doris.catalog.InternalSchema#SPM_CAPTURE_CHECKPOINT_SCHEMA}, but + * the PHYSICAL order of an upgraded table can differ: the upgrade of a pre-existing + * table APPENDS the columns it adds ({@code InternalSchemaInitializer# + * upgradeSpmCaptureCheckpointSchema}), which used to place cursor_tail after + * update_time. A positional INSERT then shifts every value behind the first + * out-of-position column - the tail JSON was written into failed_attempts, the retry + * JSON into update_time and NOW() into cursor_tail - and the checkpoint write failed + * / persisted garbage. Address the columns by NAME instead: the write must stay + * correct on every physical layout, exactly like the (by-name) CHECKPOINT_SELECT_SQL + * read. + */ + private static final String CHECKPOINT_INSERT_SQL = + "INSERT INTO " + CHECKPOINT_TABLE + + " (`id`, `last_scan_timestamp`, `pending_window_start`, `pending_window_end`," + + " `cursor_query_time`, `cursor_time`, `cursor_query_id`, `cursor_tail`," + + " `failed_attempts`, `retry_queue`, `min_query_time_ms`, `min_scan_rows`," + + " `include_pattern`, `exclude_pattern`, `scan_zone`, `update_time`)" + + " VALUES (" + CHECKPOINT_ID + ", ${lastScan}, ${pendingStart}, ${pendingEnd}," + + " ${cursorQueryTime}, '${cursorTime}', '${cursorQueryId}', '${cursorTail}'," + + " '${failedAttempts}', '${retryQueue}', ${minQueryTimeMs}, ${minScanRows}," + + " '${includePattern}', '${excludePattern}', '${scanZone}', NOW())"; + + private AuditLogScanner scanner = new AuditLogScanner(); + + /** Capture filter, refreshed from the global session variables each cycle. */ + private PlanCaptureFilter filter; + + /** Last scan window start (epoch millis); 0 means "first run, scan one interval". */ + private long lastScanTimestamp = 0; + + /** + * The session time_zone (zone ID) the most recent scan PASS rendered its window bounds + * in - the zone the audited rows' {@code time} columns are stored in. Empty = never + * scanned (a fresh process follows the global zone). audit_log keeps the WRITER's + * local rendering, so a global time_zone change makes the already published rows + * invisible to bounds rendered in the new zone: while this differs from the current + * global zone, the next pass re-renders the window in this zone first and only then in + * the new one (see {@link #resolveScanPassZone} and the exhaustion branch of + * {@link #runCaptureCycle}). Persisted with the checkpoint so a takeover resumes the + * same rendering. + */ + private String lastScanZone = ""; + + /** + * Pending scan window of a TRUNCATED cycle: the (start, end) pair the resume cursor + * below belongs to. While set, every cycle keeps scanning the SAME window - the end + * must stay fixed until the window is fully consumed, because the next + * interval-derived window would start around this window's end and leave every row the + * cursor has not reached yet permanently out of scope. + */ + private long pendingWindowStart = 0; + private long pendingWindowEnd = 0; + + /** + * The FILTER SNAPSHOT the pending window was opened with (null while none is pending). + * The audit SQL and the in-memory {@link PlanCaptureFilter#shouldCapture} stage must + * judge one window's rows by the SAME thresholds: a `SET GLOBAL + * plan_capture_min_query_time_ms` between two pages of the same window otherwise made + * the SQL return rows the stale filter rejected terminally (they were consumed, + * never captured) or pushed already-passed rows behind the cursor where a LOWERED + * threshold could no longer reach them. The whole filter is pinned, so a pattern + * change applies from the next window on. + */ + private PlanCaptureFilter pendingWindowFilter; + + /** + * Set when a cycle could not finish its pending window (page budget / failed + * checkpoint write): the daemon then reschedules the next cycle promptly instead of + * waiting the full `plan_capture_interval_seconds` (default 3h), which would grow the + * backlog by one interval's worth of eligible rows per consumed page. + */ + private volatile boolean pendingWindowNeedsPromptResume; + + /** Query ids already handled in earlier (overlapping) windows. */ + private final Map<String, Boolean> processedQueryIds = new LinkedHashMap<>(); + + /** + * Failure attempts per query id (bounded retry, see MAX_CAPTURE_ATTEMPTS). An id is + * removed here when it succeeds or is given up on; the map is capped like the + * processed-id map so a long-running failure burst cannot grow unbounded. + */ + private final Map<String, Integer> failedCaptureAttempts = new LinkedHashMap<>(); + + /** + * Candidates whose capture failed and that still have retry budget: keyset pagination + * advances the scan cursor past their audit rows and the window overlap only re-reads + * recent rows, so they are REPLAYED one attempt per cycle from here. Bounded like the + * other query-id maps; entries leave on success, on give-up, or when the id turns + * terminal elsewhere. + */ + private final Map<String, CapturedQuery> failedCaptureQueue = new LinkedHashMap<>(); + + /** + * Statement characters currently retained by {@link #failedCaptureQueue} (see + * {@link #MAX_QUEUED_FAILURE_CHARS}). Guarded by the capture daemon thread (all + * mutations happen inside a cycle) plus the test seams. + */ + private long queuedFailureChars = 0; + + /** + * Pre-page checkpoint state of the page a queued retry was FIRST seen on: the durable + * checkpoint must never move past an entry that the persisted JSON drops + * ({@link #MAX_PERSISTED_RETRIES}) - keyset pagination has already moved beyond its + * audit row, so only a cursor BEFORE that row can reach it after a restart / handoff. + * First-wins (pages only move forward) and evicted together with the queue. + */ + private final Map<String, RetryAnchor> failedCaptureAnchors = new LinkedHashMap<>(); + + /** One pre-page checkpoint state (see {@link #failedCaptureAnchors}). */ + private static final class RetryAnchor { + final long lastScanTimestamp; + final long windowStart; + final long windowEnd; + final long cursorQueryTime; + final String cursorTime; + final String cursorQueryId; + final String cursorTail; + + /** + * The zone THIS page's bounds / cursor were RENDERED in. The anchor describes a + * position of the audit stream, and that position is only reachable when it is + * rendered in the same zone again (see {@link #resolveScanPassZone}): persisting + * the CURRENT cycle's zone beside an earlier window's bounds made a takeover + * scan those bounds in a zone the rows were never written under, and the omitted + * retries (the ones the truncated queue could not carry) stayed unreachable. + */ + final String scanZone; + + /** The filter snapshot the page was judged by (see pageStartFilter). */ + final PlanCaptureFilter filter; + + RetryAnchor(long lastScanTimestamp, long windowStart, long windowEnd, + long cursorQueryTime, String cursorTime, String cursorQueryId, + String cursorTail, String scanZone, PlanCaptureFilter filter) { + this.lastScanTimestamp = lastScanTimestamp; + this.windowStart = windowStart; + this.windowEnd = windowEnd; + this.cursorQueryTime = cursorQueryTime; + this.cursorTime = cursorTime; + this.cursorQueryId = cursorQueryId; + this.cursorTail = cursorTail; + this.scanZone = scanZone; + this.filter = filter; + } + } + + /** + * Resume cursor of a TRUNCATED scan window: the FULL ORDER BY key tuple of the last + * consumed row -- (time, query_time, query_id) plus the encoded tail (client_ip, + * sql_hash, scan_rows, return_rows, statement hash) that uniquely separates audit + * rows sharing the first three keys. CURSOR_ABSENT while no partial window is + * pending - a short batch advances the watermark instead. Zero and NULL query_time + * are VALID cursors (see AuditLogScanner.CURSOR_QUERY_TIME_NULL). + */ + private long cursorQueryTime = AuditLogScanner.CURSOR_ABSENT; + private String cursorTime = ""; + private String cursorQueryId = ""; + private String cursorTail = ""; + + /** + * The pre-page state (watermark + window bounds + cursor) the CURRENT cycle's scan + * started from. It is the durable fallback persisted when the retry state is + * truncated by {@link #MAX_PERSISTED_RETRIES} (see persistCheckpoint): the durable + * cursor must never move past retries the checkpoint can no longer carry, otherwise + * a restart / leader handoff neither replays them from the queue nor re-reads their + * audit rows (the keyset cursor is beyond them and they can age outside the + * five-minute overlap), silently losing those captures. + */ + private long pageStartLastScanTimestamp = 0; + private long pageStartWindowStart = 0; + private long pageStartWindowEnd = 0; + private long pageStartCursorQueryTime = AuditLogScanner.CURSOR_ABSENT; + private String pageStartCursorTime = ""; + private String pageStartCursorQueryId = ""; + private String pageStartCursorTail = ""; + + /** + * The zone the CURRENT page's bounds / cursor were rendered in (the pass zone of the + * cycle that opened the page, see {@link #resolveScanPassZone}). It travels with + * {@link #currentPageAnchor()} and is persisted whenever the retry state rewinds the + * durable cursor to a page anchor: a takeover must re-render those bounds in the + * SAME zone, otherwise the audit rows written under the anchor's rendering are + * invisible to the re-scan (see {@link RetryAnchor#scanZone}). + */ + private String pageStartZoneId = ""; + + /** + * The FILTER SNAPSHOT the CURRENT page was scanned with (the same value handed to + * {@link AuditLogScanner#scan}). It travels with {@link #currentPageAnchor()} and is + * persisted whenever the retry state rewinds the durable cursor to an anchor: the + * rows of that page were admitted (or filtered) by THESE thresholds / patterns, so a + * takeover must re-scan the rewound range with the same eligibility - judging the + * re-scan by a configuration that changed in between could terminally filter the + * omitted oldest failure before its retry is even reachable. + */ + private PlanCaptureFilter pageStartFilter; + + /** Whether the durable checkpoint was already consulted in this process. */ + private boolean checkpointLoaded = false; + + /** + * Whether THIS process has ever seen a durable checkpoint row - read it from the store + * or written by this process. While it is false, the window a cycle consumes exists + * only in memory: runCaptureCycle records that window BEFORE scanning (see the + * initial reservation there), so a takeover can still resume it. + */ + private boolean durableCheckpointObserved = false; + + /** + * Checkpoint read / write seams. Production talks to the internal table through + * StatisticsUtil; tests replace them to simulate a failing first read and to observe + * the exact statements a persist issues. + */ + private Supplier<List<ResultRow>> checkpointReader = () -> StatisticsUtil.executeQuery( + CHECKPOINT_SELECT_SQL, Collections.emptyMap(), CHECKPOINT_IO_TIMEOUT_SECONDS); + + /** One checkpoint write statement. */ + @VisibleForTesting + public interface CheckpointWriter { + void write(String sql, Map<String, String> params) throws Exception; + } + + private CheckpointWriter checkpointWriter = (sql, params) -> StatisticsUtil.execUpdate( + sql, params, CHECKPOINT_IO_TIMEOUT_SECONDS); + + /** + * Whether a scripted checkpoint read / write seam is installed (tests only). The + * leadership fence of {@link #persistCheckpoint} guards LIVE writes: a scripted store + * stands in for the internal table, exactly like the simulator stores in + * BaselineManager.assertLeaderForWrite. + */ + private boolean checkpointSeamsForTest = false; + + /** + * Whether the cloud-mode warning was already logged (the gate fires every cycle). + */ + private boolean cloudModeWarned = false; + + /** + * The start of the window the FIRST cycle would have consumed when its checkpoint read + * FAILED (0 = none). A failed read records no window and the daemon retries promptly; + * without this floor the first cycle that succeeds on an EMPTY store would derive its + * own [now-interval, now) and permanently skip the rows of the first attempted window + * (every later window starts even later). Cleared as soon as a durable row is read or + * the reserved window becomes durable. + */ + private long firstAttemptedWindowStart = 0; + + /** + * The CLUSTER-WIDE audit publication horizon: the start time (epoch millis) of the + * oldest audit event ANY FE has accepted but not published yet (0 = nothing + * outstanding), i.e. one this FE's loader owes, one still held / queued before the + * loader of any FE, or a follower's batch whose load reported Publish Timeout. The + * capture runs on the leader alone, so only the shared view can fence a follower's + * backlog (round-36 #1). Production reads the live shared table; tests replace it. + */ + private LongSupplier auditQueueHorizon = AuditPublicationHorizon::clusterHorizon; + + // capture statistics (design doc 7.2.1 / 7.2.6) + private final AtomicLong successCount = new AtomicLong(0); + private final AtomicLong skipDuplicateCount = new AtomicLong(0); + private final AtomicLong skipSingleTableCount = new AtomicLong(0); + private final AtomicLong skipFilterCount = new AtomicLong(0); + private final AtomicLong failCount = new AtomicLong(0); + + private PlanCaptureManager() { + super("PlanCaptureManager", + Math.max(1, VariableMgr.getDefaultSessionVariable().getPlanCaptureIntervalSeconds()) + * 1000L); + this.filter = buildFilterFromGlobal(); + } + + /** The pre-page state of the CURRENT page (anchors entries queued by this page). */ + private RetryAnchor currentPageAnchor() { + return new RetryAnchor(pageStartLastScanTimestamp, pageStartWindowStart, + pageStartWindowEnd, pageStartCursorQueryTime, pageStartCursorTime, + pageStartCursorQueryId, pageStartCursorTail, pageStartZoneId, pageStartFilter); + } + + public static PlanCaptureManager getInstance() { + return INSTANCE; + } + + /** + * Builds a capture filter from the global session variables (so `SET GLOBAL` + * changes to the thresholds / table regex take effect on the next cycle). + * + * @return a new filter + */ + private static PlanCaptureFilter buildFilterFromGlobal() { + try { + SessionVariable global = VariableMgr.getDefaultSessionVariable(); + return new PlanCaptureFilter(global.getPlanCaptureIncludePattern(), + global.getPlanCaptureExcludePattern(), + global.getPlanCaptureMinQueryTimeMs(), + global.getPlanCaptureMinScanRows()); + } catch (RuntimeException e) { + // e.g. a legacy invalid regex in the global variable: never let it escape the + // singleton constructor / the daemon cycle (leader startup calls getInstance() + // before enable_plan_capture is even checked, and a PatternSyntaxException + // there would terminate the FE transition) + LOG.error("SPM plan capture disabled: invalid capture filter configuration", e); + return null; + } + } + + @Override + protected void runAfterCatalogReady() { + SessionVariable global = VariableMgr.getDefaultSessionVariable(); + // Reschedule from the cycle itself: MasterDaemon sleeps its stored intervalMs, so + // only setInterval() here makes a `SET GLOBAL plan_capture_interval_seconds` + // change affect future wakeups (rereading the variable in the cycle would only + // change the scan window). Clamp to >= 1s so a misconfiguration cannot spin. + setInterval(Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L); + // SPM baseline management (CREATE / ALTER / DROP / SHOW) explicitly rejects cloud + // mode; until the full lifecycle is supported the capture daemon must not create + // (or keep retrying to create) global baselines a cloud deployment cannot show, + // disable or drop. + if (Config.isCloudMode()) { + if (!cloudModeWarned) { + cloudModeWarned = true; + LOG.warn("SPM plan capture is not supported in cloud mode, skipping"); + } + return; + } + if (!global.isEnablePlanCapture()) { + return; + } + if (!Env.getCurrentEnv().isMaster()) { + // auto capture runs on the Leader FE only + return; + } + if (Env.isCheckpointThread()) { + return; + } + PlanCaptureFilter newFilter = buildFilterFromGlobal(); + if (newFilter == null) { + LOG.error("Plan capture filter unavailable (invalid capture regex?)," + + " skipping this capture cycle"); + return; + } + runCaptureCycle(global, newFilter); + if (pendingWindowNeedsPromptResume) { + // A window this cycle could not finish (page budget reached or its progress + // not durable) must NOT wait another full interval: the next wakeup continues + // exactly where this one stopped. The configured interval would add one + // interval's worth of eligible rows per consumed page, so a window holding + // more than one page per interval would never drain. + setInterval(Math.min( + Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L, + PENDING_WINDOW_RESUME_INTERVAL_MS)); + } + } + + /** + * One capture cycle body: everything after the runtime guards (cloud mode, enable + * flag, leader / checkpoint-thread checks, filter refresh). Split out so unit tests + * can drive a FULL cycle - checkpoint read through window derivation, scan, state + * advance and persist - without the process-global guards (master / cloud / enable) + * a test environment cannot satisfy. + * + * @param global the global session variables of this cycle + * @param newFilter the filter refreshed for this cycle + */ + @VisibleForTesting + public void runCaptureCycle(SessionVariable global, PlanCaptureFilter newFilter) { + try { + // A restarted / newly promoted leader must NOT start from a fresh + // interval-derived window: a truncated window from the previous leader is + // checkpointed here, and skipping it would permanently exclude its unconsumed + // tail (the overlap only reaches rows younger than the NEW watermark). + // A FAILED read returns false and the cycle aborts BEFORE deriving or + // persisting anything: writing a freshly derived window while the previous + // leader's unconsumed tail is still unreadable would overwrite its only + // record (the write path shares the same internal table the read failed on). + if (!loadCheckpointIfNeeded()) { + // The read failed and recorded nothing, so this process still has no window. + // Remember the window this cycle WOULD have consumed and retry promptly: the + // internal-schema initializer is asynchronous, and a later cycle deriving its + // OWN [now-interval, now) would permanently skip every eligible short row of + // this first attempted window (no later overlap reaches behind a NEW window's + // start). Nothing is written here - an unreadable checkpoint must never be + // replaced by a freshly derived one. + long attemptedStart = System.currentTimeMillis() + - Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L; + if (firstAttemptedWindowStart == 0 || attemptedStart < firstAttemptedWindowStart) { + firstAttemptedWindowStart = attemptedStart; + } + pendingWindowNeedsPromptResume = true; + LOG.warn("Plan capture cycle skipped: durable checkpoint not confirmed"); + return; + } + + // A window that is still PENDING keeps the filter snapshot it was opened with: + // its rows behind the cursor were already judged by those thresholds, and the + // audit SQL must use exactly the same values (see AuditLogScanner#scan). The + // snapshot is chosen AFTER the checkpoint load, so a TAKEOVER continues the + // restored window with the restored thresholds in its very first cycle. A NEW + // window follows the filter refreshed for this cycle, so `SET GLOBAL + // plan_capture_min_query_time_ms` takes effect from the next window on. + PlanCaptureFilter cycleFilter = pendingWindowFilter != null + ? pendingWindowFilter : newFilter; + this.filter = cycleFilter; + pendingWindowNeedsPromptResume = false; + + long currentTime = System.currentTimeMillis(); + // The CLUSTER-WIDE publication fence: the oldest audit event ANY FE has + // accepted but not published yet (round-36 #1: the local loader queue alone + // cannot see a follower's backlog - the capture runs on the leader, and the + // follower's row would land behind the advanced watermark). An unreadable + // shared table means the fence is INCOMPLETE, so the cycle is skipped and + // retried promptly instead of advancing blind. + long publicationHorizon; + try { + publicationHorizon = auditQueueHorizon.getAsLong(); + } catch (RuntimeException e) { + pendingWindowNeedsPromptResume = true; + LOG.warn("Plan capture cycle skipped: the cluster audit publication horizon" + + " could not be read", e); + return; + } + // a non-positive interval / batch size can never be written through SQL SET + // (see SessionVariable), but clamp defensively: an interval of 0 would make + // every window empty and a batch size of 0 would return LIMIT 0, mark the + // window exhausted and advance the watermark over every eligible row + long intervalMs = Math.max(1L, global.getPlanCaptureIntervalSeconds()) * 1000L; + int batchSize = Math.max(1, global.getPlanCaptureMaxBatchSize()); + // overlap the window so audit rows loaded late (published after their event + // time has passed) are still scanned; the overlap follows the audit loader's + // configured batch interval so rows written at the tail of a loader batch - + // whose event time predates the new watermark - are not lost. Duplicates are + // filtered by query id below. + long[] window = resolveScanWindow(lastScanTimestamp, pendingWindowStart, pendingWindowEnd, + currentTime, intervalMs, + scanWindowOverlapMs(GlobalVariable.auditPluginMaxBatchInternalSec, + publicationHorizon, currentTime), + firstAttemptedWindowStart); + long scanStart = window[0]; + long scanEnd = window[1]; + if (scanStart >= scanEnd) { + return; + } + + // The zone this cycle's scan RENDERS its bounds in (see resolveScanPassZone): + // remember it as the zone of the record this cycle persists, so a takeover + // resumes the same rendering and the next window can detect a change. + String currentZoneId = AuditLogScanner.auditWriteZone().getId(); + String passZoneId = resolveScanPassZone(currentZoneId); + lastScanZone = passZoneId; + + if (!durableCheckpointObserved) { + // FIRST cycle after a successful-but-EMPTY read: the store holds NO row + // describing the window this process is about to consume, so its bounds and + // page-top cursor exist only in memory. A restart / leader handoff between + // the scan and the final persist would leave the takeover with nothing to + // resume - it would derive a NEW window and permanently skip this page's + // unconsumed tail (the later overlap only reaches rows younger than the new + // watermark). Record the window TO CONSUME before consuming it: pending = + // this window, cursor = its top, watermark = the pre-page one. A failed + // write ABORTS the cycle: scanning on would advance progress no durable + // state could ever resume. + pendingWindowStart = scanStart; + pendingWindowEnd = scanEnd; + pendingWindowFilter = cycleFilter; + if (!persistCheckpointAndConfirm()) { + // either this reservation could not be confirmed readable (retried next + // cycle, idempotent UPSERT) or an earlier leader's window surfaced and + // was adopted instead (resumed next cycle) - both abort WITHOUT scanning + // a single audit row while the window waits. Resume PROMPTLY: leaving the + // daemon at the configured interval (three hours by default) delayed the + // retry of a window nothing was consumed from, exactly like the failed + // first READ above (round-35 #6; the adoption path already set the flag, + // the write / visibility failure did not). + pendingWindowNeedsPromptResume = true; + LOG.warn("Plan capture cycle skipped: the initial checkpoint row could not" + + " be confirmed VISIBLE / was superseded by an earlier window"); + return; + } + // the reservation is durable and readable: the remembered first attempted + // window is now covered by a durable record and must not widen anything + firstAttemptedWindowStart = 0; + } + + // Drain this window with a BOUNDED number of pages: one page per wakeup would + // make a window holding more than one page wait one interval per page, so a + // capture rate above `plan_capture_max_batch_size` per interval could never + // catch up with the audit stream. + Set<String> scannedQueryIds = new HashSet<>(); + AuditLogScanner.ScanBatch batch = null; + int pages = 0; + for (int page = 0; page < MAX_PAGES_PER_CYCLE; page++) { + if (queuedFailureChars > queuedFailureBudget()) { + // The retry queue holds more un-replayed failures than the leader FE + // should retain: PAUSE the drain (no page is consumed, so the window + // stays pending with its cursor and every unconsumed row remains + // reachable) and let the bounded replay below burn the queue down. + pendingWindowNeedsPromptResume = true; + break; + } + // Snapshot the PRE-PAGE state: when this page ends up with more retries + // than the durable checkpoint can carry, persistCheckpoint falls back to + // THIS state so the next leader re-scans the page instead of stepping + // over the omitted retries. Refreshed per page - a retry queued by page N + // must stay reachable from page N's top, not from the cycle's. + pageStartLastScanTimestamp = lastScanTimestamp; + pageStartWindowStart = scanStart; + pageStartWindowEnd = scanEnd; + pageStartCursorQueryTime = cursorQueryTime; + pageStartCursorTime = cursorTime; + pageStartCursorQueryId = cursorQueryId; + pageStartCursorTail = cursorTail; + pageStartZoneId = passZoneId; + pageStartFilter = cycleFilter; + + batch = scanner.scan(scanStart, scanEnd, batchSize, cycleFilter, + cursorQueryTime, cursorTime, cursorQueryId, cursorTail, passZoneId); + pages++; + for (CapturedQuery candidate : batch.getCandidates()) { + scannedQueryIds.add(retryKeyOf(candidate)); + handleCandidate(candidate); + } + if (batch.isWindowExhausted()) { + break; + } + // The batch limit truncated the window: KEEP the window BOUNDS and remember + // the full total-order cursor of the last consumed row, so the next page + // (and a later leader) resumes inside the same window. Advancing to the + // window end here would permanently skip every eligible row beyond the + // LIMIT; letting the next cycle derive a new interval window would skip + // everything the cursor has not reached yet as well. The cursor TAIL is + // what keeps rows sharing (time, query_time, query_id) - e.g. a whole page + // of NULL query ids - from looping or being skipped (see + // AuditLogScanner#ORDER_BY). + pendingWindowStart = scanStart; + pendingWindowEnd = scanEnd; + pendingWindowFilter = cycleFilter; + cursorQueryTime = batch.getCursorQueryTime(); + cursorTime = batch.getCursorTime(); + cursorQueryId = batch.getCursorQueryId(); + cursorTail = batch.getCursorTail(); + // Make this page's progress durable BEFORE consuming the next one: a page + // nothing durable describes would be re-derived as a NEW window by a + // takeover (see the reservation above). A failed write STOPS the drain - + // consuming further pages while the store is unavailable is exactly what + // the reservation exists to prevent - and the cycle resumes promptly. + if (!persistCheckpoint()) { + break; + } + } + // Rows whose capture failed stay queued: keyset pagination moved the cursor + // past their raw rows and the five-minute overlap only re-reads recent ones, + // so without this replay attempts 2..N would be unreachable for older + // failures. Replayed ONCE per cycle with the UNION of every page's keys: an id + // that appeared in ANY page of this cycle was already retried there, and + // replaying per page would burn one attempt per page for a failure the page + // loop kept failing. + replayQueuedFailures(scannedQueryIds); + if (batch == null) { + // The queue budget paused the drain before a single page: the window (a + // resumed one, or the one just reserved above) stays pending with the + // cursor it has, and the next cycle - scheduled promptly - retries. + pendingWindowNeedsPromptResume = true; + } else if (batch.isWindowExhausted()) { + if (!passZoneId.equals(currentZoneId)) { + // The window drained in the zone it was OPENED with (the cursor's + // rendering), but the global time_zone changed since: rows published + // AFTER the change were rendered in the NEW zone and the drained pass + // could not see them. Re-scan the SAME window from its top in the + // current zone instead of advancing the watermark - the reviewer's + // example: a 10:00 UTC row stored as "10:00" is invisible to a + // [17:00, 20:00) rendering, and the watermark would move past it + // forever. `lastScanZone` follows the new pass, so the re-scan itself + // advances normally once its rendering matches the global zone + // (several changes chain one pass each). + LOG.info("Plan capture: the global time_zone changed from {} to {} during" + + " window [{}, {}); re-scanning it in the new zone before advancing", + passZoneId, currentZoneId, scanStart, scanEnd); + pendingWindowStart = scanStart; + pendingWindowEnd = scanEnd; + pendingWindowFilter = cycleFilter; + cursorQueryTime = AuditLogScanner.CURSOR_ABSENT; + cursorTime = ""; + cursorQueryId = ""; + cursorTail = ""; + lastScanZone = currentZoneId; + pendingWindowNeedsPromptResume = true; + } else { + // The whole window was scanned: advance the watermark to the CONSUMED + // window end (not to `now` - rows that arrived between a resumed pending + // window's end and now would be skipped), keep the overlap so + // late-arriving audit rows stay capturable, and drop the resume state. + lastScanTimestamp = nextScanTimestamp(lastScanTimestamp, scanEnd, true); + clearPendingWindow(); + cursorQueryTime = AuditLogScanner.CURSOR_ABSENT; + cursorTime = ""; + cursorQueryId = ""; + cursorTail = ""; + lastScanZone = currentZoneId; + if (!failedCaptureQueue.isEmpty()) { + // The window is consumed but retries outlived it: this cycle replayed at + // most MAX_RETRY_REPLAY_PER_CYCLE of them, so without a prompt resume the + // remaining batches would each wait a FULL capture interval - a 25,000 + // entry outage queue would need days to burn down at the default three + // hours per 1,000 entries. Keep the per-cycle work cap and only shorten + // the WAKEUP: resume queued retries at the pending-window cadence until + // the queue is drained. + pendingWindowNeedsPromptResume = true; + } + } + } else { + // The window is still not consumed (page budget reached, the last + // checkpoint write failed, or the queue budget paused the drain): keep the + // bounds, the cursor and the pinned filter, and resume promptly instead of + // after a full interval. + pendingWindowStart = scanStart; + pendingWindowEnd = scanEnd; + pendingWindowFilter = cycleFilter; + if (batch != null) { + cursorQueryTime = batch.getCursorQueryTime(); + cursorTime = batch.getCursorTime(); + cursorQueryId = batch.getCursorQueryId(); + cursorTail = batch.getCursorTail(); + } + pendingWindowNeedsPromptResume = true; + } + // Make the progress durable for the NEXT process (leader handoff / restart). + persistCheckpoint(); + + LOG.info("PlanCapture cycle finished: pages={}, captured={}, dup={}," + + " singleTable={}, filtered={}, fail={}", + pages, successCount.get(), skipDuplicateCount.get(), skipSingleTableCount.get(), + skipFilterCount.get(), failCount.get()); + } catch (Exception e) { Review Comment: [P2] Request a prompt resume when audit scanning throws. The cycle clears `pendingWindowNeedsPromptResume`, reserves the window, then `scanner.scan` can time out; this catch only logs, so `runAfterCatalogReady` leaves the next wakeup at the default three-hour interval. Keep the durable window pending and set the prompt flag on this error path. ########## fe/fe-core/src/main/java/org/apache/doris/nereids/spm/capture/AuditLogScanner.java: ########## @@ -0,0 +1,1270 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.nereids.spm.capture; + +import org.apache.doris.common.util.TimeUtils; +import org.apache.doris.qe.SqlModeHelper; +import org.apache.doris.qe.VariableMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneId; +import java.time.ZoneOffset; +import java.time.format.DateTimeFormatter; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** + * AuditLogScanner - audit log query wrapper (Phase 2, design doc 7.2.3). + * + * Reads the __internal_schema.audit_log internal table through the internal query + * mechanism and returns the high-value query candidates for SPM auto capture. + * + * Within a capture cycle the results are deduplicated by (catalog, db, sql_digest): the + * record with the largest query_time wins, so the same query SHAPE is only processed once + * per cycle - but only within one namespace. Identical unqualified SQL executed in two + * databases is a DIFFERENT query for SPM (its eventual match key is namespace-qualified), + * so the database / catalog must take part in the dedup key. + * + * Pagination: the batch LIMIT is applied with a stable total-order cursor (see + * {@link #ORDER_BY}); the caller resumes from the returned cursor until a batch comes + * back shorter than the limit (window exhausted); advancing the window past a truncated + * batch would permanently skip every eligible row beyond the LIMIT. + */ +public class AuditLogScanner { + + /** + * Cursor sentinel: no resume cursor is pending. A valid audit query_time is + * non-negative, so the sentinel lies outside the valid domain (a zero query_time is + * a perfectly valid cursor and must not be mistaken for "no cursor"). + */ + public static final long CURSOR_ABSENT = Long.MIN_VALUE; + + /** + * Cursor sentinel: the cursor row's query_time is NULL. query_time is nullable and + * eligibility also accepts large scan_rows alone, so a full page can legitimately end + * with a NULL query_time; NULL must stay distinguishable from a zero query_time so + * the resume predicate can compare it three-valued (IS NULL). + */ + public static final long CURSOR_QUERY_TIME_NULL = Long.MIN_VALUE + 1; + + /** + * Statement timeout (seconds) of the synchronous audit read. The default + * StatisticsUtil overload assigns the ANALYZE timeout (43,200 seconds), so a stalled + * internal-table read could hold the single capture cycle for half a day and delay + * every later capture / retry. The read is latency-sensitive: fail fast and let the + * next cycle retry. + */ + static final int AUDIT_SCAN_TIMEOUT_SECONDS = 30; + + private static final Logger LOG = LogManager.getLogger(AuditLogScanner.class); + + private static final DateTimeFormatter DATETIME_FORMAT = + DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); + + /** + * Lookback floor of the completion-aware scan lower bound (see buildScanSql): a query + * that started earlier than this before the window cannot be admitted even when its + * completion reaches into the window - the trade-off that keeps the range-partitioned + * audit table prunable instead of rescanning every retained partition per page. + */ + private static final long LATE_COMPLETION_LOOKBACK_MILLIS = java.util.concurrent.TimeUnit.DAYS + .toMillis(1); + + /** audit_log SELECT columns (order must match rowToCapturedQuery / toBatch). */ + private static final String SELECT_COLUMNS = + "`stmt`, `query_time`, `scan_rows`, `return_rows`, `sql_digest`, `sql_hash`, `db`, `catalog`," + + " `query_id`, `is_internal`, `time`, `sql_mode`, `client_ip`, md5(`stmt`)"; + + /** + * Canonical name of the row-content hash pseudo column (the last ORDER BY / cursor + * tie breaker). Auditing rows that agree on EVERY ordered key are content-duplicates + * (same statement, same client, same metrics), so skipping extra copies of such a + * content class is safe: capture dedupes by (catalog, db, digest) anyway. + */ + private static final String STMT_HASH_EXPR = "md5(`stmt`)"; + + /** + * Total order of the scan / cursor. The row-EVENT time is the FIRST key (with + * query_id / client_ip / metrics / statement hash as durable tie breakers): the + * audit loader writes rows asynchronously with the ORIGINAL event time, so a row + * published after page 1 can carry an event time OLDER than the current cursor - + * under the previous query_time-first order it sorted BEFORE the cursor and every + * resumed page skipped it forever. With the event time leading, an older-event-time + * row sorts AFTER the cursor and the resumed pages reach it; a row whose event time + * is newer than the cursor is picked up by the window overlap (see + * PlanCaptureManager#scanWindowOverlapMs, which follows the loader's configured + * batch interval). + */ + private static final String ORDER_BY = + " ORDER BY `time` DESC, `query_time` DESC, `query_id` DESC, `client_ip` DESC," + + " `sql_hash` DESC, `scan_rows` DESC, `return_rows` DESC, " + STMT_HASH_EXPR + + " DESC, `catalog` DESC, `db` DESC, `sql_mode` DESC "; + + /** + * Tail of the pagination cursor AFTER (query_time, time, query_id): client_ip, + * sql_hash, scan_rows, return_rows and the statement hash. The audit table is a + * DUPLICATE KEY table whose key omits client_ip, and the raw ORDER BY has NO + * genuinely unique column: without the tail, rows sharing the first three keys made + * the resume predicate either re-select the whole (NULL query_id) group forever or + * skip the remaining duplicates after the first LIMIT. The tail is persisted with + * the checkpoint (see PlanCaptureManager#cursorTail) so a restarted / handed-off + * leader resumes exactly after the last consumed row. + */ + public static final class CursorTail { + private final String clientIp; + private final String sqlHash; + private final String scanRows; + private final String returnRows; + private final String stmtHash; + /** The row's namespace + parser mode: two NaN-id rows otherwise identical in + * client / metrics can still be SEPARATE capture identities (toBatch dedupes by + * catalog+db+sql_mode+identity), so the ordered cursor must reach them too - + * without these keys the strict after-cursor chain excluded the second row on + * every later page. */ + private final String catalog; + private final String db; + private final String sqlMode; + /** Whether catalog / db / sql_mode are PART of this tail. A legacy + * five-element tail (written by a pre-upgrade leader) has no namespace keys: + * treating its absent keys as NULL values would extend the resume chain with + * three NULL comparisons and then terminate it - skipping the group's + * remaining rows, where the old chain still reached them. */ + private final boolean hasNamespaceKeys; + /** + * The session time zone the cursor's timestamp strings were RENDERED in (the + * audit writer's zone, see {@link #auditWriteZone()}); null in a tail written + * before the element existed. A PENDING window keeps scanning in this zone so a + * global time_zone change never mixes two renderings inside one window (the + * bounds are epoch millis re-formatted every cycle, while the cursor is the + * persisted string). + */ + private final String zoneId; + + CursorTail(String clientIp, String sqlHash, String scanRows, String returnRows, + String stmtHash) { + this(clientIp, sqlHash, scanRows, returnRows, stmtHash, null, null, null, false, null); + } + + CursorTail(String clientIp, String sqlHash, String scanRows, String returnRows, + String stmtHash, String catalog, String db, String sqlMode) { + this(clientIp, sqlHash, scanRows, returnRows, stmtHash, catalog, db, sqlMode, true, + null); + } + + CursorTail(String clientIp, String sqlHash, String scanRows, String returnRows, + String stmtHash, String catalog, String db, String sqlMode, String zoneId) { + this(clientIp, sqlHash, scanRows, returnRows, stmtHash, catalog, db, sqlMode, true, + zoneId); + } + + private CursorTail(String clientIp, String sqlHash, String scanRows, String returnRows, + String stmtHash, String catalog, String db, String sqlMode, + boolean hasNamespaceKeys, String zoneId) { + this.clientIp = clientIp; + this.sqlHash = sqlHash; + this.scanRows = scanRows; + this.returnRows = returnRows; + this.stmtHash = stmtHash; + this.catalog = catalog; + this.db = db; + this.sqlMode = sqlMode; + this.hasNamespaceKeys = hasNamespaceKeys; + this.zoneId = zoneId; + } + + String getClientIp() { + return clientIp; + } + + String getSqlHash() { + return sqlHash; + } + + String getScanRows() { + return scanRows; + } + + String getReturnRows() { + return returnRows; + } + + String getStmtHash() { + return stmtHash; + } + + String getCatalog() { + return catalog; + } + + String getDb() { + return db; + } + + String getSqlMode() { + return sqlMode; + } + + /** Whether the namespace / mode keys are part of this tail (see the field). */ + boolean hasNamespaceKeys() { + return hasNamespaceKeys; + } + + /** The zone the timestamp strings were rendered in (null = not recorded). */ + String getZoneId() { + return zoneId; + } + } + + /** Encodes a cursor tail as a compact JSON list (null-safe; empty text = absent). */ + public static String encodeCursorTail(CursorTail tail) { + if (tail == null + || (tail.getClientIp() == null && tail.getSqlHash() == null + && tail.getScanRows() == null && tail.getReturnRows() == null + && tail.getStmtHash() == null && tail.getCatalog() == null + && tail.getDb() == null && tail.getSqlMode() == null)) { + // no tail information at all (a row without the appended columns - e.g. a + // pre-column audit row or a fabricated test row): treat it as a PREFIX-only + // cursor; a JSON array of nulls would otherwise extend the resume chain with + // all-NULL keys and terminate it immediately. + return ""; + } + if (!tail.hasNamespaceKeys()) { + // a legacy tail round-trips as five elements so the decoder marks it legacy + // again (the namespace keys were never observed) + return new Gson().toJson(Arrays.asList(tail.getClientIp(), tail.getSqlHash(), + tail.getScanRows(), tail.getReturnRows(), tail.getStmtHash())); + } + if (tail.getZoneId() == null) { + // a tail written before the zone element existed (or by a fixture): keep the + // eight-element form so it decodes without a zone again + return new Gson().toJson(Arrays.asList(tail.getClientIp(), tail.getSqlHash(), + tail.getScanRows(), tail.getReturnRows(), tail.getStmtHash(), + tail.getCatalog(), tail.getDb(), tail.getSqlMode())); + } + return new Gson().toJson(Arrays.asList(tail.getClientIp(), tail.getSqlHash(), + tail.getScanRows(), tail.getReturnRows(), tail.getStmtHash(), + tail.getCatalog(), tail.getDb(), tail.getSqlMode(), tail.getZoneId())); + } + + /** Decodes a cursor tail; blank / broken / all-null text decodes to null (legacy cursor). */ + static CursorTail decodeCursorTail(String text) { + if (text == null || text.trim().isEmpty()) { + return null; + } + try { + List<String> values = new Gson().fromJson(text, + new TypeToken<List<String>>() { }.getType()); + if (values == null || values.size() < 5) { + return null; + } + boolean hasAnyValue = false; + for (String value : values) { + if (value != null) { + hasAnyValue = true; + break; + } + } + if (!hasAnyValue) { + return null; + } + if (values.size() < 8) { + // five (or partially extended) element tail written before the + // namespace / mode keys existed: it carries NO information about them, + // so the resume chain keeps the legacy prefix comparison + return new CursorTail(values.get(0), values.get(1), values.get(2), + values.get(3), values.get(4)); + } + if (values.size() < 9) { + // namespace-aware tail without the zone element (pre-zone writer) + return new CursorTail(values.get(0), values.get(1), values.get(2), + values.get(3), values.get(4), values.get(5), values.get(6), + values.get(7)); + } + return new CursorTail(values.get(0), values.get(1), values.get(2), + values.get(3), values.get(4), values.get(5), values.get(6), + values.get(7), values.get(8)); + } catch (RuntimeException e) { + return null; + } + } + + /** Value of a column that may be missing (pre-column rows); null when out of range. */ + private static String valueAt(ResultRow row, int index) { + List<String> values = row.getValues(); + return index < values.size() ? values.get(index) : null; + } + + /** + * Result of one audit scan: the namespace-deduplicated candidates plus the resume + * cursor (the full ORDER BY key tuple of the last RAW row read). + */ + public static class ScanBatch { + private final List<CapturedQuery> candidates; + private final boolean windowExhausted; + private final long cursorQueryTime; + private final String cursorTime; + private final String cursorQueryId; + private final String cursorTail; + + ScanBatch(List<CapturedQuery> candidates, boolean windowExhausted, + long cursorQueryTime, String cursorTime, String cursorQueryId, + String cursorTail) { + this.candidates = candidates; + this.windowExhausted = windowExhausted; + this.cursorQueryTime = cursorQueryTime; + this.cursorTime = cursorTime == null ? "" : cursorTime; + this.cursorQueryId = cursorQueryId == null ? "" : cursorQueryId; + this.cursorTail = cursorTail == null ? "" : cursorTail; + } + + ScanBatch(List<CapturedQuery> candidates, boolean windowExhausted, + long cursorQueryTime, String cursorTime, String cursorQueryId) { + this(candidates, windowExhausted, cursorQueryTime, cursorTime, cursorQueryId, ""); + } + + public List<CapturedQuery> getCandidates() { + return candidates; + } + + /** Whether the batch returned fewer RAW rows than the limit (whole window read). */ + public boolean isWindowExhausted() { + return windowExhausted; + } + + public long getCursorQueryTime() { + return cursorQueryTime; + } + + public String getCursorTime() { + return cursorTime; + } + + public String getCursorQueryId() { + return cursorQueryId; + } + + /** Encoded tail of the cursor (see {@link CursorTail}); empty = absent. */ + public String getCursorTail() { + return cursorTail; + } + } + + /** + * Scans the audit_log table within the given time window (first page). + * + * @param startTimeMs window start (epoch millis, inclusive) + * @param endTimeMs window end (epoch millis, exclusive) + * @param maxBatchSize max number of raw rows to scan (prevents OOM) + * @return the scan batch (candidates + resume cursor) + */ + public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize) { + return scan(startTimeMs, endTimeMs, maxBatchSize, CURSOR_ABSENT, "", "", ""); + } + + /** + * Scans the audit_log table within the given time window, resuming after the cursor + * returned by the previous batch. + * + * @param startTimeMs window start (epoch millis, inclusive) + * @param endTimeMs window end (epoch millis, exclusive) + * @param maxBatchSize max number of raw rows per batch + * @param cursorQueryTime query_time of the last consumed row (CURSOR_ABSENT = start + * from the top; CURSOR_QUERY_TIME_NULL = that row's value was + * NULL; any other value - including 0 - is a real cursor) + * @param cursorTime event time of the last consumed row; empty = SQL NULL + * @param cursorQueryId query_id of the last consumed row; empty = SQL NULL + * @return the scan batch (candidates + resume cursor) + */ + public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize, + long cursorQueryTime, String cursorTime, String cursorQueryId) { + return scan(startTimeMs, endTimeMs, maxBatchSize, cursorQueryTime, cursorTime, + cursorQueryId, ""); + } + + /** + * Scans the audit_log table within the given time window, resuming after the FULL + * cursor tuple (see {@link CursorTail}). + * + * <p>The thresholds are read from the CURRENT global session variables. Only callers + * outside the capture cycle (tests, tooling) may use this overload: the cycle owns a + * pinned snapshot of them and must pass it via + * {@link #scan(long, long, int, PlanCaptureFilter, long, String, String, String)}, so + * that the SQL stage and the in-memory + * {@link PlanCaptureFilter#shouldCapture} stage never compare against two different + * threshold sets. + * + * @param startTimeMs window start (epoch millis, inclusive) + * @param endTimeMs window end (epoch millis, exclusive) + * @param maxBatchSize max number of raw rows per batch + * @param cursorQueryTime query_time of the last consumed row (CURSOR_ABSENT = start + * from the top; CURSOR_QUERY_TIME_NULL = that row's value was + * NULL; any other value - including 0 - is a real cursor) + * @param cursorTime event time of the last consumed row; empty = SQL NULL + * @param cursorQueryId query_id of the last consumed row; empty = SQL NULL + * @param cursorTail encoded tail of the last consumed row (empty = legacy cursor + * without a tail: the resume predicate falls back to the + * (time, query_time, query_id) prefix) + * @return the scan batch (candidates + resume cursor) + */ + public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize, + long cursorQueryTime, String cursorTime, String cursorQueryId, String cursorTail) { + // no pattern is used here, so only the thresholds matter: the constructor reads + // them from the current globals + return scan(startTimeMs, endTimeMs, maxBatchSize, new PlanCaptureFilter(null, null), + cursorQueryTime, cursorTime, cursorQueryId, cursorTail); + } + + /** + * Scans the audit_log table within the given time window, resuming after the FULL + * cursor tuple (see {@link CursorTail}), with the thresholds of the given filter. + * + * <p>Deriving the SQL thresholds from the SAME filter instance that later decides + * {@link PlanCaptureFilter#shouldCapture} is what keeps the two stages consistent: the + * SQL returns a row exactly when the filter would accept it, so no row the filter + * rejects is ever consumed (marked processed) and no row the filter accepts is + * unreachable behind the cursor. Reading the globals here instead made a + * `SET GLOBAL plan_capture_min_query_time_ms` between the cycle's filter construction + * and this statement return rows the stale in-memory filter then failed TERMINALLY, + * and a LOWERED threshold made already-passed rows unreachable below the cursor. + * + * @param startTimeMs window start (epoch millis, inclusive) + * @param endTimeMs window end (epoch millis, exclusive) + * @param maxBatchSize max number of raw rows per batch + * @param filter the threshold snapshot of this window (the caller's filter) + * @param cursorQueryTime query_time of the last consumed row (CURSOR_ABSENT = start + * from the top; CURSOR_QUERY_TIME_NULL = that row's value was + * NULL; any other value - including 0 - is a real cursor) + * @param cursorTime event time of the last consumed row; empty = SQL NULL + * @param cursorQueryId query_id of the last consumed row; empty = SQL NULL + * @param cursorTail encoded tail of the last consumed row (empty = legacy cursor + * without a tail: the resume predicate falls back to the + * (time, query_time, query_id) prefix) + * @return the scan batch (candidates + resume cursor) + */ + public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize, + PlanCaptureFilter filter, long cursorQueryTime, String cursorTime, + String cursorQueryId, String cursorTail) { + return scan(startTimeMs, endTimeMs, maxBatchSize, filter, cursorQueryTime, + cursorTime, cursorQueryId, cursorTail, null); + } + + /** + * As the eight-argument overload, with the zone the window must be rendered in when + * its cursor does not carry one (a window whose FIRST pass runs now). + * + * <p>audit_log.time is the audit WRITER's local rendering and the writer follows the + * global session time_zone, so after {@code SET GLOBAL time_zone} the rows published + * BEFORE the change are stored in the OLD rendering and are invisible to bounds + * rendered in the new zone - the reviewer's example: a 10:00 UTC row stored as + * "10:00" is searched as [17:00, 20:00) after the zone becomes +08, the (empty) page + * looks exhausted and the watermark moves past the row forever. The window is + * therefore opened in the zone the PREVIOUS scan used while that differs from the + * global zone: the old rendering's rows are found first, and the following pass (see + * PlanCaptureManager's exhaustion branch) revisits the SAME window in the new zone + * for the rows published after the change. Each pass is a single rendering, so the + * keyset pagination keeps walking one consistent total order. + * + * @param firstPassZoneId zone ID of the previous scan pass (empty / null = follow the + * current global time_zone) + * @return the scan batch (candidates + resume cursor) + */ + public ScanBatch scan(long startTimeMs, long endTimeMs, int maxBatchSize, + PlanCaptureFilter filter, long cursorQueryTime, String cursorTime, + String cursorQueryId, String cursorTail, String firstPassZoneId) { + // The bounds are rendered in the zone the AUDIT WRITER used (the global session + // time_zone, see auditWriteZone) - not the FE host zone - because + // __internal_schema.audit_log.time stores the writer's rendering. A PENDING + // window keeps the zone recorded in its cursor: the window's epoch bounds are + // re-rendered every cycle while the cursor is the persisted string, so a global + // time_zone change mid-window would otherwise compare two different renderings + // and skip the whole unconsumed range. Only a NEW window follows a changed + // global zone. + ZoneId auditZone = scanZoneFor(cursorTail); + if (zoneOfTail(cursorTail) == null && firstPassZoneId != null && !firstPassZoneId.isEmpty()) { + ZoneId firstPassZone = parseZone(firstPassZoneId); + if (firstPassZone != null && !firstPassZone.equals(auditZone)) { + LOG.info("SPM audit scan opens the window in zone {} (the global time_zone is" + + " now {}): its already published rows were rendered under the" + + " previous zone", firstPassZone, auditZone); + auditZone = firstPassZone; + } + } + // window bounds as MONOTONE wall-clock ranges: a UTC window crossing a DST + // transition renders as several local ranges (see localTimeRanges), never as one + // inverted range that matches nothing + List<String[]> windowRanges = localTimeRanges(startTimeMs, endTimeMs, auditZone); + // defense in depth: a non-positive batch size can no longer be written through + // SQL SET (see SessionVariable), but LIMIT 0 here would mark the window + // exhausted on an empty page and advance the watermark over every eligible row + int limit = Math.max(1, maxBatchSize); + + long minQueryTimeMs = filter.getMinQueryTimeMs(); + long minScanRows = filter.getMinScanRows(); + String sql = buildScanSql(windowRanges, limit, minQueryTimeMs, minScanRows, + cursorPredicate(cursorQueryTime, cursorTime, cursorQueryId, cursorTail), + zoneOffsetSwingSeconds(auditZone), lateCompletionFloor(startTimeMs, auditZone)); + + // bounded statement timeout: see AUDIT_SCAN_TIMEOUT_SECONDS (the no-timeout + // overload would inherit the 12h analyze timeout) + List<ResultRow> rows = StatisticsUtil.execStatisticQuery(sql, false, + AUDIT_SCAN_TIMEOUT_SECONDS); + return toBatch(rows, limit, auditZone); + } + + /** + * The zone bounds and the resume cursor are rendered in for the given cursor: a + * PENDING window keeps the zone recorded in its cursor (its epoch bounds are + * re-rendered every cycle while the cursor is the persisted string - a global + * time_zone change would otherwise mix two renderings inside one window), and a NEW + * window follows the current global zone (the audit writer's own zone). + */ + static ZoneId scanZoneFor(String cursorTail) { + ZoneId pendingZone = zoneOfTail(cursorTail); + if (pendingZone == null) { + return auditWriteZone(); + } + if (!pendingZone.equals(auditWriteZone())) { + LOG.info("SPM audit scan continues the pending window in zone {} (the global" + + " time_zone is now {}); the next window follows the new zone", + pendingZone, auditWriteZone()); + } + return pendingZone; + } + + /** + * The zone the audit WRITER rendered its timestamps in. AuditLoader formats the + * event time with {@link TimeUtils}, which on its own (context-less) worker thread + * falls back to the GLOBAL session variable time_zone; the scan bounds must use + * exactly the same zone, otherwise a non-UTC host zone makes every stored row fall + * outside the windows (or renders window bounds that match nothing). + */ + static ZoneId auditWriteZone() { + return TimeUtils.getOrSystemTimeZone( + VariableMgr.getDefaultSessionVariable().getTimeZone()).toZoneId(); + } + + /** The zone recorded in a cursor tail, or null when absent / unparsable. */ + private static ZoneId zoneOfTail(String cursorTail) { + CursorTail tail = decodeCursorTail(cursorTail); + if (tail == null || tail.getZoneId() == null || tail.getZoneId().isEmpty()) { + return null; + } + return parseZone(tail.getZoneId()); + } + + /** + * The zone ID recorded in a cursor tail, or null when the tail is absent / carries no + * (parsable) zone - a caller deciding whether the window still owes a pass in another + * zone needs exactly this distinction (see PlanCaptureManager's scan-zone handoff). + */ + static String zoneIdOfTail(String cursorTail) { + ZoneId zone = zoneOfTail(cursorTail); + return zone == null ? null : zone.getId(); + } + + /** Parses a zone ID (aliases allowed); an unusable value is null, never an error. */ + private static ZoneId parseZone(String zoneId) { + try { + return ZoneId.of(zoneId, TimeUtils.timeZoneAliasMap); + } catch (RuntimeException e) { + LOG.warn("SPM audit scan ignores an unparsable zone '{}': {}", zoneId, e.getMessage()); + return null; + } + } + + /** + * Turns one page of raw audit rows into a batch: namespace-aware dedup plus the + * resume cursor. Package-visible for tests (the SQL / pagination contract is tested + * against fabricated rows). + * + * @param rows the raw rows of one page + * @param maxBatchSize the batch limit (a shorter page exhausts the window) + * @return the scan batch + */ + static ScanBatch toBatch(List<ResultRow> rows, int maxBatchSize) { + return toBatch(rows, maxBatchSize, auditWriteZone()); + } + + /** + * Turns one page of raw audit rows into a batch with an explicit timestamp zone (the + * zone the bounds were rendered in; it travels with the cursor so a pending window + * keeps its rendering, see {@link CursorTail#getZoneId()}). + * + * @param rows the raw rows of one page + * @param maxBatchSize the batch limit (a shorter page exhausts the window) + * @param auditZone the zone the window bounds were rendered in + * @return the scan batch + */ + static ScanBatch toBatch(List<ResultRow> rows, int maxBatchSize, ZoneId auditZone) { + if (rows == null || rows.isEmpty()) { + return new ScanBatch(List.of(), true, CURSOR_ABSENT, "", ""); + } + Map<String, CapturedQuery> deduped = new LinkedHashMap<>(); + long lastQueryTime = CURSOR_ABSENT; + String lastTime = ""; + String lastQueryId = ""; + CursorTail lastTail = null; + for (ResultRow row : rows) { + // the cursor always moves to the last RAW row read, even when that row is + // unusable / filtered later: it has been consumed and must not be scanned + // again by the next page. query_time is nullable: keep NULL distinguishable + // from 0 (both are valid cursors, but the resume predicate must compare a + // NULL three-valued through IS NULL). + String rawQueryTime = row.get(1); + lastQueryTime = rawQueryTime == null + ? CURSOR_QUERY_TIME_NULL : parseLong(rawQueryTime); + lastTime = row.getWithDefault(10, ""); + lastQueryId = row.getWithDefault(8, ""); + // the full ORDER BY key tuple: without the tail a group of rows sharing + // (time, query_time, query_id) either repeated forever (NULL query_id group) + // or was skipped after the first LIMIT (duplicate non-NULL tuples) + lastTail = new CursorTail(valueAt(row, 12), valueAt(row, 5), valueAt(row, 2), + valueAt(row, 3), valueAt(row, 13), valueAt(row, 7), valueAt(row, 6), + valueAt(row, 11), auditZone == null ? null : auditZone.getId()); + CapturedQuery candidate = rowToCapturedQuery(row); + if (candidate == null || candidate.getStmt() == null || candidate.getStmt().isEmpty()) { + continue; + } + String digest = candidate.getSqlDigest(); + if (digest == null || digest.isEmpty()) { + digest = candidate.getStmt(); + } + // namespace-aware key + SPM-match identity: the database / catalog take part + // (SPM namespace-qualifies its match key), the ORIGINATING parser mode takes + // part (a || b parses differently under PIPES_AS_CONCAT), and the digest is + // refined with the CONCRETE generator arguments SPM keeps unparameterized + // (the digest masks literals, so explode(split(s,',')) and explode(split(s,';')) + // would otherwise collapse although they are different baselines). + String key = candidate.getCatalog() + '\u0001' + candidate.getDb() + '\u0001' + + candidate.getSqlMode() + '\u0001' + + dedupIdentity(candidate.getStmt(), digest, candidate.getSqlMode()); + deduped.merge(key, candidate, (a, b) -> b.getQueryTimeMs() >= a.getQueryTimeMs() ? b : a); + } + return new ScanBatch(new ArrayList<>(deduped.values()), rows.size() < maxBatchSize, + lastQueryTime, lastTime, lastQueryId, encodeCursorTail(lastTail)); + } + + /** + * SPM-match identity of one audit row used for the dedup: the audit digest when + * present (it already covers the whole logical shape), refined with the CONCRETE + * generator arguments for statements that mention a generator - SPM deliberately + * keeps LATERAL VIEW / UNNEST arguments concrete and compares them exactly, so two + * same-digest statements with different arguments are different baselines. A row + * without an audit digest falls back to its text (never coarser than SPM). + * + * The generator arguments are parsed under the row's ORIGINATING mode: the capture + * daemon's ambient mode can differ (a NO_BACKSLASH_ESCAPES session's split '\a' + * means backslash + a, while the default mode reads '\a' as 'a'), and parsing both + * rows in the daemon's mode made their fingerprints equal although SPM compares + * them concretely - one eligible row was discarded as a duplicate. + * + * <p>Package-private for tests: the identity is the only observable of the gate. + * + * @param stmt the audit statement text + * @param digest the audit digest (null / empty falls back to the statement) + * @param sqlMode the ORIGINATING parser mode of the row + * @return the dedup identity + */ + @VisibleForTesting + static String dedupIdentity(String stmt, String digest, long sqlMode) { + if (digest == null || digest.isEmpty()) { + return stmt; + } + // The digest renders every literal as "?" and every scan selector as + // PARTITION(?) / TABLET(?): two statements that differ ONLY in a concrete selector + // (PARTITION(p1) vs PARTITION(p2)) are different baselines - matching compares the + // selectors in sameScanIdentity - so the selector fingerprint joins the identity + // for statements that mention one. Generator arguments join for the same reason + // (SPM keeps LATERAL VIEW / UNNEST arguments concrete). + String generators = mentionsGenerator(stmt) ? generatorFingerprint(stmt, sqlMode) : ""; + String selectors = mentionsScanSelector(stmt) + ? scanSelectorFingerprint(stmt, sqlMode) : ""; + if (generators.isEmpty() && selectors.isEmpty()) { + return digest; + } + return digest + '\u0001' + generators + '\u0001' + selectors; + } + + /** + * Whether the statement can carry a concrete scan selector. The tokens are the ones + * the grammar actually spells out; each one is a selector the audit digest MASKS + * (PARTITION(p1) and PARTITION(p2) both render as PARTITION(?)) while SPM compares it + * concretely (sameScanIdentity / sameScanParams), so the fingerprint must join the + * dedup identity for exactly these statements: + * <ul> + * <li>PARTITION / TABLET / TABLESAMPLE / INDEX: specifiedPartition, tabletList, + * sample and index selectors;</li> + * <li>"FOR VERSION AS OF" / "FOR TIME AS OF": tableSnapshot. The formerly checked + * "FOR TIMESTAMP" is not a form the grammar accepts, so a statement using + * time travel never got a fingerprint and two same-digest variants (only the + * version / time differs) collapsed into one identity - the capture then kept + * one of them and dropped the other;</li> + * <li>'@': optScanParams, the relation-level scan parameters that SPM keeps + * concrete (sameScanParams compares type + payloads), i.e. the {@code @branch} + * / {@code @incr} / {@code @tag} / {@code @options} forms.</li> + * </ul> + * A statement mentioning none of them keeps the plain digest: the gate only has to be + * a cheap pre-filter, over-matching costs one parse, under-matching loses identity. + * + * <p>Package-private for tests. + * + * @param stmt the audit statement text + * @return whether the statement needs its concrete selectors in the dedup identity + */ + @VisibleForTesting + static boolean mentionsScanSelector(String stmt) { + if (stmt == null) { + return false; + } + String upper = stmt.toUpperCase(java.util.Locale.ROOT); + return upper.contains("PARTITION") || upper.contains("TABLET") + || upper.contains("TABLESAMPLE") || upper.contains("INDEX") + || upper.contains("FOR VERSION AS OF") || upper.contains("FOR TIME AS OF") + || upper.indexOf('@') >= 0; + } + + /** The concrete scan selectors of the statement (the full text when unparsable). */ + private static String scanSelectorFingerprint(String stmt, long sqlMode) { + try { + return org.apache.doris.qe.SqlModeHelper.withSqlMode(sqlMode, () -> { + org.apache.doris.nereids.trees.plans.Plan parsed = + new org.apache.doris.nereids.parser.NereidsParser().parseSingle(stmt); + return org.apache.doris.nereids.spm.SPMPlanTreeSupport + .scanSelectorFingerprint(parsed); + }); + } catch (Throwable t) { + // unparsable: keep the full-text identity, never a coarser one + return stmt; + } + } + + private static boolean mentionsGenerator(String stmt) { + if (stmt == null) { + return false; + } + String upper = stmt.toUpperCase(java.util.Locale.ROOT); + return upper.contains("LATERAL VIEW") || upper.contains("UNNEST"); + } + + /** The concrete generator arguments of the statement (the full text when unparsable). */ + private static String generatorFingerprint(String stmt, long sqlMode) { + try { + String fingerprint = org.apache.doris.qe.SqlModeHelper.withSqlMode(sqlMode, () -> { + org.apache.doris.nereids.trees.plans.Plan parsed = + new org.apache.doris.nereids.parser.NereidsParser().parseSingle(stmt); + StringBuilder sb = new StringBuilder(); + org.apache.doris.nereids.spm.SPMPlanTreeSupport.<RuntimeException>walkPlans( + parsed, node -> { + if (node instanceof org.apache.doris.nereids.trees.plans.logical + .LogicalGenerate) { + org.apache.doris.nereids.trees.plans.logical.LogicalGenerate<?> generate = + (org.apache.doris.nereids.trees.plans.logical.LogicalGenerate<?>) + node; + for (org.apache.doris.nereids.trees.expressions.Expression generator + : generate.getGenerators()) { + sb.append(generator.toSql()).append('|'); + } + } + }); + return sb.toString(); + }); + return fingerprint; + } catch (Throwable t) { + // unparsable: keep the full-text identity, never a coarser one + return stmt; + } + } + + /** + * Builds the audit_log scan SQL. Public for tests: the pushed-down predicate shape is + * part of the capture contract - eligibility is query time OR scanned rows (the same + * rule as PlanCaptureFilter), and internal maintenance queries are filtered in SQL + * instead of relying on a hardcoded event flag. + * + * @param start window start timestamp (formatted) + * @param end window end timestamp (formatted) + * @param maxBatchSize LIMIT for the scan + * @param minQueryTimeMs query-time threshold + * @param minScanRows scan-rows threshold + * @return the scan SQL + */ + public static String buildScanSql(String start, String end, int maxBatchSize, + long minQueryTimeMs, long minScanRows) { + return buildScanSql(start, end, maxBatchSize, minQueryTimeMs, minScanRows, ""); + } + + /** + * Builds the audit_log scan SQL with an optional resume-cursor predicate. The ORDER + * BY defines the stable total order the cursor walks (see {@link #ORDER_BY}): the + * row EVENT time first, then every remaining identity / metric key as a durable tie + * breaker. + * + * @param start window start timestamp (formatted) + * @param end window end timestamp (formatted) + * @param maxBatchSize LIMIT for the scan + * @param minQueryTimeMs query-time threshold + * @param minScanRows scan-rows threshold + * @param cursorPredicate resume-cursor predicate (empty when starting at the top) + * @return the scan SQL + */ + public static String buildScanSql(String start, String end, int maxBatchSize, + long minQueryTimeMs, long minScanRows, String cursorPredicate) { + return buildScanSql(List.<String[]>of(new String[] {start, end}), maxBatchSize, + minQueryTimeMs, minScanRows, cursorPredicate, 0L); + } + + /** + * As {@link #buildScanSql(String, String, int, long, long, String)} with an explicit + * offset swing for the completion-aware lower bound (see + * {@link #buildScanSql(List, int, long, long, String, long)}): the single-range form a + * fixed-offset zone produces. + * + * @param start window start timestamp (formatted) + * @param end window end timestamp (formatted) + * @param maxBatchSize LIMIT for the scan + * @param minQueryTimeMs query-time threshold + * @param minScanRows scan-rows threshold + * @param cursorPredicate resume-cursor predicate (empty when starting at the top) + * @param offsetSwingSeconds max offset swing of the window's zone (0 = none) + * @return the scan SQL + */ + static String buildScanSql(String start, String end, int maxBatchSize, + long minQueryTimeMs, long minScanRows, String cursorPredicate, + long offsetSwingSeconds) { + return buildScanSql(List.<String[]>of(new String[] {start, end}), maxBatchSize, + minQueryTimeMs, minScanRows, cursorPredicate, offsetSwingSeconds); + } + + /** + * Builds the audit_log scan SQL for a window rendered as one or MORE monotone + * wall-clock ranges (see {@link #localTimeRanges}) and with the zone's offset swing + * applied to the completion-aware lower bound (see + * {@link #zoneOffsetSwingSeconds}). + * + * @param windowRanges (start, end) wall-clock pairs of the window + * @param maxBatchSize LIMIT for the scan + * @param minQueryTimeMs query-time threshold + * @param minScanRows scan-rows threshold + * @param cursorPredicate resume-cursor predicate (empty when starting at the top) + * @param offsetSwingSeconds max offset swing of the window's zone (0 = none) + * @return the scan SQL + */ + static String buildScanSql(List<String[]> windowRanges, int maxBatchSize, + long minQueryTimeMs, long minScanRows, String cursorPredicate, + long offsetSwingSeconds) { + return buildScanSql(windowRanges, maxBatchSize, minQueryTimeMs, minScanRows, + cursorPredicate, offsetSwingSeconds, + completeWindowFloor(windowRanges.get(0)[0])); + } + + /** + * As the six-argument overload with an EXPLICIT completion floor (the top-level + * partition-pruning lower bound), rendered by the caller from the window-start + * INSTANT (see {@link #lateCompletionFloor}): the string form is civil arithmetic and + * therefore wrong across a DST transition (see + * {@link #completeWindowFloor(String)}). + * + * @param floor the already rendered completion floor (see {@link #lateCompletionFloor}) + * @return the scan SQL + */ + static String buildScanSql(List<String[]> windowRanges, int maxBatchSize, + long minQueryTimeMs, long minScanRows, String cursorPredicate, + long offsetSwingSeconds, String floor) { + // The window lower bound is COMPLETION-aware: audit_log.time is the query's START + // time, but its row is published only when the query FINISHES. A long-running + // query started at 11:50 is absent from the 12:00 scan; without the + // completion predicate the next (default three-hour) window starts at + // 12:00 - overlap, so its 11:50 row - now visible - would be excluded + // FOREVER. Rows are therefore also eligible while their completion + // (time + query_time) reaches into the window. + // + // The window predicate may therefore only bound the START time from ABOVE and + // split the ranges for the LOWER bound: conjoining the start-time membership + // (`time >= window start`, which windowPredicate implies for the first range) + // nullified the completion branch entirely - the earlier-start row the branch + // exists for failed the conjunct on every later scan. Only the upper bound and + // the partitionable floor are top-level conjuncts; the lower bound is the OR of + // (start-time membership in one of the ranges, completion reaching the window + // start). + // The completion branch is BOUNDED by a floor: an unbounded + // "time >= start OR completion >= start" cannot prune ANY old partition of the + // range-partitioned audit table (query_time is only known per row), so every + // keyset page would rescan retained history under the short timeout. The floor + // (start - LATE_COMPLETION_LOOKBACK_MILLIS) keeps the pruning intact for a + // three-hour window while still admitting every query whose completion reaches + // into it; a query LONGER than the lookback is the documented miss. + // + // The completion is civil arithmetic on the writer's LOCAL rendering, so it must + // be widened by the zone's offset swing: a query started 01:30 PST (09:30Z) that + // finishes 03:10:01 PDT computes as 02:10:01 without the swing, and a window + // starting 03:05 would exclude the row on EVERY later scan (no overlap reaches it + // again). Adding the swing seconds makes the bound conservative in the admitting + // direction, which is the safe side for a late-completion lookback. + String start = windowRanges.get(0)[0]; + String lastEnd = windowRanges.get(windowRanges.size() - 1)[1]; + String completionBound = "timestampadd(SECOND, CAST(`query_time` / 1000 AS BIGINT)" Review Comment: [P2] Preserve milliseconds in the late-completion predicate. `time` and `query_time` have millisecond precision, but this cast truncates duration to whole seconds. If an earlier scan advances without the row's horizon and it publishes late, a row at 11:50:00.900 lasting 299100 ms truly completes at the next 11:55 overlap start; SQL computes 11:54:59.900 and excludes it. Compare at millisecond precision or round positive durations up conservatively. ########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditPublicationHorizon.java: ########## @@ -0,0 +1,263 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.plugin.audit; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.InternalSchema; +import org.apache.doris.common.FeConstants; +import org.apache.doris.common.util.TimeUtils; +import org.apache.doris.qe.AuditEventProcessor; +import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + +/** + * The cluster-wide AUDIT PUBLICATION HORIZON: the start time (epoch millis, the + * {@code time} column of {@code audit_log}) of the oldest audit event that any FE has + * accepted but not yet PUBLISHED. The SPM capture scans the shared audit table from the + * leader, so it uses this value as a progress FENCE: its next scan window must still + * start at or before it, otherwise a row an FE still owes falls behind the advanced + * watermark and is never captured. + * + * <p>Three layers make the fence complete: + * <ul> + * <li>{@link #localHorizon()} folds THIS FE's whole audit pipeline: completed queries + * still held by the {@link WorkloadRuntimeStatusMgr} (they enter the pipeline + * before any loader sees them), the {@link AuditEventProcessor} queue and its + * in-flight event (a plugin can stall while an event is dequeued), and the + * {@link AuditLoader} queue / assembled batch / not-yet-visible batch (a stream + * load can report Publish Timeout after commit).</li> + * <li>each FE REPORTS its local horizon into the shared + * {@link InternalSchema#SPM_AUDIT_HORIZON_TBL_NAME} table, so a follower's + * backlog is visible to the leader that runs the capture.</li> + * <li>{@link #clusterHorizon()} is the MINIMUM over the local value and the FRESH + * rows of that table; a row the reporter stopped refreshing (its FE died or + * stopped reporting - the events are gone with it) is ignored.</li> + * </ul> + */ +public final class AuditPublicationHorizon { + + private static final Logger LOG = LogManager.getLogger(AuditPublicationHorizon.class); + + /** + * A reported row older than this is IGNORED: its FE stopped refreshing the fence + * (crashed / killed / its reporter thread is gone), so the events it still owed are + * lost with it and fencing progress forever would freeze the capture instead of + * protecting anything. Must be comfortably larger than the reporter's keepalive + * interval ({@link AuditLoader#HORIZON_KEEPALIVE_MILLIS}). + */ + public static final long ROW_STALE_MILLIS = 5 * 60 * 1000L; + + private static final String SELECT_ROWS_SQL = + "SELECT `horizon_ms`, `update_time` FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`"; + private static final String DELETE_OWN_ROW_SQL = "DELETE FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "` WHERE `fe_name` = '${feName}'"; + private static final String INSERT_OWN_ROW_SQL = "INSERT INTO `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`" + + " (`fe_name`, `horizon_ms`, `update_time`) VALUES ('${feName}', ${horizonMs}, '${updateTime}')"; + private static final int IO_TIMEOUT_SECONDS = 10; + + /** + * Test seam: the shared-table read (one row per FE). Null in production. + */ + @VisibleForTesting + static volatile Supplier<List<Object[]>> horizonRowsReaderForTest; + + /** + * Test seam: the shared-table write of this FE's row (delete + optional insert). + * Null in production. + */ + @VisibleForTesting + static volatile Consumer<Long> localHorizonWriterForTest; + + private AuditPublicationHorizon() { + } + + /** + * The oldest audit event THIS FE has accepted but not published, 0 when nothing is + * outstanding: the MINIMUM over every stage of the local pipeline (see the class + * javadoc). Cheap - no I/O - so callers may poll it. + */ + public static long localHorizon() { + long oldest = 0; + oldest = minPositive(oldest, AuditLoader.oldestUnpublishedEventTime()); + oldest = minPositive(oldest, preLoaderHorizon()); + return oldest; + } + + /** The stages BEFORE the audit loader: held completed queries and the processor. */ + private static long preLoaderHorizon() { + long oldest = 0; + try { + AuditEventProcessor processor = Env.getCurrentAuditEventProcessor(); + if (processor != null) { + oldest = minPositive(oldest, processor.oldestQueuedOrInFlightEventTime()); + } + } catch (Throwable t) { + // an FE without this component (tests, partial startup) has nothing to fence + LOG.debug("audit publication horizon: the audit event processor is unavailable: {}", + t.getMessage()); + } + try { + Env env = Env.getCurrentEnv(); + WorkloadRuntimeStatusMgr mgr = env == null ? null : env.getWorkloadRuntimeStatusMgr(); + if (mgr != null) { + oldest = minPositive(oldest, mgr.oldestHeldAuditEventTime()); + } + } catch (Throwable t) { + LOG.debug("audit publication horizon: the workload runtime status manager is" + + " unavailable: {}", t.getMessage()); + } + return oldest; + } + + /** + * The fence the CAPTURE uses: the minimum over this FE's own pipeline and the fresh + * rows every other FE reported. Throws {@link IllegalStateException} when the shared + * table cannot be read - the caller must NOT advance without a complete fence + * (round-36 #1: an unreadable follower row is exactly the hole this guards). + */ + public static long clusterHorizon() { + long oldest = localHorizon(); + return minPositive(oldest, remoteHorizon()); + } + + /** + * The minimum horizon over the FRESH rows of the shared table (0 when none / all + * stale). A read failure propagates as a retryable {@link IllegalStateException}. + */ + private static long remoteHorizon() { + List<Object[]> rows; + Supplier<List<Object[]>> reader = horizonRowsReaderForTest; + if (reader != null) { + rows = reader.get(); + } else if (!sharedTableAvailable()) { + return 0L; // no live FE environment (unit tests / not ready): nothing reported + } else { + try { + List<ResultRow> result = StatisticsUtil.executeQuery( + SELECT_ROWS_SQL, Collections.emptyMap(), IO_TIMEOUT_SECONDS); + rows = new ArrayList<>(); + if (result != null) { + for (ResultRow row : result) { + List<String> values = row.getValues(); + if (values == null || values.size() < 2) { + continue; + } + rows.add(new Object[] {Long.parseLong(values.get(0).trim()), + TimeUtils.timeStringToLong(values.get(1).trim())}); Review Comment: [P2] Store the horizon freshness time in a fixed zone. A follower writes `update_time` with its current global time zone, but the leader parses that zone-less DATETIME with its own zone. During `SET GLOBAL time_zone` propagation, a fresh UTC 12:00 row read by a +08 leader looks eight hours old and is discarded, so an old follower event can fall behind capture's watermark. Persist epoch milliseconds or use UTC on both sides. ########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditLoader.java: ########## @@ -329,4 +647,68 @@ public void run() { } } } + + /** + * Reports THIS FE's audit publication horizon into the shared table so the leader + * that runs the capture sees a follower's backlog (round-36 #1). It writes when the + * value CHANGED and re-reports an unchanged non-zero value on the keepalive cadence + * (the reader ignores rows whose reporter went silent). A zero horizon is reported + * once (which removes the row). + */ + private class HorizonReporter implements Runnable { + + @Override + public void run() { + long lastReported = -1; + long lastReportAt = 0; + while (!isClosed) { + try { + Thread.sleep(HORIZON_REPORT_TICK_MILLIS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + if (isClosed) { + return; + } + long horizon; + try { + horizon = AuditPublicationHorizon.localHorizon(); + } catch (Throwable t) { + LOG.warn("audit horizon reporter: cannot compute the local horizon: {}", + t.getMessage()); + continue; + } + long now = System.currentTimeMillis(); + boolean changed = horizon != lastReported; + boolean keepAlive = horizon > 0 + && now - lastReportAt >= HORIZON_KEEPALIVE_MILLIS; + if (changed || keepAlive) { + AuditPublicationHorizon.reportLocalHorizon(horizon); + lastReported = horizon; Review Comment: [P2] Acknowledge a positive horizon only after its row is readable. `reportLocalHorizon` returns silently when the table is unavailable or SQL fails, and even SQL OK can leave a COMMITTED INSERT unpublished; this reporter still records `lastReported`/`lastReportAt`. An unchanged old event can therefore have no master-visible fence until the 60-second keepalive. Confirm visibility and retry unreported values on the next tick. ########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditPublicationHorizon.java: ########## @@ -0,0 +1,263 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.plugin.audit; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.InternalSchema; +import org.apache.doris.common.FeConstants; +import org.apache.doris.common.util.TimeUtils; +import org.apache.doris.qe.AuditEventProcessor; +import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + +/** + * The cluster-wide AUDIT PUBLICATION HORIZON: the start time (epoch millis, the + * {@code time} column of {@code audit_log}) of the oldest audit event that any FE has + * accepted but not yet PUBLISHED. The SPM capture scans the shared audit table from the + * leader, so it uses this value as a progress FENCE: its next scan window must still + * start at or before it, otherwise a row an FE still owes falls behind the advanced + * watermark and is never captured. + * + * <p>Three layers make the fence complete: + * <ul> + * <li>{@link #localHorizon()} folds THIS FE's whole audit pipeline: completed queries + * still held by the {@link WorkloadRuntimeStatusMgr} (they enter the pipeline + * before any loader sees them), the {@link AuditEventProcessor} queue and its + * in-flight event (a plugin can stall while an event is dequeued), and the + * {@link AuditLoader} queue / assembled batch / not-yet-visible batch (a stream + * load can report Publish Timeout after commit).</li> + * <li>each FE REPORTS its local horizon into the shared + * {@link InternalSchema#SPM_AUDIT_HORIZON_TBL_NAME} table, so a follower's + * backlog is visible to the leader that runs the capture.</li> + * <li>{@link #clusterHorizon()} is the MINIMUM over the local value and the FRESH + * rows of that table; a row the reporter stopped refreshing (its FE died or + * stopped reporting - the events are gone with it) is ignored.</li> + * </ul> + */ +public final class AuditPublicationHorizon { + + private static final Logger LOG = LogManager.getLogger(AuditPublicationHorizon.class); + + /** + * A reported row older than this is IGNORED: its FE stopped refreshing the fence + * (crashed / killed / its reporter thread is gone), so the events it still owed are + * lost with it and fencing progress forever would freeze the capture instead of + * protecting anything. Must be comfortably larger than the reporter's keepalive + * interval ({@link AuditLoader#HORIZON_KEEPALIVE_MILLIS}). + */ + public static final long ROW_STALE_MILLIS = 5 * 60 * 1000L; + + private static final String SELECT_ROWS_SQL = + "SELECT `horizon_ms`, `update_time` FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`"; + private static final String DELETE_OWN_ROW_SQL = "DELETE FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "` WHERE `fe_name` = '${feName}'"; + private static final String INSERT_OWN_ROW_SQL = "INSERT INTO `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`" + + " (`fe_name`, `horizon_ms`, `update_time`) VALUES ('${feName}', ${horizonMs}, '${updateTime}')"; + private static final int IO_TIMEOUT_SECONDS = 10; + + /** + * Test seam: the shared-table read (one row per FE). Null in production. + */ + @VisibleForTesting + static volatile Supplier<List<Object[]>> horizonRowsReaderForTest; + + /** + * Test seam: the shared-table write of this FE's row (delete + optional insert). + * Null in production. + */ + @VisibleForTesting + static volatile Consumer<Long> localHorizonWriterForTest; + + private AuditPublicationHorizon() { + } + + /** + * The oldest audit event THIS FE has accepted but not published, 0 when nothing is + * outstanding: the MINIMUM over every stage of the local pipeline (see the class + * javadoc). Cheap - no I/O - so callers may poll it. + */ + public static long localHorizon() { + long oldest = 0; + oldest = minPositive(oldest, AuditLoader.oldestUnpublishedEventTime()); + oldest = minPositive(oldest, preLoaderHorizon()); + return oldest; + } + + /** The stages BEFORE the audit loader: held completed queries and the processor. */ + private static long preLoaderHorizon() { + long oldest = 0; + try { + AuditEventProcessor processor = Env.getCurrentAuditEventProcessor(); + if (processor != null) { + oldest = minPositive(oldest, processor.oldestQueuedOrInFlightEventTime()); + } + } catch (Throwable t) { + // an FE without this component (tests, partial startup) has nothing to fence + LOG.debug("audit publication horizon: the audit event processor is unavailable: {}", + t.getMessage()); + } + try { + Env env = Env.getCurrentEnv(); + WorkloadRuntimeStatusMgr mgr = env == null ? null : env.getWorkloadRuntimeStatusMgr(); + if (mgr != null) { + oldest = minPositive(oldest, mgr.oldestHeldAuditEventTime()); + } + } catch (Throwable t) { + LOG.debug("audit publication horizon: the workload runtime status manager is" + + " unavailable: {}", t.getMessage()); + } + return oldest; + } + + /** + * The fence the CAPTURE uses: the minimum over this FE's own pipeline and the fresh + * rows every other FE reported. Throws {@link IllegalStateException} when the shared + * table cannot be read - the caller must NOT advance without a complete fence + * (round-36 #1: an unreadable follower row is exactly the hole this guards). + */ + public static long clusterHorizon() { + long oldest = localHorizon(); + return minPositive(oldest, remoteHorizon()); + } + + /** + * The minimum horizon over the FRESH rows of the shared table (0 when none / all + * stale). A read failure propagates as a retryable {@link IllegalStateException}. + */ + private static long remoteHorizon() { + List<Object[]> rows; + Supplier<List<Object[]>> reader = horizonRowsReaderForTest; + if (reader != null) { + rows = reader.get(); + } else if (!sharedTableAvailable()) { + return 0L; // no live FE environment (unit tests / not ready): nothing reported + } else { + try { + List<ResultRow> result = StatisticsUtil.executeQuery( + SELECT_ROWS_SQL, Collections.emptyMap(), IO_TIMEOUT_SECONDS); + rows = new ArrayList<>(); + if (result != null) { + for (ResultRow row : result) { + List<String> values = row.getValues(); + if (values == null || values.size() < 2) { + continue; + } + rows.add(new Object[] {Long.parseLong(values.get(0).trim()), + TimeUtils.timeStringToLong(values.get(1).trim())}); + } + } + } catch (Exception e) { + throw new IllegalStateException("SPM capture cannot read the cluster audit" + + " publication horizon: " + e.getMessage(), e); + } + } + long oldest = 0; + long now = System.currentTimeMillis(); + for (Object[] row : rows) { + if (row == null || row.length < 2 || row[0] == null || row[1] == null) { + continue; + } + long horizon = (Long) row[0]; + long updatedAt = (Long) row[1]; + if (horizon <= 0) { + continue; + } + if (updatedAt <= 0 || now - updatedAt > ROW_STALE_MILLIS) { + continue; // the reporter stopped: its outstanding events are gone with it + } + oldest = minPositive(oldest, horizon); + } + return oldest; + } + + /** + * Publishes THIS FE's current horizon into the shared table (one row per FE). Called + * by the audit loader's reporter thread on change and on its keepalive cadence; a + * write failure only logs - the next tick retries, and the row simply goes stale if + * the FE dies. + * + * @param horizon the local horizon (0 = nothing outstanding) + */ + public static void reportLocalHorizon(long horizon) { + Consumer<Long> writer = localHorizonWriterForTest; + if (writer != null) { + writer.accept(horizon); + return; + } + if (!sharedTableAvailable()) { + return; // no live FE environment (unit tests / not ready): no shared table + } + String feName = AuditLoader.selfFeName(); + try { + Map<String, String> deleteParams = new HashMap<>(); + deleteParams.put("feName", StatisticsUtil.escapeSQL(feName)); + StatisticsUtil.execUpdate(DELETE_OWN_ROW_SQL, deleteParams, IO_TIMEOUT_SECONDS); + if (horizon > 0) { + Map<String, String> insertParams = new HashMap<>(); + insertParams.put("feName", StatisticsUtil.escapeSQL(feName)); + insertParams.put("horizonMs", String.valueOf(horizon)); + insertParams.put("updateTime", TimeUtils.longToTimeString(System.currentTimeMillis())); + StatisticsUtil.execUpdate(INSERT_OWN_ROW_SQL, insertParams, IO_TIMEOUT_SECONDS); Review Comment: [P3] Stop the horizon reporter from fencing its own writes. Even the first zero report executes a DELETE, and `StatisticsUtil.execUpdate` audits that internal statement. `localHorizon` counts the resulting event, so the next report writes DELETE/INSERT and creates more audit events; an idle FE keeps doing shared-table writes and audit loads. Exclude these internal events from the fence or suppress auditing for the reporter's SQL. ########## fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditPublicationHorizon.java: ########## @@ -0,0 +1,263 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.plugin.audit; + +import org.apache.doris.catalog.Env; +import org.apache.doris.catalog.InternalSchema; +import org.apache.doris.common.FeConstants; +import org.apache.doris.common.util.TimeUtils; +import org.apache.doris.qe.AuditEventProcessor; +import org.apache.doris.resource.workloadschedpolicy.WorkloadRuntimeStatusMgr; +import org.apache.doris.statistics.repository.ResultRow; +import org.apache.doris.statistics.util.StatisticsUtil; + +import com.google.common.annotations.VisibleForTesting; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Supplier; + +/** + * The cluster-wide AUDIT PUBLICATION HORIZON: the start time (epoch millis, the + * {@code time} column of {@code audit_log}) of the oldest audit event that any FE has + * accepted but not yet PUBLISHED. The SPM capture scans the shared audit table from the + * leader, so it uses this value as a progress FENCE: its next scan window must still + * start at or before it, otherwise a row an FE still owes falls behind the advanced + * watermark and is never captured. + * + * <p>Three layers make the fence complete: + * <ul> + * <li>{@link #localHorizon()} folds THIS FE's whole audit pipeline: completed queries + * still held by the {@link WorkloadRuntimeStatusMgr} (they enter the pipeline + * before any loader sees them), the {@link AuditEventProcessor} queue and its + * in-flight event (a plugin can stall while an event is dequeued), and the + * {@link AuditLoader} queue / assembled batch / not-yet-visible batch (a stream + * load can report Publish Timeout after commit).</li> + * <li>each FE REPORTS its local horizon into the shared + * {@link InternalSchema#SPM_AUDIT_HORIZON_TBL_NAME} table, so a follower's + * backlog is visible to the leader that runs the capture.</li> + * <li>{@link #clusterHorizon()} is the MINIMUM over the local value and the FRESH + * rows of that table; a row the reporter stopped refreshing (its FE died or + * stopped reporting - the events are gone with it) is ignored.</li> + * </ul> + */ +public final class AuditPublicationHorizon { + + private static final Logger LOG = LogManager.getLogger(AuditPublicationHorizon.class); + + /** + * A reported row older than this is IGNORED: its FE stopped refreshing the fence + * (crashed / killed / its reporter thread is gone), so the events it still owed are + * lost with it and fencing progress forever would freeze the capture instead of + * protecting anything. Must be comfortably larger than the reporter's keepalive + * interval ({@link AuditLoader#HORIZON_KEEPALIVE_MILLIS}). + */ + public static final long ROW_STALE_MILLIS = 5 * 60 * 1000L; + + private static final String SELECT_ROWS_SQL = + "SELECT `horizon_ms`, `update_time` FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`"; + private static final String DELETE_OWN_ROW_SQL = "DELETE FROM `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "` WHERE `fe_name` = '${feName}'"; + private static final String INSERT_OWN_ROW_SQL = "INSERT INTO `" + FeConstants.INTERNAL_DB_NAME + "`." + + "`" + InternalSchema.SPM_AUDIT_HORIZON_TBL_NAME + "`" + + " (`fe_name`, `horizon_ms`, `update_time`) VALUES ('${feName}', ${horizonMs}, '${updateTime}')"; + private static final int IO_TIMEOUT_SECONDS = 10; + + /** + * Test seam: the shared-table read (one row per FE). Null in production. + */ + @VisibleForTesting + static volatile Supplier<List<Object[]>> horizonRowsReaderForTest; + + /** + * Test seam: the shared-table write of this FE's row (delete + optional insert). + * Null in production. + */ + @VisibleForTesting + static volatile Consumer<Long> localHorizonWriterForTest; + + private AuditPublicationHorizon() { + } + + /** + * The oldest audit event THIS FE has accepted but not published, 0 when nothing is + * outstanding: the MINIMUM over every stage of the local pipeline (see the class + * javadoc). Cheap - no I/O - so callers may poll it. + */ + public static long localHorizon() { + long oldest = 0; + oldest = minPositive(oldest, AuditLoader.oldestUnpublishedEventTime()); + oldest = minPositive(oldest, preLoaderHorizon()); + return oldest; + } + + /** The stages BEFORE the audit loader: held completed queries and the processor. */ + private static long preLoaderHorizon() { + long oldest = 0; + try { + AuditEventProcessor processor = Env.getCurrentAuditEventProcessor(); + if (processor != null) { + oldest = minPositive(oldest, processor.oldestQueuedOrInFlightEventTime()); + } + } catch (Throwable t) { + // an FE without this component (tests, partial startup) has nothing to fence + LOG.debug("audit publication horizon: the audit event processor is unavailable: {}", + t.getMessage()); + } + try { + Env env = Env.getCurrentEnv(); + WorkloadRuntimeStatusMgr mgr = env == null ? null : env.getWorkloadRuntimeStatusMgr(); + if (mgr != null) { + oldest = minPositive(oldest, mgr.oldestHeldAuditEventTime()); + } + } catch (Throwable t) { + LOG.debug("audit publication horizon: the workload runtime status manager is" + + " unavailable: {}", t.getMessage()); + } + return oldest; + } + + /** + * The fence the CAPTURE uses: the minimum over this FE's own pipeline and the fresh + * rows every other FE reported. Throws {@link IllegalStateException} when the shared + * table cannot be read - the caller must NOT advance without a complete fence + * (round-36 #1: an unreadable follower row is exactly the hole this guards). + */ + public static long clusterHorizon() { + long oldest = localHorizon(); + return minPositive(oldest, remoteHorizon()); + } + + /** + * The minimum horizon over the FRESH rows of the shared table (0 when none / all + * stale). A read failure propagates as a retryable {@link IllegalStateException}. + */ + private static long remoteHorizon() { + List<Object[]> rows; + Supplier<List<Object[]>> reader = horizonRowsReaderForTest; + if (reader != null) { + rows = reader.get(); + } else if (!sharedTableAvailable()) { + return 0L; // no live FE environment (unit tests / not ready): nothing reported + } else { + try { + List<ResultRow> result = StatisticsUtil.executeQuery( + SELECT_ROWS_SQL, Collections.emptyMap(), IO_TIMEOUT_SECONDS); + rows = new ArrayList<>(); + if (result != null) { + for (ResultRow row : result) { + List<String> values = row.getValues(); + if (values == null || values.size() < 2) { + continue; + } + rows.add(new Object[] {Long.parseLong(values.get(0).trim()), + TimeUtils.timeStringToLong(values.get(1).trim())}); + } + } + } catch (Exception e) { + throw new IllegalStateException("SPM capture cannot read the cluster audit" + + " publication horizon: " + e.getMessage(), e); + } + } + long oldest = 0; + long now = System.currentTimeMillis(); + for (Object[] row : rows) { + if (row == null || row.length < 2 || row[0] == null || row[1] == null) { + continue; + } + long horizon = (Long) row[0]; + long updatedAt = (Long) row[1]; + if (horizon <= 0) { + continue; + } + if (updatedAt <= 0 || now - updatedAt > ROW_STALE_MILLIS) { + continue; // the reporter stopped: its outstanding events are gone with it + } + oldest = minPositive(oldest, horizon); + } + return oldest; + } + + /** + * Publishes THIS FE's current horizon into the shared table (one row per FE). Called + * by the audit loader's reporter thread on change and on its keepalive cadence; a + * write failure only logs - the next tick retries, and the row simply goes stale if + * the FE dies. + * + * @param horizon the local horizon (0 = nothing outstanding) + */ + public static void reportLocalHorizon(long horizon) { + Consumer<Long> writer = localHorizonWriterForTest; + if (writer != null) { + writer.accept(horizon); + return; + } + if (!sharedTableAvailable()) { + return; // no live FE environment (unit tests / not ready): no shared table + } + String feName = AuditLoader.selfFeName(); + try { + Map<String, String> deleteParams = new HashMap<>(); + deleteParams.put("feName", StatisticsUtil.escapeSQL(feName)); + StatisticsUtil.execUpdate(DELETE_OWN_ROW_SQL, deleteParams, IO_TIMEOUT_SECONDS); Review Comment: [P2] Preserve the follower's positive fence while refreshing it. This DELETE and the following INSERT commit separately, so the leader can read no row between them even though the follower still has an old unpublished audit event. Capture may then advance its checkpoint beyond that event's eventual publication. Use one atomic update/upsert of the per-FE row for positive values. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
