Copilot commented on code in PR #8739: URL: https://github.com/apache/hbase/pull/8739#discussion_r4200049797
########## hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/persistence/HadoopFsCachePersistenceStorage.java: ########## @@ -0,0 +1,419 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache.persistence; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.util.Objects; +import java.util.Optional; +import java.util.UUID; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.FSDataOutputStream; +import org.apache.hadoop.fs.FileSystem; +import org.apache.hadoop.fs.Path; +import org.apache.yetus.audience.InterfaceAudience; + +/** + * {@link CachePersistenceStorage} implementation backed by a Hadoop {@link FileSystem}. + * <p> + * This implementation works with filesystems available through Hadoop's filesystem abstraction, + * including HDFS and the local filesystem. + * </p> + * <p> + * Every persistence storage key is represented as a file below a configured root directory. Keys + * may contain relative path components separated by {@code /}, allowing higher-level persistence + * infrastructure to organize state by component role. + * </p> + * <p> + * Persistence writes are first written to a temporary file. The temporary file is promoted to the + * committed state path only after the caller explicitly commits the write. An abandoned write + * deletes its temporary file and leaves previously committed state unchanged. + * </p> + * <p> + * Publication uses filesystem rename operations. Rename atomicity and implementation details are + * determined by the underlying Hadoop filesystem. For filesystems where rename is implemented as + * copy-and-delete, publication has the corresponding filesystem semantics. + * </p> + * <p> + * This class does not own the supplied {@link FileSystem} and does not close it. + * </p> + */ [email protected] +public class HadoopFsCachePersistenceStorage implements CachePersistenceStorage { + + private static final String STATE_FILE_SUFFIX = ".state"; + private static final String TEMP_FILE_MARKER = ".tmp-"; + private static final String BACKUP_FILE_MARKER = ".backup-"; + + private final FileSystem fileSystem; + private final Path rootPath; + + /** + * Creates persistence storage using the filesystem associated with the supplied root path. + * @param conf Hadoop configuration used to resolve the filesystem + * @param rootPath root directory under which persisted cache state is stored + * @throws IOException if the filesystem for the root path cannot be resolved + */ + public HadoopFsCachePersistenceStorage(Configuration conf, Path rootPath) throws IOException { + this(resolveFileSystem(conf, rootPath), rootPath); + } + + /** + * Creates persistence storage using an explicitly supplied Hadoop filesystem. + * <p> + * The supplied filesystem remains owned by the caller and is not closed by this storage instance. + * </p> + * @param fileSystem Hadoop filesystem used to store cache state + * @param rootPath root directory under which persisted cache state is stored + */ + public HadoopFsCachePersistenceStorage(FileSystem fileSystem, Path rootPath) { + this.fileSystem = Objects.requireNonNull(fileSystem, "fileSystem must not be null"); + this.rootPath = Objects.requireNonNull(rootPath, "rootPath must not be null"); + } + + /** + * Opens committed persisted state for the specified storage key. + * @param key stable storage key assigned to a persistent cache component + * @return input stream containing committed state, or an empty optional if no state exists + * @throws IOException if committed state exists but cannot be opened + */ + @Override + public Optional<InputStream> open(String key) throws IOException { + Path path = getStatePath(key); + if (!fileSystem.exists(path)) { + return Optional.empty(); + } + return Optional.of(fileSystem.open(path)); + } + + /** + * Starts a new persistence write for the specified storage key. + * <p> + * Data is written to a temporary file located next to the final state file. The temporary file + * becomes committed state only when the returned handle is explicitly committed. + * </p> + * @param key stable storage key assigned to a persistent cache component + * @return handle for writing and committing persisted state + * @throws IOException if the temporary persistence file cannot be created + */ + @Override + public CachePersistenceOutput create(String key) throws IOException { + Path targetPath = getStatePath(key); + Path parent = targetPath.getParent(); + ensureDirectory(parent); + + Path temporaryPath = createTemporaryPath(targetPath); + FSDataOutputStream output = fileSystem.create(temporaryPath, false); + + return new HadoopFsCachePersistenceOutput(fileSystem, output, temporaryPath, targetPath); + } + + /** + * Creates the parent directory for a persistence state file when necessary. + * @param directory directory to create + * @throws IOException if the directory cannot be created + */ + private void ensureDirectory(Path directory) throws IOException { + if (directory == null || fileSystem.exists(directory)) { + return; + } + + if (!fileSystem.mkdirs(directory) && !fileSystem.exists(directory)) { + throw new IOException("Failed to create cache persistence directory " + directory); + } + } + + /** + * Creates a unique temporary path adjacent to the specified committed state path. + * @param targetPath committed state path + * @return unique temporary path + */ + private Path createTemporaryPath(Path targetPath) { + return new Path(targetPath.toString() + TEMP_FILE_MARKER + UUID.randomUUID()); + } + + /** + * Resolves the Hadoop filesystem associated with the supplied persistence root path. + * @param conf Hadoop configuration used to resolve the filesystem + * @param rootPath persistence root path + * @return filesystem associated with the root path + * @throws IOException if the filesystem cannot be resolved + */ + private static FileSystem resolveFileSystem(Configuration conf, Path rootPath) + throws IOException { + Objects.requireNonNull(conf, "conf must not be null"); + Objects.requireNonNull(rootPath, "rootPath must not be null"); + return rootPath.getFileSystem(conf); + } + + /** + * Returns the filesystem path containing committed state for the supplied persistence key. + * @param key persistence storage key + * @return filesystem path containing committed persisted state + */ + private Path getStatePath(String key) { + validateKey(key); + return new Path(rootPath, key + STATE_FILE_SUFFIX); + } + + /** + * Validates a persistence storage key before using it as a relative filesystem path. + * <p> + * Storage keys must be relative and may contain multiple path segments. Empty segments and the + * special {@code .} and {@code ..} path segments are not allowed. + * </p> + * @param key persistence storage key to validate + * @throws NullPointerException if the key is {@code null} + * @throws IllegalArgumentException if the key is empty or contains an invalid path component + */ + private void validateKey(String key) { + Objects.requireNonNull(key, "key must not be null"); + + if (key.isEmpty()) { + throw new IllegalArgumentException("key must not be empty"); + } + if (key.startsWith("/") || key.endsWith("/") || key.indexOf('\\') >= 0 || key.contains("://")) { + throw new IllegalArgumentException("Invalid persistence key: " + key); + } + + String[] components = key.split("/"); + for (String component : components) { + if (component.isEmpty() || ".".equals(component) || "..".equals(component)) { + throw new IllegalArgumentException("Invalid persistence key: " + key); + } + } + } + + /** + * Persistence output backed by a temporary Hadoop filesystem file. + * <p> + * A successful commit closes the temporary output stream and publishes the temporary file as the + * committed state file. Closing an uncommitted instance aborts the write. + * </p> + */ + private static final class HadoopFsCachePersistenceOutput implements CachePersistenceOutput { + + private final FileSystem fileSystem; + private final FSDataOutputStream output; + private final Path temporaryPath; + private final Path targetPath; + + private boolean streamClosed; + private boolean committed; + private boolean aborted; + + /** + * Creates a persistence output backed by a temporary filesystem file. + * @param fileSystem filesystem containing the temporary and target files + * @param output output stream writing the temporary file + * @param temporaryPath temporary persistence file + * @param targetPath committed persistence file + */ + private HadoopFsCachePersistenceOutput(FileSystem fileSystem, FSDataOutputStream output, + Path temporaryPath, Path targetPath) { + this.fileSystem = Objects.requireNonNull(fileSystem, "fileSystem must not be null"); + this.output = Objects.requireNonNull(output, "output must not be null"); + this.temporaryPath = Objects.requireNonNull(temporaryPath, "temporaryPath must not be null"); + this.targetPath = Objects.requireNonNull(targetPath, "targetPath must not be null"); + } + + /** + * Returns the stream used to write temporary persisted state. + * @return persistence output stream + */ + @Override + public OutputStream getOutputStream() { + return output; + } + + /** + * Commits the persistence write by publishing the completed temporary file. + * <p> + * If previously committed state exists, it is moved temporarily out of the way before the new + * state is published. If publication of the new state fails, this method attempts to restore + * the previous committed state. + * </p> + * @throws IOException if the output cannot be closed or the new state cannot be published + */ + @Override + public void commit() throws IOException { + ensureActive(); + + closeOutput(); + + Path backupPath = null; + boolean previousStateMoved = false; + + try { + if (fileSystem.exists(targetPath)) { + backupPath = createBackupPath(targetPath); + if (!fileSystem.rename(targetPath, backupPath)) { Review Comment: Moving the committed file to a backup leaves no readable checkpoint until the second rename. If the process exits between these renames, rollback never runs; after restart, open() checks only the .state path, so restoration silently skips the valid backup. This happens even when individual filesystem renames are atomic. Use atomic replacement where supported or a recoverable publication protocol, and add a test that simulates interruption after the backup rename. ########## hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/persistence/HadoopFsCachePersistenceStorage.java: ########## @@ -0,0 +1,419 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache.persistence; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.util.Objects; +import java.util.Optional; +import java.util.UUID; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.FSDataOutputStream; +import org.apache.hadoop.fs.FileSystem; +import org.apache.hadoop.fs.Path; +import org.apache.yetus.audience.InterfaceAudience; + +/** + * {@link CachePersistenceStorage} implementation backed by a Hadoop {@link FileSystem}. + * <p> + * This implementation works with filesystems available through Hadoop's filesystem abstraction, + * including HDFS and the local filesystem. + * </p> + * <p> + * Every persistence storage key is represented as a file below a configured root directory. Keys + * may contain relative path components separated by {@code /}, allowing higher-level persistence + * infrastructure to organize state by component role. + * </p> + * <p> + * Persistence writes are first written to a temporary file. The temporary file is promoted to the + * committed state path only after the caller explicitly commits the write. An abandoned write + * deletes its temporary file and leaves previously committed state unchanged. + * </p> + * <p> + * Publication uses filesystem rename operations. Rename atomicity and implementation details are + * determined by the underlying Hadoop filesystem. For filesystems where rename is implemented as + * copy-and-delete, publication has the corresponding filesystem semantics. + * </p> + * <p> + * This class does not own the supplied {@link FileSystem} and does not close it. + * </p> + */ [email protected] +public class HadoopFsCachePersistenceStorage implements CachePersistenceStorage { + + private static final String STATE_FILE_SUFFIX = ".state"; + private static final String TEMP_FILE_MARKER = ".tmp-"; + private static final String BACKUP_FILE_MARKER = ".backup-"; + + private final FileSystem fileSystem; + private final Path rootPath; + + /** + * Creates persistence storage using the filesystem associated with the supplied root path. + * @param conf Hadoop configuration used to resolve the filesystem + * @param rootPath root directory under which persisted cache state is stored + * @throws IOException if the filesystem for the root path cannot be resolved + */ + public HadoopFsCachePersistenceStorage(Configuration conf, Path rootPath) throws IOException { + this(resolveFileSystem(conf, rootPath), rootPath); + } + + /** + * Creates persistence storage using an explicitly supplied Hadoop filesystem. + * <p> + * The supplied filesystem remains owned by the caller and is not closed by this storage instance. + * </p> + * @param fileSystem Hadoop filesystem used to store cache state + * @param rootPath root directory under which persisted cache state is stored + */ + public HadoopFsCachePersistenceStorage(FileSystem fileSystem, Path rootPath) { + this.fileSystem = Objects.requireNonNull(fileSystem, "fileSystem must not be null"); + this.rootPath = Objects.requireNonNull(rootPath, "rootPath must not be null"); + } + + /** + * Opens committed persisted state for the specified storage key. + * @param key stable storage key assigned to a persistent cache component + * @return input stream containing committed state, or an empty optional if no state exists + * @throws IOException if committed state exists but cannot be opened + */ + @Override + public Optional<InputStream> open(String key) throws IOException { + Path path = getStatePath(key); + if (!fileSystem.exists(path)) { + return Optional.empty(); + } + return Optional.of(fileSystem.open(path)); + } + + /** + * Starts a new persistence write for the specified storage key. + * <p> + * Data is written to a temporary file located next to the final state file. The temporary file + * becomes committed state only when the returned handle is explicitly committed. + * </p> + * @param key stable storage key assigned to a persistent cache component + * @return handle for writing and committing persisted state + * @throws IOException if the temporary persistence file cannot be created + */ + @Override + public CachePersistenceOutput create(String key) throws IOException { + Path targetPath = getStatePath(key); + Path parent = targetPath.getParent(); + ensureDirectory(parent); + + Path temporaryPath = createTemporaryPath(targetPath); + FSDataOutputStream output = fileSystem.create(temporaryPath, false); + + return new HadoopFsCachePersistenceOutput(fileSystem, output, temporaryPath, targetPath); + } + + /** + * Creates the parent directory for a persistence state file when necessary. + * @param directory directory to create + * @throws IOException if the directory cannot be created + */ + private void ensureDirectory(Path directory) throws IOException { + if (directory == null || fileSystem.exists(directory)) { + return; + } + + if (!fileSystem.mkdirs(directory) && !fileSystem.exists(directory)) { + throw new IOException("Failed to create cache persistence directory " + directory); + } + } + + /** + * Creates a unique temporary path adjacent to the specified committed state path. + * @param targetPath committed state path + * @return unique temporary path + */ + private Path createTemporaryPath(Path targetPath) { + return new Path(targetPath.toString() + TEMP_FILE_MARKER + UUID.randomUUID()); + } + + /** + * Resolves the Hadoop filesystem associated with the supplied persistence root path. + * @param conf Hadoop configuration used to resolve the filesystem + * @param rootPath persistence root path + * @return filesystem associated with the root path + * @throws IOException if the filesystem cannot be resolved + */ + private static FileSystem resolveFileSystem(Configuration conf, Path rootPath) + throws IOException { + Objects.requireNonNull(conf, "conf must not be null"); + Objects.requireNonNull(rootPath, "rootPath must not be null"); + return rootPath.getFileSystem(conf); + } + + /** + * Returns the filesystem path containing committed state for the supplied persistence key. + * @param key persistence storage key + * @return filesystem path containing committed persisted state + */ + private Path getStatePath(String key) { + validateKey(key); + return new Path(rootPath, key + STATE_FILE_SUFFIX); + } + + /** + * Validates a persistence storage key before using it as a relative filesystem path. + * <p> + * Storage keys must be relative and may contain multiple path segments. Empty segments and the + * special {@code .} and {@code ..} path segments are not allowed. + * </p> + * @param key persistence storage key to validate + * @throws NullPointerException if the key is {@code null} + * @throws IllegalArgumentException if the key is empty or contains an invalid path component + */ + private void validateKey(String key) { + Objects.requireNonNull(key, "key must not be null"); + + if (key.isEmpty()) { + throw new IllegalArgumentException("key must not be empty"); + } + if (key.startsWith("/") || key.endsWith("/") || key.indexOf('\\') >= 0 || key.contains("://")) { + throw new IllegalArgumentException("Invalid persistence key: " + key); + } Review Comment: Checking for "://" still accepts URI-qualified keys such as "file:/outside". Hadoop Path interprets that child as an absolute file URI, so resolving it against rootPath can read or write outside the configured persistence directory on the local filesystem. Reject keys with a URI scheme, and test the single-slash form for both open() and create(). ########## hbase-server/src/main/java/org/apache/hadoop/hbase/io/hfile/cache/persistence/HadoopFsCachePersistenceStorage.java: ########## @@ -0,0 +1,419 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.io.hfile.cache.persistence; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.util.Objects; +import java.util.Optional; +import java.util.UUID; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.FSDataOutputStream; +import org.apache.hadoop.fs.FileSystem; +import org.apache.hadoop.fs.Path; +import org.apache.yetus.audience.InterfaceAudience; + +/** + * {@link CachePersistenceStorage} implementation backed by a Hadoop {@link FileSystem}. + * <p> + * This implementation works with filesystems available through Hadoop's filesystem abstraction, + * including HDFS and the local filesystem. + * </p> + * <p> + * Every persistence storage key is represented as a file below a configured root directory. Keys + * may contain relative path components separated by {@code /}, allowing higher-level persistence + * infrastructure to organize state by component role. + * </p> + * <p> + * Persistence writes are first written to a temporary file. The temporary file is promoted to the + * committed state path only after the caller explicitly commits the write. An abandoned write + * deletes its temporary file and leaves previously committed state unchanged. + * </p> + * <p> + * Publication uses filesystem rename operations. Rename atomicity and implementation details are + * determined by the underlying Hadoop filesystem. For filesystems where rename is implemented as + * copy-and-delete, publication has the corresponding filesystem semantics. + * </p> + * <p> + * This class does not own the supplied {@link FileSystem} and does not close it. + * </p> + */ [email protected] +public class HadoopFsCachePersistenceStorage implements CachePersistenceStorage { + + private static final String STATE_FILE_SUFFIX = ".state"; + private static final String TEMP_FILE_MARKER = ".tmp-"; + private static final String BACKUP_FILE_MARKER = ".backup-"; + + private final FileSystem fileSystem; + private final Path rootPath; + + /** + * Creates persistence storage using the filesystem associated with the supplied root path. + * @param conf Hadoop configuration used to resolve the filesystem + * @param rootPath root directory under which persisted cache state is stored + * @throws IOException if the filesystem for the root path cannot be resolved + */ + public HadoopFsCachePersistenceStorage(Configuration conf, Path rootPath) throws IOException { + this(resolveFileSystem(conf, rootPath), rootPath); + } + + /** + * Creates persistence storage using an explicitly supplied Hadoop filesystem. + * <p> + * The supplied filesystem remains owned by the caller and is not closed by this storage instance. + * </p> + * @param fileSystem Hadoop filesystem used to store cache state + * @param rootPath root directory under which persisted cache state is stored + */ + public HadoopFsCachePersistenceStorage(FileSystem fileSystem, Path rootPath) { + this.fileSystem = Objects.requireNonNull(fileSystem, "fileSystem must not be null"); + this.rootPath = Objects.requireNonNull(rootPath, "rootPath must not be null"); + } + + /** + * Opens committed persisted state for the specified storage key. + * @param key stable storage key assigned to a persistent cache component + * @return input stream containing committed state, or an empty optional if no state exists + * @throws IOException if committed state exists but cannot be opened + */ + @Override + public Optional<InputStream> open(String key) throws IOException { + Path path = getStatePath(key); + if (!fileSystem.exists(path)) { + return Optional.empty(); + } + return Optional.of(fileSystem.open(path)); + } + + /** + * Starts a new persistence write for the specified storage key. + * <p> + * Data is written to a temporary file located next to the final state file. The temporary file + * becomes committed state only when the returned handle is explicitly committed. + * </p> + * @param key stable storage key assigned to a persistent cache component + * @return handle for writing and committing persisted state + * @throws IOException if the temporary persistence file cannot be created + */ + @Override + public CachePersistenceOutput create(String key) throws IOException { + Path targetPath = getStatePath(key); + Path parent = targetPath.getParent(); + ensureDirectory(parent); + + Path temporaryPath = createTemporaryPath(targetPath); + FSDataOutputStream output = fileSystem.create(temporaryPath, false); + + return new HadoopFsCachePersistenceOutput(fileSystem, output, temporaryPath, targetPath); + } + + /** + * Creates the parent directory for a persistence state file when necessary. + * @param directory directory to create + * @throws IOException if the directory cannot be created + */ + private void ensureDirectory(Path directory) throws IOException { + if (directory == null || fileSystem.exists(directory)) { + return; + } + + if (!fileSystem.mkdirs(directory) && !fileSystem.exists(directory)) { + throw new IOException("Failed to create cache persistence directory " + directory); + } + } + + /** + * Creates a unique temporary path adjacent to the specified committed state path. + * @param targetPath committed state path + * @return unique temporary path + */ + private Path createTemporaryPath(Path targetPath) { + return new Path(targetPath.toString() + TEMP_FILE_MARKER + UUID.randomUUID()); + } + + /** + * Resolves the Hadoop filesystem associated with the supplied persistence root path. + * @param conf Hadoop configuration used to resolve the filesystem + * @param rootPath persistence root path + * @return filesystem associated with the root path + * @throws IOException if the filesystem cannot be resolved + */ + private static FileSystem resolveFileSystem(Configuration conf, Path rootPath) + throws IOException { + Objects.requireNonNull(conf, "conf must not be null"); + Objects.requireNonNull(rootPath, "rootPath must not be null"); + return rootPath.getFileSystem(conf); + } + + /** + * Returns the filesystem path containing committed state for the supplied persistence key. + * @param key persistence storage key + * @return filesystem path containing committed persisted state + */ + private Path getStatePath(String key) { + validateKey(key); + return new Path(rootPath, key + STATE_FILE_SUFFIX); + } + + /** + * Validates a persistence storage key before using it as a relative filesystem path. + * <p> + * Storage keys must be relative and may contain multiple path segments. Empty segments and the + * special {@code .} and {@code ..} path segments are not allowed. + * </p> + * @param key persistence storage key to validate + * @throws NullPointerException if the key is {@code null} + * @throws IllegalArgumentException if the key is empty or contains an invalid path component + */ + private void validateKey(String key) { + Objects.requireNonNull(key, "key must not be null"); + + if (key.isEmpty()) { + throw new IllegalArgumentException("key must not be empty"); + } + if (key.startsWith("/") || key.endsWith("/") || key.indexOf('\\') >= 0 || key.contains("://")) { + throw new IllegalArgumentException("Invalid persistence key: " + key); + } + + String[] components = key.split("/"); + for (String component : components) { + if (component.isEmpty() || ".".equals(component) || "..".equals(component)) { + throw new IllegalArgumentException("Invalid persistence key: " + key); + } + } + } + + /** + * Persistence output backed by a temporary Hadoop filesystem file. + * <p> + * A successful commit closes the temporary output stream and publishes the temporary file as the + * committed state file. Closing an uncommitted instance aborts the write. + * </p> + */ + private static final class HadoopFsCachePersistenceOutput implements CachePersistenceOutput { + + private final FileSystem fileSystem; + private final FSDataOutputStream output; + private final Path temporaryPath; + private final Path targetPath; + + private boolean streamClosed; + private boolean committed; + private boolean aborted; + + /** + * Creates a persistence output backed by a temporary filesystem file. + * @param fileSystem filesystem containing the temporary and target files + * @param output output stream writing the temporary file + * @param temporaryPath temporary persistence file + * @param targetPath committed persistence file + */ + private HadoopFsCachePersistenceOutput(FileSystem fileSystem, FSDataOutputStream output, + Path temporaryPath, Path targetPath) { + this.fileSystem = Objects.requireNonNull(fileSystem, "fileSystem must not be null"); + this.output = Objects.requireNonNull(output, "output must not be null"); + this.temporaryPath = Objects.requireNonNull(temporaryPath, "temporaryPath must not be null"); + this.targetPath = Objects.requireNonNull(targetPath, "targetPath must not be null"); + } + + /** + * Returns the stream used to write temporary persisted state. + * @return persistence output stream + */ + @Override + public OutputStream getOutputStream() { + return output; + } + + /** + * Commits the persistence write by publishing the completed temporary file. + * <p> + * If previously committed state exists, it is moved temporarily out of the way before the new + * state is published. If publication of the new state fails, this method attempts to restore + * the previous committed state. + * </p> + * @throws IOException if the output cannot be closed or the new state cannot be published + */ + @Override + public void commit() throws IOException { + ensureActive(); + + closeOutput(); + + Path backupPath = null; + boolean previousStateMoved = false; + + try { + if (fileSystem.exists(targetPath)) { + backupPath = createBackupPath(targetPath); + if (!fileSystem.rename(targetPath, backupPath)) { + throw new IOException( + "Failed to preserve existing cache persistence state " + targetPath); + } + previousStateMoved = true; + } + + if (!fileSystem.rename(temporaryPath, targetPath)) { + throw new IOException( + "Failed to publish cache persistence state " + temporaryPath + " to " + targetPath); + } + + committed = true; + + if ( + previousStateMoved && fileSystem.exists(backupPath) + && !fileSystem.delete(backupPath, false) + ) { + throw new IOException( + "Persisted cache state was committed, but failed to delete backup " + backupPath); Review Comment: Backup cleanup can throw after the replacement is published and committed is true. The coordinator then reports a failed save and skips subsequent components, but close()/abort() leave the replacement in place. This contradicts CachePersistenceStorage's successful-commit contract. Treat post-publication cleanup as non-fatal, report the cleanup failure separately, and test both delete returning false and throwing IOException. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
