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]

Reply via email to