diff --git a/changelog/unreleased/SOLR-18249.yml b/changelog/unreleased/SOLR-18249.yml new file mode 100644 index 000000000000..9fe056a8b766 --- /dev/null +++ b/changelog/unreleased/SOLR-18249.yml @@ -0,0 +1,7 @@ +title: Stage shard backup metadata before publication and fail rather than fall back to a non-atomic local move. +type: fixed +authors: + - name: Nick Shanin +links: + - name: SOLR-18249 + url: https://issues.apache.org/jira/browse/SOLR-18249 diff --git a/solr/core/src/java/org/apache/solr/cloud/api/collections/DeleteBackupCmd.java b/solr/core/src/java/org/apache/solr/cloud/api/collections/DeleteBackupCmd.java index 6c0ec870bf11..5c7fd2b43962 100644 --- a/solr/core/src/java/org/apache/solr/cloud/api/collections/DeleteBackupCmd.java +++ b/solr/core/src/java/org/apache/solr/cloud/api/collections/DeleteBackupCmd.java @@ -27,7 +27,6 @@ import java.net.URI; import java.nio.file.NoSuchFileException; import java.util.ArrayList; -import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; @@ -166,10 +165,19 @@ void deleteBackupIds( Set referencedIndexFiles = new HashSet<>(); List shardBackupIdFileDeletes = new ArrayList<>(); - List shardBackupIds = - Arrays.stream(repository.listAllOrEmpty(shardBackupMetadataDir)) - .map(sbi -> ShardBackupId.fromShardMetadataFilename(sbi)) - .collect(Collectors.toList()); + List shardBackupIds = new ArrayList<>(); + for (String filename : repository.listAllOrEmpty(shardBackupMetadataDir)) { + try { + shardBackupIds.add(ShardBackupId.fromShardMetadataFilename(filename)); + } catch (IllegalArgumentException e) { + // The directory can hold files that are not shard metadata, such as the staged temp + // file an interrupted metadata write leaves behind. Such files belong to no backup + // point, so they are ignored here instead of failing the whole deletion. + if (log.isDebugEnabled()) { + log.debug("Ignoring file [{}] in shard backup metadata directory", filename); + } + } + } for (ShardBackupId shardBackupId : shardBackupIds) { final BackupId backupId = shardBackupId.getContainingBackupId(); diff --git a/solr/core/src/java/org/apache/solr/core/backup/ShardBackupMetadata.java b/solr/core/src/java/org/apache/solr/core/backup/ShardBackupMetadata.java index 73e1a13637ad..a5256323ea19 100644 --- a/solr/core/src/java/org/apache/solr/core/backup/ShardBackupMetadata.java +++ b/solr/core/src/java/org/apache/solr/core/backup/ShardBackupMetadata.java @@ -17,6 +17,7 @@ package org.apache.solr.core.backup; +import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; @@ -28,7 +29,6 @@ import java.util.List; import java.util.Map; import java.util.Optional; -import java.util.Set; import org.apache.lucene.store.IOContext; import org.apache.lucene.store.IndexInput; import org.apache.solr.common.util.Utils; @@ -104,20 +104,16 @@ public static ShardBackupMetadata from( } /** - * Storing ShardBackupMetadata at {@code folderURI} with name {@code filename}. If a file already - * existed there, overwrite it. + * Store this metadata at {@code folderURI} under the shard backup id's filename. The JSON is + * serialized completely before the repository is asked to write it; whether publication is atomic + * depends on the repository's {@link BackupRepository#writeBytes(URI, byte[])} implementation. */ public void store(BackupRepository repository, URI folderURI, ShardBackupId shardBackupId) throws IOException { final String filename = shardBackupId.getBackupMetadataFilename(); - URI fileURI = repository.resolve(folderURI, filename); - if (repository.exists(fileURI)) { - repository.delete(folderURI, Set.of(filename)); - } - - try (OutputStream os = repository.createOutput(repository.resolve(folderURI, filename))) { - store(os); - } + ByteArrayOutputStream buffer = new ByteArrayOutputStream(); + store(buffer); + repository.writeBytes(repository.resolve(folderURI, filename), buffer.toByteArray()); } public Collection listOriginalFileNames() { diff --git a/solr/core/src/java/org/apache/solr/core/backup/repository/BackupRepository.java b/solr/core/src/java/org/apache/solr/core/backup/repository/BackupRepository.java index b8f86af5e3ed..8b0c58cc81a4 100644 --- a/solr/core/src/java/org/apache/solr/core/backup/repository/BackupRepository.java +++ b/solr/core/src/java/org/apache/solr/core/backup/repository/BackupRepository.java @@ -147,6 +147,20 @@ default URI resolveDirectory(URI baseUri, String... pathComponents) { */ OutputStream createOutput(URI path) throws IOException; + /** + * Write {@code data} to {@code path} using this repository's output semantics. + * + *

The default implementation writes directly through {@link #createOutput(URI)} and makes no + * atomicity guarantee: it does not stage the bytes, and a failed write may leave a partially + * written or replaced object. Repositories whose backing store supports it may override this + * method to stage the bytes and publish them atomically. + */ + default void writeBytes(URI path, byte[] data) throws IOException { + try (OutputStream os = createOutput(path)) { + os.write(data); + } + } + // TODO define whether this should also create any nonexistent parent directories. (i.e. is this // 'mkdir', or 'mkdir -p') /** diff --git a/solr/core/src/java/org/apache/solr/core/backup/repository/DelegatingBackupRepository.java b/solr/core/src/java/org/apache/solr/core/backup/repository/DelegatingBackupRepository.java index e3b27cb073c5..407998eec63a 100644 --- a/solr/core/src/java/org/apache/solr/core/backup/repository/DelegatingBackupRepository.java +++ b/solr/core/src/java/org/apache/solr/core/backup/repository/DelegatingBackupRepository.java @@ -96,6 +96,11 @@ public OutputStream createOutput(URI path) throws IOException { return delegate.createOutput(path); } + @Override + public void writeBytes(URI path, byte[] data) throws IOException { + delegate.writeBytes(path, data); + } + @Override public void createDirectory(URI path) throws IOException { delegate.createDirectory(path); diff --git a/solr/core/src/java/org/apache/solr/core/backup/repository/LocalFileSystemRepository.java b/solr/core/src/java/org/apache/solr/core/backup/repository/LocalFileSystemRepository.java index 77f7e1921f37..f4606fac6315 100644 --- a/solr/core/src/java/org/apache/solr/core/backup/repository/LocalFileSystemRepository.java +++ b/solr/core/src/java/org/apache/solr/core/backup/repository/LocalFileSystemRepository.java @@ -25,8 +25,11 @@ import java.nio.file.LinkOption; import java.nio.file.NoSuchFileException; import java.nio.file.Path; +import java.nio.file.StandardCopyOption; +import java.nio.file.StandardOpenOption; import java.util.Collection; import java.util.Objects; +import java.util.UUID; import org.apache.commons.io.file.PathUtils; import org.apache.lucene.store.Directory; import org.apache.lucene.store.FSDirectory; @@ -115,6 +118,37 @@ public OutputStream createOutput(URI path) throws IOException { return Files.newOutputStream(Path.of(path)); } + /** + * Stage the bytes in a sibling file and publish them by moving the staged file onto {@code path} + * atomically, replacing any existing file. + * + *

This method does not fall back to a non-atomic move. If the provider cannot perform the + * atomic move, the failure is propagated and cleanup of the staged file is attempted. + * + * @throws IOException if writing or the requested atomic move fails + */ + @Override + public void writeBytes(URI path, byte[] data) throws IOException { + Path dest = Path.of(path); + Path temp = dest.resolveSibling(dest.getFileName().toString() + ".tmp." + UUID.randomUUID()); + try { + Files.write(temp, data, StandardOpenOption.CREATE_NEW, StandardOpenOption.WRITE); + moveAtomically(temp, dest); + } catch (IOException | RuntimeException e) { + try { + Files.deleteIfExists(temp); + } catch (IOException | RuntimeException cleanupFailure) { + e.addSuppressed(cleanupFailure); + } + throw e; + } + } + + /** Performs the atomic move used by {@link #writeBytes(URI, byte[])}. */ + protected void moveAtomically(Path temp, Path dest) throws IOException { + Files.move(temp, dest, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); + } + @Override public String[] listAll(URI dirPath) throws IOException { // It is better to check the existence of the directory first since diff --git a/solr/core/src/test/org/apache/solr/cloud/api/collections/DeleteBackupCmdTest.java b/solr/core/src/test/org/apache/solr/cloud/api/collections/DeleteBackupCmdTest.java index 5d01b5ce3af3..f54f766a58a7 100644 --- a/solr/core/src/test/org/apache/solr/cloud/api/collections/DeleteBackupCmdTest.java +++ b/solr/core/src/test/org/apache/solr/cloud/api/collections/DeleteBackupCmdTest.java @@ -17,6 +17,7 @@ package org.apache.solr.cloud.api.collections; import java.io.IOException; +import java.io.OutputStream; import java.net.URI; import java.util.Set; import java.util.UUID; @@ -24,6 +25,9 @@ import org.apache.solr.common.util.NamedList; import org.apache.solr.core.backup.BackupFilePaths; import org.apache.solr.core.backup.BackupId; +import org.apache.solr.core.backup.Checksum; +import org.apache.solr.core.backup.ShardBackupId; +import org.apache.solr.core.backup.ShardBackupMetadata; import org.apache.solr.core.backup.repository.BackupRepository; import org.apache.solr.core.backup.repository.LocalFileSystemRepository; import org.junit.Before; @@ -86,6 +90,68 @@ public void deleteDirectory(URI path) throws IOException { assertEquals("simulated repository failure", thrown.getMessage()); } + @Test + public void testDeleteBackupIdsIgnoresStagedMetadataTempFile() throws Exception { + URI metadataDir = new BackupFilePaths(repository, backupUri).getShardBackupMetadataDir(); + ShardBackupId shardBackupId = new ShardBackupId("shard1", BackupId.zero()); + storeMetadata(metadataDir, shardBackupId); + String stagedFile = createStagedTempFile(metadataDir, shardBackupId); + URI metadataFile = repository.resolve(metadataDir, shardBackupId.getBackupMetadataFilename()); + assertTrue(repository.exists(metadataFile)); + assertTrue(repository.exists(repository.resolve(metadataDir, stagedFile))); + + NamedList results = new NamedList<>(); + new DeleteBackupCmd(null) + .deleteBackupIds(backupUri, repository, Set.of(BackupId.zero()), results); + + assertNotNull(results.get("deleted")); + assertFalse(repository.exists(metadataFile)); + // The staged file is not any backup point's metadata, so it is left in place. + assertTrue(repository.exists(repository.resolve(metadataDir, stagedFile))); + } + + @Test + public void testKeepNumberOfBackupIgnoresStagedMetadataTempFile() throws Exception { + URI metadataDir = new BackupFilePaths(repository, backupUri).getShardBackupMetadataDir(); + ShardBackupId oldest = new ShardBackupId("shard1", BackupId.zero()); + ShardBackupId newest = new ShardBackupId("shard1", new BackupId(1)); + storeMetadata(metadataDir, oldest); + storeMetadata(metadataDir, newest); + createStagedTempFile(metadataDir, oldest); + createBackupPropsFile(BackupId.zero()); + createBackupPropsFile(new BackupId(1)); + + new DeleteBackupCmd(null).keepNumberOfBackup(repository, backupUri, 1, new NamedList<>()); + + assertFalse( + repository.exists(repository.resolve(metadataDir, oldest.getBackupMetadataFilename()))); + assertTrue( + repository.exists(repository.resolve(metadataDir, newest.getBackupMetadataFilename()))); + } + + private void storeMetadata(URI metadataDir, ShardBackupId shardBackupId) throws IOException { + ShardBackupMetadata metadata = ShardBackupMetadata.empty(); + metadata.addBackedFile("uniq_" + shardBackupId.getIdAsString(), "orig", new Checksum(1L, 10)); + metadata.store(repository, metadataDir, shardBackupId); + } + + private String createStagedTempFile(URI metadataDir, ShardBackupId shardBackupId) + throws IOException { + String stagedName = shardBackupId.getBackupMetadataFilename() + ".tmp." + UUID.randomUUID(); + try (OutputStream out = repository.createOutput(repository.resolve(metadataDir, stagedName))) { + out.write('#'); + } + return stagedName; + } + + private void createBackupPropsFile(BackupId backupId) throws IOException { + try (OutputStream out = + repository.createOutput( + repository.resolve(backupUri, BackupFilePaths.getBackupPropsName(backupId)))) { + out.write('#'); + } + } + private URI zkStateDir(BackupId backupId) { return repository.resolveDirectory(backupUri, BackupFilePaths.getZkStateDir(backupId)); } diff --git a/solr/core/src/test/org/apache/solr/core/backup/ShardBackupMetadataOverwriteTest.java b/solr/core/src/test/org/apache/solr/core/backup/ShardBackupMetadataOverwriteTest.java new file mode 100644 index 000000000000..88a664afd591 --- /dev/null +++ b/solr/core/src/test/org/apache/solr/core/backup/ShardBackupMetadataOverwriteTest.java @@ -0,0 +1,89 @@ +/* + * 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.solr.core.backup; + +import java.io.IOException; +import java.net.URI; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import org.apache.solr.SolrTestCase; +import org.apache.solr.common.util.NamedList; +import org.apache.solr.core.backup.repository.BackupRepository; +import org.apache.solr.core.backup.repository.DelegatingBackupRepository; +import org.apache.solr.core.backup.repository.LocalFileSystemRepository; +import org.junit.Before; +import org.junit.Test; + +/** + * Verifies that overwriting shard backup metadata never deletes the previous metadata file first. + * This test deliberately uses only the {@link BackupRepository} API that predates the {@code + * writeBytes} method, so it also compiles and runs against the code from before that change, where + * the overwrite deleted the existing file and this test fails. + */ +public class ShardBackupMetadataOverwriteTest extends SolrTestCase { + + private LocalFileSystemRepository repository; + private URI folder; + private ShardBackupId shardBackupId; + + @Before + public void setUpRepo() throws Exception { + repository = new LocalFileSystemRepository(); + repository.init(new NamedList<>()); + folder = + repository.createURI(createTempDir("shard-backup-metadata").toAbsolutePath().toString()); + repository.createDirectory(folder); + shardBackupId = new ShardBackupId("shard1", BackupId.zero()); + } + + @Test + public void testStoreDoesNotDeleteExistingMetadata() throws Exception { + metadata("uniq1", "orig1", new Checksum(1L, 10)).store(repository, folder, shardBackupId); + + RecordingBackupRepository recording = new RecordingBackupRepository(repository); + metadata("uniq2", "orig2", new Checksum(2L, 20)).store(recording, folder, shardBackupId); + + assertTrue("overwrite must not delete the previous metadata file", recording.deleted.isEmpty()); + + ShardBackupMetadata loaded = ShardBackupMetadata.from(repository, folder, shardBackupId); + assertEquals(List.of("uniq2"), loaded.listUniqueFileNames()); + } + + private static ShardBackupMetadata metadata( + String uniqueFileName, String originalFileName, Checksum checksum) { + ShardBackupMetadata created = ShardBackupMetadata.empty(); + created.addBackedFile(uniqueFileName, originalFileName, checksum); + return created; + } + + private static class RecordingBackupRepository extends DelegatingBackupRepository { + final List deleted = new ArrayList<>(); + + RecordingBackupRepository(BackupRepository delegate) { + setDelegate(delegate); + } + + @Override + public void delete(URI path, Collection files) throws IOException { + for (String file : files) { + deleted.add(resolve(path, file)); + } + super.delete(path, files); + } + } +} diff --git a/solr/core/src/test/org/apache/solr/core/backup/ShardBackupMetadataTest.java b/solr/core/src/test/org/apache/solr/core/backup/ShardBackupMetadataTest.java new file mode 100644 index 000000000000..b9819149c104 --- /dev/null +++ b/solr/core/src/test/org/apache/solr/core/backup/ShardBackupMetadataTest.java @@ -0,0 +1,200 @@ +/* + * 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.solr.core.backup; + +import java.io.IOException; +import java.io.OutputStream; +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.InvocationTargetException; +import java.lang.reflect.Proxy; +import java.net.URI; +import java.nio.file.AtomicMoveNotSupportedException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.List; +import org.apache.solr.SolrTestCase; +import org.apache.solr.common.util.NamedList; +import org.apache.solr.core.backup.repository.BackupRepository; +import org.apache.solr.core.backup.repository.DelegatingBackupRepository; +import org.apache.solr.core.backup.repository.LocalFileSystemRepository; +import org.junit.Before; +import org.junit.Test; + +/** Unit tests for {@link ShardBackupMetadata} overwrite behavior. */ +public class ShardBackupMetadataTest extends SolrTestCase { + + private LocalFileSystemRepository repository; + private URI folder; + private ShardBackupId shardBackupId; + + @Before + public void setUpRepo() throws Exception { + repository = new LocalFileSystemRepository(); + repository.init(new NamedList<>()); + folder = + repository.createURI(createTempDir("shard-backup-metadata").toAbsolutePath().toString()); + repository.createDirectory(folder); + shardBackupId = new ShardBackupId("shard1", BackupId.zero()); + } + + @Test + public void testStoreOverwritesAndReadsBack() throws Exception { + metadata("uniq1", "orig1", new Checksum(1L, 10)).store(repository, folder, shardBackupId); + metadata("uniq2", "orig2", new Checksum(2L, 20)).store(repository, folder, shardBackupId); + + ShardBackupMetadata loaded = ShardBackupMetadata.from(repository, folder, shardBackupId); + assertNotNull(loaded); + assertEquals(List.of("uniq2"), loaded.listUniqueFileNames()); + assertTrue(loaded.getFile("orig2").isPresent()); + assertEquals(2L, loaded.getFile("orig2").get().fileChecksum.checksum); + assertTrue(loaded.getFile("orig1").isEmpty()); + } + + @Test + public void testFailedOverwriteKeepsPreviousMetadata() throws Exception { + metadata("uniq1", "orig1", new Checksum(1L, 10)).store(repository, folder, shardBackupId); + + FailingWriteRepository failing = new FailingWriteRepository(repository); + expectThrows( + IOException.class, + () -> + metadata("uniq2", "orig2", new Checksum(2L, 20)).store(failing, folder, shardBackupId)); + + URI dest = repository.resolve(folder, shardBackupId.getBackupMetadataFilename()); + assertTrue(repository.exists(dest)); + ShardBackupMetadata loaded = ShardBackupMetadata.from(repository, folder, shardBackupId); + assertNotNull(loaded); + assertEquals(List.of("uniq1"), loaded.listUniqueFileNames()); + assertTrue(loaded.getFile("orig1").isPresent()); + assertTrue(loaded.getFile("orig2").isEmpty()); + } + + @Test + public void testUnsupportedAtomicMovePreservesExistingMetadataAndCleansTempFile() + throws Exception { + metadata("uniq1", "orig1", new Checksum(1L, 10)).store(repository, folder, shardBackupId); + URI dest = repository.resolve(folder, shardBackupId.getBackupMetadataFilename()); + byte[] previousMetadata = Files.readAllBytes(Path.of(dest)); + + UnsupportedAtomicMoveRepository unsupported = new UnsupportedAtomicMoveRepository(); + unsupported.init(new NamedList<>()); + expectThrows( + AtomicMoveNotSupportedException.class, + () -> + metadata("uniq2", "orig2", new Checksum(2L, 20)) + .store(unsupported, folder, shardBackupId)); + + assertArrayEquals(previousMetadata, Files.readAllBytes(Path.of(dest))); + String[] files = repository.listAll(folder); + assertEquals( + "the failed publication must not leave its sibling temp file, found: " + + Arrays.toString(files), + 1, + files.length); + assertEquals(shardBackupId.getBackupMetadataFilename(), files[0]); + ShardBackupMetadata loaded = ShardBackupMetadata.from(repository, folder, shardBackupId); + assertEquals(List.of("uniq1"), loaded.listUniqueFileNames()); + } + + @Test + public void testInterfaceDefaultWriteBytesUsesCreateOutput() throws Exception { + metadata("uniq1", "orig1", new Checksum(1L, 10)).store(repository, folder, shardBackupId); + + RecordingBackupRepository recording = new RecordingBackupRepository(repository); + BackupRepository interfaceDefault = usingInterfaceDefaultWriteBytes(recording); + metadata("uniq2", "orig2", new Checksum(2L, 20)).store(interfaceDefault, folder, shardBackupId); + + assertTrue(recording.deleted.isEmpty()); + assertEquals(1, recording.created.size()); + assertEquals( + repository.resolve(folder, shardBackupId.getBackupMetadataFilename()), + recording.created.get(0)); + + ShardBackupMetadata loaded = ShardBackupMetadata.from(repository, folder, shardBackupId); + assertEquals(List.of("uniq2"), loaded.listUniqueFileNames()); + } + + private static BackupRepository usingInterfaceDefaultWriteBytes( + RecordingBackupRepository recording) { + return (BackupRepository) + Proxy.newProxyInstance( + BackupRepository.class.getClassLoader(), + new Class[] {BackupRepository.class}, + (proxy, method, args) -> { + if (method.isDefault()) { + return InvocationHandler.invokeDefault(proxy, method, args); + } + try { + return method.invoke(recording, args); + } catch (InvocationTargetException e) { + throw e.getCause(); + } + }); + } + + private static ShardBackupMetadata metadata( + String uniqueFileName, String originalFileName, Checksum checksum) { + ShardBackupMetadata created = ShardBackupMetadata.empty(); + created.addBackedFile(uniqueFileName, originalFileName, checksum); + return created; + } + + private static class RecordingBackupRepository extends DelegatingBackupRepository { + final List deleted = new ArrayList<>(); + final List created = new ArrayList<>(); + + RecordingBackupRepository(BackupRepository delegate) { + setDelegate(delegate); + } + + @Override + public OutputStream createOutput(URI path) throws IOException { + created.add(path); + return super.createOutput(path); + } + + @Override + public void delete(URI path, Collection files) throws IOException { + for (String file : files) { + deleted.add(resolve(path, file)); + } + super.delete(path, files); + } + } + + private static class FailingWriteRepository extends DelegatingBackupRepository { + FailingWriteRepository(BackupRepository delegate) { + setDelegate(delegate); + } + + @Override + public void writeBytes(URI path, byte[] data) throws IOException { + throw new IOException("injected write failure"); + } + } + + private static class UnsupportedAtomicMoveRepository extends LocalFileSystemRepository { + @Override + protected void moveAtomically(Path temp, Path dest) throws IOException { + throw new AtomicMoveNotSupportedException( + temp.toString(), dest.toString(), "injected unsupported atomic move"); + } + } +}