Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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 @@ -285,6 +285,9 @@ private record OffloadThresholds(long thresholdInBytes, long thresholdInSeconds)
// ledger metadata stored. This variable used to record the result of the latest once re-checking.
private Pair<Long, CompletableFuture<Integer>> ledgerRecheckInProgress = new ImmutablePair<>(-1L,
CompletableFuture.completedFuture(Code.OK));
// Accessed only while holding this object's monitor.
private LedgerHandle emptyLedgerMetadataCheckLedger;
private boolean emptyLedgerRolloverInProgress;

public enum State {
None, // Uninitialized
Expand Down Expand Up @@ -863,7 +866,7 @@ protected synchronized void internalAsyncAddEntry(OpAddEntry addOperation) {
if (!beforeAddEntry(addOperation)) {
return;
}
final State state = STATE_UPDATER.get(this);
State state = STATE_UPDATER.get(this);
if (state.isFenced()) {
addOperation.failed(new ManagedLedgerFencedException());
return;
Expand All @@ -879,7 +882,19 @@ protected synchronized void internalAsyncAddEntry(OpAddEntry addOperation) {
}
pendingAddEntries.add(addOperation);

if (state == State.ClosingLedger || state == State.CreatingLedger) {
if (state == State.LedgerOpened && currentLedgerEntries == 0
&& ledgers.containsKey(currentLedger.getId())) {
if (maximumRolloverTimeReached()) {
startEmptyLedgerRollover(currentLedger, null);
} else if (emptyLedgerMetadataCheckLedger == null) {
emptyLedgerMetadataCheckLedger = currentLedger;
checkEmptyLedgerMetadata(currentLedger);
}
state = STATE_UPDATER.get(this);
}

if (state == State.ClosingLedger || state == State.CreatingLedger
|| emptyLedgerMetadataCheckLedger != null) {
// We don't have a ready ledger to write into
// We are waiting for a new ledger to be created
log.debug("Queue addEntry request");
Expand Down Expand Up @@ -1536,6 +1551,11 @@ public synchronized void asyncTerminate(TerminateCallback callback, Object ctx)

log.info("Terminating managed ledger");
state = State.Terminated;
if (emptyLedgerMetadataCheckLedger != null || emptyLedgerRolloverInProgress) {
emptyLedgerMetadataCheckLedger = null;
emptyLedgerRolloverInProgress = false;
clearPendingAddEntries(new ManagedLedgerTerminatedException("Managed ledger was terminated"));
}

LedgerHandle lh = currentLedger;
log.debug().attr("ledgerId", lh.getId()).log("Closing current writing ledger");
Expand Down Expand Up @@ -1874,6 +1894,8 @@ void createNewOpAddEntryForNewLedger() {

protected synchronized void updateLedgersIdsComplete(@Nullable LedgerHandle originalCurrentLedger) {
STATE_UPDATER.set(this, State.LedgerOpened);
emptyLedgerMetadataCheckLedger = null;
emptyLedgerRolloverInProgress = false;
// Delete original "currentLedger" if it has been removed from "ledgers".
if (originalCurrentLedger != null && !ledgers.containsKey(originalCurrentLedger.getId())){
asyncDeleteLedgerWithConcurrencyLimit(originalCurrentLedger.getId(), (rc, ctx) -> {
Expand All @@ -1885,7 +1907,10 @@ protected synchronized void updateLedgersIdsComplete(@Nullable LedgerHandle orig
updateLastLedgerCreatedTimeAndScheduleRolloverTask();

log.debug().attr("count", pendingAddEntries.size()).log("Resending pending messages");
resendPendingAddEntries();
}

private void resendPendingAddEntries() {
createNewOpAddEntryForNewLedger();

// Process all the pending addEntry requests
Expand Down Expand Up @@ -2023,9 +2048,13 @@ boolean isNeededCreateNewLedgerAfterCloseLedger() {
@VisibleForTesting
@SuppressWarnings("deprecation")
@Override
public void rollCurrentLedgerIfFull() {
public synchronized void rollCurrentLedgerIfFull() {
log.info("Start checking if current ledger is full");
if (currentLedgerEntries > 0 && currentLedgerIsFull()
if (!currentLedgerIsFull()) {
return;
}

if (currentLedgerEntries > 0
&& STATE_UPDATER.compareAndSet(this, State.LedgerOpened, State.ClosingLedger)) {
currentLedger.asyncClose(new AsyncCallback.CloseCallback() {
@Override
Expand All @@ -2045,6 +2074,109 @@ public void closeComplete(int rc, LedgerHandle lh, Object o) {
ledgerClosed(lh);
}
}, null);
} else if (currentLedgerEntries == 0 && ledgers.containsKey(currentLedger.getId())) {
startEmptyLedgerRollover(currentLedger, null);
}
}

private void checkEmptyLedgerMetadata(LedgerHandle ledger) {
final CompletableFuture<ReadHandle> metadataFuture;
try {
metadataFuture = openLedgerForMetadataCheck(ledger);
} catch (Throwable cause) {
completeEmptyLedgerMetadataCheck(ledger, null, cause);
return;
}
metadataFuture.whenCompleteAsync(
(readHandle, exception) -> completeEmptyLedgerMetadataCheck(ledger, readHandle, exception), executor);
}

@VisibleForTesting
CompletableFuture<ReadHandle> openLedgerForMetadataCheck(LedgerHandle ledger) {
return bookKeeper.newOpenLedgerOp()
.withRecovery(false)
.withLedgerId(ledger.getId())
.withDigestType(config.getDigestType())
.withPassword(config.getPassword())
.withLoggerContext(log)
.execute();
}

private void completeEmptyLedgerMetadataCheck(
LedgerHandle ledger, @Nullable ReadHandle readHandle, @Nullable Throwable exception) {
try {
synchronized (this) {
if (emptyLedgerMetadataCheckLedger == ledger) {
emptyLedgerMetadataCheckLedger = null;
}
if (currentLedger != ledger || STATE_UPDATER.get(this) != State.LedgerOpened) {
return;
}
if (exception != null) {
log.warn().attr("ledgerId", ledger.getId()).exception(exception)
.log("Failed to verify empty ledger metadata before its first write; rolling it over");
if (!pendingAddEntries.isEmpty()) {
startEmptyLedgerRollover(ledger, null);
}
return;
}

LedgerMetadata metadata = readHandle.getLedgerMetadata();
if (metadata.getState() != LedgerMetadata.State.OPEN) {
startEmptyLedgerRollover(ledger, metadata.getLastEntryId());
} else if (!pendingAddEntries.isEmpty()) {
if (maximumRolloverTimeReached()) {
startEmptyLedgerRollover(ledger, null);
} else {
resendPendingAddEntries();
}
}
}
} finally {
if (readHandle != null) {
try {
readHandle.closeAsync().exceptionally(closeException -> {
log.warn().attr("ledgerId", ledger.getId()).exception(closeException)
.log("Failed to close empty ledger metadata read handle");
return null;
});
} catch (Throwable cause) {
log.warn().attr("ledgerId", ledger.getId()).exception(cause)
.log("Failed to initiate closing empty ledger metadata read handle");
}
}
}
}

private synchronized void startEmptyLedgerRollover(LedgerHandle ledger, @Nullable Long lastAddConfirmed) {
if (currentLedger != ledger || currentLedgerEntries != 0
|| !STATE_UPDATER.compareAndSet(this, State.LedgerOpened, State.ClosingLedger)) {
return;
}
emptyLedgerRolloverInProgress = true;
if (lastAddConfirmed != null) {
ledgerClosed(ledger, lastAddConfirmed);
} else {
asyncCloseEmptyLedgerBeforeFirstWrite(ledger);
}
}

@VisibleForTesting
void asyncCloseEmptyLedgerBeforeFirstWrite(LedgerHandle ledger) {
log.debug().attr("ledgerId", ledger.getId()).log("Rolling over an empty ledger before its first write");
try {
ledger.asyncClose((rc, closedLedger, ctx) -> {
if (rc != BKException.Code.OK) {
log.warn().attr("ledgerId", ledger.getId())
.attr("status", BKException.getMessage(rc))
.log("Error when closing an empty ledger before its first write");
}
ledgerClosed(ledger);
}, null);
} catch (Throwable cause) {
log.warn().attr("ledgerId", ledger.getId()).exception(cause)
.log("Failed to initiate closing an empty ledger before its first write");
ledgerClosed(ledger);
}
}

Expand Down Expand Up @@ -4465,6 +4597,11 @@ protected boolean currentLedgerIsFull() {
}
}

private boolean maximumRolloverTimeReached() {
return config.getMaximumRolloverTimeMs() > 0
&& clock.millis() - lastLedgerCreatedTimestamp >= config.getMaximumRolloverTimeMs();
}

Comment on lines +4600 to +4604

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

maximumRolloverTimeReached() returns false when the broker's metadata service is unavailable. At both call sites, a false value causes the metadata check to be skipped, allowing the write to proceed with the current handle. Since auto-recovery may have already closed the ledger via a different metadata session while this broker is disconnected, this turns a safety check into a fail-open path. Could we differentiate between "the deadline has not been reached" and "the state cannot be verified," and either queue or fail the write in the latter case?

public List<LedgerInfo> getLedgersInfoAsList() {
return Lists.newArrayList(ledgers.values());
}
Expand Down
Loading