diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java index a1f27ff62f24d..a036ee3a5089b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/AbstractWALBuffer.java @@ -31,6 +31,7 @@ import java.io.File; import java.io.IOException; +import java.nio.file.FileAlreadyExistsException; import java.nio.file.Files; import java.nio.file.StandardCopyOption; import java.util.Arrays; @@ -56,6 +57,12 @@ public abstract class AbstractWALBuffer implements IWALBuffer { @SuppressWarnings("squid:S3077") protected volatile WALWriter currentWALFileWriter; + // Only the sync thread accesses this state. Once sealed, the old WAL must never be written or + // closed again, even if renaming it or creating its successor fails because the disk is full. + private File pendingRollFile; + private WALFileStatus pendingRollStatus; + private long pendingRollSearchIndex; + protected AbstractWALBuffer( String identifier, String logDirectory, long startFileVersion, long startSearchIndex) throws IOException { @@ -99,18 +106,22 @@ public long getCurrentWALOriginalFileSize() { * @throws IOException If failing to close or open the log writer */ protected File rollLogWriter(long searchIndex, WALFileStatus fileStatus) throws IOException { - // close file - currentWALFileWriter.close(); - addDiskUsage(currentWALFileWriter.size()); - addFileNum(1); - File lastFile = currentWALFileWriter.getLogFile(); + if (!hasPendingRoll()) { + // Record the boundary only after sealing and forcing the old WAL have both succeeded. + currentWALFileWriter.close(); + pendingRollFile = currentWALFileWriter.getLogFile(); + pendingRollStatus = fileStatus; + pendingRollSearchIndex = searchIndex; + addDiskUsage(currentWALFileWriter.size()); + } + File lastFile = pendingRollFile; String lastName = lastFile.getName(); - if (WALFileUtils.parseStatusCode(lastName) != fileStatus) { + if (WALFileUtils.parseStatusCode(lastName) != pendingRollStatus) { String targetName = WALFileUtils.getLogFileName( WALFileUtils.parseVersionId(lastName), WALFileUtils.parseStartSearchIndex(lastName), - fileStatus); + pendingRollStatus); File targetFile = SystemFileFactory.INSTANCE.getFile(logDirectory, targetName); Files.move( lastFile.toPath(), @@ -118,6 +129,7 @@ protected File rollLogWriter(long searchIndex, WALFileStatus fileStatus) throws StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE); lastFile = targetFile; + pendingRollFile = targetFile; } // roll file long nextFileVersion = currentWALFileVersion + 1; @@ -125,13 +137,24 @@ protected File rollLogWriter(long searchIndex, WALFileStatus fileStatus) throws SystemFileFactory.INSTANCE.getFile( logDirectory, WALFileUtils.getLogFileName( - nextFileVersion, searchIndex, WALFileStatus.CONTAINS_SEARCH_INDEX)); + nextFileVersion, pendingRollSearchIndex, WALFileStatus.CONTAINS_SEARCH_INDEX)); + // A failed header write may leave a partial successor. Do not append to it on retry, or + // overwrite an unexpected existing WAL; either requires recovery rather than online rotation. + if (nextLogFile.length() > 0) { + throw new FileAlreadyExistsException(nextLogFile.toString()); + } currentWALFileWriter = new WALWriter(nextLogFile); currentWALFileVersion = nextFileVersion; + addFileNum(1); + pendingRollFile = null; logger.debug(StorageEngineMessages.OPEN_NEW_WAL_FILE_FOR_BUFFER, nextLogFile, identifier); return lastFile; } + protected boolean hasPendingRoll() { + return pendingRollFile != null; + } + public long getDiskUsage() { return diskUsage; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java index d67cd88a4c48f..5cee2ff26753e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBuffer.java @@ -126,6 +126,12 @@ public class WALBuffer extends AbstractWALBuffer { private final Map> memTableIdsOfWal = new ConcurrentHashMap<>(); private final BiConsumer walFileRolledListener; + // An entry may span several sync tasks. Never acknowledge its final chunk after an earlier + // chunk failed. Failures before writing a pending successor can be cleared at the batch boundary; + // failures while writing the active WAL require recovery because its record boundary is unknown. + private Exception syncFailure; + private boolean retryAfterFailedBatch; + public WALBuffer(String identifier, String logDirectory) throws IOException { this( identifier, @@ -585,6 +591,45 @@ public SyncBufferTask( public void run() { final long startTime = System.nanoTime(); + if (syncFailure != null) { + failListeners(syncFailure); + // SET SYSTEM TO RUNNING does not repair a failed batch or an unknown record boundary. + if (CommonDescriptor.getInstance().getConfig().isRunning()) { + CommonDescriptor.getInstance().getConfig().handleUnrecoverableError(); + } + if (forceFlag && retryAfterFailedBatch) { + syncFailure = null; + } + switchSyncingBufferToIdle(); + return; + } + + boolean resumedRoll = false; + final boolean hasData = syncingBuffer.position() > 0; + try { + if (hasPendingRoll()) { + // The previous task sealed the old file. Finish opening its successor before touching + // this buffer, so no bytes or metadata can be appended after the old WAL's end marker. + rollLogWriter(searchIndex, fileStatus); + resumedRoll = true; + } + } catch (IOException e) { + logger.error( + StorageEngineMessages + .STORAGE_LOG_FAIL_TO_ROLL_WAL_NODE_S_LOG_WRITER_CHANGE_SYSTEM_MODE_TO_A384AA54, + identifier, + e); + if (!forceFlag) { + syncFailure = e; + retryAfterFailedBatch = true; + } + failListeners(e); + DataNodeExceptionMetrics.getInstance().recordSuspiciousDiskException(e); + CommonDescriptor.getInstance().getConfig().handleUnrecoverableError(); + switchSyncingBufferToIdle(); + return; + } + makeMemTableCheckpoints(); long walFileVersionId = currentWALFileVersion; @@ -610,7 +655,11 @@ public void run() { .STORAGE_LOG_FAIL_TO_SYNC_WAL_NODE_S_BUFFER_CHANGE_SYSTEM_MODE_TO_ERROR_8C379D57, identifier, e); + syncFailure = e instanceof Exception exception ? exception : new IOException(e); + retryAfterFailedBatch = false; + failListeners(syncFailure); CommonDescriptor.getInstance().getConfig().handleUnrecoverableError(); + return; } finally { switchSyncingBufferToIdle(); } @@ -623,24 +672,24 @@ public void run() { boolean forceSuccess = false; // try to roll log writer - if (info.rollWALFileWriterListener != null + if ((info.rollWALFileWriterListener != null && (!resumedRoll || hasData)) // TODO: Control the wal file by the number of WALEntry || (forceFlag && currentWALFileWriter.originalSize() >= config.getWalFileSizeThresholdInByte())) { try { rollLogWriter(searchIndex, currentWALFileWriter.getWalFileStatus()); forceSuccess = true; - if (info.rollWALFileWriterListener != null) { - info.rollWALFileWriterListener.succeed(); - } } catch (IOException e) { logger.error( StorageEngineMessages .STORAGE_LOG_FAIL_TO_ROLL_WAL_NODE_S_LOG_WRITER_CHANGE_SYSTEM_MODE_TO_A384AA54, identifier, e); - if (info.rollWALFileWriterListener != null) { - info.rollWALFileWriterListener.fail(e); + failListeners(e); + if (!hasPendingRoll()) { + // A failed seal has no known durable boundary from which to resume rotation. + syncFailure = e; + retryAfterFailedBatch = false; } DataNodeExceptionMetrics.getInstance().recordSuspiciousDiskException(e); CommonDescriptor.getInstance().getConfig().handleUnrecoverableError(); @@ -657,15 +706,18 @@ public void run() { identifier, e); DataNodeExceptionMetrics.getInstance().recordSuspiciousDiskException(e); - for (WALFlushListener fsyncListener : info.fsyncListeners) { - fsyncListener.fail(e); - } + failListeners(e); + syncFailure = e; + retryAfterFailedBatch = false; CommonDescriptor.getInstance().getConfig().handleUnrecoverableError(); } } // notify all waiting listeners if (forceSuccess) { + if (info.rollWALFileWriterListener != null) { + info.rollWALFileWriterListener.succeed(); + } for (WALFlushListener fsyncListener : info.fsyncListeners) { fsyncListener.succeed(); } @@ -675,6 +727,15 @@ public void run() { WRITING_METRICS.recordSyncWALBufferCost(System.nanoTime() - startTime, forceFlag); } + private void failListeners(Exception e) { + if (info.rollWALFileWriterListener != null) { + info.rollWALFileWriterListener.fail(e); + } + for (WALFlushListener fsyncListener : info.fsyncListeners) { + fsyncListener.fail(e); + } + } + private void makeMemTableCheckpoints() { if (info.checkpoints.isEmpty()) { return; @@ -767,7 +828,7 @@ public void close() { shutdownThread(syncBufferThread, ThreadName.WAL_SYNC); } - if (currentWALFileWriter != null) { + if (currentWALFileWriter != null && !hasPendingRoll()) { try { currentWALFileWriter.close(); } catch (IOException e) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/LogWriter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/LogWriter.java index 52675dae37e7d..3652c826682dd 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/LogWriter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/LogWriter.java @@ -20,20 +20,16 @@ package org.apache.iotdb.db.storageengine.dataregion.wal.io; import org.apache.iotdb.db.conf.IoTDBDescriptor; -import org.apache.iotdb.db.i18n.StorageEngineMessages; import org.apache.iotdb.db.service.metrics.WritingMetrics; import org.apache.iotdb.db.storageengine.dataregion.wal.buffer.WALEntry; import org.apache.iotdb.db.storageengine.dataregion.wal.checkpoint.Checkpoint; import org.apache.tsfile.compress.ICompressor; import org.apache.tsfile.file.metadata.enums.CompressionType; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import java.io.File; import java.io.IOException; import java.nio.ByteBuffer; -import java.nio.channels.ClosedChannelException; import java.nio.channels.FileChannel; import java.nio.file.StandardOpenOption; @@ -42,8 +38,6 @@ * and writing {@link Checkpoint} into .checkpoint file. */ public abstract class LogWriter implements ILogWriter { - private static final Logger logger = LoggerFactory.getLogger(LogWriter.class); - protected final File logFile; protected final FileChannel logChannel; protected long originalSize = 0; @@ -73,9 +67,23 @@ protected LogWriter(File logFile, WALFileVersion version) throws IOException { StandardOpenOption.CREATE, StandardOpenOption.WRITE, StandardOpenOption.APPEND); - if ((!logFile.exists() || logFile.length() == 0) - && (version == WALFileVersion.V2 || version == WALFileVersion.V3)) { - this.logChannel.write(ByteBuffer.wrap(version.getVersionBytes())); + try { + if (logChannel.size() == 0 + && (version == WALFileVersion.V2 || version == WALFileVersion.V3)) { + ByteBuffer magic = ByteBuffer.wrap(version.getVersionBytes()); + while (magic.hasRemaining()) { + logChannel.write(magic); + } + } + } catch (IOException e) { + // A full disk can fail initialization after open() succeeds. Release the orphan channel + // before the owner retries creating the successor. + try { + logChannel.close(); + } catch (IOException closeException) { + e.addSuppressed(closeException); + } + throw e; } } @@ -124,12 +132,12 @@ public double write(ByteBuffer buffer, boolean allowCompress) throws IOException WritingMetrics.getInstance().recordCompressWALBufferCost(System.nanoTime() - startTime); } startTime = System.nanoTime(); - try { - headerBuffer.flip(); + headerBuffer.flip(); + while (headerBuffer.hasRemaining()) { logChannel.write(headerBuffer); + } + while (buffer.hasRemaining()) { logChannel.write(buffer); - } catch (ClosedChannelException e) { - logger.warn(StorageEngineMessages.CANNOT_WRITE_TO, logFile, e); } WritingMetrics.getInstance() .recordWroteWALBuffer(uncompressedSize, bufferSize, System.nanoTime() - startTime); @@ -149,9 +157,8 @@ public void force() throws IOException { @Override public void force(boolean metaData) throws IOException { - if (logChannel != null && logChannel.isOpen()) { - logChannel.force(metaData); - } + // A closed channel is a failed durability operation, not a successful no-op. + logChannel.force(metaData); } @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALWriter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALWriter.java index 10d164f3851cd..37e32d85beb96 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALWriter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALWriter.java @@ -90,7 +90,10 @@ private synchronized void endFile() throws IOException { private void writeMetadata(ByteBuffer buffer) throws IOException { buffer.flip(); - logChannel.write(buffer); + // A successful seal is the recovery boundary for switching to the next WAL file. + while (buffer.hasRemaining()) { + logChannel.write(buffer); + } } @Override diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBufferRollRecoveryTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBufferRollRecoveryTest.java new file mode 100644 index 0000000000000..ab081fa138759 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/buffer/WALBufferRollRecoveryTest.java @@ -0,0 +1,300 @@ +/* + * 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.iotdb.db.storageengine.dataregion.wal.buffer; + +import org.apache.iotdb.commons.cluster.NodeStatus; +import org.apache.iotdb.commons.conf.CommonConfig; +import org.apache.iotdb.commons.conf.CommonDescriptor; +import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode; +import org.apache.iotdb.db.storageengine.dataregion.wal.checkpoint.CheckpointManager; +import org.apache.iotdb.db.storageengine.dataregion.wal.io.WALMetaData; +import org.apache.iotdb.db.storageengine.dataregion.wal.io.WALReader; +import org.apache.iotdb.db.storageengine.dataregion.wal.io.WALWriter; +import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileStatus; +import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALFileUtils; +import org.apache.iotdb.db.storageengine.dataregion.wal.utils.listener.AbstractResultListener.Status; + +import org.apache.tsfile.common.conf.TSFileConfig; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.utils.Binary; +import org.apache.tsfile.write.schema.MeasurementSchema; +import org.awaitility.Awaitility; +import org.junit.After; +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.TemporaryFolder; + +import java.io.File; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.channels.ClosedChannelException; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; + +import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doCallRealMethod; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +public class WALBufferRollRecoveryTest { + @Rule public TemporaryFolder temporaryFolder = new TemporaryFolder(); + + private final CommonConfig commonConfig = CommonDescriptor.getInstance().getConfig(); + private final List sealedFiles = new CopyOnWriteArrayList<>(); + private NodeStatus previousStatus; + private String previousReason; + private File directory; + private WALBuffer buffer; + private WALWriter firstWriter; + + @Before + public void setUp() throws Exception { + previousStatus = commonConfig.getNodeStatus(); + previousReason = commonConfig.getStatusReason(); + commonConfig.setNodeStatus(NodeStatus.Running); + directory = temporaryFolder.newFolder("wal"); + buffer = + new WALBuffer( + "roll-recovery", + directory.getPath(), + new CheckpointManager("roll-recovery", directory.getPath()), + 0, + 0, + (sealedFile, currentFile) -> sealedFiles.add(sealedFile)); + firstWriter = spy(buffer.currentWALFileWriter); + buffer.currentWALFileWriter = firstWriter; + } + + @After + public void tearDown() throws Exception { + try { + // Tests may inject a failure before close(); restore cleanup without reopening any file. + doCallRealMethod().when(firstWriter).close(); + buffer.close(); + } finally { + commonConfig.setNodeStatus(previousStatus); + commonConfig.setStatusReason(previousReason); + } + } + + /** + * Failed successor creation must not reseal the old WAL, inflate counters, or lose later writes. + */ + @Test + public void testResumeBeforeWritingAfterRepeatedOpenFailures() throws Exception { + writeAndAwait(entry(1, 1, "before"), Status.SUCCESS); + File successor = walFile(1, 1, WALFileStatus.CONTAINS_SEARCH_INDEX); + blockPath(successor); + roll(Status.FAILURE); + awaitReadOnly(); + byte[] sealedBytes = Files.readAllBytes(firstWriter.getLogFile().toPath()); + long diskUsage = buffer.getDiskUsage(); + roll(Status.FAILURE); + assertEquals(0, buffer.getCurrentWALFileVersion()); + assertEquals(1, buffer.getFileNum()); + assertEquals(diskUsage, buffer.getDiskUsage()); + assertTrue(sealedFiles.isEmpty()); + verify(firstWriter, times(1)).close(); + + unblockPath(successor); + commonConfig.setNodeStatus(NodeStatus.Running); + writeAndAwait(entry(2, 2, "after"), Status.SUCCESS); + assertEquals(NodeStatus.Running, commonConfig.getNodeStatus()); + assertEquals(1, buffer.getCurrentWALFileVersion()); + assertEquals(successor, buffer.currentWALFileWriter.getLogFile()); + assertEquals(2, buffer.getFileNum()); + assertEquals(diskUsage, buffer.getDiskUsage()); + assertEquals(Arrays.asList(firstWriter.getLogFile()), sealedFiles); + assertArrayEquals(sealedBytes, Files.readAllBytes(firstWriter.getLogFile().toPath())); + assertEquals(Arrays.asList(1L), readTimes(firstWriter.getLogFile())); + roll(Status.SUCCESS); + assertEquals(Arrays.asList(2L), readTimes(successor)); + } + + /** A retry containing only the roll signal should open exactly one successor. */ + @Test + public void testRollSignalResumesWithoutRollingTwice() throws Exception { + writeAndAwait(entry(1, 1, "before"), Status.SUCCESS); + File successor = walFile(1, 1, WALFileStatus.CONTAINS_SEARCH_INDEX); + blockPath(successor); + roll(Status.FAILURE); + unblockPath(successor); + roll(Status.SUCCESS); + assertEquals(1, buffer.getCurrentWALFileVersion()); + assertEquals(2, buffer.getFileNum()); + assertEquals(1, sealedFiles.size()); + verify(firstWriter, times(1)).close(); + } + + /** Failed rename and then failed open must each resume at the saved stage and notify once. */ + @Test + public void testResumeRenameThenOpen() throws Exception { + writeAndAwait(entry(1, -1, "unindexed"), Status.SUCCESS); + File renamed = walFile(0, 0, WALFileStatus.CONTAINS_NONE_SEARCH_INDEX); + File successor = walFile(1, 0, WALFileStatus.CONTAINS_SEARCH_INDEX); + blockPath(renamed); + roll(Status.FAILURE); + long diskUsage = buffer.getDiskUsage(); + unblockPath(renamed); + blockPath(successor); + roll(Status.FAILURE); + assertFalse(firstWriter.getLogFile().exists()); + byte[] sealedBytes = Files.readAllBytes(renamed.toPath()); + unblockPath(successor); + roll(Status.SUCCESS); + assertEquals(Arrays.asList(renamed), sealedFiles); + assertArrayEquals(sealedBytes, Files.readAllBytes(renamed.toPath())); + assertEquals(Arrays.asList(1L), readTimes(renamed)); + assertEquals(diskUsage, buffer.getDiskUsage()); + assertEquals(2, buffer.getFileNum()); + verify(firstWriter, times(1)).close(); + } + + /** Do not append to a partial header or an unexpected nonempty successor left on disk. */ + @Test + public void testNonemptySuccessorIsNotOverwritten() throws Exception { + writeAndAwait(entry(1, 1, "before"), Status.SUCCESS); + File successor = walFile(1, 1, WALFileStatus.CONTAINS_SEARCH_INDEX); + byte[] partialHeader = new byte[] {1, 2}; + Files.write(successor.toPath(), partialHeader); + roll(Status.FAILURE); + roll(Status.FAILURE); + assertArrayEquals(partialHeader, Files.readAllBytes(successor.toPath())); + assertEquals(0, buffer.getCurrentWALFileVersion()); + verify(firstWriter, times(1)).close(); + } + + /** A failed first chunk must fail the whole large entry; the next batch can recover cleanly. */ + @Test + public void testFailedSplitEntryDoesNotLeakIntoSuccessor() throws Exception { + writeAndAwait(entry(1, 1, "before"), Status.SUCCESS); + File successor = walFile(1, 1, WALFileStatus.CONTAINS_SEARCH_INDEX); + blockPath(successor); + roll(Status.FAILURE); + buffer.setBufferSize(192); + writeAndAwait(entry(2, 2, new String(new char[4096]).replace('\0', 'x')), Status.FAILURE); + unblockPath(successor); + writeAndAwait(entry(3, 3, "after"), Status.SUCCESS); + roll(Status.SUCCESS); + assertEquals(Arrays.asList(3L), readTimes(successor)); + assertEquals(Arrays.asList(1L), readTimes(firstWriter.getLogFile())); + } + + /** A failed seal has no proven durable boundary, so an empty successor must not hide it. */ + @Test + public void testSealFailureIsNotSkipped() throws Exception { + writeAndAwait(entry(1, 1, "before"), Status.SUCCESS); + doThrow(new ClosedChannelException()).when(firstWriter).close(); + roll(Status.FAILURE); + roll(Status.FAILURE); + assertEquals(0, buffer.getCurrentWALFileVersion()); + assertEquals(1, buffer.getFileNum()); + assertTrue(sealedFiles.isEmpty()); + verify(firstWriter, times(1)).close(); + } + + /** A write failure must never be turned into success by a subsequent force or roll task. */ + @Test + public void testWriteFailureFailsSubsequentListeners() throws Exception { + doThrow(new ClosedChannelException()) + .when(firstWriter) + .write(any(ByteBuffer.class), any(WALMetaData.class)); + writeAndAwait(entry(1, 1, "failed"), Status.FAILURE); + awaitReadOnly(); + commonConfig.setNodeStatus(NodeStatus.Running); + roll(Status.FAILURE); + awaitReadOnly(); + writeAndAwait(entry(2, 2, "also failed"), Status.FAILURE); + assertEquals(0, buffer.getCurrentWALFileVersion()); + assertTrue(sealedFiles.isEmpty()); + verify(firstWriter, times(1)).write(any(ByteBuffer.class), any(WALMetaData.class)); + } + + private WALInfoEntry entry(long time, long searchIndex, String value) throws Exception { + InsertRowNode node = + new InsertRowNode( + new PlanNodeId(""), + new PartialPath("root.test.d"), + false, + new String[] {"s"}, + new TSDataType[] {TSDataType.TEXT}, + time, + new Object[] {new Binary(value, TSFileConfig.STRING_CHARSET)}, + false); + node.setMeasurementSchemas( + new MeasurementSchema[] {new MeasurementSchema("s", TSDataType.TEXT)}); + node.setSearchIndex(searchIndex); + return new WALInfoEntry(1, node, false); + } + + private void writeAndAwait(WALEntry entry, Status expected) { + buffer.write(entry); + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(expected, entry.getWalFlushListener().waitForResult())); + } + + private void roll(Status expected) { + writeAndAwait(new WALSignalEntry(WALEntryType.ROLL_WAL_LOG_WRITER_SIGNAL, false), expected); + } + + private void awaitReadOnly() { + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .untilAsserted(() -> assertEquals(NodeStatus.ReadOnly, commonConfig.getNodeStatus())); + } + + private File walFile(long version, long searchIndex, WALFileStatus status) { + return new File(directory, WALFileUtils.getLogFileName(version, searchIndex, status)); + } + + private void blockPath(File path) throws IOException { + Files.createDirectory(path.toPath()); + // A nonempty directory blocks both file creation and replacement on Windows and Unix. + Files.write(path.toPath().resolve("blocker"), new byte[] {1}); + } + + private void unblockPath(File path) throws IOException { + Files.delete(path.toPath().resolve("blocker")); + Files.delete(path.toPath()); + } + + private List readTimes(File file) throws IOException { + List times = new ArrayList<>(); + try (WALReader reader = new WALReader(file)) { + while (reader.hasNext()) { + times.add(((InsertRowNode) reader.next().getValue()).getTime()); + } + } + return times; + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALFileTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALFileTest.java index 16e9bb36f8714..766e9cb9046af 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALFileTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/wal/io/WALFileTest.java @@ -50,6 +50,7 @@ import java.io.File; import java.io.IOException; import java.nio.ByteBuffer; +import java.nio.channels.ClosedChannelException; import java.nio.channels.FileChannel; import java.nio.file.Files; import java.util.ArrayList; @@ -57,6 +58,7 @@ import java.util.Collections; import java.util.List; +import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertThrows; @@ -86,6 +88,19 @@ public void tearDown() throws Exception { } } + /** Unexpected channel closure must propagate to the buffer instead of acknowledging a write. */ + @Test + public void testClosedChannelWriteAndForceFail() throws IOException { + WALWriter writer = new WALWriter(walFile); + writer.logChannel.close(); + byte[] before = Files.readAllBytes(walFile.toPath()); + ByteBuffer buffer = ByteBuffer.allocate(1); + buffer.put((byte) 1); + assertThrows(ClosedChannelException.class, () -> writer.write(buffer)); + assertThrows(ClosedChannelException.class, writer::force); + assertArrayEquals(before, Files.readAllBytes(walFile.toPath())); + } + @Test public void testReadNormalFile() throws IOException, IllegalPathException { int fakeMemTableId = 1;