diff --git a/docs/_docs/snapshots/snapshots.adoc b/docs/_docs/snapshots/snapshots.adoc index fe61b3c49234a..a8129b4b998f1 100644 --- a/docs/_docs/snapshots/snapshots.adoc +++ b/docs/_docs/snapshots/snapshots.adoc @@ -287,6 +287,40 @@ control.(sh|bat) --snapshot restore snapshot_09062021 --groups cache-group1,cach control.(sh|bat) --snapshot restore snapshot_09062021 --increment 1 ---- +== Listing Snapshots + +To list snapshots created on the cluster, use the following command. + +[tabs] +-- +tab:Unix[] +[source,shell] +---- +# List all snapshots in the default snapshot directory. +control.sh --snapshot list + +# List all snapshots in the "/tmp/ignite/snapshots" folder. +control.sh --snapshot list --src /tmp/ignite/snapshots +---- + +tab:Windows[] +[source,shell] +---- +# List all snapshots in the default snapshot directory. +control.bat --snapshot list + +# List all snapshots in the "C:\\tmp\\ignite\\snapshots" folder. +control.bat --snapshot list --src C:\\tmp\\ignite\\snapshots +---- +-- + +=== List operation notes + +* The list operation is read-only and does not modify any snapshot data. +* If a snapshot is being modified (created or deleted), its size might not be reported. +* A snapshot is being detected by its metadata. If metadata is missing or cannot be read, snapshot is ignored. +* When listing snapshots and calculating size, no additional snapshot validation or checksum calculation is performed. + == Getting Snapshot Operation Status The status of the current snapshot operation in the cluster can be obtained using the `control.sh|bat` script or JMX interface: diff --git a/modules/control-utility/src/test/java/org/apache/ignite/testsuites/IgniteControlUtilityTestSuite6.java b/modules/control-utility/src/test/java/org/apache/ignite/testsuites/IgniteControlUtilityTestSuite6.java new file mode 100644 index 0000000000000..925394f4639b8 --- /dev/null +++ b/modules/control-utility/src/test/java/org/apache/ignite/testsuites/IgniteControlUtilityTestSuite6.java @@ -0,0 +1,32 @@ +/* + * 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.ignite.testsuites; + +import org.apache.ignite.util.GridCommandHandlerListSnapshotTest; +import org.junit.runner.RunWith; +import org.junit.runners.Suite; + +/** + * Test suite for control utility. + */ +@RunWith(Suite.class) +@Suite.SuiteClasses({ + GridCommandHandlerListSnapshotTest.class +}) +public class IgniteControlUtilityTestSuite6 { +} diff --git a/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerAbstractTest.java b/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerAbstractTest.java index 092fe78903246..7d5f2e14c289d 100644 --- a/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerAbstractTest.java +++ b/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerAbstractTest.java @@ -499,6 +499,22 @@ protected void createCacheAndPreload( ) { assert nonNull(ignite); + CacheConfiguration ccfg = testCacheConfiguration(cacheName, partitions, filter); + + ignite.createCache(ccfg); + + try (IgniteDataStreamer streamer = ignite.dataStreamer(cacheName)) { + for (int i = 0; i < countEntries; i++) + streamer.addData(i, i); + } + } + + /** */ + protected CacheConfiguration testCacheConfiguration( + String cacheName, + int partitions, + @Nullable IgnitePredicate filter + ) { CacheConfiguration ccfg = new CacheConfiguration<>(cacheName) .setAffinity(new RendezvousAffinityFunction(false, partitions)) .setBackups(1) @@ -507,12 +523,7 @@ protected void createCacheAndPreload( if (filter != null) ccfg.setNodeFilter(filter); - ignite.createCache(ccfg); - - try (IgniteDataStreamer streamer = ignite.dataStreamer(cacheName)) { - for (int i = 0; i < countEntries; i++) - streamer.addData(i, i); - } + return ccfg; } /** diff --git a/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerListSnapshotTest.java b/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerListSnapshotTest.java new file mode 100644 index 0000000000000..fea71bb75c72e --- /dev/null +++ b/modules/control-utility/src/test/java/org/apache/ignite/util/GridCommandHandlerListSnapshotTest.java @@ -0,0 +1,352 @@ +/* + * 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.ignite.util; + +import java.io.File; +import java.nio.file.DirectoryStream; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.util.Collection; +import org.apache.ignite.Ignite; +import org.apache.ignite.IgniteDataStreamer; +import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.configuration.CacheConfiguration; +import org.apache.ignite.configuration.DataStorageConfiguration; +import org.apache.ignite.configuration.IgniteConfiguration; +import org.apache.ignite.internal.IgniteEx; +import org.apache.ignite.internal.management.snapshot.SnapshotListCommand; +import org.apache.ignite.internal.processors.cache.persistence.filename.SnapshotFileTree; +import org.apache.ignite.internal.util.typedef.F; +import org.apache.ignite.internal.util.typedef.G; +import org.apache.ignite.internal.util.typedef.internal.U; +import org.apache.ignite.lang.IgnitePredicate; +import org.apache.ignite.testframework.GridTestUtils; +import org.jetbrains.annotations.Nullable; +import org.junit.Test; +import org.junit.runners.Parameterized.Parameter; +import org.junit.runners.Parameterized.Parameters; + +import static java.nio.file.Files.newDirectoryStream; +import static org.apache.ignite.cluster.ClusterState.ACTIVE; +import static org.apache.ignite.internal.commandline.CommandHandler.EXIT_CODE_OK; +import static org.apache.ignite.internal.processors.cache.persistence.snapshot.AbstractSnapshotSelfTest.snp; +import static org.junit.Assume.assumeFalse; +import static org.junit.Assume.assumeTrue; + +/** Test for the command '--snapshot list'. */ +public class GridCommandHandlerListSnapshotTest extends GridCommandHandlerAbstractTest { + /** External storage path. */ + private static final String EXT_STORAGE_PATH = "extStorage"; + + /** Flag to use {@link DataStorageConfiguration#setExtraSnapshotPaths(String...)}. */ + private boolean extStorages; + + /** Resolved external storages paths. {@code null} if {@code extStorages} is {@code false}. */ + private @Nullable String[] extStoragePaths; + + /** Consistent id postfix. */ + private @Nullable String cstIdPostfix = ""; + + /** Flag setting the usage of a custom snapshot path. */ + @Parameter(1) + public boolean customPath; + + /** Flag setting the usage of dedicated, own node working directories. */ + @Parameter(2) + public boolean separatedWorkDir; + + /** Flag to add external server node after the snapshot creation. */ + @Parameter(3) + public boolean addExtraSrvr; + + /** Number of incremental snapshots to add to the main test snapshots. */ + @Parameter(4) + public int incCnt; + + /** */ + @Parameters(name = "cmdHnd={0},customPath={1},ownWorkDir={2},addExtraSrvr={3},incCnt={4}") + public static Collection parameters() { + return GridTestUtils.cartesianProduct( + commandHandlers(), + F.asList(false, true), // Custom snapshot path + F.asList(false, true), // Separated (own) work directories + F.asList(false, true), // Add a server node after the snapshot creation + F.asList(0, 2) // Number of incremental snapshots + ); + } + + /** {@inheritDoc} */ + @Override protected void afterTest() throws Exception { + super.afterTest(); + + stopAllGrids(); + + cleanPersistenceDir(); + } + + /** {@inheritDoc} */ + @Override protected void beforeTest() throws Exception { + super.beforeTest(); + + /** Handy if test running is interrupted and {@link #afterTest()} isn't invoked. */ + cleanPersistenceDir(); + } + + /** {@inheritDoc} */ + @Override protected void cleanPersistenceDir() throws Exception { + super.cleanPersistenceDir(); + + // Also cleans separated snapshot working directories and custom snapshot paths. + try (DirectoryStream files = newDirectoryStream(Paths.get(U.defaultWorkDirectory()))) { + for (Path path : files) + U.delete(path); + } + } + + /** {@inheritDoc} */ + @Override protected IgniteConfiguration getConfiguration(String igniteInstanceName) throws Exception { + IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName); + + String workDir = separatedWorkDir + ? new File(U.defaultWorkDirectory(), igniteInstanceName).getAbsolutePath() + : U.defaultWorkDirectory(); + + cfg.setWorkDirectory(workDir); + + cfg.getDataStorageConfiguration().setWalCompactionEnabled(incCnt > 0); + + if (extStorages) { + cfg.getDataStorageConfiguration().setExtraStoragePaths( + workDir + File.separator, + workDir + File.separator + EXT_STORAGE_PATH + ); + + extStoragePaths = cfg.getDataStorageConfiguration().getExtraStoragePaths(); + + cfg.getDataStorageConfiguration().setExtraSnapshotPaths("", EXT_STORAGE_PATH); + } + + return cfg; + } + + /** {@inheritDoc} */ + @Override public String getTestIgniteInstanceName() { + return super.getTestIgniteInstanceName() + cstIdPostfix; + } + + /** */ + @Test + public void testNoSnapshots() throws Exception { + assumeFalse(addExtraSrvr || incCnt > 0); + + doTestSnapshotsLists(0, false, false); + } + + /** */ + @Test + public void testSingleSnapshot() throws Exception { + doTestSnapshotsLists(1, false, false); + } + + /** */ + @Test + public void testSeveralSnapshots() throws Exception { + doTestSnapshotsLists(3, false, false); + } + + /** */ + @Test + public void testOneNodeMisses() throws Exception { + // Doesn't matter here. + assumeFalse(incCnt > 0); + // Let's keep just one node not seeing the snapshot. + assumeTrue(separatedWorkDir); + // Almost the same tests. + assumeFalse(addExtraSrvr); + + doTestSnapshotsLists(2, true, false); + } + + /** */ + @Test + public void testExternalStorages() throws Exception { + // External storages are required to be the same as configured in the node's PDS storages. Thus, we skip different work folders. + // Also, extra snapshot storages aren't used if snapshot is created with a custom path. + assumeFalse(separatedWorkDir || customPath); + + extStorages = true; + + doTestSnapshotsLists(2, false, false); + } + + /** */ + @Test + public void testChangedConsistentId() throws Exception { + // In new created working directories there will be obviously no snapshots. + assumeFalse(separatedWorkDir); + // Fastens the tests + assumeFalse(incCnt > 0); + + doTestSnapshotsLists(1, false, true); + } + + /** */ + private void doTestSnapshotsLists(int snpCnt, boolean deleteOnOneNode, boolean restartWithChangedCstIds) throws Exception { + // A custom snapshot path actually puts snapshots in a shared directory. This skews the results when dedicated + // work directories are set. + assumeFalse(customPath && separatedWorkDir); + + int srvrsCnt = 3; + int entriesCnt = 10; + int partitions = 4; + + IgniteEx ig = (IgniteEx)startGridsMultiThreaded(srvrsCnt); + + startGrid(CLIENT_NODE_NAME_PREFIX); + + ig.cluster().state(ACTIVE); + + File snpsRootFile = customPath + ? new File(ig.context().pdsFolderResolver().fileTree().snapshotsRoot(), "ex_snapshots") + : null; + + // Flag if 'testSnapshot0' deleted on node0. + boolean grid0HasNoSnapshot0 = false; + + // Create snapshots. + if (snpCnt > 0) { + createCacheAndPreload(ig, DEFAULT_CACHE_NAME, entriesCnt, partitions, null); + + String absPathStr = null; + + for (int snpIdx = 0; snpIdx < snpCnt; snpIdx++) { + absPathStr = customPath ? snpsRootFile.getAbsolutePath() : null; + + snp(ig).createSnapshot("testSnapshot" + snpIdx, absPathStr, false, false).get(getTestTimeout()); + + for (int incIdx = 0; incIdx < incCnt; incIdx++) { + try (IgniteDataStreamer ds = ig.dataStreamer(DEFAULT_CACHE_NAME)) { + for (int val = (snpIdx + incIdx + 1) * entriesCnt; val < (snpIdx + incIdx + 2) * entriesCnt; val++) + ds.addData(val, val); + } + + snp(ig).createSnapshot("testSnapshot" + snpIdx, absPathStr, true, false).get(getTestTimeout()); + } + } + + if (deleteOnOneNode) { + SnapshotFileTree sft = new SnapshotFileTree(ig.context(), "testSnapshot0", absPathStr); + + assertTrue(sft.root().exists()); + assertTrue(U.delete(sft.root())); + assertFalse(sft.root().exists()); + + grid0HasNoSnapshot0 = true; + } + } + + if (restartWithChangedCstIds) { + stopAllGrids(); + + cstIdPostfix = "_changed"; + + startGridsMultiThreaded(srvrsCnt); + + startGrid(CLIENT_NODE_NAME_PREFIX); + } + + // Add a server. + if (addExtraSrvr) + startGrid(G.allGrids().size()); + + injectTestSystemOut(); + + // Requests snapshots. + if (customPath) { + assertEquals(EXIT_CODE_OK, execute(newCommandHandler(createTestLogger()), "--snapshot", "list", "--src", + snpsRootFile.getAbsolutePath())); + } + else + assertEquals(EXIT_CODE_OK, execute(newCommandHandler(createTestLogger()), "--snapshot", "list")); + + String out = testOut.toString(); + + // Find the nodes in the output. + for (Ignite g : G.allGrids()) { + ClusterNode n = g.cluster().localNode(); + + assertEquals(n.isClient() ? 0 : 1, countEntries(out, "Node '%s'".formatted(n.consistentId().toString()))); + } + + // Ensure that there are no snapshots. + if (snpCnt == 0) { + assertFalse(out.contains("Snapshot '")); + + assertEquals(srvrsCnt + (addExtraSrvr ? 1 : 0), countEntries(out, SnapshotListCommand.NO_SNAPSHOTS)); + + return; + } + + assertEquals((addExtraSrvr && separatedWorkDir ? 1 : 0), countEntries(out, SnapshotListCommand.NO_SNAPSHOTS)); + + // The additional server node doesn't have snapshots. But it can see them if shared the snapshot directory. + int snpsRecordsCnt = srvrsCnt + (addExtraSrvr + ? (separatedWorkDir ? 0 : 1) + : 0 + ); + + // Find the snapshots in the output. + for (int snpIdx = 0; snpIdx < snpCnt; snpIdx++) { + // 'testSnapshot0' has fewer records if was deleted on one node. + int certainSnpRecordsCnt = snpIdx == 0 && grid0HasNoSnapshot0 + ? snpsRecordsCnt - 1 + : snpsRecordsCnt; + + assertEquals(certainSnpRecordsCnt, countEntries(out, "Snapshot 'testSnapshot" + snpIdx + "'")); + + if (incCnt > 0) + assertEquals(certainSnpRecordsCnt * snpCnt, countEntries(out, "incremental snapshots: cnt=" + incCnt)); + + if (extStorages) + assertEquals(certainSnpRecordsCnt * snpCnt, countEntries(out, "external storages: cnt=1, size=")); + } + } + + /** */ + private static int countEntries(String txt, String entry) { + String prev = txt; + + txt = txt.replace(entry, ""); + + return (prev.length() - txt.length()) / entry.length(); + } + + /** */ + @Override protected CacheConfiguration testCacheConfiguration( + String cacheName, + int partitions, + @Nullable IgnitePredicate filter + ) { + CacheConfiguration ccfg = super.testCacheConfiguration(cacheName, partitions, filter); + + if (extStorages) + ccfg.setStoragePaths(extStoragePaths); + + return ccfg; + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotCommand.java b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotCommand.java index deb5416b8e28b..e9d4e7504351d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotCommand.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotCommand.java @@ -28,7 +28,8 @@ public SnapshotCommand() { new SnapshotCancelCommand(), new SnapshotCheckCommand(), new SnapshotRestoreCommand(), - new SnapshotStatusCommand() + new SnapshotStatusCommand(), + new SnapshotListCommand() ); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListCommand.java b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListCommand.java new file mode 100644 index 0000000000000..cfc6c6c0e0386 --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListCommand.java @@ -0,0 +1,113 @@ +/* + * 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.ignite.internal.management.snapshot; + +import java.time.Instant; +import java.time.ZoneId; +import java.time.format.DateTimeFormatter; +import java.util.function.Consumer; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.SnapshotListJobResult; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.SnapshotListTaskResult; +import org.apache.ignite.internal.util.typedef.internal.U; + +/** Snapshot list command. */ +public class SnapshotListCommand extends AbstractSnapshotCommand { + /** */ + public static final String HEADER = "Snapshots found on the following nodes:"; + + /** */ + public static final String DESC = "Lists all snapshots on all online server nodes with their sizes"; + + /** */ + public static final String NO_SNAPSHOTS = "No snapshots found."; + + /** */ + private static final String PATTERN_FORMAT = "yyyy-MM-dd HH:mm:ss Z"; + + /** */ + private static final DateTimeFormatter DATE_FORMATTER = DateTimeFormatter.ofPattern(PATTERN_FORMAT).withZone(ZoneId.systemDefault()); + + /** {@inheritDoc} */ + @Override public String description() { + return DESC; + } + + /** {@inheritDoc} */ + @Override public Class argClass() { + return SnapshotListCommandArg.class; + } + + /** {@inheritDoc} */ + @Override public Class taskClass() { + return SnapshotListTask.class; + } + + /** {@inheritDoc} */ + @Override public void printResult(SnapshotListCommandArg arg, SnapshotListTaskResult res, Consumer printer) { + printer.accept(HEADER); + + for (int nodeIdx = 0; nodeIdx < res.nodesIds().length; nodeIdx++) { + // Skip line before node. + printer.accept(""); + + printer.accept("\tNode '%s' [uuid=%s]:".formatted(res.consistentIds()[nodeIdx], res.nodesIds()[nodeIdx])); + + SnapshotListJobResult nodeResult = res.nodesSnapshots()[nodeIdx]; + + if (nodeResult.snapshots().isEmpty()) { + printer.accept("\t\t" + NO_SNAPSHOTS); + + continue; + } + + nodeResult.snapshots().forEach((snpName, snpInfo) -> { + printer.accept("\t\tSnapshot '%s': totalSize=%s (%db), created='%s' (epoch=%d)".formatted( + snpName, + U.humanReadableByteCount(snpInfo.size()), + snpInfo.size(), + DATE_FORMATTER.format(Instant.ofEpochMilli(snpInfo.date())), + snpInfo.date() + )); + + SnapshotListJobResult.SnapshotInfo extStors = snpInfo.externalStorages(); + SnapshotListJobResult.SnapshotInfo incs = snpInfo.incrementals(); + + if (extStors != null) { + printer.accept("\t\t\texternal storages: cnt=%d, size=%s (%db)".formatted( + extStors.number(), + U.humanReadableByteCount(extStors.size()), + extStors.size() + )); + } + + if (incs != null) { + printer.accept("\t\t\tincremental snapshots: cnt=%d, size=%s (%db), modified='%s' (epoch=%d)".formatted( + incs.number(), + U.humanReadableByteCount(incs.size()), + incs.size(), + DATE_FORMATTER.format(Instant.ofEpochMilli(incs.date())), + incs.date() + )); + } + }); + } + + // Drop a line. + printer.accept(""); + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListCommandArg.java b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListCommandArg.java new file mode 100644 index 0000000000000..af94efaab54ff --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListCommandArg.java @@ -0,0 +1,45 @@ +/* + * 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.ignite.internal.management.snapshot; + +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.dto.IgniteDataTransferObject; +import org.apache.ignite.internal.management.api.Argument; +import org.jetbrains.annotations.Nullable; + +/** */ +public class SnapshotListCommandArg extends IgniteDataTransferObject { + /** */ + private static final long serialVersionUID = 0; + + /** */ + @Order(0) + @Argument(example = "path/to/snapshots", optional = true, description = "Path to snapshot location directory. " + + "If not specified, the default configured snapshot directory will be used") + @Nullable String src; + + /** */ + public @Nullable String src() { + return src; + } + + /** */ + public void src(@Nullable String src) { + this.src = src; + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListTask.java b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListTask.java new file mode 100644 index 0000000000000..17bbbe30c6bbe --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/management/snapshot/SnapshotListTask.java @@ -0,0 +1,441 @@ +/* + * 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.ignite.internal.management.snapshot; + +import java.io.File; +import java.io.IOException; +import java.nio.file.FileVisitResult; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.SimpleFileVisitor; +import java.nio.file.attribute.BasicFileAttributes; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.stream.Collectors; +import org.apache.ignite.IgniteException; +import org.apache.ignite.IgniteLogger; +import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.compute.ComputeJobResult; +import org.apache.ignite.internal.NodeStoppingException; +import org.apache.ignite.internal.processors.cache.persistence.filename.SnapshotFileTree; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.IgniteSnapshotManager; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.IncrementalSnapshotMetadata; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.SnapshotListJobResult; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.SnapshotListTaskResult; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.SnapshotMetadata; +import org.apache.ignite.internal.processors.rollingupgrade.feature.CoreFeatureRegistry; +import org.apache.ignite.internal.processors.task.GridInternal; +import org.apache.ignite.internal.thread.pool.IgniteThreadPoolExecutor; +import org.apache.ignite.internal.util.typedef.F; +import org.apache.ignite.internal.util.typedef.T2; +import org.apache.ignite.internal.util.typedef.X; +import org.apache.ignite.internal.visor.VisorJob; +import org.apache.ignite.internal.visor.VisorMultiNodeTask; +import org.apache.ignite.internal.visor.VisorTaskArgument; +import org.apache.ignite.resources.LoggerResource; +import org.jetbrains.annotations.Nullable; + +/** */ +@GridInternal +public class SnapshotListTask extends VisorMultiNodeTask { + /** Serial version uid. */ + private static final long serialVersionUID = 0L; + + /** {@inheritDoc} */ + @Override protected VisorJob job(SnapshotListCommandArg arg) { + return new SnapshotListJob(arg, debug); + } + + /** {@inheritDoc} */ + @Override protected Collection jobNodes(VisorTaskArgument arg) { + if (!ignite.context().rollingUpgrade().features().isActive(CoreFeatureRegistry.SNAPSHOT_LIST_FEATURE)) + throw new IgniteException("Won't search for local snapshots. The snapshot list feature isn't activated yet."); + + /** Allows {@link #map0(List, VisorTaskArgument)} to use the entire subgrid. */ + return ignite.cluster().forServers().nodes().stream().map(ClusterNode::id).collect(Collectors.toList()); + } + + /** {@inheritDoc} */ + @Override protected SnapshotListTaskResult reduce0(List nodesJobsResults) throws IgniteException { + String[] cstIds = new String[nodesJobsResults.size()]; + UUID[] nodesIds = new UUID[nodesJobsResults.size()]; + SnapshotListJobResult[] nodesResults = new SnapshotListJobResult[nodesJobsResults.size()]; + + // Sorting the results by consistent id for better reading. + nodesJobsResults = nodesJobsResults.stream() + .sorted((jr0, jr1) -> nodeConsistentId(jr0.getNode()).compareTo(nodeConsistentId(jr1.getNode()))) + .toList(); + + for (int i = 0; i < nodesJobsResults.size(); i++) { + ComputeJobResult nodeJobRes = nodesJobsResults.get(i); + + if (nodeJobRes.getException() != null) { + throw new IgniteException("Failed to execute snapshot list job on node [uuid=" + nodeJobRes.getNode().id() + ']', + nodeJobRes.getException()); + } + + assert nodeJobRes.getData() != null; + + cstIds[i] = nodeConsistentId(nodeJobRes.getNode()); + nodesIds[i] = nodeJobRes.getNode().id(); + nodesResults[i] = nodeJobRes.getData(); + } + + return new SnapshotListTaskResult(cstIds, nodesIds, nodesResults); + } + + /** */ + private String nodeConsistentId(ClusterNode n) { + UUID nodeId = n.id(); + + n = ignite.context().discovery().node(nodeId); + + if (n == null) + n = ignite.context().discovery().historicalNode(nodeId); + + return n == null ? "" : n.consistentId().toString(); + } + + /** + * Walk through a directory. Doesn't lock it or its content. Tries to find files and summarize their size. + * Tolerates and skips concurrent modification errors. + */ + public static long calculateDirectorySize(File path) throws IOException { + AtomicLong size = new AtomicLong(0); + AtomicBoolean entered = new AtomicBoolean(); + + Files.walkFileTree(path.toPath(), new SimpleFileVisitor<>() { + @Override public FileVisitResult visitFile(Path file, BasicFileAttributes attrs) { + entered.compareAndSet(false, true); + + // Use attrs instead of Files.size() for efficiency. + if (attrs.isRegularFile()) + size.addAndGet(attrs.size()); + + return FileVisitResult.CONTINUE; + } + + @Override public FileVisitResult visitFileFailed(Path file, IOException err) throws IOException { + if (!entered.get() && file.toFile().equals(path)) { + // Cannot even start snapshot size calculation - can't enter snapshot directory. + throw err; + } + + return FileVisitResult.CONTINUE; + } + }); + + return size.get(); + } + + /** */ + private static class SnapshotListJob extends SnapshotJob { + /** Serial version uid. */ + private static final long serialVersionUID = 0L; + + /** */ + @LoggerResource + private IgniteLogger log; + + /** + * @param arg Snapshot list task argument. + * @param debug Flag indicating whether debug information should be printed into node log. + */ + protected SnapshotListJob(SnapshotListCommandArg arg, boolean debug) { + super(arg, debug); + } + + /** {@inheritDoc} */ + @Override protected SnapshotListJobResult run(SnapshotListCommandArg arg) { + assert !ignite.localNode().isClient(); + + if (ignite.context().isStopping()) + throw new IgniteException("Won't search for local snapshots.", new NodeStoppingException("Node is stopping.")); + + // Read local snapshots. + List> locSnps = findLocalSnapshots(arg.src()); + + if (locSnps.isEmpty()) + return new SnapshotListJobResult(Collections.emptyMap()); + + // Optional external storages descriptions (per snapshot name). + Map extStors = new ConcurrentHashMap<>(locSnps.size(), 1.0f); + // Optional incremental parts descriptions (per snapshot name). + Map incs = new ConcurrentHashMap<>(locSnps.size(), 1.0f); + + // Incremental parts and external storages futures. + List> futs = new ArrayList<>(locSnps.size() * 2); + // Names of snapshots with read failures. + Set failedSnps = ConcurrentHashMap.newKeySet(locSnps.size() / 2); + + IgniteThreadPoolExecutor exec = ignite.context().pools().getSnapshotExecutorService(); + + // Optional descriptions of snapshot external storages and incremental parts. + for (T2 snpPair : locSnps) { + SnapshotFileTree sft = snpPair.get1(); + String snpName = sft.name(); + + // Future for optional external storages. + Future fut = exec.submit(() -> { + if (ignite.context().isStopping()) + throw new IgniteException("Won't search for local snapshots.", new NodeStoppingException("Node is stopping.")); + + if (failedSnps.contains(snpName)) + return; + + try { + SnapshotListJobResult.SnapshotInfo extDesc = externalStorages(sft); + + if (extDesc != null) + extStors.put(snpName, extDesc); + } + catch (Exception e) { + failedSnps.add(snpName); + + log.warning("Failed to read snapshot's external storages, snapshot ignored [snpName=" + snpName + ']', e); + } + }); + + futs.add(fut); + + // Future for optional incremental parts. + fut = exec.submit(() -> { + if (ignite.context().isStopping()) + throw new IgniteException("Won't search for local snapshots.", new NodeStoppingException("Node is stopping.")); + + if (failedSnps.contains(snpName)) + return; + + try { + SnapshotListJobResult.SnapshotInfo incDesc = incrementals(sft); + + if (incDesc != null) + incs.put(snpName, incDesc); + } + catch (Exception e) { + failedSnps.add(snpName); + + log.warning("Failed to read snapshot's incremental parts, snapshot ignored [snpName=" + snpName + ']', e); + } + }); + + futs.add(fut); + } + + // Wait for the snapshot optional description futures. + for (Future fut : futs) { + try { + fut.get(); + } + catch (ExecutionException e) { + if (X.hasCause(e, NodeStoppingException.class)) + throw new IgniteException("Won't search for local snapshots.", e.getCause()); + + // All the futures have internal exceptions being logged. No errors expected. + throw new IgniteException("Failed to read local nodes' snapshots.", e); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + + throw new IgniteException("Interrupted while reading local snapshots.", e); + } + } + + // Result snapshot descriptions. + Map resMap = new ConcurrentHashMap<>(locSnps.size() - failedSnps.size(), 1.0f); + + // Reduce results. + for (T2 snpPair : locSnps) { + SnapshotFileTree sft = snpPair.get1(); + String snpName = sft.name(); + + if (failedSnps.contains(snpName)) + continue; + + long size; + + try { + size = calculateDirectorySize(sft.root()); + } + catch (IOException e) { + log.warning("Failed to calculate snapshot's size, snapshot ignored [snpName=" + snpName + ']', e); + + continue; + } + + SnapshotListJobResult.SnapshotInfo snpDesc = new SnapshotListJobResult.SnapshotInfo( + size, + snpPair.get2(), + extStors.get(snpName), + incs.get(snpName) + ); + + resMap.put(snpName, snpDesc); + } + + return new SnapshotListJobResult(resMap); + } + + /** */ + private @Nullable SnapshotListJobResult.SnapshotInfo externalStorages(SnapshotFileTree sft) { + int extStoragesCnt = 0; + long extStoragesSize = 0; + + for (File extraStorage : sft.allStorages().toList()) { + if (sft.nodeStorage().equals(extraStorage)) + continue; + + try { + extStoragesSize += calculateDirectorySize(extraStorage); + } + catch (IOException e) { + log.warning("Failed to calculate snapshot's external storage size, storage ignored [extraStorage=" + + extraStorage + ']', e); + + continue; + } + + extStoragesCnt++; + } + + return extStoragesCnt == 0 ? null : new SnapshotListJobResult.SnapshotInfo(extStoragesCnt, extStoragesSize); + } + + /** @return Snapshot file tree and creation time from the snapshot metadata. */ + private List> findLocalSnapshots(@Nullable String snpPath) { + // The tree is used only to extract the snapshots root directory. The snapshot name isn't used. + File[] dirsToParse = new SnapshotFileTree(ignite.context(), "snp", snpPath).root().getParentFile().listFiles(); + + if (F.isEmpty(dirsToParse)) + return Collections.emptyList(); + + List>> futs = new ArrayList<>(dirsToParse.length); + + IgniteThreadPoolExecutor exec = ignite.context().pools().getSnapshotExecutorService(); + + for (File snpDir : dirsToParse) { + Future> snpDirFut = exec.submit(() -> { + if (ignite.context().isStopping()) + throw new IgniteException("Won't search for local snapshots.", new NodeStoppingException("Node is stopping.")); + + String snpName = snpDir.getName(); + + // Snapshot tree being used as a path, to read the metas only. + SnapshotFileTree sft = new SnapshotFileTree(ignite.context(), snpName, snpPath); + + List metas = ignite.context().cache().context().snapshotMgr().readSnapshotMetadatas(sft, false); + + if (metas.isEmpty()) + return null; + + // Get the last-created time snapshot metadata. + SnapshotMetadata snpMeta = metas.stream().max(Comparator.comparingLong(SnapshotMetadata::snapshotTime)).get(); + + // Real, meta-based snapshot file tree. Can belong to other cluster, other consistent id. + sft = new SnapshotFileTree( + ignite.configuration(), + ignite.context().pdsFolderResolver().fileTree(), + snpName, + snpPath, + snpMeta.folderName(), + snpMeta.consistentId() + ); + + return new T2<>(sft, snpMeta.snapshotTime()); + }); + + futs.add(snpDirFut); + } + + List> res = new ArrayList<>(dirsToParse.length); + + futs.forEach(f -> { + try { + T2 snpDirRes = f.get(); + + if (snpDirRes != null) + res.add(snpDirRes); + } + catch (ExecutionException e) { + if (X.hasCause(e, NodeStoppingException.class)) + throw new IgniteException("Won't search for local snapshots.", e.getCause()); + + log.warning("Failed to read snapshot, snapshot ignored [path=" + snpPath + ']', e); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + + throw new IgniteException("Interrupted while reading local snapshots.", e); + } + }); + + return res; + } + + /** @return Number, total size and last creation time of incremental snapshots. */ + private @Nullable SnapshotListJobResult.SnapshotInfo incrementals(SnapshotFileTree sft) { + File[] incs = sft.incrementsRoot().listFiles(); + + if (F.isEmpty(incs)) + return null; + + int cnt = 0; + long size = 0L; + long createTime = 0L; + + IgniteSnapshotManager snpMgr = ignite.context().cache().context().snapshotMgr(); + + int incIdx; + + for (File incDir : incs) { + if (!SnapshotFileTree.incrementSnapshotDir(incDir) || !incDir.exists()) + continue; + + try { + incIdx = Integer.parseInt(incDir.getName()); + + SnapshotFileTree.IncrementalSnapshotFileTree incTree = sft.incrementalSnapshotFileTree(incIdx); + + IncrementalSnapshotMetadata incMeta = snpMgr.readIncrementalSnapshotMetadata(incTree.meta()); + + size += calculateDirectorySize(incDir); + + createTime = Math.max(createTime, incMeta.snapshotTime()); + + cnt++; + } + catch (Exception e) { + log.warning("Failed to read incremental snapshot, skipped [dir=" + incDir + ']', e); + } + } + + return cnt == 0 ? null : new SnapshotListJobResult.SnapshotInfo(cnt, size, createTime); + } + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsSingleMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsSingleMessage.java index a9f1f42ca872b..1857aabd1f965 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsSingleMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsSingleMessage.java @@ -34,7 +34,7 @@ /** * Information about partitions of a single node. Sent in response to {@link GridDhtPartitionsSingleRequest} and during * processing partitions exchange future.
- * Has to be completelly restored after receiving from another node. + * Has to be completely restored after receiving from another node. * * @see #afterReceive() */ diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/filename/SnapshotFileTree.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/filename/SnapshotFileTree.java index f53abc94ba480..d686e6ade95a7 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/filename/SnapshotFileTree.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/filename/SnapshotFileTree.java @@ -332,6 +332,7 @@ public static String snapshotMetaFileName(String consId) { public static File root(SharedFileTree ft, String name, @Nullable String path) { assert name != null : "Snapshot name cannot be empty or null."; + // TODO: use the default snapshots directory for relative {@code path} https://issues.apache.org/jira/browse/IGNITE-29126 return path == null ? new File(ft.snapshotsRoot(), name) : new File(path, name); } @@ -373,6 +374,8 @@ public File meta() { } /** + * TODO: support the extra storages for incremental snapshots https://issues.apache.org/jira/browse/IGNITE-29128 + * * Modifies {@link #extraStorages} for this tree to reflect snapshot options. * In case {@link IgniteConfiguration#getSnapshotPath()} points to absolute directory or {@link #path} for snapshot provided * then all snapshot files must be stored inside one folder. @@ -457,7 +460,6 @@ private NodeFileTree tempFileTree(GridKernalContext ctx) { return res; } - /** {@inheritDoc} */ @Override public String toString() { return S.toString(SnapshotFileTree.class, this); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteSnapshotManager.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteSnapshotManager.java index 4d32bc1d38c09..34d6167271cf6 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteSnapshotManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteSnapshotManager.java @@ -1830,6 +1830,21 @@ public T readFromFile(File smf) throws IgniteCheckedException, IOException { * local node will be placed on the first place. */ public List readSnapshotMetadatas(SnapshotFileTree sft) { + return readSnapshotMetadatas(sft, true); + } + + /** + * Note, there can be snapshots from other nodes. + * This method will read all metadata. + * Some instances can return {@link SnapshotMetadata#folderName()} and {@link SnapshotMetadata#consistentId()} that differs from local. + * + * @param sft Snapshot file tree. + * @param failIfCantRead If {@code true}, throws exception if cannot read a metadata file. + * @return List of snapshot metadata for the given snapshot name on local node. + * If snapshot has been taken from local node the snapshot metadata for given + * local node will be placed on the first place. + */ + public List readSnapshotMetadatas(SnapshotFileTree sft, boolean failIfCantRead) { if (!(sft.root().exists() && sft.root().isDirectory())) return Collections.emptyList(); @@ -1841,8 +1856,8 @@ public List readSnapshotMetadatas(SnapshotFileTree sft) { Map metasMap = new HashMap<>(); SnapshotMetadata prev = null; - try { - for (File smf : smfs) { + for (File smf : smfs) { + try { SnapshotMetadata curr = readSnapshotMetadata(smf); if (prev != null && !prev.sameSnapshot(curr)) { @@ -1854,9 +1869,12 @@ public List readSnapshotMetadatas(SnapshotFileTree sft) { prev = curr; } - } - catch (IgniteCheckedException | IOException e) { - throw new IgniteException(e); + catch (Exception e) { + if (failIfCantRead) + throw new IgniteException("Failed to read snapshot metadata [meta=" + smf + ']', e); + else + log.warning("Failed to read snapshot metadata, snapshot skipped [meta=" + smf + ']', e); + } } SnapshotMetadata currNodeSmf = metasMap.remove(cctx.localNode().consistentId().toString()); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotCheckProcess.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotCheckProcess.java index 533f3e1a44d60..9ca188fa4c076 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotCheckProcess.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotCheckProcess.java @@ -772,7 +772,7 @@ private static final class SnapshotCheckContext { */ @Nullable private volatile List metas; - /** Map of snapshot pathes per consistent id for {@link #metas}. */ + /** Map of snapshot paths per consistent id for {@link #metas}. */ @GridToStringInclude @Nullable private Map locFileTree; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotListJobResult.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotListJobResult.java new file mode 100644 index 0000000000000..83e839e2c7179 --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotListJobResult.java @@ -0,0 +1,139 @@ +/* + * 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.ignite.internal.processors.cache.persistence.snapshot; + +import java.util.Map; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.dto.IgniteDataTransferObject; +import org.jetbrains.annotations.Nullable; + +/** Per-node result of the snapshot list command. Contains information of the snapshots found on the current node. */ +public final class SnapshotListJobResult extends IgniteDataTransferObject { + /** Serial version uid. */ + private static final long serialVersionUID = 0L; + + /** Local node's snapshot descriptions. */ + @Order(0) + Map snapshots; + + /** Default constructor for serialization purposes. */ + public SnapshotListJobResult() { + // No-op. + } + + /** */ + public SnapshotListJobResult(Map snapshots) { + this.snapshots = snapshots; + } + + /** @return The snapshots descriptions. */ + public Map snapshots() { + return snapshots; + } + + /** Holds combined snapshot data information: size, creation time, number of incremental parts or external storages. */ + public static class SnapshotInfo extends IgniteDataTransferObject { + /** Serial version uid. */ + private static final long serialVersionUID = 0L; + + /** Total size of a snapshot. Or size of its external storages or incremental parts. */ + @Order(0) + long size; + + /** Creation date of snapshot or of its incremental parts. Is {@code null} for snapshot external storages description. */ + @Order(1) + @Nullable Long date; + + /** Snapshot external storages description if exists. Is always {@code null} for not the main snapshot description. */ + @Order(2) + @Nullable SnapshotInfo extStors; + + /** Snapshot incremental parts description if exists. Is always {@code null} for not the main snapshot description. */ + @Order(3) + @Nullable SnapshotInfo incs; + + /** Number of snapshot external storages or incremental parts. Is {@code null} for the main snapshot description. */ + @Order(4) + @Nullable Integer cnt; + + /** Empty constructor for serialization purposes. */ + public SnapshotInfo() { + // No-op. + } + + /** Creates snapshot main description. */ + public SnapshotInfo( + long mainSize, + long date, + @Nullable SnapshotInfo extStors, + @Nullable SnapshotInfo incs + ) { + size = mainSize + (extStors == null ? 0L : extStors.size()); + this.date = date; + this.extStors = extStors; + this.incs = incs; + } + + /** Creates external storages description. */ + public SnapshotInfo(int cnt, long size) { + this.cnt = cnt; + this.size = size; + } + + /** Creates incremental parts description. */ + public SnapshotInfo(int cnt, long size, long date) { + this.cnt = cnt; + this.size = size; + this.date = date; + } + + /** @return Total size of a snapshot. Or size of its external storages or incremental parts. */ + public long size() { + return size; + } + + /** + * @return Creation date of a snapshot or of its incremental parts. {@code Null} for snapshot + * external storages description. + */ + public @Nullable Long date() { + return date; + } + + /** + * @return Snapshot external storages description if exists. Is always {@code null} for anything other than + * the main snapshot description. + */ + public @Nullable SnapshotInfo externalStorages() { + return extStors; + } + + /** + * @return Snapshot incremental parts description if exists. Is always {@code null} for anything other than + * the main snapshot description. + */ + public @Nullable SnapshotInfo incrementals() { + return incs; + } + + /** @return Number of snapshot external storages or incremental parts. Or {@code null} for the main snapshot data. */ + public @Nullable Integer number() { + return cnt; + } + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotListTaskResult.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotListTaskResult.java new file mode 100644 index 0000000000000..f00ba2543aacc --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotListTaskResult.java @@ -0,0 +1,68 @@ +/* + * 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.ignite.internal.processors.cache.persistence.snapshot; + +import java.util.UUID; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.dto.IgniteDataTransferObject; +import org.apache.ignite.internal.management.snapshot.SnapshotListTask; + +/** Accumulated result of {@link SnapshotListTask}. */ +public final class SnapshotListTaskResult extends IgniteDataTransferObject { + /** Serial version uid. */ + private static final long serialVersionUID = 0L; + + /** Nodes consistent ids. */ + @Order(0) + String[] cstIds; + + /** Nodes UUIDs. */ + @Order(1) + UUID[] nodesIds; + + /** Results. */ + @Order(2) + SnapshotListJobResult[] nodesSnapshots; + + /** Default constructor for serialization purposes. */ + public SnapshotListTaskResult() { + // No-op. + } + + /** */ + public SnapshotListTaskResult(String[] consistentIds, UUID[] nodesIds, SnapshotListJobResult[] nodesSnapshots) { + this.cstIds = consistentIds; + this.nodesIds = nodesIds; + this.nodesSnapshots = nodesSnapshots; + } + + /** @return Nodes consistent ids. */ + public String[] consistentIds() { + return cstIds; + } + + /** @return Nodes UUIDs. */ + public UUID[] nodesIds() { + return nodesIds; + } + + /** @return The results. */ + public SnapshotListJobResult[] nodesSnapshots() { + return nodesSnapshots; + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotMetadataVerificationTask.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotMetadataVerificationTask.java index 643838ac3c445..4a8fea2b43b73 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotMetadataVerificationTask.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/SnapshotMetadataVerificationTask.java @@ -35,7 +35,10 @@ import org.apache.ignite.internal.processors.cache.persistence.wal.reader.IgniteWalIteratorFactory; import org.apache.ignite.internal.util.typedef.F; -/** Snapshot task to verify snapshot metadata on the baseline nodes for given snapshot name. */ +/** + * Snapshot task to verify snapshot metadata on the baseline nodes for given snapshot name. + * TODO : Revise in https://issues.apache.org/jira/browse/IGNITE-29062 + */ public class SnapshotMetadataVerificationTask implements Supplier> { /** */ private final IgniteEx ignite; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 575300772a905..bf7f6f34c2566 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -103,7 +103,7 @@ *

*

Whole updates history until some point in the past is stored along with the data, so when an outdated node * connects to the cluster it will receive all the missing data and apply it locally. Listeners will also be invoked - * after such updates. If there's not enough history stored or joining node is clear then it'll receive shapshot of + * after such updates. If there's not enough history stored or joining node is clear then it'll receive snapshot of * distributed metastorage (usually called {@code fullData} in code) so there won't be inconsistencies. *

* diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java index 92b7eeb1a8c5b..b6d79ee0f6d84 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/rollingupgrade/feature/CoreFeatureRegistry.java @@ -93,4 +93,7 @@ public class CoreFeatureRegistry { /** */ public static final IgniteFeature ROLLING_UPGRADE_FEATURE = new IgniteCoreFeature(0); + + /** */ + public static final IgniteFeature SNAPSHOT_LIST_FEATURE = new IgniteCoreFeature(1); } diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/AbstractSnapshotSelfTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/AbstractSnapshotSelfTest.java index 8857a4ddd656e..8aca648f0e019 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/AbstractSnapshotSelfTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/AbstractSnapshotSelfTest.java @@ -701,8 +701,8 @@ protected void checkSnapshot(String snpName, String snpPath) { * @param ignite Ignite instance. * @return Snapshot manager related to given ignite instance. */ - public static IgniteSnapshotManager snp(IgniteEx ignite) { - return ignite.context().cache().context().snapshotMgr(); + public static IgniteSnapshotManager snp(Ignite ignite) { + return ((IgniteEx)ignite).context().cache().context().snapshotMgr(); } /** @@ -745,7 +745,7 @@ protected static List setBlockingSnapshotExecutor(List execs = new ArrayList<>(); for (Ignite grid : grids) { - IgniteSnapshotManager mgr = snp((IgniteEx)grid); + IgniteSnapshotManager mgr = snp(grid); Function old = mgr.localSnapshotSenderFactory(); BlockingExecutor block = new BlockingExecutor(mgr.snapshotExecutorService()); diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteClusterSnapshotListRollingUpgradeTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteClusterSnapshotListRollingUpgradeTest.java new file mode 100644 index 0000000000000..e335c8989711a --- /dev/null +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteClusterSnapshotListRollingUpgradeTest.java @@ -0,0 +1,245 @@ +/* + * 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.ignite.internal.processors.cache.persistence.snapshot; + +import java.io.File; +import java.nio.file.DirectoryStream; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.util.Arrays; +import java.util.Collection; +import java.util.HashSet; +import java.util.UUID; +import org.apache.ignite.IgniteDataStreamer; +import org.apache.ignite.IgniteException; +import org.apache.ignite.cache.CacheAtomicityMode; +import org.apache.ignite.cache.CacheMode; +import org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction; +import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.configuration.CacheConfiguration; +import org.apache.ignite.configuration.DataRegionConfiguration; +import org.apache.ignite.configuration.DataStorageConfiguration; +import org.apache.ignite.configuration.IgniteConfiguration; +import org.apache.ignite.internal.IgniteEx; +import org.apache.ignite.internal.IgniteInternalFuture; +import org.apache.ignite.internal.management.snapshot.SnapshotListCommandArg; +import org.apache.ignite.internal.management.snapshot.SnapshotListTask; +import org.apache.ignite.internal.processors.rollingupgrade.AbstractRollingUpgradeTest; +import org.apache.ignite.internal.util.distributed.SingleNodeMessage; +import org.apache.ignite.internal.util.typedef.internal.U; +import org.apache.ignite.internal.visor.VisorTaskArgument; +import org.apache.ignite.testframework.GridTestUtils; +import org.junit.Test; + +import static java.nio.file.Files.newDirectoryStream; +import static org.apache.ignite.internal.TestRecordingCommunicationSpi.spi; +import static org.apache.ignite.internal.util.distributed.DistributedProcess.DistributedProcessType.RU_PREPARE_VERSION_FINALIZATION; +import static org.apache.ignite.testframework.GridTestUtils.assertThrowsAnyCause; +import static org.apache.ignite.testframework.GridTestUtils.waitForCondition; + +/** */ +public class IgniteClusterSnapshotListRollingUpgradeTest extends AbstractRollingUpgradeTest { + /** */ + private static final int ALL_GRIDS = 4; + + /** */ + private static final int CLIENTS = 1; + + /** */ + private static final String SNP_NAME = "testSnapshot"; + + /** {@inheritDoc} */ + @Override protected void afterTest() throws Exception { + super.afterTest(); + + cleanPersistenceDir(); + + // Clean all: also separated snapshot working directories. + try (DirectoryStream files = newDirectoryStream(Paths.get(U.defaultWorkDirectory()))) { + for (Path path : files) + U.delete(path); + } + } + + /** {@inheritDoc} */ + @Override protected IgniteConfiguration getConfiguration(String igniteInstanceName, String ver) throws Exception { + IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName, ver); + + cfg.setDataStorageConfiguration( + new DataStorageConfiguration() + .setDefaultDataRegionConfiguration( + new DataRegionConfiguration() + .setPersistenceEnabled(true) + .setMaxSize(DataStorageConfiguration.DFLT_DATA_REGION_INITIAL_SIZE) + ) + ); + + cfg.setWorkDirectory(new File(U.defaultWorkDirectory(), igniteInstanceName).getAbsolutePath()); + + return cfg; + } + + /** */ + @Test + public void testParallelRollingUpgradeInProgress() throws Exception { + for (int i = 0; i < ALL_GRIDS; i++) + startGrid(i, "2.19.4", i >= ALL_GRIDS - CLIENTS); + + grid(0).cluster().active(true); + + int testNodeIx = ALL_GRIDS - CLIENTS - 1; + + createCacheAndSnapshot(testNodeIx); + + ru(grid(testNodeIx)).enableVersionUpgrade(); + + for (int i = 0; i < ALL_GRIDS; i++) { + assertTrue(ru(grid(i)).isVersionUpgradeEnabled()); + + upgradeNodeVersion(i, "2.19.4", "2.19.5"); + } + + spi(grid(testNodeIx)).blockMessages((node, msg) -> msg instanceof SingleNodeMessage snm && + snm.type() == RU_PREPARE_VERSION_FINALIZATION.ordinal()); + + IgniteInternalFuture finalizeFut = GridTestUtils.runAsync(() -> ru(testNodeIx).finalizeClusterVersion()); + + assertTrue(spi(grid(testNodeIx)).waitForBlocked(1, getTestTimeout())); + + ensureSnapshotListFailed(); + + spi(grid(testNodeIx)).stopBlock(); + + assertFalse(spi(grid(testNodeIx)).hasBlockedMessages()); + + finalizeFut.get(getTestTimeout()); + + for (int i = 0; i < ALL_GRIDS; i++) { + int i0 = i; + + assertTrue(waitForCondition(() -> !ru(grid(i0)).isVersionUpgradeEnabled(), getTestTimeout())); + } + + ensureSnapshotListSucceeds(); + } + + /** */ + @Test + public void testSnapshotListFeature() throws Exception { + for (int i = 0; i < ALL_GRIDS; i++) + startGrid(i, "2.19.4", i >= ALL_GRIDS - CLIENTS); + + grid(0).cluster().active(true); + + createCacheAndSnapshot(1); + + ensureSnapshotListFailed(); + + ru(grid(0)).enableVersionUpgrade(); + + for (int i = 0; i < ALL_GRIDS; i++) { + assertTrue(ru(grid(i)).isVersionUpgradeEnabled()); + + upgradeNodeVersion(i, "2.19.4", "2.19.5"); + + ensureSnapshotListFailed(); + } + + ru(grid(1)).finalizeClusterVersion(); + + for (int i = 0; i < ALL_GRIDS; i++) { + int i0 = i; + + assertTrue(waitForCondition(() -> !ru(grid(i0)).isVersionUpgradeEnabled(), getTestTimeout())); + } + + ensureSnapshotListSucceeds(); + } + + /** */ + private void ensureSnapshotListSucceeds() throws Exception { + for (int g = 0; g < ALL_GRIDS; g++) { + SnapshotListTaskResult listOpRes = listSnapshots(grid(g)); + + assertEquals(3, listOpRes.nodesSnapshots().length); + + // Tests client exclusion. + if (CLIENTS > 0) { + UUID cliId = grid(ALL_GRIDS - CLIENTS).cluster().localNode().id(); + + assertFalse(new HashSet<>(Arrays.asList(listOpRes.nodesIds())).contains(cliId)); + } + + for (SnapshotListJobResult nodeRes : listOpRes.nodesSnapshots()) + assertEquals(1, nodeRes.snapshots().size()); + } + } + + /** */ + private void createCacheAndSnapshot(int gridIdx) { + int partsCnt = 4; + int keysCnt = partsCnt * 5; + + grid(gridIdx).createCache(new CacheConfiguration<>(DEFAULT_CACHE_NAME) + .setCacheMode(CacheMode.REPLICATED) + .setAffinity(new RendezvousAffinityFunction().setPartitions(partsCnt)) + .setAtomicityMode(CacheAtomicityMode.ATOMIC)); + + try (IgniteDataStreamer ds = grid(gridIdx).dataStreamer(DEFAULT_CACHE_NAME)) { + for (int i = 0; i < keysCnt; i++) + ds.addData(i, i); + } + + createSnapshot(gridIdx); + } + + /** */ + private void createSnapshot(int gridIdx) { + snp(gridIdx).createSnapshot(SNP_NAME).get(getTestTimeout()); + } + + /** */ + private void ensureSnapshotListFailed() { + String err = "The snapshot list feature isn't activated yet"; + + for (int i = 0; i < ALL_GRIDS; i++) { + int i0 = i; + + assertThrowsAnyCause( + null, + () -> listSnapshots(grid(i0)), + IgniteException.class, + err + ); + } + } + + /** */ + private IgniteSnapshotManager snp(int gridIdx) { + return grid(gridIdx).context().cache().context().snapshotMgr(); + } + + /** */ + private static SnapshotListTaskResult listSnapshots(IgniteEx grid) throws Exception { + SnapshotListCommandArg arg = new SnapshotListCommandArg(); + + Collection nodes = grid.cluster().forServers().nodes().stream().map(ClusterNode::id).toList(); + + return grid.compute().execute(SnapshotListTask.class, new VisorTaskArgument<>(nodes, arg, false)).result(); + } +} diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteClusterSnapshotListTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteClusterSnapshotListTest.java new file mode 100644 index 0000000000000..8fc18e641cd96 --- /dev/null +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/IgniteClusterSnapshotListTest.java @@ -0,0 +1,1004 @@ +/* + * 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.ignite.internal.processors.cache.persistence.snapshot; + +import java.io.File; +import java.io.IOException; +import java.io.RandomAccessFile; +import java.io.Serializable; +import java.nio.file.DirectoryStream; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.nio.file.attribute.PosixFilePermission; +import java.nio.file.attribute.PosixFilePermissions; +import java.util.Collection; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Stream; +import org.apache.ignite.IgniteDataStreamer; +import org.apache.ignite.IgniteException; +import org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction; +import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.configuration.CacheConfiguration; +import org.apache.ignite.configuration.IgniteConfiguration; +import org.apache.ignite.internal.IgniteEx; +import org.apache.ignite.internal.IgniteInternalFuture; +import org.apache.ignite.internal.NodeStoppingException; +import org.apache.ignite.internal.management.snapshot.SnapshotListCommandArg; +import org.apache.ignite.internal.management.snapshot.SnapshotListTask; +import org.apache.ignite.internal.processors.cache.persistence.filename.SnapshotFileTree; +import org.apache.ignite.internal.util.typedef.F; +import org.apache.ignite.internal.util.typedef.internal.U; +import org.apache.ignite.internal.visor.VisorTaskArgument; +import org.apache.ignite.lang.IgniteFuture; +import org.apache.ignite.plugin.AbstractTestPluginProvider; +import org.apache.ignite.plugin.PluginConfiguration; +import org.apache.ignite.plugin.PluginContext; +import org.apache.ignite.plugin.PluginProvider; +import org.jetbrains.annotations.Nullable; +import org.junit.Ignore; +import org.junit.Test; +import org.junit.runners.Parameterized; + +import static java.nio.file.Files.newDirectoryStream; +import static org.apache.ignite.configuration.IgniteConfiguration.DFLT_SNAPSHOT_THREAD_POOL_SIZE; +import static org.apache.ignite.testframework.GridTestUtils.assertThrowsAnyCause; +import static org.apache.ignite.testframework.GridTestUtils.cartesianProduct; +import static org.apache.ignite.testframework.GridTestUtils.runAsync; +import static org.apache.ignite.testframework.GridTestUtils.waitForCondition; +import static org.junit.Assume.assumeFalse; +import static org.junit.Assume.assumeTrue; + +/** Cluster-wide snapshot list procedure tests. */ +public class IgniteClusterSnapshotListTest extends AbstractSnapshotSelfTest { + /** Number of cache keys to pre-create at node start. */ + private static final int CACHE_KEYS_RANGE = 15; + + /** Number of partitions within a snapshot cache group. */ + private static final int CACHE_PARTITIONS_COUNT = 4; + + /** */ + private static boolean posixPermissions; + + /** Size of the snapshot utility thread pool. */ + @Parameterized.Parameter(2) + public int snpThrdPoolSz; + + /** */ + private PluginProvider pluginProvider; + + /** Flag to spread the test cache data over the external storages. */ + private boolean extStorages; + + /** Resolved external storages paths. {@code null} if {@link #extStorages} is {@code false}. */ + private @Nullable String[] extStoragePaths; + + /** Parameters. */ + @Parameterized.Parameters(name = "encryption={0}, onlyPrimary={1}, snpThrdPoolSz={2}") + public static Collection params() { + return cartesianProduct( + encryptionParameters(), // Encryption + F.asList(false, true), // Only primary + F.asList(DFLT_SNAPSHOT_THREAD_POOL_SIZE, 1) // Snapshots thread pool size + ); + } + + /** {@inheritDoc} */ + @Override protected IgniteConfiguration getConfiguration(String igniteInstanceName) throws Exception { + IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName) + .setWorkDirectory(new File(U.defaultWorkDirectory(), igniteInstanceName).getAbsolutePath()); + + if (pluginProvider != null) + cfg.setPluginProviders(pluginProvider); + + if (extStorages) { + // External storage paths must be identical on all the nodes (a cache has the single storage paths + // setting). Thus, a shared work directory is used, like in GridCommandHandlerListSnapshotTest. + cfg.setWorkDirectory(U.defaultWorkDirectory()); + + cfg.getDataStorageConfiguration().setExtraStoragePaths( + U.defaultWorkDirectory() + File.separator, + U.defaultWorkDirectory() + File.separator + "extStorage" + ); + + extStoragePaths = cfg.getDataStorageConfiguration().getExtraStoragePaths(); + + cfg.getDataStorageConfiguration().setExtraSnapshotPaths("", "extStorage"); + } + + cfg.setSnapshotThreadPoolSize(snpThrdPoolSz); + + return cfg; + } + + /** {@inheritDoc} */ + @Override protected void cleanPersistenceDir() throws Exception { + super.cleanPersistenceDir(); + + // Also cleans separated snapshot working directories and custom snapshot paths. + try (DirectoryStream files = newDirectoryStream(Paths.get(U.defaultWorkDirectory()))) { + for (Path path : files) + U.delete(path); + } + } + + /** {@inheritDoc} */ + @Override protected void beforeTestsStarted() throws Exception { + super.beforeTestsStarted(); + + File workDir = new File(U.defaultWorkDirectory()); + + assertTrue(workDir.exists()); + + try { + Files.getPosixFilePermissions(workDir.toPath()); + + posixPermissions = true; + } + catch (UnsupportedOperationException ignored) { + // No-op. + } + } + + /** {@inheritDoc} */ + @Override protected CacheConfiguration txCacheConfig(CacheConfiguration ccfg) { + // Speeds up the tests. + ccfg = super.txCacheConfig(ccfg) + .setAffinity(new RendezvousAffinityFunction(false, CACHE_PARTITIONS_COUNT)) + .setBackups(1); + + if (extStorages) { + assert !F.isEmpty(extStoragePaths); + + ccfg.setStoragePaths(extStoragePaths); + } + + return ccfg; + } + + /** */ + @Test + public void testConcurrentCreation() throws Exception { + // The test uses the thread blocking and conditional waitings. Won't proceed with 1 thread. + assumeTrue(snpThrdPoolSz > 1); + + int grids = 3; + int testNodeIdx = 1; + int testNodeOrder = testNodeIdx + 1; + + CountDownLatch beginSnpCreation = new CountDownLatch(grids); + CountDownLatch proceedSnpCreation = new CountDownLatch(1); + + // Delays snapshot creation after its metadata is written. + pluginProvider = new AbstractTestPluginProvider() { + @Override public String name() { + return "TestSnpMgrProvider"; + } + + @Override public T createComponent(PluginContext ctx, Class cls) { + if (IgniteSnapshotManager.class.isAssignableFrom(cls)) { + return (T)new IgniteSnapshotManager(((IgniteEx)ctx.grid()).context()) { + @Override public void storeSnapshotMeta(M meta, File smf) { + super.storeSnapshotMeta(meta, smf); + + beginSnpCreation.countDown(); + + if (((IgniteEx)ctx.grid()).localNode().order() == testNodeOrder) { + try { + assertTrue(proceedSnpCreation.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + } + catch (InterruptedException e) { + throw new IllegalStateException(e); + } + } + } + }; + } + + return super.createComponent(ctx, cls); + } + }; + + startGridsWithCache(grids, txCacheConfig(defaultCacheConfiguration()), CACHE_KEYS_RANGE); + + IgniteFuture createSnpFut = snp(grid(0)).createSnapshot(SNAPSHOT_NAME, null, false, onlyPrimary); + + assertTrue(beginSnpCreation.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + + SnapshotListTaskResult lstOpRes = listSnapshots(grid(2)); + + int snpsCnt = Stream.of(lstOpRes.nodesSnapshots()).mapToInt(jr -> jr.snapshots().size()).sum(); + + assertEquals(grids, snpsCnt); + + SnapshotFileTree testSft = new SnapshotFileTree(grid(testNodeIdx).context(), SNAPSHOT_NAME, null); + + assertTrue(SnapshotListTask.calculateDirectorySize(testSft.root()) > 0L); + + proceedSnpCreation.countDown(); + + createSnpFut.get(getTestTimeout()); + } + + /** */ + @Test + public void testCompletelyDeletedAfterMetaRead() throws Exception { + doTestDeletionAfterMetaRead(true); + } + + /** */ + @Test + public void testPartlyDeletedAfterMetaRead() throws Exception { + doTestDeletionAfterMetaRead(false); + } + + /** */ + private void doTestDeletionAfterMetaRead(boolean completeDeletion) throws Exception { + int grids = 3; + int testGridIdx = 1; + + CountDownLatch metaReadBeginLatch = new CountDownLatch(grids); + CountDownLatch metaReadProceedLatch = new CountDownLatch(1); + + pluginProvider = new AbstractTestPluginProvider() { + @Override public String name() { + return "TestSnpMgrProvider"; + } + + @Override public T createComponent(PluginContext ctx, Class cls) { + if (IgniteSnapshotManager.class.isAssignableFrom(cls)) { + return (T)new IgniteSnapshotManager(((IgniteEx)ctx.grid()).context()) { + @Override public List readSnapshotMetadatas(SnapshotFileTree sft, boolean failIfCantRead) { + metaReadBeginLatch.countDown(); + + if (cctx.localNode().order() == testGridIdx + 1) { + try { + assertTrue(metaReadProceedLatch.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + } + catch (InterruptedException e) { + throw new IllegalStateException(e); + } + } + + return super.readSnapshotMetadatas(sft, failIfCantRead); + } + }; + } + + return super.createComponent(ctx, cls); + } + }; + + startGridsWithSnapshot(grids, CACHE_KEYS_RANGE, true, false); + + SnapshotFileTree testSnpFt = new SnapshotFileTree(grid(testGridIdx).context(), SNAPSHOT_NAME, null); + + long testSnpSz = SnapshotListTask.calculateDirectorySize(testSnpFt.root()); + + assertTrue(testSnpSz > 0); + + IgniteInternalFuture lstOpFut = runAsync(() -> listSnapshots(grid(0))); + + assertTrue(metaReadBeginLatch.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + + if (completeDeletion) + assertTrue(testSnpFt.root().exists() && U.delete(testSnpFt.root()) && !testSnpFt.root().exists()); + else + assertTrue(testSnpFt.nodeStorage().exists() && U.delete(testSnpFt.nodeStorage()) && !testSnpFt.nodeStorage().exists()); + + metaReadProceedLatch.countDown(); + + SnapshotListTaskResult lstOpRes = lstOpFut.get(); + + if (completeDeletion) { + int snpsCnt = Stream.of(lstOpRes.nodesSnapshots()).mapToInt(jr -> jr.snapshots().size()).sum(); + + assertEquals(grids - 1, snpsCnt); + } + else { + int snpsCnt = 0; + boolean victimNodeFound = false; + + for (int i = 0; i < lstOpRes.nodesIds().length; i++) { + UUID nid = lstOpRes.nodesIds()[i]; + + Map nodeSnps = lstOpRes.nodesSnapshots()[i].snapshots(); + + snpsCnt += nodeSnps.size(); + + if (nid.equals(grid(testGridIdx).localNode().id())) { + victimNodeFound = true; + + assertEquals(1, nodeSnps.size()); + + long curSnpSz = SnapshotListTask.calculateDirectorySize(testSnpFt.root()); + + assertTrue(curSnpSz < testSnpSz); + } + } + + assertTrue(victimNodeFound); + assertEquals(grids, snpsCnt); + } + } + + /** + * Ensures that a file can be read without locking or other issues while being concurrently written on current OS and + * file system. This behavior is important for the case when snapshots are being read while creation. + */ + @Test + public void testReadWhileCreating() throws Exception { + // Doesn't matter here. speeds up the tests. + assumeFalse(encryption || onlyPrimary); + + assertTrue(new File(U.defaultWorkDirectory()).exists()); + + File testF = new File(U.defaultWorkDirectory(), "test.out"); + + assumeFalse(testF.exists()); + + CountDownLatch proceedWriteLatch = new CountDownLatch(1); + AtomicBoolean readFlag = new AtomicBoolean(); + + try { + try (RandomAccessFile raf = new RandomAccessFile(testF, "rw");) { + raf.write(1); + raf.write(2); + + Thread t = new Thread(() -> { + try (RandomAccessFile raf0 = new RandomAccessFile(testF, "r")) { + assertEquals(1, raf0.read()); + assertEquals(2, raf0.read()); + + readFlag.set(true); + + proceedWriteLatch.countDown(); + } + catch (Throwable e) { + log.error("Failed to concurrently read the test file.", e); + } + }); + + t.setDaemon(true); + t.start(); + + assertTrue(proceedWriteLatch.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + + raf.write(3); + } + + try (RandomAccessFile raf = new RandomAccessFile(testF, "r");) { + assertEquals(1, raf.read()); + assertEquals(2, raf.read()); + assertEquals(3, raf.read()); + } + } + finally { + assertTrue(!testF.exists() || testF.delete()); + } + + assertTrue(readFlag.get()); + } + + /** */ + @Test + public void testMissingSnapshotPath() throws Exception { + doTestWrongSnapshotPath(true); + } + + /** */ + @Test + public void testEmptySnapshotPath() throws Exception { + doTestWrongSnapshotPath(false); + } + + /** */ + private void doTestWrongSnapshotPath(boolean missing) throws Exception { + // Doesn't matter here, speeds up the tests. + assumeFalse(encryption || onlyPrimary || snpThrdPoolSz < 2); + + doTestSnapshotPath(null); + + File path = new File(U.defaultWorkDirectory(), "not_snapshots"); + + if (!missing) + assertTrue(new File(path, SNAPSHOT_NAME).mkdirs()); + + SnapshotListJobResult[] lstOpRes = listSnapshots(grid(0), path.getAbsolutePath()).nodesSnapshots(); + + int cnt = Stream.of(lstOpRes).mapToInt(nodeRes -> nodeRes.snapshots().size()).sum(); + + assertEquals(0, cnt); + } + + /** */ + @Test + public void testDefaultSnapshotPath() throws Exception { + doTestSnapshotPath(null); + } + + /** */ + @Test + @Ignore("https://issues.apache.org/jira/browse/IGNITE-29126") + public void testRelativeSnapshotPath() throws Exception { + doTestSnapshotPath("ex_snapshots"); + } + + /** */ + @Test + public void testAbsoluteSnapshotPath() throws Exception { + doTestSnapshotPath(new File(U.defaultWorkDirectory(), "ex_snapshots").getAbsolutePath()); + } + + /** */ + private void doTestSnapshotPath(@Nullable String path) throws Exception { + int grids = 3; + + startGridsWithCache(grids, txCacheConfig(defaultCacheConfiguration()), CACHE_KEYS_RANGE); + + snp(grid(0)).createSnapshot(SNAPSHOT_NAME, path, false, onlyPrimary).get(getTestTimeout()); + + SnapshotListJobResult[] lstOpRes = listSnapshots(grid(0), path).nodesSnapshots(); + + int cnt = Stream.of(lstOpRes).mapToInt(nodeRes -> nodeRes.snapshots().size()).sum(); + + assertEquals(grids, cnt); + } + + /** */ + @Test + public void testMissingMeta() throws Exception { + doTestWithWrongMeta(false); + } + + /** */ + @Test + public void testCorruptedMeta() throws Exception { + doTestWithWrongMeta(true); + } + + /** + * Tests snapshot lists when one snapshot metadata cannot be read. + * + * @param corruptFile If {@code true}, corrupts metadata. Otherwise, deletes metadata. + */ + private void doTestWithWrongMeta(boolean corruptFile) throws Exception { + // Speeds up the tests. + assumeFalse(encryption); + + int grids = 3; + int testGridIdx = 1; + + startGridsWithSnapshot(grids, CACHE_KEYS_RANGE, true, false); + + SnapshotFileTree testSnpFt = new SnapshotFileTree(grid(testGridIdx).context(), SNAPSHOT_NAME, null); + + assertTrue(testSnpFt.meta().exists()); + + if (corruptFile) { + try (RandomAccessFile raf = new RandomAccessFile(testSnpFt.meta(), "rw")) { + raf.write(UUID.randomUUID().toString().getBytes()); + } + } + else + assertTrue(testSnpFt.meta().delete() && !testSnpFt.meta().exists()); + + SnapshotListTaskResult res = listSnapshots(grid(0)); + + int foundSnpsCnt = 0; + boolean victimNodeFound = false; + + for (int i = 0; i < res.nodesIds().length; i++) { + UUID nid = res.nodesIds()[i]; + + foundSnpsCnt += res.nodesSnapshots()[i].snapshots().size(); + + if (nid.equals(grid(testGridIdx).localNode().id())) { + victimNodeFound = true; + + assertTrue(res.nodesSnapshots()[i].snapshots().isEmpty()); + } + } + + assertTrue(victimNodeFound); + assertEquals(grids - 1, foundSnpsCnt); + } + + /** */ + @Test + public void testMissingIncrementalMeta() throws Exception { + doTestWithWrongIncrementalMeta(false); + } + + /** */ + @Test + public void testCorruptedIncrementalMeta() throws Exception { + doTestWithWrongIncrementalMeta(true); + } + + /** + * Tests snapshot list when incremental snapshot metadata cannot be read. + * The main snapshot should still be listed, but without incremental info on the affected node. + * + * @param corruptFile If {@code true}, corrupts metadata. Otherwise, deletes metadata. + */ + private void doTestWithWrongIncrementalMeta(boolean corruptFile) throws Exception { + // Incremental snapshots do not support the only-primary mode or encryption. + assumeFalse(onlyPrimary || encryption); + + int grids = 3; + int testGridIdx = 1; + int incsCnt = 3; + + IgniteEx ig = startGridsWithCache(grids, txCacheConfig(defaultCacheConfiguration()), CACHE_KEYS_RANGE); + + snp(ig).createSnapshot(SNAPSHOT_NAME, null, false, onlyPrimary).get(getTestTimeout()); + + for (int i = 0; i < incsCnt; ++i) { + try (IgniteDataStreamer ds = grid(0).dataStreamer(DEFAULT_CACHE_NAME)) { + for (int kv = (i + 1) * CACHE_KEYS_RANGE; kv < (i + 1) * CACHE_KEYS_RANGE * 2; ++kv) + ds.addData(kv, kv); + } + + snp(ig).createSnapshot(SNAPSHOT_NAME, null, true, onlyPrimary).get(getTestTimeout()); + } + + SnapshotFileTree testSnpFt = new SnapshotFileTree(grid(testGridIdx).context(), SNAPSHOT_NAME, null); + SnapshotFileTree.IncrementalSnapshotFileTree incFt = testSnpFt.incrementalSnapshotFileTree(1); + File incMeta = incFt.meta(); + + assertTrue(incMeta.exists()); + + if (corruptFile) { + try (RandomAccessFile raf = new RandomAccessFile(incMeta, "rw")) { + raf.write(UUID.randomUUID().toString().getBytes()); + } + } + else + assertTrue(incMeta.delete() && !incMeta.exists()); + + SnapshotListTaskResult res = listSnapshots(grid(0)); + + for (int i = 0; i < res.nodesIds().length; i++) { + UUID nid = res.nodesIds()[i]; + + Map snps = res.nodesSnapshots()[i].snapshots(); + + assertEquals(1, snps.size()); + + SnapshotListJobResult.SnapshotInfo info = snps.get(SNAPSHOT_NAME); + + assertNotNull(info); + assertNotNull(info.incrementals()); + + if (nid.equals(grid(testGridIdx).localNode().id())) + assertEquals(incsCnt - 1, info.incrementals().number().intValue()); + else + assertEquals(incsCnt, info.incrementals().number().intValue()); + } + } + + /** + * Test snapshot list operation when a node can't read some snapshot part due to insufficient permissions. + * I.e. a test node is able to read snapshot meta but can't read some of the snapshot's data. + */ + @Test + public void testDeniedPermissions() throws Exception { + assumeTrue(posixPermissions); + // We rely on sizes here. Better to avoid empty data nodes not to become flaky. + assumeFalse(onlyPrimary); + + int grids = 3; + int testGridIdx = 1; + + // Permissions to restore. + Map> oldPerms = new ConcurrentHashMap<>(); + // The 'change permissions' flag. + AtomicBoolean changePermissions = new AtomicBoolean(true); + + // Deny reading on a couple of incremental snapshot metadata files on one node. + pluginProvider = new AbstractTestPluginProvider() { + @Override public String name() { + return "TestSnpMgrProvider"; + } + + @Override public T createComponent(PluginContext ctx, Class cls) { + if (IgniteSnapshotManager.class.isAssignableFrom(cls)) { + return (T)new IgniteSnapshotManager(((IgniteEx)ctx.grid()).context()) { + @Override public List readSnapshotMetadatas(SnapshotFileTree sft, + boolean failIfCantRead) { + if (changePermissions.get() && ctx.localNode().order() == testGridIdx + 1) { + File victimDir = sft.nodeStorage(); + + assertTrue(victimDir.exists()); + assertTrue(victimDir.isDirectory()); + + // Ensure that blocked snapshot part has some size. We use sizes to compare later. + try { + assertTrue(SnapshotListTask.calculateDirectorySize(victimDir) > 0L); + } + catch (IOException e) { + throw new IllegalStateException(e); + } + + Path victimDirPath = victimDir.toPath(); + + try { + Set perms = Files.getPosixFilePermissions(victimDirPath); + + assertFalse(perms.isEmpty()); + + // Deny reading. + Files.setPosixFilePermissions(victimDirPath, PosixFilePermissions.fromString("---------")); + + // Save actual permissions to restore. + oldPerms.put(victimDirPath, perms); + } + catch (Exception e) { + throw new IgniteException("Unable to set the test posix permissions.", e); + } + } + + return super.readSnapshotMetadatas(sft, failIfCantRead); + } + }; + } + + return super.createComponent(ctx, cls); + } + }; + + startGridsWithSnapshot(grids, CACHE_KEYS_RANGE, true, true); + + UUID testNodeId = grid(testGridIdx).localNode().id(); + + // Snapshot size with restricted permissions. + long testSize0 = 0L; + // Snapshot size with normal permissions. + long testSize1 = 0L; + + SnapshotListTaskResult snpLstOpRes; + + try { + // First run. + snpLstOpRes = listSnapshots(grid(0)); + + assertFalse(oldPerms.isEmpty()); + + for (int i = 0; i < snpLstOpRes.nodesIds().length; i++) { + UUID nodeId = snpLstOpRes.nodesIds()[i]; + + // Store size of the partly read snapshot. + if (nodeId.equals(testNodeId)) { + Map snps = snpLstOpRes.nodesSnapshots()[i].snapshots(); + + assertEquals(1, snps.size()); + + testSize0 = snps.get(SNAPSHOT_NAME).size(); + } + } + + // Ensure that we've found and read snapshot on the test node. + assertTrue(testSize0 > 0L); + } + finally { + // Restore the permissions in any case. + oldPerms.forEach((path, perms) -> { + try { + Files.setPosixFilePermissions(path, perms); + } + catch (IOException e) { + throw new IllegalStateException(e); + } + }); + } + + // Relaunch the operation. + changePermissions.set(false); + + snpLstOpRes = listSnapshots(grid(0)); + + for (int i = 0; i < snpLstOpRes.nodesIds().length; i++) { + UUID nodeId = snpLstOpRes.nodesIds()[i]; + + // Store size of the partly read snapshot. + if (nodeId.equals(testNodeId)) { + Map snps = snpLstOpRes.nodesSnapshots()[i].snapshots(); + + assertEquals(1, snps.size()); + + testSize1 = snps.get(SNAPSHOT_NAME).size(); + } + } + + // Ensure that the calculated anew size is bigger than in the previous run. + assertTrue(testSize1 > 0L); + assertTrue(testSize1 > testSize0); + } + + /** */ + @Test + public void testSnapshotListsDates() throws Exception { + // Incremental snapshots do not support the only-primary mode or encryption. + assumeFalse(onlyPrimary || encryption); + + int grids = 3; + + IgniteEx ig = startGridsWithCache(grids, txCacheConfig(defaultCacheConfiguration()), CACHE_KEYS_RANGE); + + long time0 = U.currentTimeMillis(); + + // Wait for a while, spend some time. + U.sleep(300L); + + snp(ig).createSnapshot(SNAPSHOT_NAME, null, false, onlyPrimary).get(getTestTimeout()); + + SnapshotListTaskResult res = listSnapshots(grid(0)); + + for (int i = 0; i < res.nodesIds().length; i++) { + Map snps = res.nodesSnapshots()[i].snapshots(); + + assertEquals(1, snps.size()); + + SnapshotListJobResult.SnapshotInfo info = snps.get(SNAPSHOT_NAME); + + assertTrue(info.date() > time0); + } + + // Wait for a while, spend some time. + U.sleep(300L); + + long time1 = U.currentTimeMillis(); + + try (IgniteDataStreamer ds = ig.dataStreamer(DEFAULT_CACHE_NAME)) { + for (int kv = CACHE_KEYS_RANGE; kv < CACHE_KEYS_RANGE * 2; kv++) + ds.addData(kv, kv); + } + + snp(ig).createSnapshot(SNAPSHOT_NAME, null, true, onlyPrimary).get(getTestTimeout()); + + SnapshotListTaskResult res2 = listSnapshots(grid(0)); + + for (int i = 0; i < res2.nodesIds().length; i++) { + Map snps = res2.nodesSnapshots()[i].snapshots(); + + assertEquals(1, snps.size()); + + SnapshotListJobResult.SnapshotInfo info = snps.get(SNAPSHOT_NAME); + + assertNotNull(info.incrementals()); + + assertTrue(info.date() < time1); + assertTrue(info.date() < info.incrementals().date()); + assertTrue(info.incrementals().date() > time1); + } + } + + /** */ + @Test + public void testSnapshotListsSizes() throws Exception { + // Incremental snapshots do not support the only-primary mode or encryption. + assumeFalse(onlyPrimary || encryption); + // Speeds up the tests. + assumeTrue(snpThrdPoolSz > 1); + + int grids = 3; + + IgniteEx ig = startGridsWithCache(grids, txCacheConfig(defaultCacheConfiguration()), CACHE_KEYS_RANGE); + + snp(ig).createSnapshot(SNAPSHOT_NAME, null, false, onlyPrimary).get(getTestTimeout()); + + SnapshotListTaskResult res = listSnapshots(grid(0)); + + Map sizes0 = collectSizes(res); + + // Check the sizes. + for (int g = 0; g < grids; g++) { + long sz = SnapshotListTask.calculateDirectorySize(new SnapshotFileTree(grid(g).context(), SNAPSHOT_NAME, null).root()); + + assertEquals(sz, sizes0.get(grid(g).localNode().id()).longValue()); + + assertNull(snapshotInfo(res, grid(g)).externalStorages()); + assertNull(snapshotInfo(res, grid(g)).incrementals()); + } + + // Add some data and create an incremental snapshot. + try (IgniteDataStreamer ds = ig.dataStreamer(DEFAULT_CACHE_NAME)) { + for (int kv = CACHE_KEYS_RANGE; kv < CACHE_KEYS_RANGE * 2; ++kv) + ds.addData(kv, kv); + } + + snp(ig).createSnapshot(SNAPSHOT_NAME, null, true, onlyPrimary).get(getTestTimeout()); + + // Repeat the operation. + res = listSnapshots(grid(0)); + + Map sizes1 = collectSizes(res); + + for (int g = 0; g < grids; g++) { + UUID nodeId = grid(g).localNode().id(); + + SnapshotListJobResult.SnapshotInfo info = snapshotInfo(res, grid(g)); + + assertNotNull(info.incrementals()); + assertEquals(1, info.incrementals().number().intValue()); + + // The size must grow, but exactly to the actual directory size: the incremental part + // is placed inside the snapshot root and must not be counted twice. + assertTrue(sizes1.get(nodeId) > sizes0.get(nodeId)); + + long sz = SnapshotListTask.calculateDirectorySize(new SnapshotFileTree(grid(g).context(), SNAPSHOT_NAME, null).root()); + + assertEquals(sz, sizes1.get(nodeId).longValue()); + } + } + + /** */ + @Test + public void testSnapshotListsSizesExternalStorages() throws Exception { + // Speeds up the tests. + assumeTrue(snpThrdPoolSz > 1); + assumeTrue(onlyPrimary); + + int grids = 3; + + extStorages = true; + + // Properly delays the test cache creation with the configured external storages. + dfltCacheCfg = null; + + startGridsMultiThreaded(grids); + + grid(0).createCache(txCacheConfig(defaultCacheConfiguration())); + + try (IgniteDataStreamer ds = grid(0).dataStreamer(DEFAULT_CACHE_NAME)) { + for (int i = 0; i < CACHE_KEYS_RANGE; i++) + ds.addData(i, i); + } + + snp(grid(0)).createSnapshot(SNAPSHOT_NAME, null, false, onlyPrimary).get(getTestTimeout()); + + SnapshotListTaskResult res = listSnapshots(grid(0)); + + for (int g = 0; g < grids; g++) { + SnapshotListJobResult.SnapshotInfo info = snapshotInfo(res, grid(g)); + + assertNotNull(info); + assertNotNull(info.externalStorages()); + + long rootSize = SnapshotListTask.calculateDirectorySize(new SnapshotFileTree(grid(g).context(), SNAPSHOT_NAME, null).root()); + + // If there is a data withing the snapshot's external storage, its size must be added to the total size. + assertTrue(info.externalStorages().size() > 0L ? info.size() > rootSize : info.size() == rootSize); + assertEquals(rootSize + info.externalStorages().size(), info.size()); + } + } + + /** */ + @Test + public void testNodeStopDuringSnapshotList() throws Exception { + // Doesn't matter here, speeds up the tests. + assumeFalse(encryption || onlyPrimary); + + int grids = 3; + int testGridIdx = 1; + + CountDownLatch snpLstBeginLatch = new CountDownLatch(grids); + CountDownLatch snpLstProceedLatch = new CountDownLatch(1); + + // Delays snapshot reading. + pluginProvider = new AbstractTestPluginProvider() { + @Override public String name() { + return "TestSnpMgrProvider"; + } + + @Override public T createComponent(PluginContext ctx, Class cls) { + if (IgniteSnapshotManager.class.isAssignableFrom(cls)) { + return (T)new IgniteSnapshotManager(((IgniteEx)ctx.grid()).context()) { + @Override public List readSnapshotMetadatas(SnapshotFileTree sft, boolean failIfCantRead) { + snpLstBeginLatch.countDown(); + + if (((IgniteEx)ctx.grid()).localNode().order() == testGridIdx + 1) { + try { + assertTrue(snpLstProceedLatch.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + } + catch (InterruptedException e) { + throw new IllegalStateException(e); + } + } + + return super.readSnapshotMetadatas(sft, failIfCantRead); + } + }; + } + + return super.createComponent(ctx, cls); + } + }; + + startGridsWithSnapshot(grids, CACHE_KEYS_RANGE, true); + + IgniteInternalFuture lstOpFut = runAsync(() -> listSnapshots(grid(0))); + + assertTrue(snpLstBeginLatch.await(getTestTimeout(), TimeUnit.MILLISECONDS)); + + IgniteInternalFuture stopFut = runAsync(() -> stopGrid(testGridIdx)); + + assertTrue(waitForCondition(() -> grid(testGridIdx).context().isStopping(), getTestTimeout())); + + snpLstProceedLatch.countDown(); + + assertThrowsAnyCause( + null, + () -> lstOpFut.get(getTestTimeout()), + NodeStoppingException.class, + "Node is stopping" + ); + + stopFut.get(getTestTimeout()); + } + + /** */ + private static SnapshotListTaskResult listSnapshots(IgniteEx grid, @Nullable String src) throws Exception { + SnapshotListCommandArg arg = new SnapshotListCommandArg(); + + arg.src(src); + + Collection nodes = grid.cluster().forServers().nodes().stream().map(ClusterNode::id).toList(); + + return grid.compute().execute(SnapshotListTask.class, new VisorTaskArgument<>(nodes, arg, false)).result(); + } + + /** */ + private static SnapshotListTaskResult listSnapshots(IgniteEx grid) throws Exception { + return listSnapshots(grid, null); + } + + /** @return The test snapshot info reported for the node. */ + private static @Nullable SnapshotListJobResult.SnapshotInfo snapshotInfo(SnapshotListTaskResult res, IgniteEx node) { + for (int i = 0; i < res.nodesIds().length; i++) { + if (res.nodesIds()[i].equals(node.localNode().id())) + return res.nodesSnapshots()[i].snapshots().get(SNAPSHOT_NAME); + } + + return null; + } + + /** @return The test snapshot sizes reported by the list operation per node id. */ + private static Map collectSizes(SnapshotListTaskResult res) { + Map sizes = new HashMap<>(); + + for (int i = 0; i < res.nodesIds().length; i++) { + SnapshotListJobResult.SnapshotInfo info = res.nodesSnapshots()[i].snapshots().get(SNAPSHOT_NAME); + + assertNotNull(info); + + sizes.put(res.nodesIds()[i], info.size()); + } + + return sizes; + } +} diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/AbstractIncrementalSnapshotTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/AbstractIncrementalSnapshotTest.java index fcf65fa9b30ed..8055e7795e396 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/AbstractIncrementalSnapshotTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/AbstractIncrementalSnapshotTest.java @@ -132,7 +132,7 @@ protected void awaitSnapshotResourcesCleaned() { try { assertTrue(GridTestUtils.waitForCondition(() -> { for (Ignite g: G.allGrids()) { - if (snp((IgniteEx)g).currentCreateRequest() != null) + if (snp(g).currentCreateRequest() != null) return false; } diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/IncrementalSnapshotNodeFailureTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/IncrementalSnapshotNodeFailureTest.java index 88aad62a4985f..68ff0e443740f 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/IncrementalSnapshotNodeFailureTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/persistence/snapshot/incremental/IncrementalSnapshotNodeFailureTest.java @@ -109,7 +109,7 @@ private void runIncrementalSnapshotAndBreak(Supplier breakSnpWithExcp) t awaitSnapshotResourcesCleaned(); for (Ignite g: G.allGrids()) - assertNull(snp((IgniteEx)g).incrementalSnapshotId()); + assertNull(snp(g).incrementalSnapshotId()); stopAllGrids(); diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/AbstractRollingUpgradeTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/AbstractRollingUpgradeTest.java index 3d7f76477c251..4371c77ae4c17 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/AbstractRollingUpgradeTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/AbstractRollingUpgradeTest.java @@ -91,6 +91,8 @@ * 2.19.1 0-1 * 2.19.2 0-2 * 2.19.3 0-2,6 + * 2.19.4 0 + * 2.19.5 0-1 * 2.20.0 2-5 * 2.20.1 2-6 * 2.21.0 6 diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/feature/TestIgniteReleaseFeatures_2_19_4.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/feature/TestIgniteReleaseFeatures_2_19_4.java new file mode 100644 index 0000000000000..30a8a8ae030ee --- /dev/null +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/feature/TestIgniteReleaseFeatures_2_19_4.java @@ -0,0 +1,24 @@ +/* + * 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.ignite.internal.processors.rollingupgrade.feature; + +/** */ +public class TestIgniteReleaseFeatures_2_19_4 { + /** */ + public static final IgniteFeature ROLLING_UPGRADE_FEATURE = new IgniteCoreFeature(0); +} diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/feature/TestIgniteReleaseFeatures_2_19_5.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/feature/TestIgniteReleaseFeatures_2_19_5.java new file mode 100644 index 0000000000000..aeccf2e733eda --- /dev/null +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/rollingupgrade/feature/TestIgniteReleaseFeatures_2_19_5.java @@ -0,0 +1,27 @@ +/* + * 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.ignite.internal.processors.rollingupgrade.feature; + +/** */ +public class TestIgniteReleaseFeatures_2_19_5 { + /** */ + public static final IgniteFeature ROLLING_UPGRADE_FEATURE = new IgniteCoreFeature(0); + + /** */ + public static final IgniteFeature SNAPSHOT_LIST_FEATURE = CoreFeatureRegistry.SNAPSHOT_LIST_FEATURE; +} diff --git a/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteSnapshotTestSuite9.java b/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteSnapshotTestSuite9.java new file mode 100644 index 0000000000000..1ad9acf650aa6 --- /dev/null +++ b/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteSnapshotTestSuite9.java @@ -0,0 +1,51 @@ +/* + * 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.ignite.testsuites; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.IgniteClusterSnapshotListRollingUpgradeTest; +import org.apache.ignite.internal.processors.cache.persistence.snapshot.IgniteClusterSnapshotListTest; +import org.apache.ignite.testframework.GridTestUtils; +import org.apache.ignite.testframework.junits.DynamicSuite; +import org.junit.runner.RunWith; + +/** Split off from {@link IgniteSnapshotTestSuite8} to reduce the single-suite runtime in CI. */ +@RunWith(DynamicSuite.class) +public class IgniteSnapshotTestSuite9 { + /** + * @return Test suite. + */ + public static List> suite() { + return suite(null); + } + + /** + * @param ignoredTests Tests to ignore. + * @return Test suite. + */ + public static List> suite(Collection ignoredTests) { + List> suite = new ArrayList<>(); + + GridTestUtils.addTestIfNeeded(suite, IgniteClusterSnapshotListTest.class, ignoredTests); + GridTestUtils.addTestIfNeeded(suite, IgniteClusterSnapshotListRollingUpgradeTest.class, ignoredTests); + + return suite; + } +} diff --git a/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassTest_help.output b/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassTest_help.output index dd6ff4e43caf5..b86b79a73834f 100644 --- a/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassTest_help.output +++ b/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassTest_help.output @@ -269,6 +269,12 @@ This utility can do the following commands: Get the status of the current snapshot operation: control.(sh|bat) --snapshot status + Lists all snapshots on all online server nodes with their sizes: + control.(sh|bat) --snapshot list [--src path/to/snapshots] + + Parameters: + --src path/to/snapshots - Path to snapshot location directory. If not specified, the default configured snapshot directory will be used. + Change cluster tag to new value: control.(sh|bat) --change-tag newTagValue [--yes] diff --git a/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassWithSSLTest_help.output b/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassWithSSLTest_help.output index 76f0f8c7e0e37..54b21e4ef88bf 100644 --- a/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassWithSSLTest_help.output +++ b/modules/core/src/test/resources/org.apache.ignite.util/GridCommandHandlerClusterByClassWithSSLTest_help.output @@ -269,6 +269,12 @@ This utility can do the following commands: Get the status of the current snapshot operation: control.(sh|bat) --snapshot status + Lists all snapshots on all online server nodes with their sizes: + control.(sh|bat) --snapshot list [--src path/to/snapshots] + + Parameters: + --src path/to/snapshots - Path to snapshot location directory. If not specified, the default configured snapshot directory will be used. + Change cluster tag to new value: control.(sh|bat) --change-tag newTagValue [--yes] diff --git a/modules/ducktests/src/main/java/org/apache/ignite/internal/ducktest/tests/ContinuousDataLoadApplication.java b/modules/ducktests/src/main/java/org/apache/ignite/internal/ducktest/tests/ContinuousDataLoadApplication.java index bff6187c71532..5df1a59847a60 100644 --- a/modules/ducktests/src/main/java/org/apache/ignite/internal/ducktest/tests/ContinuousDataLoadApplication.java +++ b/modules/ducktests/src/main/java/org/apache/ignite/internal/ducktest/tests/ContinuousDataLoadApplication.java @@ -78,7 +78,7 @@ public class ContinuousDataLoadApplication extends IgniteAwareApplication { if (notifyTime + TimeUnit.MILLISECONDS.toNanos(1500) < System.nanoTime()) notifyTime = System.nanoTime(); - // Delayed notify of the initialization to make sure the data load has completelly began and + // Delayed notify of the initialization to make sure the data load has completely begun and // has produced some valuable amount of data. if (!inited() && warmUpCnt == loaded) markInitialized();