From 644d9af5090966c37506acb7beabcc6ebb4b780b Mon Sep 17 00:00:00 2001 From: Samuel Date: Thu, 18 Jun 2026 00:39:32 +0100 Subject: [PATCH] fix: improve Raft log compaction safety, refine election logic, and update repository documentation --- .gitignore | 5 +- README.md | 2 + .../java/com/drmq/broker/raft/RaftLog.java | 81 ++++++++---- .../java/com/drmq/broker/raft/RaftNode.java | 121 +++++++++++------- 4 files changed, 143 insertions(+), 66 deletions(-) diff --git a/.gitignore b/.gitignore index 5a4f1c0..85a734e 100644 --- a/.gitignore +++ b/.gitignore @@ -19,4 +19,7 @@ drmq-python-client/*.pyc # TypeScript SDK drmq-ts-client/node_modules/ drmq-ts-client/dist/ -drmq-ts-client/build/ \ No newline at end of file +drmq-ts-client/build/ + +.agent/ +drmq-docs/ diff --git a/README.md b/README.md index 9c5b6c1..4a8370d 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/drmq-broker/src/main/java/com/drmq/broker/raft/RaftLog.java b/drmq-broker/src/main/java/com/drmq/broker/raft/RaftLog.java index ff09b61..465b527 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/raft/RaftLog.java +++ b/drmq-broker/src/main/java/com/drmq/broker/raft/RaftLog.java @@ -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 remainingEntries; - List 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 newPositions = new ArrayList<>(remainingEntries.size()); @@ -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 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); } /** diff --git a/drmq-broker/src/main/java/com/drmq/broker/raft/RaftNode.java b/drmq-broker/src/main/java/com/drmq/broker/raft/RaftNode.java index d914d90..6c44a49 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/raft/RaftNode.java +++ b/drmq-broker/src/main/java/com/drmq/broker/raft/RaftNode.java @@ -58,6 +58,7 @@ public class RaftNode { // Leader-only volatile state private final Map nextIndex; private final Map matchIndex; + private final Map snapshotInProgress = new ConcurrentHashMap<>(); private final String nodeId; private final int port; @@ -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> voteRpcHandlers = new ConcurrentHashMap<>(); private final Map> appendRpcHandlers = new ConcurrentHashMap<>(); private final Map> preVoteRpcHandlers = new ConcurrentHashMap<>(); @@ -290,19 +293,24 @@ public void registerInstallSnapshotHandler(String peerId, Function= votesNeeded) { if (electionStarted.compareAndSet(false, true)) { logger.info("[{}] Pre-vote succeeded (single-node quorum), starting real election", nodeId); - startElection(); + startElection(proposedTerm); } } @@ -382,7 +390,7 @@ private void startPreVote() { } // Release lock before startElection() to avoid holding it during disk I/O if (shouldStartElection) { - startElection(); + startElection(proposedTerm); } }); } @@ -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; @@ -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; @@ -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 @@ -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++; @@ -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 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) { @@ -634,6 +646,18 @@ private void replicateTo(PeerAddress peer) { } if (needsSnapshot) return; + long prevLogTerm = raftLog.getTermAt(prevLogIndex); + List entries = raftLog.getEntriesFrom(peerNextIndex); + + AppendEntriesRequest request = AppendEntriesRequest.newBuilder() + .setTerm(currentTermLocal) + .setLeaderId(leaderIdLocal) + .setPrevLogIndex(prevLogIndex) + .setPrevLogTerm(prevLogTerm) + .addAllEntries(entries) + .setLeaderCommit(commitIndexLocal) + .build(); + Function handler = appendRpcHandlers.get(peer.id()); if (handler == null) return; @@ -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 handler = installSnapshotRpcHandlers.get(peer.id()); if (handler == null) return; @@ -722,6 +748,7 @@ private void sendInstallSnapshotToPeer(PeerAddress peer) { .build(); InstallSnapshotResponse response = handler.apply(request); + lastContactTime.put(peer.id(), System.currentTimeMillis()); lock.lock(); try { @@ -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--) { @@ -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); + } + }); } } }