apache / apache/rocketmq

[Bug] AutoSwitchHA slave truncation is not crash-safe because stale CommitLog tail remains valid on disk

Open
#10,666 7 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
22.6k
Forks
12k
Avg merge
3d 1h
Merged PRs (30d)
27

Description

### Before Creating the Bug Report

- [x] I found a bug, not just asking a question, which should be created in [GitHub Discussions](https://github.com/apache/rocketmq/discussions).

- [x] I have searched the [GitHub Issues](https://github.com/apache/rocketmq/issues) and [GitHub Discussions](https://github.com/apache/rocketmq/discussions) of this repository and believe that this is not a duplicate.

- [x] I have confirmed that this bug belongs to the current repository, not other repositories of RocketMQ.

### Runtime platform environment

ubuntu 24.04

### RocketMQ version

branch: develop version: 5.5.0 Git commit id: a6fb9e2fa0d4e446c88b0051cd4fac233dcee9ec

### JDK Version

JDK 8u202

### Describe the Bug

In Controller mode, RocketMQ uses `AutoSwitchHAService`.

When a slave reconnects to the master, `AutoSwitchHAClient#doTruncate()` calculates a consistent point and truncates the local CommitLog:

```java
final long truncateOffset =
localEpochCache.findConsistentPoint(masterEpochCache);

if (!this.messageStore.truncateFiles(truncateOffset)) {
LOGGER.error("Failed to truncate slave log to {}", truncateOffset);
return false;
}
```

The truncation offset is validated as a message boundary, but the actual truncation only resets the in-memory positions of the target `MappedFile`.

The physical bytes after the truncation offset remain unchanged on disk.

If the slave crashes again after truncation but before the master has durably overwritten the stale tail, CommitLog recovery may scan and recover the stale messages again.

### 1. Controller mode uses AutoSwitchHAService

```java
if (brokerConfig.isEnableControllerMode()) {
this.haService = new AutoSwitchHAService();
}
```

Source:

https://github.com/apache/rocketmq/blob/a6fb9e2fa0d4e446c88b0051cd4fac233dcee9ec/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java#L1978-L1988

### 2. AutoSwitchHAClient truncates the slave CommitLog

```java
final long truncateOffset =
localEpochCache.findConsistentPoint(masterEpochCache);

if (!this.messageStore.truncateFiles(truncateOffset)) {
LOGGER.error("Failed to truncate slave log to {}", truncateOffset);
return false;
}

changeCurrentState(HAConnectionState.TRANSFER);
this.currentReportedOffset = truncateOffset;
```

Source:

https://github.com/apache/rocketmq/blob/a6fb9e2fa0d4e446c88b0051cd4fac233dcee9ec/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAClient.java#L450-L477

### 3. DefaultMessageStore only checks whether the offset is aligned

```java
public boolean truncateFiles(long offsetToTruncate) {
if (offsetToTruncate >= this.getMaxPhyOffset()) {
return true;
}

if (!isOffsetAligned(offsetToTruncate)) {
return false;
}

truncateDirtyFiles(offsetToTruncate);
return true;
}
```

Source:

https://github.com/apache/rocketmq/blob/a6fb9e2fa0d4e446c88b0051cd4fac233dcee9ec/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java#L828-L851

The alignment check prevents truncation in the middle of a message, but it does not make the truncation durable.

### 4. MappedFileQueue only resets in-memory positions

```java
public void truncateDirtyFiles(long offset) {
List willRemoveFiles = new ArrayList<>();

for (MappedFile file : this.mappedFiles) {
long fileTailOffset =
file.getFileFromOffset() + this.mappedFileSize;

if (fileTailOffset > offset) {
if (offset >= file.getFileFromOffset()) {
file.setWrotePosition(
(int) (offset % this.mappedFileSize));
file.setCommittedPosition(
(int) (offset % this.mappedFileSize));
file.setFlushedPosition(
(int) (offset % this.mappedFileSize));
} else {
file.destroy(1000);
willRemoveFiles.add(file);
}
}
}

this.deleteExpiredFile(willRemoveFiles);
}
```

Source:

https://github.com/apache/rocketmq/blob/a6fb9e2fa0d4e446c88b0051cd4fac233dcee9ec/store/src/main/java/org/apache/rocketmq/store/MappedFileQueue.java#L217-L234

After truncating to offset `T`, the state is:

```text
In-memory wrotePosition = T
In-memory committedPosition = T
In-memory flushedPosition = T

Physical bytes in [T, oldMaxOffset) remain unchanged
```

This also applies when `T` is exactly the beginning of a `MappedFile`.

Because the condition is:

```java
offset >= file.getFileFromOffset()
```

the target file is retained with position zero instead of being deleted.

### 5. Existing MappedFiles are loaded as fully written after restart

```java
mappedFile.setWrotePosition(this.mappedFileSize);
mappedFile.setFlushedPosition(this.mappedFileSize);
mappedFile.setCommittedPosition(this.mappedFileSize);
```

Source:

https://github.com/apache/rocketmq/blob/a6fb9e2fa0d4e446c88b0051cd4fac233dcee9ec/store/src/main/java/org/apache/rocketmq/store/MappedFileQueue.java#L270-L298

CommitLog recovery then scans the physical contents to determine the valid tail.

Because the discarded bytes were not erased or invalidated, the stale records after `T` may still contain:

- a valid total size;
- a valid magic code;
- a correct physical offset;
- a valid CRC.

As a result, recovery may advance beyond the intended HA truncation offset.

## Failure Scenario

A possible failure sequence is:

1. The slave has CommitLog data up to offset `S`.
2. After a master switch or epoch divergence, `AutoSwitchHAClient` calculates an aligned consistent point `T`.
3. `T` is smaller than `S`:

```text
T < S
```

4. The slave executes:

```java
messageStore.truncateFiles(T);
```

5. The in-memory CommitLog maximum offset becomes `T`.
6. The physical bytes in `[T, S)` remain unchanged.
7. Before the master has durably overwritten the stale tail, the slave crashes again.
8. On restart, the MappedFile is initially loaded as fully written.
9. CommitLog recovery scans the stale records after `T`.
10. The recovered maximum physical offset may become greater than `T` again.

### Steps to Reproduce

The storage behavior can be reproduced without running a complete Controller cluster.

1. Create a CommitLog MappedFile.
2. Append and flush at least two valid messages.
3. Record the physical start offset of the second message as `T`.
4. Call:

```java
messageStore.truncateFiles(T);
```

5. Verify:

```text
messageStore.getMaxPhyOffset() == T
```

6. Do not append new data from `T`.
7. Simulate an abrupt process termination.
8. Reload the same store directory.
9. Execute CommitLog recovery.
10. Check the recovered maximum physical offset.

Expected result:

```text
recoveredMaxPhyOffset == T
```

Potential current result:

```text
recoveredMaxPhyOffset > T
```

The old second message may be recovered because its physical bytes remain valid on disk.

### What Did You Expect to See?

## Expected Behavior

After `AutoSwitchHAClient` successfully truncates the slave to offset `T`, a subsequent restart must not recover any CommitLog record at or after `T` unless that record has been received again from the current master.

The HA truncation should provide a durable recovery boundary.

### What Did You See Instead?

Slave commitLog keeps the stale tail.

### Additional Context

_No response_

Contributor guide

Open the contributing guide

Research direction

Start with AutoSwitchHAClient#doTruncate(), DefaultMessageStore#truncateFiles(), and MappedFileQueue#truncateDirtyFiles() to trace the truncation and restart path. Reproduce the two-message scenario described in the issue, then inspect MappedFileQueue loading and CommitLog recovery. Done means a restart cannot recover records at or after T unless they are received again from the current master.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
distributed-systems
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
48/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.