This is an automated email from the ASF dual-hosted git repository.

smengcl pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/master by this push:
     new 37e75cb7787 HDDS-16092. Stop takeSnapshot from lowering the persisted 
transaction index (#10953)
37e75cb7787 is described below

commit 37e75cb7787b8bef29307437f6074d44ef23c4e0
Author: Ritesh H Shukla <[email protected]>
AuthorDate: Wed Sep 9 17:40:29 2026 -0700

    HDDS-16092. Stop takeSnapshot from lowering the persisted transaction index 
(#10953)
    
    Generated-by: Claude Code (Fable 5)
---
 .../ozone/om/ratis/OzoneManagerDoubleBuffer.java   |  61 +++++-
 .../ozone/om/ratis/OzoneManagerStateMachine.java   |  14 +-
 ...estOzoneManagerDoubleBufferTransactionInfo.java | 235 +++++++++++++++++++++
 .../om/ratis/TestOzoneManagerStateMachine.java     | 140 +++++++++++-
 4 files changed, 435 insertions(+), 15 deletions(-)

diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerDoubleBuffer.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerDoubleBuffer.java
index 1f77f6f5b49..c822e0b9fe2 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerDoubleBuffer.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerDoubleBuffer.java
@@ -41,6 +41,7 @@
 import org.apache.hadoop.hdds.tracing.TracingUtil;
 import org.apache.hadoop.hdds.utils.TransactionInfo;
 import org.apache.hadoop.hdds.utils.db.BatchOperation;
+import org.apache.hadoop.hdds.utils.db.Table;
 import org.apache.hadoop.ozone.om.OMMetadataManager;
 import org.apache.hadoop.ozone.om.S3SecretManager;
 import org.apache.hadoop.ozone.om.codec.OMDBDefinition;
@@ -102,6 +103,23 @@ public final class OzoneManagerDoubleBuffer {
   /** The number of flush iterations (for testing and debug only). */
   private final AtomicLong flushIterations = new AtomicLong();
 
+  /**
+   * Orders the two writers of {@link 
org.apache.hadoop.ozone.OzoneConsts#TRANSACTION_INFO_KEY}:
+   * the batch commit in {@link #flushBatch} and {@link #persistIfNewer}. The 
commit publishes a
+   * batch of transactions together with the index describing them, while the 
state machine
+   * derives its own index from a field this class only advances once that 
commit has returned.
+   * Unordered, a snapshot taken in between can store an index lower than the 
data already in the
+   * DB, leaving a watermark that disclaims committed transactions.
+   */
+  private final Object transactionInfoLock = new Object();
+
+  /**
+   * Runs inside {@link #persistIfNewer} between reading the stored 
transaction info and writing
+   * the candidate. Lets a test attempt a batch commit in exactly that window, 
which is the one
+   * {@link #transactionInfoLock} has to close. No-op in production.
+   */
+  private volatile Runnable afterTransactionInfoRead = () -> { };
+
   /** Entry for {@link #currentBuffer} and {@link #readyBuffer}. */
   private static class Entry {
     private final TermIndex termIndex;
@@ -351,6 +369,38 @@ void flushCurrentBuffer() {
     }
   }
 
+  /**
+   * Stores {@code candidate} as the transaction info unless what is already 
stored is at or ahead
+   * of it, and returns the value stored when this method returns.
+   * <p>
+   * Every write of this key that is not part of a batch commit must come 
through here. The stored
+   * index states which transactions the DB contains, so lowering it would 
make a restart replay
+   * transactions that are already applied.
+   *
+   * @return the stored transaction info, never null
+   */
+  TransactionInfo persistIfNewer(TransactionInfo candidate) throws IOException 
{
+    synchronized (transactionInfoLock) {
+      final Table<String, TransactionInfo> table = 
omMetadataManager.getTransactionInfoTable();
+      // Skip the table cache, as TransactionInfo.readTransactionInfo does for 
this same key. The
+      // value this compares against is written by a batch commit, which does 
not populate the
+      // cache; reading through the cache would risk comparing against a stale 
value and writing
+      // the lower index anyway, which is the whole defect.
+      final TransactionInfo stored = table.getSkipCache(TRANSACTION_INFO_KEY);
+      afterTransactionInfoRead.run();
+      if (stored != null && stored.compareTo(candidate) >= 0) {
+        return stored;
+      }
+      table.put(TRANSACTION_INFO_KEY, candidate);
+      return candidate;
+    }
+  }
+
+  @VisibleForTesting
+  void setAfterTransactionInfoRead(Runnable hook) {
+    this.afterTransactionInfoRead = hook;
+  }
+
   private void flushBatch(Queue<Entry> buffer) throws IOException {
     Map<String, List<Long>> cleanupEpochs = new HashMap<>();
     // Commit transaction info to DB.
@@ -376,9 +426,14 @@ private void flushBatch(Queue<Entry> buffer) throws 
IOException {
               batchOperation, TRANSACTION_INFO_KEY, 
TransactionInfo.valueOf(lastTransaction)));
 
       long startTime = Time.monotonicNow();
-      flushBatchWithTrace(lastTraceId, buffer.size(),
-          () -> omMetadataManager.getStore()
-              .commitBatchOperation(batchOperation));
+      // The commit is the point where this batch's transactions and the index 
describing them
+      // both become visible, so it is what persistIfNewer has to be ordered 
against. Only the
+      // flush daemon commits and snapshots are rare, so the lock is all but 
uncontended.
+      synchronized (transactionInfoLock) {
+        flushBatchWithTrace(lastTraceId, buffer.size(),
+            () -> omMetadataManager.getStore()
+                .commitBatchOperation(batchOperation));
+      }
 
       metrics.updateFlushTime(Time.monotonicNow() - startTime);
     }
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
index 50bb1f3b376..c80a7977a1f 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java
@@ -17,7 +17,6 @@
 
 package org.apache.hadoop.ozone.om.ratis;
 
-import static org.apache.hadoop.ozone.OzoneConsts.TRANSACTION_INFO_KEY;
 import static 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Status.INTERNAL_ERROR;
 import static 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Status.METADATA_ERROR;
 
@@ -626,14 +625,19 @@ private synchronized long takeSnapshotImpl() throws 
IOException {
       final TermIndex snapshot = applied.compareTo(notified) > 0 ? applied : 
notified;
 
       long startTime = Time.monotonicNow();
-      final TransactionInfo transactionInfo = 
TransactionInfo.valueOf(snapshot);
+      // The double buffer may have committed transactions past the index 
computed above, since it
+      // advances lastAppliedTermIndex only after its batch commit returns. 
Persisting through it
+      // keeps the stored index from moving backwards over that data, and 
yields whichever value is
+      // actually stored so the in-memory copy Ratis reads cannot disagree 
with the DB.
+      final TransactionInfo transactionInfo =
+          
ozoneManagerDoubleBuffer.persistIfNewer(TransactionInfo.valueOf(snapshot));
       ozoneManager.setTransactionInfo(transactionInfo);
-      
ozoneManager.getMetadataManager().getTransactionInfoTable().put(TRANSACTION_INFO_KEY,
 transactionInfo);
       ozoneManager.getMetadataManager().getStore().flushDB();
       LOG.info("{}: taking snapshot. applied = {}, skipped = {}, " +
           "notified = {}, current snapshot index = {}, took {} ms",
-              getId(), applied, lastSkippedIndex, notified, snapshot, 
Time.monotonicNow() - startTime);
-      return snapshot.getIndex();
+              getId(), applied, lastSkippedIndex, notified, 
transactionInfo.getTermIndex(),
+              Time.monotonicNow() - startTime);
+      return transactionInfo.getTermIndex().getIndex();
     }
   }
 
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerDoubleBufferTransactionInfo.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerDoubleBufferTransactionInfo.java
new file mode 100644
index 00000000000..2ff1f15e24c
--- /dev/null
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerDoubleBufferTransactionInfo.java
@@ -0,0 +1,235 @@
+/*
+ * 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.hadoop.ozone.om.ratis;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.any;
+import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.utils.TransactionInfo;
+import org.apache.hadoop.hdds.utils.db.Table;
+import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.audit.AuditLogger;
+import org.apache.hadoop.ozone.audit.AuditMessage;
+import org.apache.hadoop.ozone.om.OMConfigKeys;
+import org.apache.hadoop.ozone.om.OMMetrics;
+import org.apache.hadoop.ozone.om.OmMetadataManagerImpl;
+import org.apache.hadoop.ozone.om.OzoneManager;
+import org.apache.hadoop.ozone.om.response.OMClientResponse;
+import org.apache.hadoop.ozone.om.response.key.OMKeyCreateResponse;
+import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
+import 
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
+import org.apache.ratis.server.protocol.TermIndex;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/**
+ * Tests that the transaction index persisted under TRANSACTION_INFO_KEY only 
ever moves forward.
+ * <p>
+ * Two writers reach that key: the batch commit inside the double buffer, 
which publishes a batch
+ * of transactions together with the index describing them, and {@code 
persistIfNewer}, used when
+ * the state machine takes a snapshot from an index the buffer has not yet 
caught up to. If the
+ * second can land on top of a commit it did not observe, the DB is left 
holding transactions its
+ * own watermark disclaims and a restart replays them. See HDDS-16092.
+ * <p>
+ * These tests keep the flush daemon stopped so the test thread is the only 
flusher, matching
+ * production, where {@code flushCurrentBuffer} is reached from the daemon 
alone. Two threads
+ * committing concurrently would write indexes out of order for reasons that 
have nothing to do
+ * with the ordering under test.
+ */
+class TestOzoneManagerDoubleBufferTransactionInfo {
+
+  private OzoneManagerDoubleBuffer doubleBuffer;
+  private OmMetadataManagerImpl omMetadataManager;
+  private OMClientResponse keyCreateResponse;
+
+  @TempDir
+  private File tempDir;
+
+  @BeforeEach
+  public void setup() throws IOException {
+    OzoneConfiguration conf = new OzoneConfiguration();
+    conf.set(OMConfigKeys.OZONE_OM_DB_DIRS, tempDir.getAbsolutePath());
+
+    OzoneManager ozoneManager = mock(OzoneManager.class);
+    when(ozoneManager.getMetrics()).thenReturn(OMMetrics.create(conf));
+    AuditLogger auditLogger = mock(AuditLogger.class);
+    doNothing().when(auditLogger).logWrite(any(AuditMessage.class));
+    when(ozoneManager.getAuditLogger()).thenReturn(auditLogger);
+
+    omMetadataManager = new OmMetadataManagerImpl(conf, ozoneManager);
+    when(ozoneManager.getMetadataManager()).thenReturn(omMetadataManager);
+
+    // Not started: the test thread drives every flush, so the ordering under 
test is the only
+    // concurrency in play.
+    doubleBuffer = OzoneManagerDoubleBuffer.newBuilder()
+        .setOmMetadataManager(omMetadataManager)
+        .setMaxUnFlushedTransactionCount(1000)
+        .build();
+
+    OMResponse omResponse = mock(OMResponse.class);
+    when(omResponse.getTraceID()).thenReturn("traceId");
+    
when(omResponse.getCmdType()).thenReturn(OzoneManagerProtocolProtos.Type.CreateKey);
+    keyCreateResponse = mock(OMKeyCreateResponse.class);
+    when(keyCreateResponse.getOMResponse()).thenReturn(omResponse);
+    doNothing().when(keyCreateResponse).checkAndUpdateDB(any(), any());
+  }
+
+  @AfterEach
+  public void tearDown() throws IOException {
+    if (doubleBuffer != null) {
+      doubleBuffer.stop();
+    }
+    if (omMetadataManager != null) {
+      omMetadataManager.stop();
+    }
+  }
+
+  private Table<String, TransactionInfo> transactionInfoTable() {
+    return omMetadataManager.getTransactionInfoTable();
+  }
+
+  private TermIndex storedTermIndex() throws IOException {
+    final TransactionInfo stored = 
transactionInfoTable().get(OzoneConsts.TRANSACTION_INFO_KEY);
+    return stored == null ? null : stored.getTermIndex();
+  }
+
+  private void commit(long index) throws IOException {
+    doubleBuffer.add(keyCreateResponse, TransactionInfo.getTermIndex(index));
+    doubleBuffer.flushCurrentBuffer();
+  }
+
+  /**
+   * Drives the exact interleaving the ordering exists to exclude: a snapshot 
is held between
+   * reading the stored index and writing its own, while a real batch commit 
is attempted. The
+   * commit has to be excluded from that window, otherwise the snapshot's 
older index lands on
+   * top of it.
+   * <p>
+   * This is the only test that covers the locking half of the fix. Removing 
the lock while
+   * leaving the comparison in persistIfNewer fails here and nowhere else; the 
window it needs
+   * is far too narrow to hit by chance.
+   */
+  @Test
+  public void testPersistIfNewerIsOrderedAgainstBatchCommit() throws Exception 
{
+    transactionInfoTable().put(OzoneConsts.TRANSACTION_INFO_KEY,
+        TransactionInfo.valueOf(TransactionInfo.getTermIndex(50)));
+
+    final CountDownLatch snapshotRead = new CountDownLatch(1);
+    final CountDownLatch commitDone = new CountDownLatch(1);
+    final ExecutorService committer = Executors.newSingleThreadExecutor();
+    try {
+      doubleBuffer.setAfterTransactionInfoRead(() -> {
+        snapshotRead.countDown();
+        try {
+          // Long enough that an unordered commit would have finished inside 
the window. When it
+          // is ordered this simply times out and the snapshot carries on.
+          commitDone.await(5, TimeUnit.SECONDS);
+        } catch (InterruptedException e) {
+          Thread.currentThread().interrupt();
+        }
+      });
+
+      final Future<?> commit = committer.submit(() -> {
+        snapshotRead.await();
+        commit(105);
+        commitDone.countDown();
+        return null;
+      });
+
+      
doubleBuffer.persistIfNewer(TransactionInfo.valueOf(TransactionInfo.getTermIndex(100)));
+      commit.get(30, TimeUnit.SECONDS);
+
+      assertEquals(TransactionInfo.getTermIndex(105), storedTermIndex(),
+          "a batch commit must not be overwritten by a snapshot that read 
before it");
+    } finally {
+      doubleBuffer.setAfterTransactionInfoRead(() -> { });
+      committer.shutdownNow();
+    }
+  }
+
+  /**
+   * The same invariant without a staged interleaving: with commits and 
snapshots running
+   * concurrently, a reader must never observe the stored index go backwards.
+   * <p>
+   * This covers the comparison in persistIfNewer, which it fails within a 
handful of rounds when
+   * removed. It does not cover the lock -- the read-to-write window is too 
narrow to lose the
+   * race by chance -- so it complements the staged test above rather than 
repeating it.
+   */
+  @Test
+  public void testPersistedTransactionInfoNeverMovesBackwards() throws 
Exception {
+    final int rounds = 300;
+    final AtomicBoolean running = new AtomicBoolean(true);
+    final AtomicReference<String> regression = new AtomicReference<>();
+    final ExecutorService threads = Executors.newFixedThreadPool(2);
+    try {
+      final Future<?> snapshots = threads.submit(() -> {
+        while (running.get()) {
+          final TermIndex before = storedTermIndex();
+          if (before != null) {
+            doubleBuffer.persistIfNewer(TransactionInfo.valueOf(before));
+          }
+        }
+        return null;
+      });
+
+      final Future<?> reader = threads.submit(() -> {
+        long highest = -1;
+        while (running.get()) {
+          final TermIndex seen = storedTermIndex();
+          if (seen != null) {
+            if (seen.getIndex() < highest) {
+              regression.compareAndSet(null,
+                  "stored index went backwards: " + highest + " -> " + 
seen.getIndex());
+              return null;
+            }
+            highest = seen.getIndex();
+          }
+        }
+        return null;
+      });
+
+      for (int i = 1; i <= rounds; i++) {
+        commit(i);
+      }
+      running.set(false);
+      snapshots.get(30, TimeUnit.SECONDS);
+      reader.get(30, TimeUnit.SECONDS);
+
+      assertNull(regression.get(), regression.get());
+      assertEquals(TransactionInfo.getTermIndex(rounds), storedTermIndex());
+    } finally {
+      running.set(false);
+      threads.shutdownNow();
+    }
+  }
+}
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
index 5cc0b52cfc5..072bb84188e 100644
--- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java
@@ -52,7 +52,6 @@
 import org.apache.hadoop.hdds.conf.OzoneConfiguration;
 import org.apache.hadoop.hdds.utils.TransactionInfo;
 import org.apache.hadoop.hdds.utils.db.DBStore;
-import org.apache.hadoop.hdds.utils.db.Table;
 import org.apache.hadoop.ozone.OzoneConsts;
 import org.apache.hadoop.ozone.audit.AuditEventStatus;
 import org.apache.hadoop.ozone.audit.AuditMessage;
@@ -700,12 +699,10 @@ public void testTakeSnapshotAlreadyCaughtUp() throws 
Exception {
     assertTermIndex(1, 1, sm.getLastAppliedTermIndex());
 
     OMMetadataManager metaMgr = mock(OMMetadataManager.class);
-    @SuppressWarnings("unchecked")
-    Table<String, TransactionInfo> txnTable = mock(Table.class);
     DBStore store = mock(DBStore.class);
     when(om.getMetadataManager()).thenReturn(metaMgr);
-    when(metaMgr.getTransactionInfoTable()).thenReturn(txnTable);
     when(metaMgr.getStore()).thenReturn(store);
+    stubPersistIfNewerAsAccepting();
 
     long snapshotIndex = sm.takeSnapshot();
 
@@ -714,6 +711,137 @@ public void testTakeSnapshotAlreadyCaughtUp() throws 
Exception {
     verify(om).setTransactionInfo(any(TransactionInfo.class));
   }
 
+  /**
+   * Makes the mocked double buffer behave as if nothing newer were stored, so 
the candidate the
+   * state machine computes is the one that gets persisted and returned.
+   */
+  private void stubPersistIfNewerAsAccepting() throws IOException {
+    when(doubleBuffer.persistIfNewer(any(TransactionInfo.class)))
+        .thenAnswer(invocation -> invocation.getArgument(0));
+  }
+
+  /**
+   * The double buffer commits the transaction data and TRANSACTION_INFO_KEY 
in one batch and
+   * only advances the state machine's applied index afterwards, so a snapshot 
taken in that
+   * gap computes its index from a value that lags what is already durable. 
Writing it must not
+   * move the stored index backwards: the DB would then hold data its own 
watermark disclaims,
+   * and a restart would replay transactions the DB already contains. See 
HDDS-16092.
+   */
+  @Test
+  public void testTakeSnapshotDoesNotLowerPersistedTransactionInfo(@TempDir 
Path tmpDir)
+      throws Exception {
+    OzoneConfiguration conf = new OzoneConfiguration();
+    conf.set(OMConfigKeys.OZONE_OM_DB_DIRS, 
tmpDir.toAbsolutePath().toString());
+    OzoneManager realishOm = mock(OzoneManager.class);
+    when(realishOm.getConfiguration()).thenReturn(conf);
+    when(realishOm.getConfig()).thenReturn(conf.getObject(OmConfig.class));
+    OMMetadataManager realMetaMgr = new OmMetadataManagerImpl(conf, realishOm);
+    try {
+      when(realishOm.getMetadataManager()).thenReturn(realMetaMgr);
+      // Not started: the flush daemon plays no part here, only the persist 
path does.
+      OzoneManagerDoubleBuffer realBuffer = 
OzoneManagerDoubleBuffer.newBuilder()
+          .setOmMetadataManager(realMetaMgr)
+          .setMaxUnFlushedTransactionCount(1)
+          .build();
+
+      // A batch committed transactions up to index 105 and is blocked before 
advancing the
+      // state machine's applied index, which still reads 100.
+      TransactionInfo flushed = TransactionInfo.valueOf(TermIndex.valueOf(1, 
105));
+      
realMetaMgr.getTransactionInfoTable().put(OzoneConsts.TRANSACTION_INFO_KEY, 
flushed);
+
+      OzoneManagerStateMachine testSm =
+          new OzoneManagerStateMachine(realishOm, realBuffer, handler, 
executor, null);
+      testSm.updateLastAppliedTermIndex(TermIndex.valueOf(1, 100));
+
+      long snapshotIndex = testSm.takeSnapshot();
+
+      assertEquals(flushed,
+          
realMetaMgr.getTransactionInfoTable().get(OzoneConsts.TRANSACTION_INFO_KEY),
+          "takeSnapshot must not lower the persisted transaction info");
+
+      // The in-memory copy is what getLatestSnapshot() reports to Ratis, so 
it must not
+      // disagree with the DB either.
+      ArgumentCaptor<TransactionInfo> published =
+          ArgumentCaptor.forClass(TransactionInfo.class);
+      verify(realishOm).setTransactionInfo(published.capture());
+      assertEquals(flushed, published.getValue());
+      assertEquals(105, snapshotIndex);
+    } finally {
+      realMetaMgr.stop();
+    }
+  }
+
+  /**
+   * The guard added for HDDS-16092 must not stop a snapshot from advancing 
the stored index
+   * when its own value is the newer one -- that is the whole point of the 
write.
+   */
+  @Test
+  public void testTakeSnapshotAdvancesPersistedTransactionInfo(@TempDir Path 
tmpDir)
+      throws Exception {
+    OzoneConfiguration conf = new OzoneConfiguration();
+    conf.set(OMConfigKeys.OZONE_OM_DB_DIRS, 
tmpDir.toAbsolutePath().toString());
+    OzoneManager realishOm = mock(OzoneManager.class);
+    when(realishOm.getConfiguration()).thenReturn(conf);
+    when(realishOm.getConfig()).thenReturn(conf.getObject(OmConfig.class));
+    OMMetadataManager realMetaMgr = new OmMetadataManagerImpl(conf, realishOm);
+    try {
+      when(realishOm.getMetadataManager()).thenReturn(realMetaMgr);
+      // Not started: the flush daemon plays no part here, only the persist 
path does.
+      OzoneManagerDoubleBuffer realBuffer = 
OzoneManagerDoubleBuffer.newBuilder()
+          .setOmMetadataManager(realMetaMgr)
+          .setMaxUnFlushedTransactionCount(1)
+          .build();
+
+      
realMetaMgr.getTransactionInfoTable().put(OzoneConsts.TRANSACTION_INFO_KEY,
+          TransactionInfo.valueOf(TermIndex.valueOf(1, 50)));
+
+      OzoneManagerStateMachine testSm =
+          new OzoneManagerStateMachine(realishOm, realBuffer, handler, 
executor, null);
+      testSm.updateLastAppliedTermIndex(TermIndex.valueOf(1, 100));
+
+      long snapshotIndex = testSm.takeSnapshot();
+
+      assertEquals(TransactionInfo.valueOf(TermIndex.valueOf(1, 100)),
+          
realMetaMgr.getTransactionInfoTable().get(OzoneConsts.TRANSACTION_INFO_KEY));
+      assertEquals(100, snapshotIndex);
+    } finally {
+      realMetaMgr.stop();
+    }
+  }
+
+  /**
+   * With nothing stored yet the snapshot's value is the only candidate and 
must be written.
+   */
+  @Test
+  public void testTakeSnapshotWritesWhenNothingPersisted(@TempDir Path tmpDir) 
throws Exception {
+    OzoneConfiguration conf = new OzoneConfiguration();
+    conf.set(OMConfigKeys.OZONE_OM_DB_DIRS, 
tmpDir.toAbsolutePath().toString());
+    OzoneManager realishOm = mock(OzoneManager.class);
+    when(realishOm.getConfiguration()).thenReturn(conf);
+    when(realishOm.getConfig()).thenReturn(conf.getObject(OmConfig.class));
+    OMMetadataManager realMetaMgr = new OmMetadataManagerImpl(conf, realishOm);
+    try {
+      when(realishOm.getMetadataManager()).thenReturn(realMetaMgr);
+      // Not started: the flush daemon plays no part here, only the persist 
path does.
+      OzoneManagerDoubleBuffer realBuffer = 
OzoneManagerDoubleBuffer.newBuilder()
+          .setOmMetadataManager(realMetaMgr)
+          .setMaxUnFlushedTransactionCount(1)
+          .build();
+
+      OzoneManagerStateMachine testSm =
+          new OzoneManagerStateMachine(realishOm, realBuffer, handler, 
executor, null);
+      testSm.updateLastAppliedTermIndex(TermIndex.valueOf(2, 7));
+
+      long snapshotIndex = testSm.takeSnapshot();
+
+      assertEquals(TransactionInfo.valueOf(TermIndex.valueOf(2, 7)),
+          
realMetaMgr.getTransactionInfoTable().get(OzoneConsts.TRANSACTION_INFO_KEY));
+      assertEquals(7, snapshotIndex);
+    } finally {
+      realMetaMgr.stop();
+    }
+  }
+
   @Test
   public void testTakeSnapshotWaitsForFlush() throws Exception {
     // Create a skip gap: notify 0, then skip to 3
@@ -732,12 +860,10 @@ public void testTakeSnapshotWaitsForFlush() throws 
Exception {
     }).when(doubleBuffer).awaitFlush();
 
     OMMetadataManager metaMgr = mock(OMMetadataManager.class);
-    @SuppressWarnings("unchecked")
-    Table<String, TransactionInfo> txnTable = mock(Table.class);
     DBStore store = mock(DBStore.class);
     when(om.getMetadataManager()).thenReturn(metaMgr);
-    when(metaMgr.getTransactionInfoTable()).thenReturn(txnTable);
     when(metaMgr.getStore()).thenReturn(store);
+    stubPersistIfNewerAsAccepting();
 
     long snapshotIndex = sm.takeSnapshot();
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to