From 59f41e0d4b2dd964c9fdddf98d79b05d57226181 Mon Sep 17 00:00:00 2001 From: Samuel Date: Sun, 19 Jul 2026 22:25:57 +0100 Subject: [PATCH 1/3] feat: add topic and consumer management pages to dashboard, improve throughput formatting, and implement broker admin HTTP server. --- .../java/com/drmq/broker/AdminHttpServer.java | 138 ++++++++++++++++++ .../java/com/drmq/broker/BrokerServer.java | 8 + .../drmq/broker/ConsumerGroupCoordinator.java | 5 - .../java/com/drmq/broker/MessageStore.java | 24 +++ .../java/com/drmq/broker/OffsetManager.java | 8 + .../drmq/broker/TelemetryWebSocketServer.java | 70 ++++++--- .../java/com/drmq/broker/raft/RaftNode.java | 14 +- drmq-dashboard/src/App.tsx | 8 +- drmq-dashboard/src/pages/Consumers.tsx | 134 +++++++++++++++++ drmq-dashboard/src/pages/Dashboard.tsx | 24 +-- drmq-dashboard/src/pages/Topics.tsx | 102 +++++++++++++ drmq-python-client/test_time_seek.py | 72 +++++++++ 12 files changed, 562 insertions(+), 45 deletions(-) create mode 100644 drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java create mode 100644 drmq-dashboard/src/pages/Consumers.tsx create mode 100644 drmq-dashboard/src/pages/Topics.tsx create mode 100644 drmq-python-client/test_time_seek.py diff --git a/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java b/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java new file mode 100644 index 0000000..2db1e62 --- /dev/null +++ b/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java @@ -0,0 +1,138 @@ +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); + + // 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 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 allOffsets = offsetManager.getAllOffsets(); + + // Group by consumer group name + java.util.Map groupsMap = new java.util.HashMap<>(); + + for (java.util.Map.Entry 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 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 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"); + } + + 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); + } + } +} diff --git a/drmq-broker/src/main/java/com/drmq/broker/BrokerServer.java b/drmq-broker/src/main/java/com/drmq/broker/BrokerServer.java index 7534ec9..8a9b627 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/BrokerServer.java +++ b/drmq-broker/src/main/java/com/drmq/broker/BrokerServer.java @@ -46,6 +46,7 @@ public class BrokerServer { private final List raftPeers; private final BrokerMetrics metrics; private TelemetryWebSocketServer telemetryServer; + private AdminHttpServer adminHttpServer; private volatile boolean running = false; @@ -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()); @@ -214,6 +219,9 @@ public void shutdown() { if (telemetryServer != null) { telemetryServer.shutdown(); } + if (adminHttpServer != null) { + adminHttpServer.stop(); + } if (activeChannels != null) { activeChannels.close().awaitUninterruptibly(); diff --git a/drmq-broker/src/main/java/com/drmq/broker/ConsumerGroupCoordinator.java b/drmq-broker/src/main/java/com/drmq/broker/ConsumerGroupCoordinator.java index 66be5cc..f4aed93 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/ConsumerGroupCoordinator.java +++ b/drmq-broker/src/main/java/com/drmq/broker/ConsumerGroupCoordinator.java @@ -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. @@ -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); @@ -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; @@ -379,7 +376,6 @@ 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); @@ -387,7 +383,6 @@ private void routeToDlq(GroupTopicState state, String group, String topic, long state.dispatchOffset = nextOffset; } - // Clean up the delivery counter for this offset state.deliveryCounts.remove(badOffset); } diff --git a/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java b/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java index 47d18f1..eff9c4e 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java +++ b/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java @@ -48,6 +48,9 @@ public class MessageStore implements Closeable { // Topic -> Total number of messages private final ConcurrentHashMap topicMessageCounts = new ConcurrentHashMap<>(); + + // Topic -> Highest offset appended + private final ConcurrentHashMap topicHeadOffsets = new ConcurrentHashMap<>(); // In-memory cache for recent messages (Topic -> BoundedMessageCache) private final ConcurrentHashMap messageCache = new ConcurrentHashMap<>(); @@ -124,6 +127,11 @@ private void recoverInternal() throws IOException { addToCache(topic, message); topicMessageCounts.computeIfAbsent(topic, k -> new AtomicLong()).incrementAndGet(); + AtomicLong head = topicHeadOffsets.computeIfAbsent(topic, k -> new AtomicLong(-1)); + if (offset > head.get()) { + head.set(offset); + } + if (offset > maxOffset) { maxOffset = offset; } @@ -156,6 +164,7 @@ public void reload() throws IOException { logger.info("Reloading MessageStore state from disk..."); topicIndex.clear(); topicMessageCounts.clear(); + topicHeadOffsets.clear(); messageCache.clear(); topicWriteLocks.clear(); @@ -222,6 +231,11 @@ public long append(String topic, byte[] payload, String key, long clientTimestam addToCache(topic, message); topicMessageCounts.computeIfAbsent(topic, k -> new AtomicLong()).incrementAndGet(); + + AtomicLong head = topicHeadOffsets.computeIfAbsent(topic, k -> new AtomicLong(-1)); + if (offset > head.get()) { + head.set(offset); + } logger.debug("Persisted and indexed message: topic={}, offset={}, position={}, segment={}", topic, offset, position, segment.getFilePath().getFileName()); @@ -292,11 +306,15 @@ public long appendBatch(String topic, List entri } AtomicLong counter = topicMessageCounts.computeIfAbsent(topic, k -> new AtomicLong()); + AtomicLong head = topicHeadOffsets.computeIfAbsent(topic, k -> new AtomicLong(-1)); for (int i = 0; i < messages.size(); i++) { StoredMessage msg = messages.get(i); indexMessage(topic, msg.getOffset(), positions.get(i)); addToCache(topic, msg); counter.incrementAndGet(); + if (msg.getOffset() > head.get()) { + head.set(msg.getOffset()); + } } logger.debug("Batch persisted: topic={}, baseOffset={}, count={}", topic, baseOffset, batchSize); @@ -523,6 +541,11 @@ public int getMessageCount(String topic) { return count == null ? 0 : count.intValue(); } + public long getHeadOffset(String topic) { + AtomicLong head = topicHeadOffsets.get(topic); + return head == null ? -1 : head.get(); + } + public List getTopics() { return new ArrayList<>(topicMessageCounts.keySet()); } @@ -542,6 +565,7 @@ public long getCachedMessageCount() { public void clear() { topicIndex.clear(); topicMessageCounts.clear(); + topicHeadOffsets.clear(); messageCache.clear(); globalOffset.set(0); logger.info("Message store memory state cleared"); diff --git a/drmq-broker/src/main/java/com/drmq/broker/OffsetManager.java b/drmq-broker/src/main/java/com/drmq/broker/OffsetManager.java index 6c49d01..28b7e77 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/OffsetManager.java +++ b/drmq-broker/src/main/java/com/drmq/broker/OffsetManager.java @@ -74,6 +74,14 @@ public int getOffsetEntryCount() { return offsets.size(); } + /** + * Get a snapshot of all committed offsets across all groups and topics. + * @return map of "group/topic" to offset + */ + public Map getAllOffsets() { + return new java.util.HashMap<>(offsets); + } + /** Load all offsets from disk on startup. */ private void load() throws IOException { if (!Files.exists(offsetsFile)) { diff --git a/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java b/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java index 642e786..fc5b8cb 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java +++ b/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java @@ -16,8 +16,8 @@ import java.util.List; import java.util.Timer; import java.util.TimerTask; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CopyOnWriteArraySet; -import java.util.concurrent.TimeUnit; public class TelemetryWebSocketServer extends WebSocketServer { private static final Logger logger = LoggerFactory.getLogger(TelemetryWebSocketServer.class); @@ -28,7 +28,10 @@ public class TelemetryWebSocketServer extends WebSocketServer { private final Gson gson; // Throughput history (0-100 scaled for chart) - private final List throughputHistory = new ArrayList<>(); + private final List throughputHistory = new CopyOnWriteArrayList<>(); + private final List produceHistory = new CopyOnWriteArrayList<>(); + private final List consumeHistory = new CopyOnWriteArrayList<>(); + private final List errorHistory = new CopyOnWriteArrayList<>(); // Rolling byte snapshots to compute per-second rates private double lastProduceBytes = 0; @@ -53,14 +56,21 @@ public TelemetryWebSocketServer(int port, BrokerServer brokerServer) { for (int i = 0; i < 30; i++) { throughputHistory.add(0.0); + produceHistory.add(0.0); + consumeHistory.add(0.0); + errorHistory.add(0.0); } } @Override public void onOpen(WebSocket conn, ClientHandshake handshake) { - connections.add(conn); - logger.info("New telemetry dashboard connected: {}", conn.getRemoteSocketAddress()); - conn.send(buildTelemetryPayload()); + try { + connections.add(conn); + logger.info("New telemetry dashboard connected: {}", conn.getRemoteSocketAddress()); + conn.send(buildTelemetryPayload()); + } catch (Exception e) { + logger.error("Error during onOpen, preventing SelectorThread crash", e); + } } @Override @@ -138,14 +148,19 @@ private void updateRates() { lastConsumeRecords = consumeRecords; lastErrorCount = errors; - // Throughput chart history: scale (produce + consume) MB/s to 0-100 double totalMBps = currentProduceMBps + currentConsumeMBps; - // Scale: 50 MB/s = 100 on chart - double chartVal = Math.min(100, (totalMBps / 50.0) * 100); - if (chartVal < 3 && totalMBps > 0) chartVal = 3; // minimum visibility synchronized (throughputHistory) { throughputHistory.remove(0); - throughputHistory.add(chartVal); + throughputHistory.add(totalMBps); + + produceHistory.remove(0); + produceHistory.add(currentProduceMBps); + + consumeHistory.remove(0); + consumeHistory.add(currentConsumeMBps); + + errorHistory.remove(0); + errorHistory.add(currentErrorRate); } } @@ -159,9 +174,9 @@ private String buildTelemetryPayload() { JsonObject metrics = new JsonObject(); double totalMBps = currentProduceMBps + currentConsumeMBps; - metrics.addProperty("totalThroughputMB", round2(totalMBps)); - metrics.addProperty("produceThroughputMB", round2(currentProduceMBps)); - metrics.addProperty("consumeThroughputMB", round2(currentConsumeMBps)); + metrics.addProperty("totalThroughputMB", round4(totalMBps)); + metrics.addProperty("produceThroughputMB", round4(currentProduceMBps)); + metrics.addProperty("consumeThroughputMB", round4(currentConsumeMBps)); metrics.addProperty("produceRate", Math.round(currentProduceRecordRate)); // msgs/s metrics.addProperty("consumeRate", Math.round(currentConsumeRecordRate)); // msgs/s metrics.addProperty("errorRate", Math.round(currentErrorRate)); // errors/s @@ -175,8 +190,8 @@ private String buildTelemetryPayload() { produceLatencyMs = Math.max(singleLatency, batchLatency); // use whichever path is active consumeLatencyMs = bm.getTimerMeanMs("drmq.broker.request.latency", "consume"); } - metrics.addProperty("produceLatencyMs", round2(produceLatencyMs)); - metrics.addProperty("consumeLatencyMs", round2(consumeLatencyMs)); + metrics.addProperty("produceLatencyMs", round4(produceLatencyMs)); + metrics.addProperty("consumeLatencyMs", round4(consumeLatencyMs)); // Active connections: real handler count from Netty int totalHandlers = brokerServer.getActiveChannelsCount(); @@ -265,12 +280,21 @@ private String buildTelemetryPayload() { // Throughput chart history JsonArray history = new JsonArray(); + JsonArray pHistory = new JsonArray(); + JsonArray cHistory = new JsonArray(); + JsonArray eHistory = new JsonArray(); synchronized (throughputHistory) { - for (Double val : throughputHistory) { - history.add(round2(val)); + for (int i = 0; i < throughputHistory.size(); i++) { + history.add(round4(throughputHistory.get(i))); + pHistory.add(round4(produceHistory.get(i))); + cHistory.add(round4(consumeHistory.get(i))); + eHistory.add(round4(errorHistory.get(i))); } } metrics.add("throughputHistory", history); + metrics.add("produceHistory", pHistory); + metrics.add("consumeHistory", cHistory); + metrics.add("errorHistory", eHistory); // Produce-rate history (msgs/s scaled similarly) state.add("metrics", metrics); @@ -286,7 +310,7 @@ private String buildTelemetryPayload() { localNode.addProperty("name", "Broker-" + localId.toUpperCase()); localNode.addProperty("status", localStatus); // throughput in bytes/s for this node (real produce traffic it handled) - localNode.addProperty("throughputMBps", round2(currentProduceMBps + currentConsumeMBps)); + localNode.addProperty("throughputMBps", round4(currentProduceMBps + currentConsumeMBps)); localNode.addProperty("produceRate", Math.round(currentProduceRecordRate)); localNode.addProperty("consumeRate", Math.round(currentConsumeRecordRate)); localNode.addProperty("commitIndex", commitIndex); @@ -366,9 +390,9 @@ private String buildTelemetryPayload() { } JsonObject latencies = new JsonObject(); // Inter-node Raft RPC latency (same for all links since we only know our own outbound) - latencies.addProperty("alphaBeta", round2(raftRpcLatencyMs > 0 ? raftRpcLatencyMs : 0)); - latencies.addProperty("betaGamma", round2(raftRpcLatencyMs > 0 ? raftRpcLatencyMs : 0)); - latencies.addProperty("raftRpcMs", round2(raftRpcLatencyMs)); + latencies.addProperty("alphaBeta", round4(raftRpcLatencyMs > 0 ? raftRpcLatencyMs : 0)); + latencies.addProperty("betaGamma", round4(raftRpcLatencyMs > 0 ? raftRpcLatencyMs : 0)); + latencies.addProperty("raftRpcMs", round4(raftRpcLatencyMs)); state.add("latencies", latencies); // ── EVENTS ──────────────────────────────────────────────────────────── @@ -377,7 +401,7 @@ private String buildTelemetryPayload() { return gson.toJson(state); } - private static double round2(double v) { - return Math.round(v * 100.0) / 100.0; + private static double round4(double v) { + return Math.round(v * 10000.0) / 10000.0; } } 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 d9329bc..d0074c9 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 @@ -45,16 +45,16 @@ public class RaftNode { private static final int MAX_PENDING_PROPOSALS = 10000; // Safety limit // Persistent state (survives restart) - private long currentTerm; - private String votedFor; + private volatile long currentTerm; + private volatile String votedFor; private final RaftLog raftLog; // Volatile state - private RaftState state; - private long commitIndex; - private long lastApplied; - private long lastAppliedTerm; - private String leaderId; + private volatile RaftState state; + private volatile long commitIndex; + private volatile long lastApplied; + private volatile long lastAppliedTerm; + private volatile String leaderId; // Leader-only volatile state private final Map nextIndex; diff --git a/drmq-dashboard/src/App.tsx b/drmq-dashboard/src/App.tsx index fa43a3c..a7b0c28 100644 --- a/drmq-dashboard/src/App.tsx +++ b/drmq-dashboard/src/App.tsx @@ -1,8 +1,10 @@ import { BrowserRouter, Routes, Route, Link, useLocation } from 'react-router-dom'; -import { LayoutDashboard, Book, Server, ChevronRight } from 'lucide-react'; +import { LayoutDashboard, Book, Server, ChevronRight, Layers, Users } from 'lucide-react'; import { useState, useCallback } from 'react'; import Dashboard from './pages/Dashboard'; import { useClusterTelemetry } from './useClusterTelemetry'; +import Topics from './pages/Topics'; +import Consumers from './pages/Consumers'; function Sidebar({ telemetryState }: { telemetryState: any }) { const location = useLocation(); @@ -11,6 +13,8 @@ function Sidebar({ telemetryState }: { telemetryState: any }) { const navItems = [ { to: '/', icon: LayoutDashboard, label: 'Telemetry' }, + { to: '/topics', icon: Layers, label: 'Topics' }, + { to: '/consumers', icon: Users, label: 'Consumers' }, { to: '/docs', icon: Book, label: 'Documentation' }, ]; @@ -130,6 +134,8 @@ export default function App() { style={{ minHeight: '100vh' }} /> } /> + } /> + } /> diff --git a/drmq-dashboard/src/pages/Consumers.tsx b/drmq-dashboard/src/pages/Consumers.tsx new file mode 100644 index 0000000..b044c78 --- /dev/null +++ b/drmq-dashboard/src/pages/Consumers.tsx @@ -0,0 +1,134 @@ +import React, { useState, useEffect } from 'react'; +import { Panel } from '../components/DashboardWidgets'; +import { RefreshCw, Users, AlertCircle, CheckCircle2 } from 'lucide-react'; +import { DitherButton } from '../components/dither-kit'; + +interface ConsumerTopic { + topic: string; + committedOffset: number; + lag: number; + activeMembers: number; +} + +interface ConsumerGroup { + groupId: string; + topics: ConsumerTopic[]; +} + +export default function Consumers() { + const [groups, setGroups] = useState([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + + const fetchConsumers = async () => { + setLoading(true); + setError(null); + try { + const res = await fetch('http://localhost:9392/api/consumers'); + if (!res.ok) throw new Error(`HTTP error! status: ${res.status}`); + const data = await res.json(); + setGroups(data); + } catch (err: any) { + setError(err.message || 'Failed to fetch consumers'); + } finally { + setLoading(false); + } + }; + + useEffect(() => { + fetchConsumers(); + const interval = setInterval(fetchConsumers, 5000); + return () => clearInterval(interval); + }, []); + + const getLagColor = (lag: number) => { + if (lag === 0) return 'text-emerald-400'; + if (lag < 100) return 'text-amber-400'; + return 'text-red-400'; + }; + + return ( +
+
+
+

+ + Consumer Groups & Lag +

+

Monitor message processing lag across all consumer groups

+
+ + + Refresh + +
+ + + {error && ( +
+ {error} +
+ )} + +
+ + + + + + + + + + + + {groups.length === 0 && !loading && !error && ( + + + + )} + {groups.flatMap((group) => + group.topics.map((topic, i) => ( + + {i === 0 ? ( + + ) : null} + + + + + + )) + )} + +
Group IDTopicOffsetConsumer LagMembers
+ No consumer groups registered. +
+ + {group.groupId} + + + {topic.topic} + + {topic.committedOffset.toLocaleString()} + +
+ {topic.lag === 0 ? ( + + ) : ( + 100 ? 'text-red-500/70' : 'text-amber-500/70'}`} /> + )} + + {topic.lag.toLocaleString()} + +
+
+ + {topic.activeMembers} + +
+
+
+
+ ); +} diff --git a/drmq-dashboard/src/pages/Dashboard.tsx b/drmq-dashboard/src/pages/Dashboard.tsx index 81ba2a1..ab1cd76 100644 --- a/drmq-dashboard/src/pages/Dashboard.tsx +++ b/drmq-dashboard/src/pages/Dashboard.tsx @@ -92,7 +92,7 @@ export default function Dashboard({ const { metrics, latencies } = telemetryState; setLatencyHistory(prev => { const next = [...prev, { - t: new Date().toLocaleTimeString([], { second: '2-digit', minute: '2-digit' }), + t: new Date().toLocaleTimeString([], { hour: '2-digit', minute: '2-digit', second: '2-digit' }), produce: metrics.produceLatencyMs, consume: metrics.consumeLatencyMs, rpc: latencies.raftRpcMs @@ -139,18 +139,24 @@ export default function Dashboard({ /* ── Throughput chart data ─────────────────────────────────────── */ const totalHist = metrics.throughputHistory ?? []; + const produceHist = metrics.produceHistory ?? []; const consumeHist = metrics.consumeHistory ?? []; const errorHist = metrics.errorHistory ?? []; const throughputData = totalHist.map((v: number, i: number) => { const consumeVal = consumeHist[i] ?? 0; - const produceVal = Math.max(0, v - consumeVal); + let produceVal = produceHist[i]; + if (produceVal === undefined) { + produceVal = Math.max(0, v - consumeVal); + } return { t: `-${totalHist.length - i}s`, - produce: parseFloat(((produceVal / 100) * 50).toFixed(2)), - consume: parseFloat(((consumeVal / 100) * 50).toFixed(2)), + produce: parseFloat(produceVal.toFixed(2)), + consume: parseFloat(consumeVal.toFixed(2)), }; }); + + const formatMBps = (val: number) => val > 0 && val < 0.005 ? '< 0.01' : val.toFixed(2); /* ── Radar data for selected/leader node ──────────────────────── */ const leaderNode = nodes.find((n: any) => n.status === 'LEADER'); @@ -219,11 +225,11 @@ export default function Dashboard({
} /> } /> } className="shrink-0">
- {metrics.totalThroughputMB.toFixed(2)} MB/s + {formatMBps(metrics.totalThroughputMB)} MB/s
- ↑ {metrics.produceThroughputMB.toFixed(2)} · ↓ {metrics.consumeThroughputMB.toFixed(2)} + ↑ {formatMBps(metrics.produceThroughputMB)} · ↓ {formatMBps(metrics.consumeThroughputMB)}
- + diff --git a/drmq-dashboard/src/pages/Topics.tsx b/drmq-dashboard/src/pages/Topics.tsx new file mode 100644 index 0000000..e6e63a6 --- /dev/null +++ b/drmq-dashboard/src/pages/Topics.tsx @@ -0,0 +1,102 @@ +import React, { useState, useEffect } from 'react'; +import { Panel } from '../components/DashboardWidgets'; +import { RefreshCw, Database } from 'lucide-react'; +import { DitherButton } from '../components/dither-kit'; + +interface TopicData { + name: string; + messageCount: number; +} + +export default function Topics() { + const [topics, setTopics] = useState([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + + const fetchTopics = async () => { + setLoading(true); + setError(null); + try { + // For now, default to the local admin port of Broker-1 (9392) + const res = await fetch('http://localhost:9392/api/topics'); + if (!res.ok) throw new Error(`HTTP error! status: ${res.status}`); + const data = await res.json(); + setTopics(data); + } catch (err: any) { + setError(err.message || 'Failed to fetch topics'); + } finally { + setLoading(false); + } + }; + + useEffect(() => { + fetchTopics(); + const interval = setInterval(fetchTopics, 5000); + return () => clearInterval(interval); + }, []); + + return ( +
+
+
+

+ + Topic Explorer +

+

Live visibility into DRMQ active topics

+
+ + + Refresh + +
+ + + {error && ( +
+ {error} +
+ )} + +
+ + + + + + + + + + {topics.length === 0 && !loading && !error && ( + + + + )} + {topics.map(topic => ( + + + + + + ))} + +
Topic NameMessage CountStatus
+ No active topics found. +
+
+
+ {topic.name} +
+
+ {topic.messageCount.toLocaleString()} + + + Active + +
+
+
+
+ ); +} diff --git a/drmq-python-client/test_time_seek.py b/drmq-python-client/test_time_seek.py new file mode 100644 index 0000000..1832520 --- /dev/null +++ b/drmq-python-client/test_time_seek.py @@ -0,0 +1,72 @@ +import time +import logging +from drmq_client import DRMQProducer, DRMQConsumer + +logging.basicConfig(level=logging.INFO) + +def run_test(): + servers = "localhost:9092,localhost:9093,localhost:9094" + topic = "time-seek-topic-" + str(int(time.time())) + + print(f"--- Starting Time-Based Seek Test on {topic} ---") + producer = DRMQProducer(servers) + try: + producer.connect() + + # Produce 1st message + producer.send(topic, b"Message 1 - Early").result() + print("Sent Message 1") + + time.sleep(1) + # Capture timestamp before Message 2 + target_time = int(time.time() * 1000) + print(f"Captured target timestamp: {target_time}") + time.sleep(1) + + # Produce 2nd and 3rd messages + producer.send(topic, b"Message 2 - Target").result() + print("Sent Message 2") + producer.send(topic, b"Message 3 - Late").result() + print("Sent Message 3") + + finally: + producer.close() + + print("\n--- Starting Consumer ---") + # Single mode consumer (no group_id) + consumer = DRMQConsumer(servers, group_id=None) + try: + consumer.connect() + + print(f"Seeking to timestamp {target_time}...") + consumer.seek_by_time(topic, target_time) + print(f"Consumer local offsets: {consumer.local_offsets}") + + print("Polling...") + messages = consumer.poll(max_messages=10, timeout_ms=2000) + print(f"Polled {len(messages)} messages") + + if len(messages) == 0: + print("TEST FAILED: No messages received.") + # Let's see what is actually in the broker if we start from 0 + consumer.local_offsets[topic] = 0 + msgs_from_0 = consumer.poll(max_messages=100, timeout_ms=2000) + print(f"Polled from 0: {len(msgs_from_0)} messages") + for m in msgs_from_0: + print(f" - offset={m.offset}, ts={m.timestamp}, payload={m.payload.decode('utf-8')}") + exit(1) + + first_msg = messages[0].payload.decode('utf-8') + print(f"First message received: {first_msg}") + + if "Message 2 - Target" in first_msg: + print("TEST PASSED! Seek successfully skipped Message 1 and started at Message 2.") + else: + print(f"TEST FAILED! Expected Message 2, but got: {first_msg}") + exit(1) + + finally: + consumer.close() + +if __name__ == "__main__": + run_test() From d3345677b48c8b51df7ca797533ddc38607e5de1 Mon Sep 17 00:00:00 2001 From: Samuel Date: Sun, 19 Jul 2026 22:59:45 +0100 Subject: [PATCH 2/3] feat: add message inspection UI to dashboard and introduce broker configuration and stress testing scripts --- .../java/com/drmq/broker/AdminHttpServer.java | 58 +++++++ .../drmq/broker/TelemetryWebSocketServer.java | 2 +- drmq-dashboard/src/pages/Dashboard.tsx | 30 +++- drmq-dashboard/src/pages/Topics.tsx | 158 +++++++++++++++++- 4 files changed, 242 insertions(+), 6 deletions(-) diff --git a/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java b/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java index 2db1e62..ac02d76 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java +++ b/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java @@ -34,6 +34,7 @@ public AdminHttpServer(int port, MessageStore messageStore, OffsetManager offset 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); @@ -121,6 +122,63 @@ private void handleConsumers(HttpExchange exchange) throws IOException { 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; + + 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]); + } + } + + if (topic == null) { + sendJsonResponse(exchange, 400, "{\"error\":\"Missing 'topic' parameter\"}"); + return; + } + + // Limit bounds to avoid OOM + limit = Math.min(100, Math.max(1, limit)); + + List 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() + "\"}"); + } + } + private void addCorsHeaders(HttpExchange exchange) { exchange.getResponseHeaders().add("Access-Control-Allow-Origin", "*"); exchange.getResponseHeaders().add("Access-Control-Allow-Methods", "GET, POST, DELETE, OPTIONS"); diff --git a/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java b/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java index fc5b8cb..15bff02 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java +++ b/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java @@ -54,7 +54,7 @@ public TelemetryWebSocketServer(int port, BrokerServer brokerServer) { this.broadcastTimer = new Timer("TelemetryBroadcastTimer", true); this.gson = new Gson(); - for (int i = 0; i < 30; i++) { + for (int i = 0; i < 300; i++) { throughputHistory.add(0.0); produceHistory.add(0.0); consumeHistory.add(0.0); diff --git a/drmq-dashboard/src/pages/Dashboard.tsx b/drmq-dashboard/src/pages/Dashboard.tsx index ab1cd76..d030d5e 100644 --- a/drmq-dashboard/src/pages/Dashboard.tsx +++ b/drmq-dashboard/src/pages/Dashboard.tsx @@ -78,6 +78,7 @@ export default function Dashboard({ }) { /* ── Hooks (must come before any early returns) ───────────────── */ + const [historyWindowSeconds, setHistoryWindowSeconds] = useState(30); const prevCommit = useRef(0); const [latencyHistory, setLatencyHistory] = useState([]); @@ -138,10 +139,15 @@ export default function Dashboard({ const leaderName = nodes.find((n: any) => n.status === 'LEADER')?.name ?? 'No Leader'; /* ── Throughput chart data ─────────────────────────────────────── */ - const totalHist = metrics.throughputHistory ?? []; - const produceHist = metrics.produceHistory ?? []; - const consumeHist = metrics.consumeHistory ?? []; - const errorHist = metrics.errorHistory ?? []; + const totalHistFull = metrics.throughputHistory ?? []; + const produceHistFull = metrics.produceHistory ?? []; + const consumeHistFull = metrics.consumeHistory ?? []; + const errorHistFull = metrics.errorHistory ?? []; + + const totalHist = totalHistFull.slice(-historyWindowSeconds); + const produceHist = produceHistFull.slice(-historyWindowSeconds); + const consumeHist = consumeHistFull.slice(-historyWindowSeconds); + const errorHist = errorHistFull.slice(-historyWindowSeconds); const throughputData = totalHist.map((v: number, i: number) => { const consumeVal = consumeHist[i] ?? 0; @@ -322,6 +328,22 @@ export default function Dashboard({
+
+ 30s + setHistoryWindowSeconds(Number(e.target.value))} + className="flex-1 h-1.5 bg-white/10 rounded-full appearance-none [&::-webkit-slider-thumb]:appearance-none [&::-webkit-slider-thumb]:w-3 [&::-webkit-slider-thumb]:h-3 [&::-webkit-slider-thumb]:bg-cyan-400 [&::-webkit-slider-thumb]:rounded-full cursor-pointer" + /> + 5min +
+
+ VIEWING LAST {historyWindowSeconds >= 60 ? `${historyWindowSeconds / 60} MIN` : `${historyWindowSeconds} SEC`} +
{/* ── Node Health Radar ──────────────────────────────── */} diff --git a/drmq-dashboard/src/pages/Topics.tsx b/drmq-dashboard/src/pages/Topics.tsx index e6e63a6..f4e6ab2 100644 --- a/drmq-dashboard/src/pages/Topics.tsx +++ b/drmq-dashboard/src/pages/Topics.tsx @@ -2,6 +2,7 @@ import React, { useState, useEffect } from 'react'; import { Panel } from '../components/DashboardWidgets'; import { RefreshCw, Database } from 'lucide-react'; import { DitherButton } from '../components/dither-kit'; +import { motion, AnimatePresence } from 'framer-motion'; interface TopicData { name: string; @@ -13,6 +14,42 @@ export default function Topics() { const [loading, setLoading] = useState(true); const [error, setError] = useState(null); + const [inspectorTopic, setInspectorTopic] = useState(null); + const [inspectOffset, setInspectOffset] = useState(0); + const [inspectLimit, setInspectLimit] = useState(10); + const [messages, setMessages] = useState([]); + const [inspectLoading, setInspectLoading] = useState(false); + const [inspectError, setInspectError] = useState(null); + + const fetchMessages = async () => { + if (!inspectorTopic) return; + setInspectLoading(true); + setInspectError(null); + try { + const offsetParam = inspectOffset === '' ? 0 : inspectOffset; + const limitParam = inspectLimit === '' ? 1 : inspectLimit; + const res = await fetch(`http://localhost:9392/api/messages?topic=${inspectorTopic}&offset=${offsetParam}&limit=${limitParam}`); + if (!res.ok) throw new Error(`HTTP error! status: ${res.status}`); + const data = await res.json(); + if (data.error) throw new Error(data.error); + setMessages(data); + } catch (err: any) { + setInspectError(err.message || 'Failed to fetch messages'); + } finally { + setInspectLoading(false); + } + }; + + // Automatically fetch messages when inspector opens + useEffect(() => { + if (inspectorTopic) { + setInspectOffset(0); // Reset to 0 when opening new topic + fetchMessages(); + } else { + setMessages([]); + } + }, [inspectorTopic]); + const fetchTopics = async () => { setLoading(true); setError(null); @@ -65,12 +102,13 @@ export default function Topics() { Topic Name Message Count Status + Actions {topics.length === 0 && !loading && !error && ( - + No active topics found. @@ -91,12 +129,130 @@ export default function Topics() { Active + + + ))}
+ + {/* Live Message Inspector Drawer */} + + {inspectorTopic && ( + +
+
+

+ + Inspect Topic +

+
{inspectorTopic}
+
+ +
+ +
+
+ + setInspectOffset(e.target.value === '' ? '' : parseInt(e.target.value))} + className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors" + /> +
+
+ + setInspectLimit(e.target.value === '' ? '' : parseInt(e.target.value))} + className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors" + /> +
+
+ + + + Fetch + +
+
+ +
+ {inspectError && ( +
+ ⚠️ +
{inspectError}
+
+ )} + + {messages.length === 0 && !inspectLoading && !inspectError && ( +
+ +
NO MESSAGES FOUND
+
+ )} + +
+ {messages.map((msg, idx) => { + // Attempt to prettify JSON payloads + let formattedPayload = msg.payload; + try { + const parsed = JSON.parse(msg.payload); + formattedPayload = JSON.stringify(parsed, null, 2); + } catch (e) { /* ignore */ } + + return ( +
+
+
+ + OFFSET {msg.offset} + +
+ {new Date(msg.storedAt).toISOString()} +
+ + {msg.key && ( +
+ KEY + {msg.key} +
+ )} + +
+
+                          {formattedPayload}
+                        
+
+
+ ); + })} +
+
+
+ )} +
); } From 68b9fc38e54c683cf6b55fb903c5f16bbb92ddb4 Mon Sep 17 00:00:00 2001 From: Samuel Date: Mon, 20 Jul 2026 00:02:41 +0100 Subject: [PATCH 3/3] feat: add message timestamp seeking support to broker and update dashboard inspection UI --- .../java/com/drmq/broker/AdminHttpServer.java | 10 ++ .../java/com/drmq/broker/MessageStore.java | 9 +- drmq-dashboard/src/pages/Topics.tsx | 95 +++++++++++++------ 3 files changed, 85 insertions(+), 29 deletions(-) diff --git a/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java b/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java index ac02d76..72ef596 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java +++ b/drmq-broker/src/main/java/com/drmq/broker/AdminHttpServer.java @@ -139,6 +139,7 @@ private void handleMessages(HttpExchange exchange) throws IOException { String topic = null; long offset = 0; int limit = 10; + Long timestamp = null; for (String param : query.split("&")) { String[] pair = param.split("="); @@ -146,6 +147,7 @@ private void handleMessages(HttpExchange exchange) throws IOException { 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]); } } @@ -157,6 +159,14 @@ private void handleMessages(HttpExchange exchange) throws IOException { // 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 messages = messageStore.getMessages(topic, offset, limit); JsonArray msgsArray = new JsonArray(); diff --git a/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java b/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java index eff9c4e..e558b8e 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java +++ b/drmq-broker/src/main/java/com/drmq/broker/MessageStore.java @@ -706,9 +706,16 @@ public void removeRange(long fromOffset, long toOffset) { public List getMessagesFrom(long fromOffset, int maxCount) { lock.readLock().lock(); try { + if (!cache.containsKey(fromOffset)) { + return Collections.emptyList(); + } List result = new ArrayList<>(); + boolean found = false; for (Map.Entry entry : cache.entrySet()) { - if (entry.getKey() >= fromOffset) { + if (entry.getKey() == fromOffset) { + found = true; + } + if (found) { result.add(entry.getValue()); if (result.size() >= maxCount) { break; diff --git a/drmq-dashboard/src/pages/Topics.tsx b/drmq-dashboard/src/pages/Topics.tsx index f4e6ab2..a9d628b 100644 --- a/drmq-dashboard/src/pages/Topics.tsx +++ b/drmq-dashboard/src/pages/Topics.tsx @@ -16,6 +16,7 @@ export default function Topics() { const [inspectorTopic, setInspectorTopic] = useState(null); const [inspectOffset, setInspectOffset] = useState(0); + const [inspectTimestamp, setInspectTimestamp] = useState(''); const [inspectLimit, setInspectLimit] = useState(10); const [messages, setMessages] = useState([]); const [inspectLoading, setInspectLoading] = useState(false); @@ -28,11 +29,29 @@ export default function Topics() { try { const offsetParam = inspectOffset === '' ? 0 : inspectOffset; const limitParam = inspectLimit === '' ? 1 : inspectLimit; - const res = await fetch(`http://localhost:9392/api/messages?topic=${inspectorTopic}&offset=${offsetParam}&limit=${limitParam}`); + + let url = `http://localhost:9392/api/messages?topic=${inspectorTopic}&limit=${limitParam}`; + if (inspectTimestamp) { + const ts = new Date(inspectTimestamp).getTime(); + if (!isNaN(ts)) { + url += `×tamp=${ts}`; + } else { + url += `&offset=${offsetParam}`; + } + } else { + url += `&offset=${offsetParam}`; + } + + const res = await fetch(url); if (!res.ok) throw new Error(`HTTP error! status: ${res.status}`); const data = await res.json(); if (data.error) throw new Error(data.error); setMessages(data); + + // Snap the input offset to the actual returned offset to show compaction jumps + if (data && data.length > 0) { + setInspectOffset(data[0].offset); + } } catch (err: any) { setInspectError(err.message || 'Failed to fetch messages'); } finally { @@ -82,10 +101,13 @@ export default function Topics() {

Live visibility into DRMQ active topics

- + @@ -132,7 +154,7 @@ export default function Topics() { @@ -170,31 +192,48 @@ export default function Topics() { -
-
- - setInspectOffset(e.target.value === '' ? '' : parseInt(e.target.value))} - className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors" - /> -
-
- - setInspectLimit(e.target.value === '' ? '' : parseInt(e.target.value))} - className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors" - /> +
+
+
+ + setInspectOffset(e.target.value === '' ? '' : parseInt(e.target.value))} + className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors disabled:opacity-30" + placeholder="e.g. 100" + /> +
+
+ + setInspectTimestamp(e.target.value)} + className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors [&::-webkit-calendar-picker-indicator]:invert" + /> +
-
- - - - Fetch - +
+
+ + setInspectLimit(e.target.value === '' ? '' : parseInt(e.target.value))} + className="bg-black/40 border border-white/10 rounded px-3 py-2 text-sm text-white font-mono focus:outline-none focus:border-cyan-500/50 transition-colors" + /> +
+
+ +