github-actions[bot] commented on code in PR #68499:
URL: https://github.com/apache/doris/pull/68499#discussion_r4110391953


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/SPMPlanner.java:
##########
@@ -0,0 +1,687 @@
+// 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;
+
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Pair;
+import org.apache.doris.common.UserException;
+import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.spm.builder.SPMPlan2SQLBuilder;
+import org.apache.doris.nereids.spm.manager.BaselineManager;
+import org.apache.doris.nereids.spm.manager.SessionBaselineStore;
+import org.apache.doris.nereids.spm.matcher.SPMFrozenTreeReplacer;
+import org.apache.doris.nereids.spm.matcher.SPMPlaceholderReplacer;
+import org.apache.doris.nereids.spm.placeholder.SPMPlaceholderBuilder;
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.plans.Plan;
+import org.apache.doris.nereids.trees.plans.commands.Command;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.SqlModeHelper;
+
+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.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * SPMPlanner - the controller of the whole-query SPM bind and rewrite.
+ *
+ * SPM works on the WHOLE parsed (still unbound) SELECT plan tree, not only on 
the
+ * per-query-block WHERE predicates:
+ *
+ * 1. Build (CREATE BASELINE PLAN / auto capture): parse bindSql, replace 
EVERY literal
+ *    of the whole tree (filters, having, projections, aggregate 
group-by/output and
+ *    every subquery inside an expression, recursively) with a placeholder 
through one
+ *    shared SPMPlaceholderBuilder, and do the same for planSql so both trees 
share the
+ *    same placeholder ids. The baseline keeps the parameterized bind tree 
(used for the
+ *    Level 3 structural match) and the parameterized plan tree (used to 
produce the
+ *    rewritten query).
+ * 2. Rewrite (user query, called from StmtExecutor / EXPLAIN): compute the 
value-free
+ *    full-query digest / hash, find candidate baselines (Level 1 hash + Level 
2 digest),
+ *    structurally compare the candidate's parameterized bind tree with the 
user's tree
+ *    over the whole tree (Level 3, extracting the user values for every 
placeholder id),
+ *    and on a match replay the baseline's FROZEN planSql (M3): the frozen 
text - the
+ *    decompiled optimal plan carrying placeholder ids, join distribution 
hints and the
+ *    pushed-down structure - is re-parsed and the extracted values are 
substituted by
+ *    placeholder id, so the rewrite reproduces the frozen optimal structure 
(SR-aligned).
+ *    Baselines without a frozen placeholder text fall back to substituting 
the candidate's
+ *    parameterized plan tree. The substituted tree (still unbound, no 
placeholders left)
+ *    is returned and planned normally by the caller.
+ *
+ * A query without any query-block WHERE - e.g. one whose SELECT-list scalar 
subquery
+ * carries its own WHERE - is handled naturally: the subquery's literals are 
part of the
+ * whole tree and take part in parameterization, matching, value extraction 
and rewrite,
+ * so no special case or rejection is needed.
+ *
+ * A failed rewrite, a timeout or a no-match returns null (the caller executes 
the
+ * original query normally; SPM degrades transparently).
+ */
+public class SPMPlanner {
+
+    private static final Logger LOG = LogManager.getLogger(SPMPlanner.class);
+
+    /** Id of the baseline used by the last successful rewrite (-1 when none). 
*/
+    private long usedBaselineId = -1;
+
+    /**
+     * Returns the id of the baseline used by the last successful rewrite.
+     *
+     * @return the baseline id, or -1 when the last rewrite did not hit a 
baseline
+     */
+    public long getUsedBaselineId() {
+        return usedBaselineId;
+    }
+
+    // ==================== query rewrite (called from StmtExecutor / EXPLAIN) 
====================
+
+    /**
+     * Rewrites a user query (parsed, still unbound) via SPM.
+     *
+     * Runs the three-level match against every candidate baseline and, on a 
hit,
+     * substitutes the user's actual literal values into the candidate's 
parameterized
+     * plan tree and returns the resulting (unbound) plan. The caller 
(StmtExecutor)
+     * plans the returned tree normally, so the rewritten query goes through 
the regular
+     * analyze / optimize pipeline.
+     *
+     * @param userPlan the parsed (unbound) plan of the user query
+     * @param deadline rewrite deadline (epoch millis) for the 
spm_rewrite_timeout_ms budget
+     * @return the rewritten unbound plan, or null when no baseline matched / 
timed out /
+     *         the substituted plan is unusable
+     */
+    public LogicalPlan tryRewritePlan(LogicalPlan userPlan, long deadline) {
+        if (System.currentTimeMillis() > deadline) {
+            LOG.info("SPM tryRewritePlan: timeout before matching, degrade to 
the original plan");
+            return null;
+        }
+        // SELECT ... INTO OUTFILE writes to a destination that lives OUTSIDE 
the plan
+        // expressions: the generic match cannot see a different path / format 
and the
+        // rewritten tree keeps the CAPTURED sink fields, so a replay would 
export to the
+        // baseline's destination. SPM refuses to rewrite such statements (the 
original
+        // query runs normally, which is always correct).
+        if (SPMPlanTreeSupport.containsFileSink(userPlan)) {
+            LOG.info("SPM tryRewritePlan: the statement writes to a file sink 
(OUTFILE); keeping"
+                    + " the original plan");
+            return null;
+        }
+        ConnectContext ctx = ConnectContext.get();
+        SessionBaselineStore sessionStore = ctx == null ? null : 
ctx.getSessionBaselineStore();
+        // Fast path: with no baseline anywhere (global or session) there is 
nothing to
+        // match, so skip the whole-tree digest rendering below - it walks 
every literal
+        // of the whole query and is the most expensive part of a rewrite 
attempt.
+        if (!BaselineManager.getInstance().hasBaselines()
+                && (sessionStore == null || sessionStore.isEmpty())) {
+            return null;
+        }
+        // Level 1/2 use the FULL-QUERY value-free digest (Plan.toSpmDigest() 
renders every
+        // literal as "?", so the digest is value-independent). The digest is 
computed on a
+        // namespace-qualified COPY: "FROM t" is keyed as the CURRENT 
catalog/db's t, so a
+        // baseline captured under db1 can never match the same text executed 
under db2
+        // (the frozen planSql is fully qualified and would silently keep 
running against
+        // db1.t). Explicitly qualified references stay verbatim on both 
sides, and so do
+        // references to a CTE alias: those bind inside the query's own WITH 
clause and are
+        // therefore namespace-independent.
+        LogicalPlan matchPlan = SPMPlanTreeSupport.namespaceQualified(userPlan,
+                captureCatalogName(ctx), captureDatabaseName(ctx));
+        String queryDigest = matchPlan.toSpmDigest();
+        long queryHash = SPMUtils.hashOf(queryDigest);
+        // SESSION-scope baselines of the current connection are consulted 
BEFORE the
+        // global ones, so a session baseline can override a global baseline 
for this
+        // session only; each candidate list is already priority-ordered.
+        List<BaselinePlan> candidates = new ArrayList<>();
+        if (sessionStore != null) {
+            candidates.addAll(sessionStore.findCandidateBaselines(queryDigest, 
queryHash));
+        }
+        candidates.addAll(
+                
BaselineManager.getInstance().findCandidateBaselines(queryDigest, queryHash));
+        if (candidates.isEmpty()) {
+            return null;
+        }
+        // View guard (computed lazily on the first structural match, so the 
per-query
+        // hot path pays nothing): the replay of a matched baseline is planned 
before the
+        // normal authorization pass and resolves a view to its base tables, so
+        // authorization would check the base tables instead of the view - a 
view-only
+        // user is denied on the base tables, while a base-table user passes 
the same
+        // view query without any view check. View queries keep the original 
plan.
+        boolean viewChecked = false;
+        boolean viewReferenced = false;
+        for (BaselinePlan candidate : candidates) {
+            if (System.currentTimeMillis() > deadline) {
+                LOG.info("SPM tryRewritePlan: timeout before matching baseline 
{}, "
+                        + "degrade to the original plan", candidate.getId());
+                return null;
+            }
+            LogicalPlan bindTree = candidate.getParameterizedBindPlan();
+            if (bindTree == null) {
+                continue;
+            }
+            // Level 3: whole-tree structural match + value extraction
+            Map<Long, Expression> placeholderValues = new HashMap<>();
+            if (!SPMPlanTreeSupport.check(bindTree, matchPlan, 
placeholderValues)) {
+                continue;
+            }
+            LOG.info("SPM tryRewritePlan: baseline {} matched, extracted {} 
placeholder values",
+                    candidate.getId(), placeholderValues.size());
+            if (!viewChecked) {
+                viewChecked = true;
+                viewReferenced = SPMPlanTreeSupport.referencesView(ctx, 
userPlan);
+            }
+            if (viewReferenced) {
+                LOG.info("SPM tryRewritePlan: the query references a view; 
keeping the original"
+                        + " plan so view authorization is preserved");
+                return null;
+            }
+            // The expensive part (value-free digest, candidate lookup, 
whole-tree check) is
+            // already done and a baseline has matched: finish the (cheap) 
value substitution
+            // even when the budget has been consumed. Previously this path 
re-checked the
+            // deadline here and returned null on expiry, silently throwing 
the completed
+            // match away (rewrite degraded to the original plan).
+            // M3: replay the FROZEN optimal plan. A baseline whose frozen 
planSql carries
+            // placeholder calls (decompiled from the SPM-optimized physical 
plan, structure
+            // + join distribution hints preserved) is re-parsed and the user 
values are
+            // substituted by placeholder id - the rewrite reproduces the 
frozen optimal
+            // structure exactly (SR-aligned). Baselines without a frozen 
placeholder text
+            // (in-memory engine, decompile-fallback planSql, legacy "?" 
frozen text) fall
+            // back to substituting the parameterized plan tree.
+            LogicalPlan rewritten = rewriteFromFrozenTree(candidate, 
placeholderValues);
+            if (rewritten != null) {
+                // non-expression literals (LIMIT / OFFSET) are long fields, 
not
+                // placeholders; adopt the user's values so a structurally 
identical query
+                // with a different limit is rewritten with the USER limit
+                usedBaselineId = candidate.getId();
+                return SPMPlanTreeSupport.mergeLimits(rewritten, matchPlan);
+            }
+            LOG.info("SPM tryRewritePlan: baseline {} frozen planSql replay 
unavailable, "
+                    + "falling back to parameterized plan tree", 
candidate.getId());
+            LogicalPlan planTree = 
stripSelectHints(candidate.getParameterizedPlanPlan());
+            if (planTree == null) {
+                continue;
+            }
+            // fallback: substitute the extracted user values into the 
parameterized plan
+            // tree
+            SPMPlaceholderReplacer replacer = new SPMPlaceholderReplacer();
+            rewritten = SPMPlanTreeSupport.transform(
+                    planTree, expr -> expr.accept(replacer, 
placeholderValues));
+            // safety net: a plan-only placeholder (no value extracted from 
the user query)
+            // would reach the analyzer - never rewrite with this candidate; 
keep trying the
+            // remaining candidates (they are validated independently) and 
degrade to the
+            // original plan when none works
+            if (SPMPlanTreeSupport.containsPlaceholder(rewritten)) {
+                continue;
+            }
+            usedBaselineId = candidate.getId();
+            return SPMPlanTreeSupport.mergeLimits(rewritten, matchPlan);
+        }
+        return null;
+    }
+
+    /**
+     * M3: replays a baseline from its frozen planSql text.
+     *
+     * When the baseline's planSql carries placeholder calls (the decompiled 
frozen
+     * optimal plan), the text is re-parsed into an unbound logical plan - 
preserving the
+     * frozen structure (join order / subquery nesting / pushed-down filters) 
and the join
+     * distribution hints ([BROADCAST] / [SHUFFLE]) - and the user values are 
substituted
+     * by placeholder id (SPMFrozenTreeReplacer). The caller plans the 
returned tree
+     * normally, so the optimizer starts from the frozen optimal structure 
instead of the
+     * user's raw SQL structure (aligned with SR's PlaceholderReplacer flow).
+     *
+     * @param candidate         the matched baseline
+     * @param placeholderValues placeholder id -> user value (extracted by the 
Level 3
+     *                          structural check on the bind tree)
+     * @return the substituted unbound plan, or null when the baseline has no 
frozen
+     *         placeholder text / re-parsing fails / a placeholder has no user 
value
+     */
+    private LogicalPlan rewriteFromFrozenTree(BaselinePlan candidate,
+            Map<Long, Expression> placeholderValues) {
+        String planSql = candidate.getPlanSql();
+        if (planSql == null
+                || (!planSql.contains(SPMFrozenTreeReplacer.CONST_VAR_FUNC)

Review Comment:
   [P1] Classify actual placeholder calls, not raw text. When the decompiler 
rejects a physical node such as `PhysicalAssertNumRows`, this PR deliberately 
stores the original `planSql`. If that ordinary SQL merely contains 
`_spm_const_var` or `_spm_const_list` in a string literal, identifier, or 
comment, this branch treats it as frozen; parsing/replacement changes nothing, 
the residue check passes, and replay returns the captured literals before 
reaching the correctly parameterized fallback. The same substring test in 
`BaselineManager.parsePersistedRow` also discards that fallback after reload. 
Traverse the parsed tree for real placeholder functions (or persist an explicit 
frozen flag), and cover a fallback query whose literal contains one of these 
names.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/capture/AuditLogScanner.java:
##########
@@ -0,0 +1,355 @@
+// 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.qe.SessionVariable;
+import org.apache.doris.qe.VariableMgr;
+import org.apache.doris.statistics.repository.ResultRow;
+import org.apache.doris.statistics.util.StatisticsUtil;
+
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.format.DateTimeFormatter;
+import java.util.ArrayList;
+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 (query_time, time, 
query_id)
+ * cursor. 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;
+
+    private static final DateTimeFormatter DATETIME_FORMAT =
+            DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+
+    /** audit_log SELECT columns (order must match rowToCapturedQuery / 
toBatch). */
+    private static final String SELECT_COLUMNS =

Review Comment:
   [P2] Auto capture needs the originating parser mode. `audit_log` already 
stores decoded `sql_mode`, but this projection and `CapturedQuery` omit it; 
`processCandidate` then reparses the statement in a fresh internal context 
after restoring only catalog/database, so both the bind tree and persisted 
`creatorSqlMode` use the internal default. A `PIPES_AS_CONCAT` statement 
containing `a || b` is captured as boolean OR (or fails) and cannot produce a 
usable baseline for later CONCAT-mode executions. Select/decode the audit mode, 
carry it through retry-queue persistence, and build under that mode, with an 
audit-to-reload test.



##########
fe/fe-core/src/main/java/org/apache/doris/catalog/InternalSchemaInitializer.java:
##########
@@ -405,6 +426,49 @@ public static void createTbl() throws UserException {
          * )
          */
         createTable(getAuditLogCreateSql());
+        createTable(getSpmBaselinesCreateSql());
+        createTable(getSpmCaptureCheckpointCreateSql());
+    }
+
+    /**
+     * Adds the `sql_mode` column to a PRE-EXISTING spm_baselines table (new 
clusters get
+     * it from the create SQL). The column carries the parser mode of the 
creating session:
+     * without it a PIPES_AS_CONCAT baseline is re-parsed under the default 
mode after a
+     * restart, so the stored digest still finds the row while the structural 
match rejects
+     * every CONCAT-mode query and the baseline silently stops applying. 
Idempotent: a
+     * table that already carries the column is left untouched.
+     */
+    private static void upgradeSpmBaselinesSchema() {
+        try {
+            Optional<Database> dbOpt =
+                    
Env.getCurrentEnv().getInternalCatalog().getDb(FeConstants.INTERNAL_DB_NAME);
+            if (!dbOpt.isPresent()) {
+                return;
+            }
+            Table table = 
dbOpt.get().getTable(InternalSchema.SPM_BASELINES_TBL_NAME).orElse(null);
+            if (table == null) {
+                return;
+            }
+            if (table.getBaseSchema().stream()
+                    .anyMatch(column -> 
"sql_mode".equalsIgnoreCase(column.getName()))) {
+                return;
+            }
+            ColumnDefinition definition = new ColumnDefinition("sql_mode",
+                    
DataType.fromCatalogType(ScalarType.createType(PrimitiveType.BIGINT)),
+                    true, null, ColumnNullableType.NULLABLE, -1, 
Optional.empty(),
+                    Optional.empty(), "", true, Optional.empty());
+            AddColumnOp addColumnOp = new AddColumnOp(definition, null, null, 
null);
+            
addColumnOp.setColumn(definition.translateToCatalogStyleForSchemaChange());
+            TableNameInfo tableNameInfo = new 
TableNameInfo(InternalCatalog.INTERNAL_CATALOG_NAME,
+                    FeConstants.INTERNAL_DB_NAME, 
InternalSchema.SPM_BASELINES_TBL_NAME);
+            Env.getCurrentEnv().alterTable(
+                    new AlterTableCommand(tableNameInfo, 
Lists.newArrayList(addColumnOp)));
+            LOG.info("SPM: added the sql_mode column to {}", 
InternalSchema.SPM_BASELINES_TBL_NAME);
+        } catch (Throwable t) {
+            // Retried on the next initializer iteration / FE start; the 
baseline load
+            // fails fast (bounded timeout) and retries until the column 
exists.
+            LOG.warn("SPM: failed to add the spm_baselines sql_mode column, 
will retry", t);

Review Comment:
   [P2] A transient ALTER failure is not retried in this FE process. `run()` 
calls this method only once after its creation loop has already exited, and the 
following replica loop never comes back here. On an upgraded cluster the table 
then remains without `sql_mode`, while `BaselineManager` unconditionally 
selects/writes that column, so baseline loading and global DDL stay broken 
until restart despite this `will retry` message. Keep the initializer alive 
until the column is observed (or invoke the idempotent upgrade from a real 
retry path), and cover a first ALTER failure followed by success without 
restart.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/matcher/SPMFrozenTreeReplacer.java:
##########
@@ -0,0 +1,196 @@
+// 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.matcher;
+
+import org.apache.doris.nereids.analyzer.UnboundFunction;
+import org.apache.doris.nereids.trees.expressions.Expression;
+import org.apache.doris.nereids.trees.expressions.InPredicate;
+import org.apache.doris.nereids.trees.expressions.literal.Literal;
+import org.apache.doris.nereids.trees.expressions.visitor.ExpressionVisitor;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * SPMFrozenTreeReplacer - placeholder replacer for the RE-PARSED frozen plan 
tree (M3).
+ *
+ * At rewrite time the frozen planSql - the decompiled optimal plan
+ * carrying the placeholder ids, join distribution hints ([BROADCAST] / 
[SHUFFLE]) and the
+ * pushed-down structure - is re-parsed into an unbound logical plan. In that 
tree a
+ * placeholder is NOT an in-memory SpmConstVar / SpmConstList node (those
+ * only exist inside the parameterized bind / plan trees) but a raw parsed 
function call:
+ *
+ *   _spm_const_var(id) - a scalar placeholder, possibly wrapped by an outer
+ *       CAST (the optimizer typed it at CREATE time), e.g.
+ *       a > CAST(_spm_const_var(1) AS INT);
+ *   _spm_const_list(id) - the single option of an IN predicate, e.g.
+ *       x IN (_spm_const_list(3)).
+ *
+ * This replacer identifies those calls by name + leading integer literal (the 
id) and
+ * substitutes the actual user value (extracted by SPMAstCheckVisitor from the 
bind tree)
+ * - mirroring SR's PlaceholderReplacer which visits FunctionCallExpr and 
InPredicate.
+ */
+public class SPMFrozenTreeReplacer extends ExpressionVisitor<Expression, 
Map<Long, Expression>> {
+
+    /** Function name of the scalar placeholder as printed into the frozen 
planSql. */
+    public static final String CONST_VAR_FUNC = "_spm_const_var";
+
+    /** Function name of the IN-list placeholder as printed into the frozen 
planSql. */
+    public static final String CONST_LIST_FUNC = "_spm_const_list";
+
+    /**
+     * Substitutes every placeholder call of the frozen (re-parsed) tree 
expression with
+     * the user value registered under the placeholder id.
+     *
+     * @param expr              the frozen-tree expression (may contain 
placeholders)
+     * @param placeholderValues placeholder id -> user actual value (extracted 
from the
+     *                          bind-side structural check)
+     * @return the substituted expression
+     */
+    public Expression replace(Expression expr, Map<Long, Expression> 
placeholderValues) {
+        return expr.accept(this, placeholderValues);
+    }
+
+    /**
+     * Returns whether the expression is a frozen-tree placeholder call
+     * (_spm_const_var(id) / _spm_const_list(id)) that was NOT substituted
+     * - such an expression must never reach the analyzer (the function is not 
registered),
+     * so a rewritten tree containing one must be rejected by the caller.
+     */
+    public static boolean isUnsubstitutedPlaceholder(Expression expr) {
+        if (expr instanceof UnboundFunction) {
+            String name = ((UnboundFunction) expr).getName();
+            return CONST_VAR_FUNC.equals(name) || CONST_LIST_FUNC.equals(name);
+        }
+        return false;
+    }
+
+    @Override

Review Comment:
   [P2] Recurse into `SubqueryExpr.queryPlan` here. The generic visitor only 
walks `Expression.children()`, while the builder and in-memory replacer both 
need explicit `visitSubqueryExpr` overrides. A residual `LEFT NULL_AWARE ANTI 
JOIN` is frozen as `NOT IN (SELECT ... WHERE ... _spm_const_var(...))`; 
replacement misses that inner plan, the residue scan finds the call and rejects 
replay, and after refresh/restart a frozen row has no parameterized fallback 
tree. Mirror the existing subquery rebuild logic and cover persisted/reloaded 
IN, EXISTS, scalar, and residual NOT IN cases.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/spm/manager/BaselineManager.java:
##########
@@ -0,0 +1,1345 @@
+// 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.manager;
+
+import org.apache.doris.catalog.InternalSchema;
+import org.apache.doris.common.FeConstants;
+import org.apache.doris.common.Pair;
+import org.apache.doris.nereids.spm.BaselinePlan;
+import org.apache.doris.nereids.spm.BaselineScope;
+import org.apache.doris.nereids.spm.BaselineSource;
+import org.apache.doris.nereids.spm.BaselineStatus;
+import org.apache.doris.nereids.spm.SPMPlanner;
+import org.apache.doris.nereids.spm.matcher.SPMFrozenTreeReplacer;
+import org.apache.doris.nereids.trees.plans.logical.LogicalPlan;
+import org.apache.doris.qe.SqlModeHelper;
+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.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.format.DateTimeFormatter;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
+import java.util.stream.Collectors;
+
+/**
+ * BaselineManager - baseline storage, cache and index management (M3).
+ *
+ * Corresponds to design doc section 6.6. Manages the CRUD of baselines and 
maintains
+ * two query structures:
+ *
+ * - hashIndex: {@code Map<Long, List<Long>>}, bindSqlHash -> baseline id 
list. Level 1
+ *   coarse filtering with O(1) lookup.
+ * - baselines: id -> BaselinePlan in-memory storage (Phase 1 MVP). Phase 2 
persists it
+ *   to the __internal_schema.spm_baselines internal table (see design doc 
6.14).
+ *
+ * Candidate baseline lookup (the first two of the three-level filter):
+ *
+ * 1. Level 1: hashIndex.get(queryHash) -> candidate id list
+ * 2. Level 2: exact digest.equals(baseline.bindSqlDigest) matching
+ * 3. Sort by priority (see the comparator below)
+ *
+ * Id source (the GLOBAL id invariant): GLOBAL writes have a single writer 
(user DDL is
+ * forwarded to the master, auto capture runs on the Leader), so "read the 
persistence
+ * watermark, then allocate" is sufficient to keep the local generator from 
ever handing
+ * out an id that is already used in the shared table: before EVERY allocation
+ * createBaseline reads MAX(id) from the internal table (one light 
aggregation) and
+ * advances idGenerator past it; when that read fails the CREATE fails visibly 
with a
+ * retryable error and no id is allocated. This closes the failover lag hole 
(a new
+ * master may start with a generator behind the table) and the load-failure 
hole (a
+ * broken startup load leaves the generator at 1 while the table is full); it 
stays
+ * correct even when a row was inserted out of contract (e.g. a manual table 
write). The
+ * startup load / periodic refresh keep advancing the generator as a second 
safety net.
+ * Nothing that does not allocate (ALTER / DROP / matching) reads the 
watermark, and
+ * creates only come from the DDL / capture cycle, so the read is off the 
query hot path.
+ * GLOBAL ids start at 1 and stay in [1, 2^62); SESSION ids live in [2^62, 
2^63) (see
+ * BaselineScope.ofId), so the scope of an id is always exact.
+ *
+ * Concurrency: the in-memory store (baselines / hashIndex / stateVersion / 
load state) is
+ * guarded by one read-write lock. Lookups take the read lock and never 
serialize on each
+ * other; create / drop / status / load / refresh take the write lock; 
internal-table I/O
+ * runs outside the lock wherever correctness allows (see the individual 
methods).
+ */
+public class BaselineManager {
+
+    private static final Logger LOG = 
LogManager.getLogger(BaselineManager.class);
+
+    /** Singleton. */
+    private static final BaselineManager INSTANCE = new BaselineManager();
+
+    // ==================== internal-table persistence 
(__internal_schema.spm_baselines) ==========
+
+    /** Fully qualified internal table (FeConstants.INTERNAL_DB_NAME == 
"__internal_schema"). */
+    private static final String SPM_BASELINES_TABLE =
+            FeConstants.INTERNAL_DB_NAME + "." + 
InternalSchema.SPM_BASELINES_TBL_NAME;
+
+    /** Column order follows InternalSchema.SPM_BASELINES_SCHEMA. */
+    private static final String SELECT_ALL_SQL = "SELECT `id`, `bind_sql`, 
`bind_sql_digest`,"
+            + " `bind_sql_hash`, `plan_sql`, `query_id`, `cost`, 
`query_time_ms`, `source`,"
+            + " `status`, `create_time`, `update_time`, `sql_mode` FROM " + 
SPM_BASELINES_TABLE;
+
+    /** The persistence-layer id watermark (see the class javadoc "Id 
source"): read
+     *  before every id allocation. MAX over an aggregate is a light 
single-row query. */
+    private static final String SELECT_MAX_ID_SQL = "SELECT MAX(`id`) FROM " + 
SPM_BASELINES_TABLE;
+
+    /**
+     * Durable-key lookup used by the create-time dedup: the in-memory index 
can be stale
+     * (a follower that loaded=true before becoming master missed rows written 
afterwards),
+     * so the authoritative duplicate check reads the (bind_sql_digest, 
plan_sql) key back
+     * from the table before a new row is inserted.
+     */
+    private static final String SELECT_BY_KEY_SQL = "SELECT `id`, `bind_sql`, 
`bind_sql_digest`,"
+            + " `bind_sql_hash`, `plan_sql`, `query_id`, `cost`, 
`query_time_ms`, `source`,"
+            + " `status`, `create_time`, `update_time`, `sql_mode` FROM " + 
SPM_BASELINES_TABLE
+            + " WHERE `bind_sql_digest` = '${bindSqlDigest}' AND `plan_sql` = 
'${planSql}'";
+
+    private static final String INSERT_SQL = "INSERT INTO " + 
SPM_BASELINES_TABLE
+            + " VALUES (${id}, '${bindSql}', '${bindSqlDigest}', 
${bindSqlHash},"
+            + " '${planSql}', '${queryId}', ${cost}, ${queryTimeMs}, 
'${source}', '${status}',"
+            + " '${createTime}', '${updateTime}', ${sqlMode})";
+
+    /**
+     * Deletes one row by id AND content key (bind_sql_digest + plan_sql): a 
delete
+     * triggered by a stale id must never remove a row that happens to carry 
the same id
+     * but different content (e.g. after the table was edited out of band, or 
the id was
+     * rewound by an FE started with stale metadata).
+     */
+    private static final String DELETE_BY_IDENTITY_SQL = "DELETE FROM " + 
SPM_BASELINES_TABLE
+            + " WHERE `id` = ${id} AND `bind_sql_digest` = '${bindSqlDigest}'"
+            + " AND `plan_sql` = '${planSql}'";
+
+    /** Removes one baseline row by id + its previous status (status UPDATE 
support). */
+    private static final String DELETE_BY_ID_AND_STATUS_SQL = "DELETE FROM " + 
SPM_BASELINES_TABLE
+            + " WHERE `id` = ${id} AND `status` = '${status}'";
+
+    /** DATETIME column format (internal table create_time / update_time). */
+    private static final DateTimeFormatter TS_FORMAT =
+            DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
+
+    /**
+     * Statement timeout (seconds) of the SPM internal-table reads. The 
temporary context
+     * StatisticsUtil builds otherwise inherits the analyze timeout (12h by 
default): an
+     * unavailable tablet / BE would stall the caller (master readiness, 
background load)
+     * far beyond the advertised SPM budget.
+     */
+    private static final int INTERNAL_QUERY_TIMEOUT_SECONDS = 10;
+
+    /** How long a management caller waits for an in-flight background load. */
+    private static final long MANAGEMENT_LOAD_WAIT_MILLIS = 5_000L;
+
+    // ==================== priority ordering ====================
+
+    /**
+     * Candidate baseline priority ordering:
+     *
+     * 1. both have queryMs (>= 0) -> the one with shorter time wins
+     * 2. neither has queryMs (-1) -> the one with lower cost wins
+     * 3. only one has queryMs -> the one without time wins (manual baselines 
win over
+     *    auto-captured ones)
+     *
+     * queryMs = -1 means no actual execution time (manually created); >= 0 
means known
+     * execution stats (0 is a valid measured sub-millisecond time). The 
comparison is a
+     * total order: equal known times compare equal and a known time never 
compares -1 in
+     * both directions.
+     */
+    private static final Comparator<BaselinePlan> comparator = (o1, o2) -> {
+        // -1 is "unknown" (manual create); 0 is a VALID measured time, so 
"known" is
+        // >= 0. compare(0, 0) must be 0 and 0 vs -1 must be ordered 
consistently, or
+        // stream sorting can misorder candidates / throw a 
comparator-contract error.
+        long time1 = o1.getQueryTimeMs();
+        long time2 = o2.getQueryTimeMs();
+        boolean known1 = time1 >= 0;
+        boolean known2 = time2 >= 0;
+        if (known1 && known2) {
+            return Long.compare(time1, time2);
+        } else if (!known1 && !known2) {
+            return Double.compare(o1.getCost(), o2.getCost());
+        } else {
+            // the one with time is ordered last (the one without time wins)
+            return known1 ? 1 : -1;
+        }
+    };
+
+    /** Auto-increment id counter. */
+    private final AtomicLong idGenerator = new AtomicLong(1);
+
+    /** Whether the persisted baselines have been loaded from the internal 
table (or the
+     *  internal schema db is disabled). Set by loadFromInternalTable() and
+     *  clearForTest(). */
+    private volatile boolean loaded = false;
+
+    /**
+     * Coalesces the internal-table loads: at most ONE read may be in flight. 
The query
+     * path (ensureLoaded) never blocks on the table - an unavailable tablet / 
BE used to
+     * stall every rewrite attempt for the temporary context's full analyze 
timeout - it
+     * just requests a background load; a failed read keeps {@code 
loaded=false} and is
+     * retried by the next access / refresh cycle.
+     */
+    private final AtomicBoolean loadInProgress = new AtomicBoolean(false);
+
+    /** Notified when an in-flight load finishes (management callers wait on 
it). */
+    private final Object loadMonitor = new Object();
+
+    /** Whether CRUD writes to the internal table (disabled by clearForTest 
for tests). */
+    private volatile boolean persistToTable = true;
+
+    /**
+     * Guards the in-memory store: {@code baselines} / {@code hashIndex} /
+     * {@code stateVersion} and the {@code loaded} state machine. Only 
accesses to those
+     * structures are critical sections - internal-table I/O and the CPU work 
derived
+     * from a snapshot (digest filter, priority sort) run OUTSIDE the lock. 
Rewrite
+     * lookups (hasBaselines + findCandidateBaselines on every query) take the 
read lock
+     * and therefore never serialize on each other; create / drop / status / 
load /
+     * refresh take the write lock for their in-memory validation / 
publication only.
+     */
+    private final ReentrantReadWriteLock stateLock = new 
ReentrantReadWriteLock();
+
+    /**
+     * Serializes WRITERS (create / drop / status) against each other for their
+     * internal-table read-modify-write sequences. Deliberately separate from 
stateLock:
+     * every SPM query takes stateLock.readLock() in hasBaselines / 
findCandidateBaselines
+     * BEFORE its rewrite timeout can help, so ONE slow persistence statement 
(each
+     * auto-capture candidate issues one) would stall planning for every such 
query on the
+     * FE. The state lock only protects in-memory validation and publication - 
see the
+     * two-phase structure of {@link #createBaseline} / {@link #dropBaseline} /
+     * {@link #updateStatus}. Readers never touch this lock.
+     */
+    private final Object writerLock = new Object();
+
+    /**
+     * Monotonic version of the in-memory state, bumped by every local 
mutation (create /
+     * drop / status change / load) and by every refresh that changed 
something. The periodic
+     * refresh compares it to skip a table snapshot whose read overlapped a 
local mutation
+     * (see refreshFromInternalTable). Guarded by stateLock.
+     */
+    private long stateVersion = 0;
+
+    /** id -> BaselinePlan (Phase 1 in-memory storage). */
+    private final Map<Long, BaselinePlan> baselines = new HashMap<>();
+
+    /** bindSqlHash -> baseline id list (Level 1 coarse filter index). */
+    private final Map<Long, List<Long>> hashIndex = new HashMap<>();
+
+    private BaselineManager() {
+    }
+
+    public static BaselineManager getInstance() {
+        return INSTANCE;
+    }
+
+    /**
+     * Candidate baseline ordering shared with the SESSION-scope store: it 
applies the
+     * exact same priority rules as this manager (see the comparator field).
+     * Public so tests can verify the total-order contract directly.
+     *
+     * @param o1 first baseline
+     * @param o2 second baseline
+     * @return the comparison result
+     */
+    public static int compareCandidates(BaselinePlan o1, BaselinePlan o2) {
+        return comparator.compare(o1, o2);
+    }
+
+    // ==================== CRUD ====================
+
+    /**
+     * Creates a baseline.
+     *
+     * Duplicate detection: identical (bindSqlHash, bindSqlDigest) with the 
exact same
+     * planSql is skipped; a different planSql is allowed to coexist (one 
query can have
+     * multiple plan baselines).
+     *
+     * @param plan the baseline (id is assigned here if unset)
+     * @return the id of the created baseline (or the existing id when 
duplicated)
+     */
+    public long createBaseline(BaselinePlan plan) {
+        ensureLoadedOrThrow();
+        // Two-phase create. I/O (watermark read, durable-key dedup, repair 
deletes, INSERT)
+        // runs under writerLock but NEVER under stateLock: the state lock 
only protects the
+        // in-memory duplicate validation (phase 1) and the publication (phase 
2). Holding
+        // the state lock across the I/O would stall every SPM query's rewrite 
lookup.
+        synchronized (writerLock) {
+            // Id watermark first (see the class javadoc "Id source"): the 
generator must be
+            // advanced past the persistence layer BEFORE an id is handed out. 
A create whose
+            // watermark read fails fails visibly and allocates nothing, 
instead of silently
+            // colliding with a row written by a newer master.
+            final long watermark = readPersistedWatermark();
+            // Phase 1: exact-duplicate validation against the in-memory index 
(read lock).
+            stateLock.readLock().lock();
+            try {
+                if (plan.getBindSqlHash() != 0) {
+                    for (BaselinePlan existing : 
findByHash(plan.getBindSqlHash())) {
+                        if 
(existing.getBindSqlDigest().equals(plan.getBindSqlDigest())
+                                && 
existing.getPlanSql().equals(plan.getPlanSql())) {
+                            return existing.getId(); // exact duplicate -> skip
+                        }
+                    }
+                }
+            } finally {
+                stateLock.readLock().unlock();
+            }
+            // Durable-key check: the in-memory index can be stale (e.g. a 
follower that
+            // loaded=true before becoming master missed rows written after 
its last
+            // refresh). Without this check the INSERT below would REPLACE a 
durable
+            // baseline - changing its id and, on an INSERT failure, losing 
the old row.
+            // A durable duplicate returns its id and is adopted into memory 
instead;
+            // extra same-key rows (partial-state survivors) are repaired away 
idempotently.
+            // writerLock keeps another writer's INSERT/DELETE pair out of 
this window.
+            if (persistenceEnabled() && plan.getBindSqlDigest() != null) {
+                List<BaselinePlan> durable =
+                        readPersistedByKey(plan.getBindSqlDigest(), 
plan.getPlanSql());
+                if (!durable.isEmpty()) {
+                    BaselinePlan winner = durable.get(0);
+                    for (int i = 1; i < durable.size(); i++) {
+                        winner = pickDurableWinner(winner, durable.get(i));
+                    }
+                    for (BaselinePlan row : durable) {
+                        if (row.getId() != winner.getId()) {
+                            try {
+                                persistDeleteByIdentity(row);
+                            } catch (RuntimeException e) {
+                                // best-effort repair: the key state is 
deterministic either
+                                // way (the winner above), the leftover row is 
warned below
+                                LOG.warn("SPM failed to repair a duplicate 
baseline row (id={}): {}",
+                                        row.getId(), e.getMessage());
+                            }
+                        }
+                    }
+                    publishBaseline(winner);
+                    LOG.info("SPM baseline create deduplicated against the 
durable key: id={}",
+                            winner.getId());
+                    return winner.getId();
+                }
+            }
+            // every baseline owned by the global manager is GLOBAL-scope 
(same value the
+            // rows loaded from the internal table and the auto capturer get); 
the id stays
+            // in the GLOBAL range [1, 2^62), so BaselineScope.ofId(id) is 
exact
+            plan.setScope(BaselineScope.GLOBAL);
+            if (watermark >= idGenerator.get()) {
+                // watermark + 1 is an id the table has never seen; the 
generator is
+                // monotonic and the watermark was read at the top of this 
create, so
+                // concurrent local creates can only jump the generator 
further up, never
+                // below a watermark any create has observed
+                idGenerator.set(watermark + 1);
+            }
+            long id = idGenerator.getAndIncrement();
+            plan.setId(id);
+            long now = System.currentTimeMillis();
+            plan.setCreateTime(now);
+            plan.setUpdateTime(now);
+            // persist first so a persist failure leaves the in-memory state 
untouched and
+            // fails the DDL visibly; no same-key row can exist here (the 
durable-key check
+            // above returned any), so the INSERT cannot overwrite an existing 
baseline
+            persistInsert(plan);
+            // Phase 2: publish (the only state-lock section of a create).
+            publishBaseline(plan);
+            return id;
+        }
+    }
+
+    /**
+     * Phase 2 of a writer: publishes a baseline into the in-memory store. No 
I/O - only
+     * the hash index / map / version are touched, under the write lock. 
Replaces a row a
+     * concurrent refresh may have loaded for the same id (its index entry is 
dropped
+     * first, so no id is indexed twice).
+     */
+    private void publishBaseline(BaselinePlan plan) {
+        stateLock.writeLock().lock();
+        try {
+            BaselinePlan replaced = baselines.put(plan.getId(), plan);
+            if (replaced != null) {
+                removeFromHashIndex(replaced);
+            }
+            addToHashIndex(plan);
+            stateVersion++;
+        } finally {
+            stateLock.writeLock().unlock();
+        }
+    }
+
+    /**
+     * Drops a baseline.
+     *
+     * @param id the baseline id
+     * @return whether the drop succeeded
+     */
+    public boolean dropBaseline(long id) {
+        ensureLoadedOrThrow();
+        // Two-phase drop: the DELETE I/O runs under writerLock but NOT under 
stateLock;
+        // the state lock only removes the in-memory entry afterwards.
+        synchronized (writerLock) {
+            BaselinePlan removed;
+            stateLock.readLock().lock();
+            try {
+                removed = baselines.get(id);
+            } finally {
+                stateLock.readLock().unlock();
+            }
+            if (removed == null) {
+                return false;
+            }
+            // persist first so a failure keeps both the in-memory state and 
the table row;
+            // the delete is keyed by id + content, so a stale id can never 
remove an
+            // unrelated row that reused the id
+            persistDeleteByIdentity(removed);
+            stateLock.writeLock().lock();
+            try {
+                BaselinePlan gone = baselines.remove(id);
+                if (gone == null) {
+                    // raced with a concurrent drop / refresh: the row is gone 
either way
+                    return true;
+                }
+                removeFromHashIndex(gone);
+                stateVersion++;
+                return true;
+            } finally {
+                stateLock.writeLock().unlock();
+            }
+        }
+    }
+
+    /**
+     * Updates the status of a baseline (ENABLE / DISABLE).
+     *
+     * @param id     the baseline id
+     * @param status the target status
+     * @return whether the update succeeded
+     */
+    public boolean updateStatus(long id, BaselineStatus status) {
+        ensureLoadedOrThrow();
+        // Two-phase status change: the INSERT(new) + DELETE(old) pair is 
internal-table I/O
+        // and runs under writerLock, never under stateLock. Matching 
tolerates a status
+        // flip racing with a lookup exactly like an ALTER landing right after 
the lookup
+        // (see findCandidateBaselines); on failure the in-memory flip is 
reverted.
+        synchronized (writerLock) {
+            BaselinePlan plan;
+            stateLock.readLock().lock();
+            try {
+                plan = baselines.get(id);
+            } finally {
+                stateLock.readLock().unlock();
+            }
+            if (plan == null) {
+                return false;
+            }
+            BaselineStatus previousStatus = plan.getStatus();
+            if (previousStatus == status) {
+                // ALTER to the already-set status: nothing to persist. 
(Inserting + "deleting
+                // the old row by its status" would delete the freshly 
inserted row as well,
+                // because both rows carry the same status.)
+                return true;
+            }
+            // The internal table is a DUPLICATE-key table on which UPDATE is 
not supported, so a
+            // status change is persisted as INSERT (new status) + DELETE (old 
status). The INSERT
+            // runs FIRST so the durable new row exists before any delete: a 
failure can never
+            // leave the in-memory state "old" while the only table row was 
already removed (the
+            // delete-then-insert gap, where the next refresh / restart 
silently dropped the
+            // baseline). Deleting by the PREVIOUS status can never touch the 
freshly inserted
+            // row (the statuses differ).
+            long previousUpdateTime = plan.getUpdateTime();
+            plan.setStatus(status);
+            plan.setUpdateTime(System.currentTimeMillis());
+            try {
+                persistInsert(plan);
+                persistDeleteByIdAndStatus(id, previousStatus);
+            } catch (RuntimeException e) {
+                // repair: delete the freshly inserted row by its (new) status 
- the old row
+                // was not touched yet, so the durable state is the old row 
again - then
+                // revert memory so memory and the table agree
+                try {
+                    persistDeleteByIdAndStatus(id, status);

Review Comment:
   [P1] Do not blindly compensate this multi-transaction status change. 
DELETE(old) can commit but report `KV_TXN_MAYBE_COMMITTED`; this catch then 
DELETEs(new) too and reverts memory, leaving no durable baseline for the next 
refresh/restart. Reconcile ambiguous outcomes (or use an atomic/versioned 
protocol) while preserving at least one version, and add a 
committed-delete/error fault-injection case.



-- 
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