Skip to content
Merged
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
5 changes: 4 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -19,4 +19,7 @@ drmq-python-client/*.pyc
# TypeScript SDK
drmq-ts-client/node_modules/
drmq-ts-client/dist/
drmq-ts-client/build/
drmq-ts-client/build/

.agent/
drmq-docs/
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
# DRMQ - Distributed Reliable Message Queue

**Official Documentation:** [https://drmq-web.vercel.app](https://drmq-web.vercel.app)

DRMQ is a fault-tolerant, high-performance distributed message broker built from first principles. It provides guaranteed message delivery, strict ordering, and high availability through the implementation of the Raft consensus algorithm for log replication and leader election. DRMQ supports scalable consumption via multi-consumer groups — multiple consumers can share a group to load-balance message processing without partitions.

## Overview
Expand Down
81 changes: 59 additions & 22 deletions drmq-broker/src/main/java/com/drmq/broker/raft/RaftLog.java
Original file line number Diff line number Diff line change
Expand Up @@ -210,13 +210,17 @@ public synchronized void truncateFrom(long fromIndex) throws IOException {
* Compact the Raft log by removing entries from memory up to the given index,
* and rewriting the log file on disk to reclaim space.
*/
public synchronized void compact(long upToIndex) throws IOException {
if (upToIndex <= startIndex || upToIndex > getLastIndex()) {
return;
}
int removeCount = (int) (upToIndex - startIndex + 1);
public void compact(long upToIndex) throws IOException {
int removeCount;
List<RaftEntry> remainingEntries;

List<RaftEntry> remainingEntries = new ArrayList<>(entries.subList(removeCount, entries.size()));
synchronized(this) {
if (upToIndex <= startIndex || upToIndex > getLastIndex()) {
return;
}
removeCount = (int) (upToIndex - startIndex + 1);
remainingEntries = new ArrayList<>(entries.subList(removeCount, entries.size()));
}

java.io.File tempFile = new java.io.File(logPath.getParent().toFile(), logPath.getFileName().toString() + ".tmp");
List<Long> newPositions = new ArrayList<>(remainingEntries.size());
Expand All @@ -231,23 +235,56 @@ public synchronized void compact(long upToIndex) throws IOException {
tempRaf.getFD().sync();
}

// Swap files
raf.close();
java.nio.file.Files.move(tempFile.toPath(), logPath, java.nio.file.StandardCopyOption.REPLACE_EXISTING, java.nio.file.StandardCopyOption.ATOMIC_MOVE);
try {
raf = new RandomAccessFile(logPath.toFile(), "rw");
} catch (IOException e) {
throw new IOException("Failed to reopen compacted log: " + logPath, e);
synchronized(this) {
int currentExpectedSize = remainingEntries.size() + removeCount;
if (entries.size() < currentExpectedSize) {
// Truncation happened during compaction! Abort to avoid restoring truncated entries.
tempFile.delete();
return;
}

int addedCount = entries.size() - currentExpectedSize;
if (addedCount > 0) {
List<RaftEntry> newlyAdded = entries.subList(entries.size() - addedCount, entries.size());
try (RandomAccessFile tempRaf = new RandomAccessFile(tempFile, "rw")) {
tempRaf.seek(tempRaf.length());
for (RaftEntry entry : newlyAdded) {
newPositions.add(tempRaf.getFilePointer());
byte[] data = entry.toByteArray();
tempRaf.writeInt(data.length);
tempRaf.write(data);
}
tempRaf.getFD().sync();
}
}

// Swap files
raf.close();
try {
java.nio.file.Files.move(tempFile.toPath(), logPath, java.nio.file.StandardCopyOption.REPLACE_EXISTING, java.nio.file.StandardCopyOption.ATOMIC_MOVE);
} catch (Exception e) {
try {
raf = new RandomAccessFile(logPath.toFile(), "rw");
raf.seek(raf.length());
} catch (Exception ignore) {
}
throw e;
}

try {
raf = new RandomAccessFile(logPath.toFile(), "rw");
} catch (IOException e) {
throw new IOException("Failed to reopen compacted log: " + logPath, e);
}
raf.seek(raf.length());

entries.subList(0, removeCount).clear();
filePositions.clear();
filePositions.addAll(newPositions);
startIndex = upToIndex + 1;

logger.debug("Compacted Raft log on disk up to index {}", upToIndex);
}
raf.seek(raf.length());

entries.clear();
entries.addAll(remainingEntries);
filePositions.clear();
filePositions.addAll(newPositions);
startIndex = upToIndex + 1;

logger.debug("Compacted Raft log on disk up to index {}", upToIndex);
}

/**
Expand Down
121 changes: 78 additions & 43 deletions drmq-broker/src/main/java/com/drmq/broker/raft/RaftNode.java
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ public class RaftNode {
// Leader-only volatile state
private final Map<String, Long> nextIndex;
private final Map<String, Long> matchIndex;
private final Map<String, Boolean> snapshotInProgress = new ConcurrentHashMap<>();

private final String nodeId;
private final int port;
Expand All @@ -69,6 +70,8 @@ public class RaftNode {
private final Path stateFilePath;
private final long raftCompactThreshold;

private final AtomicBoolean isCompacting = new AtomicBoolean(false);

private final Map<String, Function<RequestVoteRequest, RequestVoteResponse>> voteRpcHandlers = new ConcurrentHashMap<>();
private final Map<String, Function<AppendEntriesRequest, AppendEntriesResponse>> appendRpcHandlers = new ConcurrentHashMap<>();
private final Map<String, Function<PreVoteRequest, PreVoteResponse>> preVoteRpcHandlers = new ConcurrentHashMap<>();
Expand Down Expand Up @@ -290,19 +293,24 @@ public void registerInstallSnapshotHandler(String peerId, Function<InstallSnapsh
* If the timer fires, the node starts an election.
*/
private void resetElectionTimer() {
if (electionTimer != null) {
electionTimer.cancel(false);
}
long timeout;
if (startupGrace) {
timeout = ELECTION_TIMEOUT_MAX_MS * 3;
startupGrace = false;
logger.info("[{}] Startup grace: election timeout set to {}ms", nodeId, timeout);
} else {
timeout = ELECTION_TIMEOUT_MIN_MS +
ThreadLocalRandom.current().nextLong(ELECTION_TIMEOUT_MAX_MS - ELECTION_TIMEOUT_MIN_MS);
lock.lock();
try {
if (electionTimer != null) {
electionTimer.cancel(false);
}
long timeout;
if (startupGrace) {
timeout = ELECTION_TIMEOUT_MAX_MS * 3;
startupGrace = false;
logger.info("[{}] Startup grace: election timeout set to {}ms", nodeId, timeout);
} else {
timeout = ELECTION_TIMEOUT_MIN_MS +
ThreadLocalRandom.current().nextLong(ELECTION_TIMEOUT_MAX_MS - ELECTION_TIMEOUT_MIN_MS);
}
electionTimer = scheduler.schedule(this::startPreVote, timeout, TimeUnit.MILLISECONDS);
} finally {
lock.unlock();
}
electionTimer = scheduler.schedule(this::startPreVote, timeout, TimeUnit.MILLISECONDS);
}

/**
Expand Down Expand Up @@ -341,7 +349,7 @@ private void startPreVote() {
if (votesReceived.get() >= votesNeeded) {
if (electionStarted.compareAndSet(false, true)) {
logger.info("[{}] Pre-vote succeeded (single-node quorum), starting real election", nodeId);
startElection();
startElection(proposedTerm);
}
}

Expand Down Expand Up @@ -382,7 +390,7 @@ private void startPreVote() {
}
// Release lock before startElection() to avoid holding it during disk I/O
if (shouldStartElection) {
startElection();
startElection(proposedTerm);
}
});
}
Expand All @@ -394,7 +402,7 @@ private void startPreVote() {
* Start a real election: transition to CANDIDATE, increment term, vote for self,
* request votes from peers. Only called after a successful pre-vote.
*/
private void startElection() {
private void startElection(long proposedTerm) {
if (!running) return;

long myTerm;
Expand All @@ -404,8 +412,12 @@ private void startElection() {
try {
if (!running || state == RaftState.LEADER) return;

if (currentTerm >= proposedTerm) {
return;
}

electionStartNanos = System.nanoTime();
currentTerm++;
currentTerm = proposedTerm;
state = RaftState.CANDIDATE;
votedFor = nodeId;
leaderId = null;
Expand All @@ -427,8 +439,8 @@ private void startElection() {
lock.unlock();
}


savePersistentState();
resetElectionTimer();

int votesNeeded = (peers.size() + 1) / 2 + 1;
AtomicLong votesReceived = new AtomicLong(1); // self-vote
Expand Down Expand Up @@ -558,6 +570,10 @@ private void checkQuorum() {
int activePeers = 1;

for (PeerAddress peer : peers) {
if (snapshotInProgress.getOrDefault(peer.id(), false)) {
activePeers++;
continue;
}
Long lastContact = lastContactTime.get(peer.id());
if (lastContact != null && (now - lastContact) <= quorumWindow) {
activePeers++;
Expand Down Expand Up @@ -601,31 +617,27 @@ private void sendHeartbeats() {
*/
private void replicateTo(PeerAddress peer) {
boolean needsSnapshot = false;
AppendEntriesRequest request;
long peerNextIndex;
long prevLogIndex;
long currentTermLocal;
long commitIndexLocal;
String leaderIdLocal;

lock.lock();
try {
if (state != RaftState.LEADER) return;

long peerNextIndex = nextIndex.getOrDefault(peer.id(), raftLog.getLastIndex() + 1);
long prevLogIndex = peerNextIndex - 1;
peerNextIndex = nextIndex.getOrDefault(peer.id(), raftLog.getLastIndex() + 1);
prevLogIndex = peerNextIndex - 1;

if (prevLogIndex > 0 && prevLogIndex < raftLog.getStartIndex()) {
needsSnapshot = true;
return;
}

long prevLogTerm = raftLog.getTermAt(prevLogIndex);

List<RaftEntry> entries = raftLog.getEntriesFrom(peerNextIndex);

request = AppendEntriesRequest.newBuilder()
.setTerm(currentTerm)
.setLeaderId(nodeId)
.setPrevLogIndex(prevLogIndex)
.setPrevLogTerm(prevLogTerm)
.addAllEntries(entries)
.setLeaderCommit(commitIndex)
.build();
currentTermLocal = currentTerm;
commitIndexLocal = commitIndex;
leaderIdLocal = nodeId;
} finally {
lock.unlock();
if (needsSnapshot) {
Expand All @@ -634,6 +646,18 @@ private void replicateTo(PeerAddress peer) {
}
if (needsSnapshot) return;

long prevLogTerm = raftLog.getTermAt(prevLogIndex);
List<RaftEntry> entries = raftLog.getEntriesFrom(peerNextIndex);

AppendEntriesRequest request = AppendEntriesRequest.newBuilder()
.setTerm(currentTermLocal)
.setLeaderId(leaderIdLocal)
.setPrevLogIndex(prevLogIndex)
.setPrevLogTerm(prevLogTerm)
.addAllEntries(entries)
.setLeaderCommit(commitIndexLocal)
.build();

Function<AppendEntriesRequest, AppendEntriesResponse> handler = appendRpcHandlers.get(peer.id());
if (handler == null) return;

Expand Down Expand Up @@ -689,13 +713,15 @@ private void sendInstallSnapshotToPeer(PeerAddress peer) {
lock.unlock();
}

Path snapshotZip;
snapshotInProgress.put(peer.id(), true);
try {
snapshotZip = snapshotManager.createSnapshot(snapshotIndex);
} catch (IOException e) {
logger.error("[{}] Failed to create snapshot for peer {}", nodeId, peer.id(), e);
return;
}
Path snapshotZip;
try {
snapshotZip = snapshotManager.createSnapshot(snapshotIndex);
} catch (IOException e) {
logger.error("[{}] Failed to create snapshot for peer {}", nodeId, peer.id(), e);
return;
}

Function<InstallSnapshotRequest, InstallSnapshotResponse> handler = installSnapshotRpcHandlers.get(peer.id());
if (handler == null) return;
Expand All @@ -722,6 +748,7 @@ private void sendInstallSnapshotToPeer(PeerAddress peer) {
.build();

InstallSnapshotResponse response = handler.apply(request);
lastContactTime.put(peer.id(), System.currentTimeMillis());

lock.lock();
try {
Expand Down Expand Up @@ -753,9 +780,11 @@ private void sendInstallSnapshotToPeer(PeerAddress peer) {
} catch (Exception e) {
logger.debug("[{}] InstallSnapshot to {} failed: {}", nodeId, peer.id(), e.getMessage());
}
} finally {
snapshotInProgress.put(peer.id(), false);
}
}


private void advanceCommitIndex() {
long lastIndex = raftLog.getLastIndex();
for (long n = lastIndex; n > commitIndex; n--) {
Expand Down Expand Up @@ -875,10 +904,16 @@ private void applyCommitted() {
long finalCompactIndex = Math.min(safeCompactIndex, lastApplied - raftCompactThreshold);

if (finalCompactIndex > 0) {
try {
raftLog.compact(finalCompactIndex);
} catch (IOException e) {
logger.error("Failed to compact Raft log", e);
if (isCompacting.compareAndSet(false, true)) {
raftExecutor.execute(() -> {
try {
raftLog.compact(finalCompactIndex);
} catch (IOException e) {
logger.error("Failed to compact Raft log", e);
} finally {
isCompacting.set(false);
}
});
}
}
}
Expand Down