Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
12b1a94
fix: fail closed after buffered entrylog write failure
Aug 4, 2026
0b74a31
fix: fail closed on entry log flush failure
Aug 4, 2026
0a18155
test: add e2e coverage for entry log flush failure
Aug 4, 2026
3f42820
fix: close entrylog failure propagation gaps
yangxianjungree Aug 4, 2026
5e3674b
fix: keep fatal entrylog notifications compatible
yangxianjungree Aug 4, 2026
d761317
fix: adapt fatal entrylog logging to master
yangxianjungree Aug 5, 2026
8988984
test: release entrylog test buffers
yangxianjungree Aug 5, 2026
9088500
fix: make entrylog failure shutdown cleanup idempotent
yangxianjungree Aug 5, 2026
acff463
fix: continue db storage shutdown cleanup after flush failure
yangxianjungree Aug 5, 2026
198fa0f
style: order new imports to satisfy checkstyle
Aug 21, 2026
c7f0106
fix: clean up after entry log flush failure
yangxianjungree Sep 12, 2026
b650b47
fix: preserve checkpoint boundary after entry log failure
yangxianjungree Sep 13, 2026
2b41514
fix: preserve snapshots after parallel flush errors
yangxianjungree Sep 13, 2026
65562fa
fix: serialize entry log flush lifecycle
yangxianjungree Sep 13, 2026
50ad5d2
fix: recheck terminal entry log failure under flush lock
yangxianjungree Sep 13, 2026
0d12b32
fix: make per-ledger log rotation handoff atomic
yangxianjungree Sep 13, 2026
77e4a10
fix: propagate fatal entry log failures from compaction
yangxianjungree Sep 13, 2026
77712e6
test: exercise real storage checkpoint failure path
yangxianjungree Sep 13, 2026
ea95365
style: order checkpoint test imports
yangxianjungree Sep 13, 2026
4d98daf
fix: guard checkpoint fast path under flush lock
yangxianjungree Sep 13, 2026
cf66cc4
fix: close checkpoint completion race after entry log failure
yangxianjungree Sep 13, 2026
4dc68d5
fix: persist background entry log failure state
yangxianjungree Sep 13, 2026
c88e087
fix: retain fatal state from entry logger callbacks
yangxianjungree Sep 13, 2026
0853b77
fix: use lock protocol for fatal flush state
yangxianjungree Sep 14, 2026
747ae77
test: skip parallel flusher case when unsupported
yangxianjungree Sep 14, 2026
9640c20
fix: avoid lock inversion when recording fatal flush
yangxianjungree Sep 14, 2026
6fd21f7
fix: serialize fatal state with checkpoint completion
yangxianjungree Sep 14, 2026
33b04bd
fix: drain active flush before storage shutdown
yangxianjungree Sep 14, 2026
28d00cd
test: skip closed executor case when unsupported
yangxianjungree Sep 14, 2026
c8cc4b2
fix: wait for flushes before shutdown cleanup
yangxianjungree Sep 14, 2026
d59116f
fix: abort storage cleanup when gc shutdown is interrupted
yangxianjungree Sep 15, 2026
37ad080
fix: make entry log channel cleanup retryable
yangxianjungree Sep 15, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import java.time.Duration;
import java.util.Iterator;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Delayed;
import java.util.concurrent.ExecutionException;
Expand Down Expand Up @@ -120,6 +121,7 @@ void run() {

public MockExecutorController controlSubmit(ScheduledExecutorService service) {
doAnswer(answerNow()).when(service).submit(any(Runnable.class));
doAnswer(answerNowCallable()).when(service).submit(org.mockito.ArgumentMatchers.<Callable<Object>>any());
return this;
}

Expand Down Expand Up @@ -183,7 +185,23 @@ private static Answer<Future<?>> answerNow() {
SettableFuture<Void> future = SettableFuture.create();
future.set(null);
return future;
};
};
}

private static Answer<Future<?>> answerNowCallable() {
return invocationOnMock -> {
// Keep Callable submission semantics: task failures complete the Future exceptionally.
ThreadRegistry.forceClearRegistrationForTests(Thread.currentThread().getId());

Callable<?> task = invocationOnMock.getArgument(0);
SettableFuture<Object> future = SettableFuture.create();
try {
future.set(task.call());
} catch (Throwable t) {
future.setException(t);
}
return future;
};
}

private DeferredTask addDelayedTask(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,18 @@
*/
package org.apache.bookkeeper.common.testing.executors;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.fail;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;

import java.time.Duration;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import org.junit.Before;
Expand Down Expand Up @@ -57,6 +63,41 @@ public void testSubmit() {
verify(task, times(1)).run();
}

@Test
public void testSubmitCallable() throws Exception {
Callable<String> task = () -> "done";

Future<String> future = executor.submit(task);

assertEquals("done", future.get());
}

@Test
public void testSubmitCallableWithNullResult() throws Exception {
Callable<Object> task = () -> null;

Future<Object> future = executor.submit(task);

assertNull(future.get());
}

@Test
public void testSubmitCallableFailure() throws Exception {
RuntimeException failure = new RuntimeException("failure");
Callable<String> task = () -> {
throw failure;
};

Future<String> future = executor.submit(task);

try {
future.get();
fail("Expected the submitted Callable failure to be visible from Future.get()");
} catch (ExecutionException e) {
assertEquals(failure, e.getCause());
}
}

@Test
public void testExecute() {
Runnable task = mock(Runnable.class);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -500,6 +500,9 @@ public void ledgerDeleted(long ledgerId) {
ledgerStorage.setStateManager(stateManager);
ledgerStorage.setCheckpointSource(checkpointSource);
ledgerStorage.setCheckpointer(syncThread);
if (isDbLedgerStorage) {
((DbLedgerStorage) ledgerStorage).setFatalErrorListener(getLedgerDirsListener());
}
ledgerStorage.registerLedgerDeletionListener(ledgerDeletionListener);
handles = new HandleFactoryImpl(ledgerStorage);

Expand Down Expand Up @@ -871,6 +874,9 @@ public void run() {
// because shutdown can be called from sync thread which would be
// interrupted by shutdown call.
AtomicBoolean shutdownTriggered = new AtomicBoolean(false);
// Startup flush runs before stateManager.initState(), so isRunning() is false if it fails.
// Track shutdown independently to ensure the cleanup path still runs exactly once.
private final AtomicBoolean shutdownStarted = new AtomicBoolean(false);
void triggerBookieShutdown(final int exitCode) {
if (!shutdownTriggered.compareAndSet(false, true)) {
return;
Expand Down Expand Up @@ -899,7 +905,7 @@ public int shutdown() {
int shutdown(int exitCode) {
lock.lock();
try {
if (isRunning()) {
if (shutdownStarted.compareAndSet(false, true)) {
// the exitCode only set when first shutdown usually due to exception found
log.info()
.attr("bookiePort", conf.getBookiePort())
Expand Down Expand Up @@ -1033,6 +1039,9 @@ public void recoveryAddEntry(ByteBuf entry, WriteCallback cb, Object ctx, byte[]
addEntryInternal(handle, entry, false /* ackBeforeSync */, cb, ctx, masterKey);
}
success = true;
} catch (EntryLogWriteException e) {
triggerBookieShutdown(ExitCode.BOOKIE_EXCEPTION);
throw e;
} catch (NoWritableLedgerDirException e) {
stateManager.transitionToReadOnlyMode();
throw new IOException(e);
Expand Down Expand Up @@ -1125,6 +1134,9 @@ public void addEntry(ByteBuf entry, boolean ackBeforeSync, WriteCallback cb, Obj
addEntryInternal(handle, entry, ackBeforeSync, cb, ctx, masterKey);
}
success = true;
} catch (EntryLogWriteException e) {
triggerBookieShutdown(ExitCode.BOOKIE_EXCEPTION);
throw e;
} catch (NoWritableLedgerDirException e) {
stateManager.transitionToReadOnlyMode();
throw new IOException(e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,10 @@ public class BufferedChannel extends BufferedReadChannel implements Closeable {
* The buffer used to write operations.
*/
protected final ByteBuf writeBuffer;
// The file channel may fail to close after the buffer has already been released.
// Track the two resources independently so a later force-close can retry the file
// channel without releasing the buffer twice.
private volatile boolean writeBufferReleased;
/**
* The absolute position of the next write operation.
*/
Expand All @@ -71,6 +75,7 @@ public class BufferedChannel extends BufferedReadChannel implements Closeable {
protected final AtomicLong unpersistedBytes;

private boolean closed = false;
private volatile IOException writeFailure;

// make constructor to be public for unit test
public BufferedChannel(ByteBufAllocator allocator, FileChannel fc, int capacity) throws IOException {
Expand Down Expand Up @@ -101,7 +106,10 @@ public synchronized void close() throws IOException {
if (closed) {
return;
}
ReferenceCountUtil.release(writeBuffer);
if (!writeBufferReleased) {
ReferenceCountUtil.release(writeBuffer);
writeBufferReleased = true;
}
fileChannel.close();
closed = true;
}
Expand All @@ -117,8 +125,14 @@ public synchronized void close() throws IOException {
public void write(ByteBuf src) throws IOException {
boolean shouldForceWrite = false;
synchronized (this) {
int copied = copyIntoWriteBuffer(src);
shouldForceWrite = updatePositionAndFlushIfNeeded(copied);
checkWritable();
try {
int copied = copyIntoWriteBuffer(src);
shouldForceWrite = updatePositionAndFlushIfNeeded(copied);
} catch (IOException e) {
markWriteFailure(e);
throw e;
}
}
if (shouldForceWrite) {
forceWrite(false);
Expand All @@ -137,9 +151,15 @@ public void write(ByteBuf src) throws IOException {
public void write(ByteBuf src1, ByteBuf src2) throws IOException {
boolean shouldForceWrite = false;
synchronized (this) {
int copied = copyIntoWriteBuffer(src1);
copied += copyIntoWriteBuffer(src2);
shouldForceWrite = updatePositionAndFlushIfNeeded(copied);
checkWritable();
try {
int copied = copyIntoWriteBuffer(src1);
copied += copyIntoWriteBuffer(src2);
shouldForceWrite = updatePositionAndFlushIfNeeded(copied);
} catch (IOException e) {
markWriteFailure(e);
throw e;
}
}
if (shouldForceWrite) {
forceWrite(false);
Expand Down Expand Up @@ -203,6 +223,7 @@ public long getFileChannelPosition() {
* @throws IOException
*/
public void flushAndForceWrite(boolean forceMetadata) throws IOException {
checkWritable();
flush();
forceWrite(forceMetadata);
}
Expand All @@ -217,6 +238,7 @@ public void flushAndForceWrite(boolean forceMetadata) throws IOException {
* @throws IOException
*/
public void flushAndForceWriteIfRegularFlush(boolean forceMetadata) throws IOException {
checkWritable();
if (doRegularFlushes) {
flushAndForceWrite(forceMetadata);
}
Expand All @@ -229,10 +251,19 @@ public void flushAndForceWriteIfRegularFlush(boolean forceMetadata) throws IOExc
* @throws IOException if the write fails.
*/
public synchronized void flush() throws IOException {
checkWritable();
ByteBuffer toWrite = writeBuffer.internalNioBuffer(0, writeBuffer.writerIndex());
do {
fileChannel.write(toWrite);
} while (toWrite.hasRemaining());
try {
while (toWrite.hasRemaining()) {
int written = fileChannel.write(toWrite);
if (written <= 0) {
throw new IOException("Unable to make progress while flushing buffered channel");
}
}
} catch (IOException e) {
markWriteFailure(e);
throw e;
}
writeBuffer.clear();
writeBufferStartPosition.set(fileChannel.position());
}
Expand All @@ -244,6 +275,7 @@ public synchronized void flush() throws IOException {
* @throws IOException
*/
public long forceWrite(boolean forceMetadata) throws IOException {
checkWritable();
// This is the point up to which we had flushed to the file system page cache
// before issuing this force write hence is guaranteed to be made durable by
// the force write, any flush that happens after this may or may
Expand All @@ -270,7 +302,12 @@ public long forceWrite(boolean forceMetadata) throws IOException {
}
}

fileChannel.force(forceMetadata);
try {
fileChannel.force(forceMetadata);
} catch (IOException e) {
markWriteFailure(e);
throw e;
}
return positionForceWrite;
}

Expand Down Expand Up @@ -333,4 +370,23 @@ public synchronized int getNumOfBytesInWriteBuffer() {
long getUnpersistedBytes() {
return unpersistedBytes.get();
}

final void checkWritable() throws IOException {
IOException failure = writeFailure;
if (failure != null) {
throw new IOException("BufferedChannel is in failed state", failure);
}
// close() releases the write buffer before attempting the file-channel close.
// If that close fails, forceClose() must be able to retry the file channel, but
// no caller may continue writing through the already released buffer.
if (writeBufferReleased) {
throw new IOException("BufferedChannel is closed");
}
}

final void markWriteFailure(IOException e) {
if (writeFailure == null) {
writeFailure = e;
}
}
Comment on lines +374 to +391
}
Original file line number Diff line number Diff line change
Expand Up @@ -56,8 +56,10 @@
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.regex.Pattern;
import lombok.CustomLog;
import org.apache.bookkeeper.bookie.LedgerDirsManager.LedgerDirsListener;
import org.apache.bookkeeper.bookie.storage.CompactionEntryLog;
import org.apache.bookkeeper.bookie.storage.EntryLogScanner;
import org.apache.bookkeeper.bookie.storage.EntryLogger;
Expand Down Expand Up @@ -136,6 +138,7 @@ public String toString() {
* Updates the entry log file header with the offset and size of the map.
*/
void appendLedgersMap() throws IOException {
checkWritable();

long ledgerMapOffset = this.position();

Expand Down Expand Up @@ -204,7 +207,23 @@ public void accept(long ledgerId, long size) {
mapInfo.putLong(ledgerMapOffset);
mapInfo.putInt(numberOfLedgers);
mapInfo.flip();
this.fileChannel.write(mapInfo, LEDGERS_MAP_OFFSET_POSITION);
try {
writeFully(this.fileChannel, mapInfo, LEDGERS_MAP_OFFSET_POSITION);
} catch (IOException e) {
markWriteFailure(e);
throw e;
}
}

private static void writeFully(FileChannel fileChannel, ByteBuffer buffer, long position) throws IOException {
long writePosition = position;
while (buffer.hasRemaining()) {
int written = fileChannel.write(buffer, writePosition);
if (written <= 0) {
throw new IOException("Unable to make progress while updating entry log header");
}
writePosition += written;
}
}
Comment on lines +218 to 227
}

Expand All @@ -222,6 +241,7 @@ public void accept(long ledgerId, long size) {

final EntryLoggerAllocator entryLoggerAllocator;
private final EntryLogManager entryLogManager;
private final AtomicBoolean closed = new AtomicBoolean(false);

private final CopyOnWriteArrayList<EntryLogListener> listeners = new CopyOnWriteArrayList<EntryLogListener>();

Expand Down Expand Up @@ -365,6 +385,10 @@ EntryLogManager getEntryLogManager() {
return entryLogManager;
}

public void setFatalErrorListener(LedgerDirsListener fatalErrorListener) {
entryLogManager.setFatalErrorListener(fatalErrorListener);
}

void addListener(EntryLogListener listener) {
if (null != listener) {
listeners.add(listener);
Expand Down Expand Up @@ -1208,6 +1232,10 @@ public boolean accept(long ledgerId) {
*/
@Override
public void close() {
if (!closed.compareAndSet(false, true)) {
log.debug("EntryLogger is already stopped");
return;
}
// since logChannel is buffered channel, do flush when shutting down
log.info("Stopping EntryLogger");
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,17 @@ public boolean compact(EntryLogMetadata entryLogMeta) {
} catch (LedgerDirsManager.NoWritableLedgerDirException nwlde) {
log.warn().exception(nwlde).log("No writable ledger directory available, aborting compaction");
return false;
} catch (EntryLogWriteException elwe) {
// A poisoned live writer is a bookie-fatal condition, not a
// recoverable compaction failure.
if (entryLogger instanceof DefaultEntryLogger) {
EntryLogManager manager = ((DefaultEntryLogger) entryLogger).getEntryLogManager();
if (manager instanceof EntryLogManagerBase) {
((EntryLogManagerBase) manager).notifyFatalEntryLogWriteFailure(
"Fatal entry log write failure while compacting entry log", elwe);
}
}
return false;
} catch (IOException ioe) {
// if compact entry log throws IOException, we don't want to remove that
// entry log. however, if some entries from that log have been re-added
Expand Down
Loading