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]