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
206 changes: 206 additions & 0 deletions drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,206 @@
package com.drmq.broker;

import com.google.gson.Gson;
import com.google.gson.JsonArray;
import com.google.gson.JsonObject;
import com.sun.net.httpserver.HttpExchange;
import com.sun.net.httpserver.HttpHandler;
import com.sun.net.httpserver.HttpServer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.IOException;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.util.List;

/**
* Lightweight HTTP server for administrative REST APIs.
*/
public class AdminHttpServer {
private static final Logger logger = LoggerFactory.getLogger(AdminHttpServer.class);
private final HttpServer server;
private final MessageStore messageStore;
private final OffsetManager offsetManager;
private final ConsumerGroupCoordinator groupCoordinator;
private final Gson gson = new Gson();

public AdminHttpServer(int port, MessageStore messageStore, OffsetManager offsetManager, ConsumerGroupCoordinator groupCoordinator) throws IOException {
this.messageStore = messageStore;
this.offsetManager = offsetManager;
this.groupCoordinator = groupCoordinator;

this.server = HttpServer.create(new InetSocketAddress(port), 0);

this.server.createContext("/api/topics", this::handleTopics);
this.server.createContext("/api/consumers", this::handleConsumers);
this.server.createContext("/api/messages", this::handleMessages);

// CORS and standard executor
this.server.setExecutor(null);
}

public void start() {
server.start();
logger.info("Admin HTTP Server started on port {}", server.getAddress().getPort());
}

public void stop() {
server.stop(1);
logger.info("Admin HTTP Server stopped.");
}

private void handleTopics(HttpExchange exchange) throws IOException {
addCorsHeaders(exchange);
if ("OPTIONS".equals(exchange.getRequestMethod())) {
exchange.sendResponseHeaders(204, -1);
return;
}

JsonArray topicsArray = new JsonArray();
List<String> topics = messageStore.getTopics();

for (String topic : topics) {
JsonObject obj = new JsonObject();
obj.addProperty("name", topic);
obj.addProperty("messageCount", messageStore.getMessageCount(topic));
// The global offset is global, not per topic.
// We'll just return the message count for now.
topicsArray.add(obj);
}

sendJsonResponse(exchange, 200, gson.toJson(topicsArray));
}

private void handleConsumers(HttpExchange exchange) throws IOException {
addCorsHeaders(exchange);
if ("OPTIONS".equals(exchange.getRequestMethod())) {
exchange.sendResponseHeaders(204, -1);
return;
}

JsonArray groupsArray = new JsonArray();
java.util.Map<String, Long> allOffsets = offsetManager.getAllOffsets();

// Group by consumer group name
java.util.Map<String, JsonArray> groupsMap = new java.util.HashMap<>();

for (java.util.Map.Entry<String, Long> entry : allOffsets.entrySet()) {
String[] parts = entry.getKey().split("/");
if (parts.length != 2) continue;

String groupName = parts[0];
String topicName = parts[1];
long committedOffset = entry.getValue();

// Calculate lag using true topic head offset rather than message count
// Since DRMQ uses global offsets, messageCount does not correlate to the offset values.
long headOffset = messageStore.getHeadOffset(topicName);
long lag = 0;
if (headOffset >= 0) {
long effectiveCommitted = Math.max(0, committedOffset); // -1 means none committed
lag = Math.max(0, (headOffset + 1) - effectiveCommitted);
}

JsonObject topicObj = new JsonObject();
topicObj.addProperty("topic", topicName);
topicObj.addProperty("headOffset", headOffset);
topicObj.addProperty("committedOffset", committedOffset);
topicObj.addProperty("lag", lag);
topicObj.addProperty("activeMembers", groupCoordinator.getConsumerCount(groupName, topicName));

groupsMap.computeIfAbsent(groupName, k -> new JsonArray()).add(topicObj);
}

for (java.util.Map.Entry<String, JsonArray> entry : groupsMap.entrySet()) {
JsonObject groupObj = new JsonObject();
groupObj.addProperty("groupId", entry.getKey());
groupObj.add("topics", entry.getValue());
groupsArray.add(groupObj);
}

sendJsonResponse(exchange, 200, gson.toJson(groupsArray));
}

private void handleMessages(HttpExchange exchange) throws IOException {
addCorsHeaders(exchange);
if ("OPTIONS".equals(exchange.getRequestMethod())) {
exchange.sendResponseHeaders(204, -1);
return;
}

try {
String query = exchange.getRequestURI().getQuery();
if (query == null) {
sendJsonResponse(exchange, 400, "{\"error\":\"Missing query parameters\"}");
return;
}

String topic = null;
long offset = 0;
int limit = 10;
Long timestamp = null;

for (String param : query.split("&")) {
String[] pair = param.split("=");
if (pair.length == 2) {
if ("topic".equals(pair[0])) topic = pair[1];
else if ("offset".equals(pair[0])) offset = Long.parseLong(pair[1]);
else if ("limit".equals(pair[0])) limit = Integer.parseInt(pair[1]);
else if ("timestamp".equals(pair[0])) timestamp = Long.parseLong(pair[1]);
}
}

if (topic == null) {
sendJsonResponse(exchange, 400, "{\"error\":\"Missing 'topic' parameter\"}");
return;
}

// Limit bounds to avoid OOM
limit = Math.min(100, Math.max(1, limit));

if (timestamp != null && timestamp > 0) {
offset = messageStore.findOffsetByTimestamp(topic, timestamp);
if (offset == -1) {
sendJsonResponse(exchange, 200, "[]");
return;
}
}

List<com.drmq.protocol.DRMQProtocol.StoredMessage> messages = messageStore.getMessages(topic, offset, limit);
JsonArray msgsArray = new JsonArray();

for (com.drmq.protocol.DRMQProtocol.StoredMessage msg : messages) {
JsonObject obj = new JsonObject();
obj.addProperty("offset", msg.getOffset());
obj.addProperty("timestamp", msg.getTimestamp());
obj.addProperty("storedAt", msg.getStoredAt());
if (msg.hasKey()) {
obj.addProperty("key", msg.getKey());
}
obj.addProperty("payload", msg.getPayload().toStringUtf8());
msgsArray.add(obj);
}

sendJsonResponse(exchange, 200, gson.toJson(msgsArray));
} catch (Exception e) {
logger.error("Error handling messages request", e);
sendJsonResponse(exchange, 500, "{\"error\":\"" + e.getMessage() + "\"}");
}
Comment on lines +186 to +189

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Build the error JSON with Gson instead of string concatenation.

e.getMessage() is interpolated raw into the JSON body. If the message contains a ", newline, or backslash the response becomes malformed JSON and the dashboard's res.json() fails; if it is null the body reads {"error":"null"}. It also leaks internal exception detail verbatim.

🐛 Proposed fix
         } catch (Exception e) {
             logger.error("Error handling messages request", e);
-            sendJsonResponse(exchange, 500, "{\"error\":\"" + e.getMessage() + "\"}");
+            JsonObject err = new JsonObject();
+            err.addProperty("error", e.getMessage() != null ? e.getMessage() : "Internal error");
+            sendJsonResponse(exchange, 500, gson.toJson(err));
         }
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
} catch (Exception e) {
logger.error("Error handling messages request", e);
sendJsonResponse(exchange, 500, "{\"error\":\"" + e.getMessage() + "\"}");
}
} catch (Exception e) {
logger.error("Error handling messages request", e);
JsonObject err = new JsonObject();
err.addProperty("error", e.getMessage() != null ? e.getMessage() : "Internal error");
sendJsonResponse(exchange, 500, gson.toJson(err));
}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java` around lines
186 - 189, Update the exception response handling in AdminHttpServer’s
message-request catch block to construct the error payload with the existing
Gson serializer instead of concatenating e.getMessage() into JSON. Serialize a
structured error value so quotes, newlines, backslashes, and null messages
remain valid JSON, while preserving the 500 response through sendJsonResponse
and avoiding raw internal exception details in the response.

}

private void addCorsHeaders(HttpExchange exchange) {
exchange.getResponseHeaders().add("Access-Control-Allow-Origin", "*");
exchange.getResponseHeaders().add("Access-Control-Allow-Methods", "GET, POST, DELETE, OPTIONS");
exchange.getResponseHeaders().add("Access-Control-Allow-Headers", "Content-Type, Authorization");
}
Comment on lines +192 to +196

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🔒 Security & Privacy | 🟠 Major | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
# Confirm whether any auth/token check gates these admin endpoints elsewhere.
rg -nP 'Authorization|Bearer|token|apiKey|AdminHttpServer' drmq-broker/src/main/java --type=java -C2

Repository: samuel025/DRMQ

Length of output: 3093


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Inspect AdminHttpServer handlers and related message endpoints/cors behavior.
wc -l drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java
cat -n drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java | sed -n '1,280p'

# Search for /api/messages/path, getPayload, or raw message serialization in Java sources.
rg -n '/api/messages|messages|getPayload|payload|Access-Control-Allow-Origin|addCorsHeaders|sendResponse' drmq-broker/src/main/java drmq-broker/src/main -g '*.java' -C2

Repository: samuel025/DRMQ

Length of output: 50371


Restrict admin CORS or require authentication.

/api/messages returns raw message payloads, and the admin server sets Access-Control-Allow-Origin: "*". Any page an operator visits can treat localhost:<adminPort> as same-origin for data reads; also /api/consumers exposes consumer/lag state. Require auth/token or restrict CORS to the dashboard origin instead of allowing cross-origin reads from every origin.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java` around lines
192 - 196, Update addCorsHeaders in AdminHttpServer to remove the wildcard
Access-Control-Allow-Origin policy and enforce authentication/token validation
or allow only the configured dashboard origin. Ensure the restriction applies to
sensitive endpoints including /api/messages and /api/consumers, while preserving
CORS support only for authorized dashboard requests.


private void sendJsonResponse(HttpExchange exchange, int statusCode, String response) throws IOException {
byte[] bytes = response.getBytes("UTF-8");
exchange.getResponseHeaders().add("Content-Type", "application/json");
exchange.sendResponseHeaders(statusCode, bytes.length);
try (OutputStream os = exchange.getResponseBody()) {
os.write(bytes);
}
}
}
6 changes: 3 additions & 3 deletions drmq-broker/src/main/java/com/drmq/broker/BrokerConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -141,7 +141,7 @@ public static BrokerConfig fromArgs(String[] args) {
String metricsPath = "/metrics";
long logSegmentBytes = 100 * 1024 * 1024L; // 100MB
long logRetentionMs = 7L * 24 * 60 * 60 * 1000; // 7 days
long raftCompactThreshold = 50000L; // Keep 50,000 entries to buffer followers during short outages
long raftCompactThreshold = 1000L;
int maxDeliveries = 5;
String dlqTopicPrefix = "dlq.";
boolean logSegmentFsync = true;
Expand Down Expand Up @@ -188,12 +188,12 @@ public static BrokerConfig fromArgs(String[] args) {

for (int i = 0; i < args.length; i++) {
switch (args[i]) {
case "--config" -> i++; // skip the value, already handled
case "--config" -> i++;
case "--id", "--node-id" -> nodeId = args[++i];
case "--port" -> port = Integer.parseInt(args[++i]);
case "--data-dir" -> dataDir = args[++i];
case "--peers" -> {
peers.clear(); // Override config file peers
peers.clear();
String[] peerStrs = args[++i].split(",");
for (String peerStr : peerStrs) {
peers.add(PeerAddress.parse(peerStr));
Expand Down
8 changes: 8 additions & 0 deletions drmq-broker/src/main/java/com/drmq/broker/BrokerServer.java
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ public class BrokerServer {
private final List<RaftPeer> raftPeers;
private final BrokerMetrics metrics;
private TelemetryWebSocketServer telemetryServer;
private AdminHttpServer adminHttpServer;

private volatile boolean running = false;

Expand Down Expand Up @@ -171,6 +172,10 @@ public void initChannel(SocketChannel ch) {
telemetryServer = new TelemetryWebSocketServer(wsPort, this);
telemetryServer.start();

int adminPort = config.getPort() + 300;
adminHttpServer = new AdminHttpServer(adminPort, messageStore, offsetManager, groupCoordinator);
adminHttpServer.start();

logger.info("DRMQ Broker started on port {} with data directory {}",
config.getPort(), config.getDataDir());

Expand Down Expand Up @@ -214,6 +219,9 @@ public void shutdown() {
if (telemetryServer != null) {
telemetryServer.shutdown();
}
if (adminHttpServer != null) {
adminHttpServer.stop();
}

if (activeChannels != null) {
activeChannels.close().awaitUninterruptibly();
Expand Down
70 changes: 70 additions & 0 deletions drmq-broker/src/main/java/com/drmq/broker/ClientHandler.java
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ private MessageEnvelope handleMessage(MessageEnvelope envelope) throws IOExcepti
case APPEND_ENTRIES_REQUEST -> handleAppendEntriesRequest(envelope);
case INSTALL_SNAPSHOT_REQUEST -> handleInstallSnapshotRequest(envelope);
case SEARCH_OFFSET_BY_TIME_REQUEST -> handleSearchOffsetByTimeRequest(envelope);
case ATOMIC_PRODUCE_REQUEST -> handleAtomicProduceRequest(envelope);
default -> createErrorResponse("Unknown message type: " + envelope.getType());
};
}
Expand Down Expand Up @@ -214,6 +215,74 @@ private MessageEnvelope createProduceBatchErrorResponse(String errorMessage, Err
.build();
}

private MessageEnvelope handleAtomicProduceRequest(MessageEnvelope envelope) throws IOException {
long startNanos = System.nanoTime();
long totalPayloadBytes = 0;
int batchCount = 0;
try {
com.drmq.protocol.DRMQProtocol.AtomicProduceRequest request = com.drmq.protocol.DRMQProtocol.AtomicProduceRequest.parseFrom(envelope.getPayload());

for (var slice : request.getSlicesList()) {
batchCount += slice.getEntriesCount();
for (var entry : slice.getEntriesList()) {
totalPayloadBytes += entry.getPayload().size();
}
}

if (batchCount == 0) {
return createAtomicProduceErrorResponse("Atomic batch must contain at least one message", ErrorCode.UNKNOWN_ERROR);
}
if (totalPayloadBytes > MAX_PAYLOAD_BYTES) {
return createAtomicProduceErrorResponse("Batch payload exceeds maximum size of " + MAX_PAYLOAD_BYTES + " bytes", ErrorCode.UNKNOWN_ERROR);
}

java.util.Map<String, Long> offsets;
if (raftNode != null) {
if (!raftNode.isLeader()) {
String leaderAddr = raftNode.getLeaderAddress();
return createAtomicProduceErrorResponse("NOT_LEADER:" +
(leaderAddr != null ? leaderAddr : "UNKNOWN"), ErrorCode.NOT_LEADER);
}
offsets = raftNode.proposeAtomicBatch(request.getSlicesList());
} else {
offsets = messageStore.appendAtomicBatch(request.getSlicesList());
}
Comment on lines +232 to +249

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Inconsistent handling of single-slice atomic requests between cluster and single-node paths.

There is no explicit "at least 2 topics" check here. In cluster mode raftNode.proposeAtomicBatch throws IllegalArgumentException for <2 slices, which surfaces as a generic UNKNOWN_ERROR; in single-node mode messageStore.appendAtomicBatch accepts a single slice. Add an up-front validation returning a clear error so both paths behave the same regardless of Raft mode. Also note MAX_BATCH_MESSAGES is enforced for PRODUCE_BATCH_REQUEST but not here — confirm whether the atomic path should bound message count too.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@drmq-broker/src/main/java/com/drmq/broker/ClientHandler.java` around lines
232 - 249, Update the atomic produce validation in ClientHandler to reject
requests with fewer than two slices before selecting the Raft or message-store
path, returning a clear client error consistently in both modes. Also enforce
MAX_BATCH_MESSAGES for the atomic request’s total message count if that limit is
intended to apply, preserving the existing payload-size and empty-batch checks.


logger.debug("Produced atomic batch: topics={}, count={}", offsets.keySet(), batchCount);

com.drmq.protocol.DRMQProtocol.AtomicProduceResponse response = com.drmq.protocol.DRMQProtocol.AtomicProduceResponse.newBuilder()
.setSuccess(true)
.putAllBaseOffsets(offsets)
.build();

BrokerMetrics.get().recordRequest("atomic_produce", true,
System.nanoTime() - startNanos, totalPayloadBytes, batchCount);

return MessageEnvelope.newBuilder()
.setType(MessageType.ATOMIC_PRODUCE_RESPONSE)
.setPayload(response.toByteString())
.build();

} catch (Exception e) {
logger.error("Error processing atomic produce request", e);
BrokerMetrics.get().recordRequest("atomic_produce", false,
System.nanoTime() - startNanos, totalPayloadBytes, batchCount);
return createAtomicProduceErrorResponse(e.getMessage(), ErrorCode.UNKNOWN_ERROR);
}
}

private MessageEnvelope createAtomicProduceErrorResponse(String errorMessage, ErrorCode errorCode) {
com.drmq.protocol.DRMQProtocol.AtomicProduceResponse response = com.drmq.protocol.DRMQProtocol.AtomicProduceResponse.newBuilder()
.setSuccess(false)
.setErrorMessage(errorMessage != null ? errorMessage : "Unknown error")
.setErrorCode(errorCode)
.build();
return MessageEnvelope.newBuilder()
.setType(MessageType.ATOMIC_PRODUCE_RESPONSE)
.setPayload(response.toByteString())
.build();
}

private MessageEnvelope handleConsumeRequest(MessageEnvelope envelope) throws IOException {
long startNanos = System.nanoTime();
try {
Expand Down Expand Up @@ -532,6 +601,7 @@ private MessageEnvelope createErrorResponse(String errorMessage, MessageType mes
case COMMIT_OFFSET_RESPONSE -> createCommitOffsetErrorResponse(errorMessage);
case FETCH_OFFSET_RESPONSE -> createFetchOffsetErrorResponse(errorMessage);
case NACK_RESPONSE -> createNackErrorResponse(errorMessage);
case ATOMIC_PRODUCE_RESPONSE -> createAtomicProduceErrorResponse(errorMessage, ErrorCode.UNKNOWN_ERROR);
default -> createProduceErrorResponse(errorMessage, ErrorCode.UNKNOWN_ERROR);
};
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -293,7 +293,6 @@ public int getActiveLeasesCount() {
return total;
}

// ---- Dead-Letter Queue (DLQ) Support ----

/**
* Explicitly reject (NACK) a message offset for a consumer within a group.
Expand All @@ -319,7 +318,6 @@ public boolean nackOffset(String group, String topic, String consumerId, long of

state.lock.lock();
try {
// Remove the consumer's active lease
Lease lease = state.activeLeases.remove(consumerId);
if (lease != null) {
state.members.remove(consumerId);
Expand All @@ -333,7 +331,6 @@ public boolean nackOffset(String group, String topic, String consumerId, long of
routeToDlq(state, group, topic, offset);
return true;
} else {
// Rewind for redelivery using the lease's fromOffset to not skip messages
long rewindOffset = (lease != null) ? lease.fromOffset : offset;
if (rewindOffset < state.dispatchOffset) {
state.dispatchOffset = rewindOffset;
Expand Down Expand Up @@ -379,15 +376,13 @@ private void routeToDlq(GroupTopicState state, String group, String topic, long
}
});

// Advance past the bad offset regardless — don't let a DLQ write failure block progress
long nextOffset = badOffset + 1;
state.committedRanges.add(new CommittedRange(badOffset, nextOffset));
advanceCommittedOffset(state, group, topic);
if (state.dispatchOffset <= badOffset) {
state.dispatchOffset = nextOffset;
}

// Clean up the delivery counter for this offset
state.deliveryCounts.remove(badOffset);
}

Expand Down
Loading