From 955e5350da218781bd8026076125b79a2d99d1c1 Mon Sep 17 00:00:00 2001 From: Samuel Date: Mon, 13 Jul 2026 19:12:16 +0100 Subject: [PATCH] refactor: improve replication lag reporting accuracy and optimize Raft log compaction logic --- README.md | 4 +- .../java/com/drmq/broker/BrokerConfig.java | 2 +- .../drmq/broker/TelemetryWebSocketServer.java | 35 +++- .../java/com/drmq/broker/raft/RaftNode.java | 59 ++----- .../java/com/drmq/client/DRMQConsumer.java | 2 +- drmq-ts-client/src/client.ts | 2 +- drmq-ts-client/src/example.ts | 41 ----- drmq-ts-client/src/messages.ts | 167 ++++++++++++++---- 8 files changed, 185 insertions(+), 127 deletions(-) delete mode 100644 drmq-ts-client/src/example.ts diff --git a/README.md b/README.md index 7ef6b4b..d9d0bea 100644 --- a/README.md +++ b/README.md @@ -18,9 +18,11 @@ The project is structured as a multi-module Maven build, separating the core bro - **Persistent Storage:** Custom Write-Ahead Log (WAL) and segment-based message storage ensure messages are durably persisted to disk. Features thread-safe, atomic consumer offset management with bounds locking designed to minimize data loss during concurrent background writes and handle shutdowns gracefully. - **Graceful Teardown Coordination:** Orchestrated, safe termination of Netty EventLoops, RPC executors, and disk storage guaranteeing state integrity without resource leaks during node shutdowns. - **High Performance:** + - **Client-Side Batching:** Producers feature high-throughput, latency-optimized message batching via a configurable `linger.ms` window. This groups thousands of messages into a single network round-trip and Raft log flush, massively increasing throughput. + - **Configurable Disk Durability:** By default, DRMQ guarantees strict flush-before-ack durability (`fsync`). However, administrators can explicitly disable this (`--log-segment-fsync false`) for extreme throughput scenarios where hardware page-cache flushing is acceptable. - **Follower-based Reads:** Scalable read operations allowing consumers to fetch messages from follower nodes, distributing the load across the cluster. - **Dedicated Thread Pools:** Separated executor services for Raft tasks and client handling prevent thread starvation and ensure consistent performance. Independent scheduler threads prevent I/O blocking from stunting cluster heartbeats. -- **Robust Client Ecosystem:** Includes a producer and consumer client with automatic reconnects and randomized bootstrap load balancing for seamless failover operations. +- **Robust Client Ecosystem:** Includes Java, Python, and TypeScript SDKs featuring automatic reconnects, randomized bootstrap load balancing, typed Error Code handling, and seamless leader failovers. - **Metrics & Observability:** Integrated with Micrometer and Prometheus to provide deep visibility into broker health, log replication lag, and throughput metrics. ## Architecture & Modules diff --git a/drmq-broker/src/main/java/com/drmq/broker/BrokerConfig.java b/drmq-broker/src/main/java/com/drmq/broker/BrokerConfig.java index 3b2de22..dfeef2e 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/BrokerConfig.java +++ b/drmq-broker/src/main/java/com/drmq/broker/BrokerConfig.java @@ -125,7 +125,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 = 1000L; + long raftCompactThreshold = 50000L; // Keep 50,000 entries to buffer followers during short outages int maxDeliveries = 5; String dlqTopicPrefix = "dlq."; boolean logSegmentFsync = true; 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 22d7777..642e786 100644 --- a/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java +++ b/drmq-broker/src/main/java/com/drmq/broker/TelemetryWebSocketServer.java @@ -296,7 +296,6 @@ private String buildTelemetryPayload() { localNode.addProperty("y", 200); nodes.add(localNode); - // Peer nodes — real Raft state, honest about what we don't know if (raftNode != null) { int[] positions = { 0, 1 }; int i = 0; @@ -306,15 +305,37 @@ private String buildTelemetryPayload() { Long matchIdx = raftNode.getMatchIndexMap().get(peerId); long peerApplied = matchIdx != null ? matchIdx : 0; - // Replication lag for this peer. - // If matchIndex is unconfirmed (null/0) and commitIndex is large, - // the new leader simply hasn't heard from this peer since election. - // Report -1 to signal "unreachable" rather than a fictitious huge lag. - long replicationLag; + long replicationLag = 0; if ((matchIdx == null || matchIdx == 0) && commitIndex > 100) { replicationLag = -1; // unreachable / unknown } else { - replicationLag = Math.max(0, commitIndex - peerApplied); + long lagEntries = Math.max(0, commitIndex - peerApplied); + if (lagEntries == 0) { + replicationLag = 0; + } else if (raftNode.getRaftLog() != null && lagEntries <= 2000) { + // Calculate exact message lag by inspecting uncommitted Raft entries + long msgLag = 0; + for (long idx = peerApplied + 1; idx <= commitIndex; idx++) { + com.drmq.protocol.DRMQProtocol.RaftEntry entry = raftNode.getRaftLog().getEntry(idx); + if (entry != null && entry.getCommandType() == com.drmq.protocol.DRMQProtocol.RaftCommandType.BATCH_MESSAGE) { + try { + com.drmq.protocol.DRMQProtocol.ProduceBatchRequest batch = + com.drmq.protocol.DRMQProtocol.ProduceBatchRequest.parseFrom(entry.getPayload()); + msgLag += batch.getEntriesCount(); + } catch (Exception e) { + msgLag += 1; // Fallback + } + } else if (entry != null && entry.getCommandType() == com.drmq.protocol.DRMQProtocol.RaftCommandType.MESSAGE) { + msgLag += 1; + } + } + replicationLag = msgLag; + } else if (lagEntries > 2000) { + // Fallback estimate if lag is extremely large to avoid blocking telemetry thread + double avgBatchSize = currentProduceRecordRate > 0 ? + Math.max(1.0, currentProduceRecordRate / 10.0) : 50.0; + replicationLag = (long) (lagEntries * avgBatchSize); + } } peerNode.addProperty("id", peerId); 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 c4005a3..459e476 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 @@ -79,9 +79,6 @@ public class RaftNode { private final Map> installSnapshotRpcHandlers = new ConcurrentHashMap<>(); private final ReentrantLock lock = new ReentrantLock(); - // Using a cached thread pool or larger scheduled pool size - // We have 5 timed tasks (election, heartbeat, quorum check, proposal cleanup, state save). - // Given they are short-lived, 4 threads minimizes scheduling collision while preserving efficiency. private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(4); private final ExecutorService raftExecutor; private ScheduledFuture electionTimer; @@ -112,7 +109,6 @@ private static class ProposalState { private final Map lastContactTime = new ConcurrentHashMap<>(); private static final long LOG_RATE_LIMIT_MS = 1000; - // Snapshot receive state private long snapshotReceiveOffset = 0; private java.io.OutputStream snapshotReceiveStream = null; private Path snapshotTempFile = null; @@ -152,9 +148,6 @@ public RaftNode(String nodeId, int port, List peers, loadPersistentState(); } - /** - * Backward-compatible constructor using default compaction threshold (1000). - */ public RaftNode(String nodeId, int port, List peers, MessageStore messageStore, OffsetManager offsetManager, Path dataDir) throws IOException { this(nodeId, port, peers, messageStore, offsetManager, dataDir, 1000L); @@ -254,6 +247,11 @@ public void stop() { logger.info("[{}] Raft node stopped", nodeId); } + public RaftLog getRaftLog() { + return raftLog; + } + + // Peer RPC Registration @@ -346,13 +344,6 @@ private void startPreVote() { AtomicLong votesReceived = new AtomicLong(1); // self-vote AtomicBoolean electionStarted = new AtomicBoolean(false); - // Single-node cluster: self-vote alone satisfies quorum, start election immediately - if (votesReceived.get() >= votesNeeded) { - if (electionStarted.compareAndSet(false, true)) { - logger.info("[{}] Pre-vote succeeded (single-node quorum), starting real election", nodeId); - startElection(proposedTerm); - } - } for (PeerAddress peer : peers) { CompletableFuture.supplyAsync(() -> { @@ -389,7 +380,6 @@ private void startPreVote() { } finally { lock.unlock(); } - // Release lock before startElection() to avoid holding it during disk I/O if (shouldStartElection) { startElection(proposedTerm); } @@ -399,10 +389,7 @@ private void startPreVote() { resetElectionTimer(); } - /** - * 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(long proposedTerm) { if (!running) return; @@ -447,19 +434,6 @@ private void startElection(long proposedTerm) { AtomicLong votesReceived = new AtomicLong(1); // self-vote AtomicBoolean electionWon = new AtomicBoolean(false); - // Single-node cluster: self-vote alone satisfies quorum, become leader immediately - if (votesReceived.get() >= votesNeeded) { - if (electionWon.compareAndSet(false, true)) { - lock.lock(); - try { - if (currentTerm == myTerm && state == RaftState.CANDIDATE) { - becomeLeader(); - } - } finally { - lock.unlock(); - } - } - } for (PeerAddress peer : peers) { CompletableFuture.supplyAsync(() -> { @@ -647,7 +621,7 @@ private void replicateTo(PeerAddress peer) { sendInstallSnapshotToPeer(peer); } } - if (needsSnapshot) return; + long prevLogTerm = raftLog.getTermAt(prevLogIndex); List entries = raftLog.getEntriesFrom(peerNextIndex); @@ -856,7 +830,6 @@ private void applyCommitted() { nodeId, lastApplied, entry.getTopic(), batchRequest.getEntriesCount()); } default -> { - // MESSAGE (default): apply to MessageStore long msgOffset = messageStore.append( entry.getTopic(), entry.getPayload().toByteArray(), @@ -896,11 +869,10 @@ private void applyCommitted() { if (applied) { stateSaveNeeded = true; - - // Keep at least 100x the compact threshold as a retention buffer + long retentionLimit = lastApplied - (raftCompactThreshold * 100); long safeCompactIndex; - + if (isLeader()) { long minMatchIndex = lastApplied; for (long idx : matchIndex.values()) { @@ -910,15 +882,20 @@ private void applyCommitted() { } else { safeCompactIndex = retentionLimit; } - - // Always keep at least raftCompactThreshold entries after lastApplied + long finalCompactIndex = Math.min(safeCompactIndex, lastApplied - raftCompactThreshold); - - if (finalCompactIndex > 0) { + + long currentLogStart = raftLog.getStartIndex(); + boolean compactionDue = finalCompactIndex > 0 + && (finalCompactIndex - currentLogStart) >= raftCompactThreshold; + + if (compactionDue) { if (isCompacting.compareAndSet(false, true)) { raftExecutor.execute(() -> { try { raftLog.compact(finalCompactIndex); + logger.debug("[{}] Chunked compaction complete: log now starts at {}", + nodeId, finalCompactIndex + 1); } catch (IOException e) { logger.error("Failed to compact Raft log", e); } finally { diff --git a/drmq-client/src/main/java/com/drmq/client/DRMQConsumer.java b/drmq-client/src/main/java/com/drmq/client/DRMQConsumer.java index 05227f3..750751f 100644 --- a/drmq-client/src/main/java/com/drmq/client/DRMQConsumer.java +++ b/drmq-client/src/main/java/com/drmq/client/DRMQConsumer.java @@ -35,7 +35,7 @@ public class DRMQConsumer implements AutoCloseable { private static final int DEFAULT_MAX_MESSAGES = 100; private static final long DEFAULT_POLL_TIMEOUT_MS = 1000; private static final int MAX_RETRIES = 5; - private static final long RECONNECT_DELAY_MS = 500; // Brief pause between retries to allow leader election + private static final long RECONNECT_DELAY_MS = 500; private String host; private int port; diff --git a/drmq-ts-client/src/client.ts b/drmq-ts-client/src/client.ts index 72a1b2a..9ad4242 100644 --- a/drmq-ts-client/src/client.ts +++ b/drmq-ts-client/src/client.ts @@ -274,8 +274,8 @@ export class DRMQProducer extends DRMQClient { batch.forEach((pm, i) => pm.resolve({ success: true, offset: baseOffset + i })); return; } else { + const errorMsg = resp.errorMessage; if (resp.errorCode === ErrorCode.NOT_LEADER) { - const errorMsg = resp.errorMessage; const leaderAddr = errorMsg && errorMsg.startsWith("NOT_LEADER:") ? errorMsg.substring("NOT_LEADER:".length) : "UNKNOWN"; diff --git a/drmq-ts-client/src/example.ts b/drmq-ts-client/src/example.ts deleted file mode 100644 index 227fa97..0000000 --- a/drmq-ts-client/src/example.ts +++ /dev/null @@ -1,41 +0,0 @@ -import { DRMQProducer, DRMQConsumer } from './client'; - -async function main() { - const servers = "localhost:9092,localhost:9093"; - - console.log("--- Testing TS Producer ---"); - const producer = new DRMQProducer(servers); - try { - await producer.connect(); - const payload = Buffer.from("Hello from TypeScript!"); - const res = await producer.send("ts-topic", payload); - if (res.success) { - console.log(`Message sent successfully at offset ${res.offset}`); - } else { - console.log(`Failed to send: ${res.errorMessage}`); - } - } catch (err) { - console.error(`Producer error: ${(err as Error).message}`); - } finally { - producer.close(); - } - - console.log("\n--- Testing TS Consumer ---"); - const consumer = new DRMQConsumer(servers, "ts-workers"); - try { - await consumer.connect(); - await consumer.subscribe("ts-topic"); - - console.log("Polling for messages..."); - const messages = await consumer.poll(10, 5000); - for (const msg of messages) { - console.log(`Received (offset ${msg.offset}): ${Buffer.from(msg.payload).toString('utf-8')}`); - } - } catch (err) { - console.error(`Consumer error: ${(err as Error).message}`); - } finally { - consumer.close(); - } -} - -main().catch(console.error); diff --git a/drmq-ts-client/src/messages.ts b/drmq-ts-client/src/messages.ts index 47f9dee..090c3fa 100644 --- a/drmq-ts-client/src/messages.ts +++ b/drmq-ts-client/src/messages.ts @@ -205,10 +205,50 @@ export function messageTypeToJSON(object: MessageType): string { } } +/** --- Standard Error Codes --- */ +export enum ErrorCode { + NONE = 0, + NOT_LEADER = 1, + UNKNOWN_ERROR = 99, + UNRECOGNIZED = -1, +} + +export function errorCodeFromJSON(object: any): ErrorCode { + switch (object) { + case 0: + case "NONE": + return ErrorCode.NONE; + case 1: + case "NOT_LEADER": + return ErrorCode.NOT_LEADER; + case 99: + case "UNKNOWN_ERROR": + return ErrorCode.UNKNOWN_ERROR; + case -1: + case "UNRECOGNIZED": + default: + return ErrorCode.UNRECOGNIZED; + } +} + +export function errorCodeToJSON(object: ErrorCode): string { + switch (object) { + case ErrorCode.NONE: + return "NONE"; + case ErrorCode.NOT_LEADER: + return "NOT_LEADER"; + case ErrorCode.UNKNOWN_ERROR: + return "UNKNOWN_ERROR"; + case ErrorCode.UNRECOGNIZED: + default: + return "UNRECOGNIZED"; + } +} + /** Request to produce a message to a topic */ export interface ProduceRequest { topic: string; - payload: Buffer; + payload: Uint8Array; /** Optional message key for partitioning (future use) */ key?: | string @@ -224,6 +264,7 @@ export interface ProduceResponse { offset: number; /** Error details if failed */ errorMessage: string; + errorCode: ErrorCode; } /** Request to produce a batch of messages to a topic */ @@ -233,7 +274,7 @@ export interface ProduceBatchRequest { } export interface ProduceBatchRequest_BatchEntry { - payload: Buffer; + payload: Uint8Array; key?: string | undefined; clientTimestamp: number; } @@ -247,13 +288,14 @@ export interface ProduceBatchResponse { count: number; /** Error details if failed */ errorMessage: string; + errorCode: ErrorCode; } /** Internal message representation stored in the broker */ export interface StoredMessage { offset: number; topic: string; - payload: Buffer; + payload: Uint8Array; key?: | string | undefined; @@ -341,7 +383,7 @@ export interface NackResponse { export interface MessageEnvelope { type: MessageType; /** Serialized inner message */ - payload: Buffer; + payload: Uint8Array; } /** --- Raft log entry (wraps a user message for Raft replication) --- */ @@ -353,7 +395,7 @@ export interface RaftEntry { /** Target topic for the message */ topic: string; /** Message payload */ - payload: Buffer; + payload: Uint8Array; /** Optional message key */ key?: | string @@ -447,7 +489,7 @@ export interface InstallSnapshotRequest { /** Byte offset of this chunk in the file */ offset: number; /** Raw zip chunk payload */ - data: Buffer; + data: Uint8Array; /** True if this is the final chunk */ done: boolean; } @@ -457,7 +499,7 @@ export interface InstallSnapshotResponse { } function createBaseProduceRequest(): ProduceRequest { - return { topic: "", payload: Buffer.alloc(0), key: undefined, timestamp: 0 }; + return { topic: "", payload: new Uint8Array(0), key: undefined, timestamp: 0 }; } export const ProduceRequest: MessageFns = { @@ -497,7 +539,7 @@ export const ProduceRequest: MessageFns = { break; } - message.payload = Buffer.from(reader.bytes()); + message.payload = reader.bytes(); continue; } case 3: { @@ -528,7 +570,7 @@ export const ProduceRequest: MessageFns = { fromJSON(object: any): ProduceRequest { return { topic: isSet(object.topic) ? globalThis.String(object.topic) : "", - payload: isSet(object.payload) ? Buffer.from(bytesFromBase64(object.payload)) : Buffer.alloc(0), + payload: isSet(object.payload) ? bytesFromBase64(object.payload) : new Uint8Array(0), key: isSet(object.key) ? globalThis.String(object.key) : undefined, timestamp: isSet(object.timestamp) ? globalThis.Number(object.timestamp) : 0, }; @@ -557,7 +599,7 @@ export const ProduceRequest: MessageFns = { fromPartial, I>>(object: I): ProduceRequest { const message = createBaseProduceRequest(); message.topic = object.topic ?? ""; - message.payload = object.payload ?? Buffer.alloc(0); + message.payload = object.payload ?? new Uint8Array(0); message.key = object.key ?? undefined; message.timestamp = object.timestamp ?? 0; return message; @@ -565,7 +607,7 @@ export const ProduceRequest: MessageFns = { }; function createBaseProduceResponse(): ProduceResponse { - return { success: false, offset: 0, errorMessage: "" }; + return { success: false, offset: 0, errorMessage: "", errorCode: 0 }; } export const ProduceResponse: MessageFns = { @@ -579,6 +621,9 @@ export const ProduceResponse: MessageFns = { if (message.errorMessage !== "") { writer.uint32(26).string(message.errorMessage); } + if (message.errorCode !== 0) { + writer.uint32(32).int32(message.errorCode); + } return writer; }, @@ -613,6 +658,14 @@ export const ProduceResponse: MessageFns = { message.errorMessage = reader.string(); continue; } + case 4: { + if (tag !== 32) { + break; + } + + message.errorCode = reader.int32() as any; + continue; + } } if ((tag & 7) === 4 || tag === 0) { break; @@ -631,6 +684,11 @@ export const ProduceResponse: MessageFns = { : isSet(object.error_message) ? globalThis.String(object.error_message) : "", + errorCode: isSet(object.errorCode) + ? errorCodeFromJSON(object.errorCode) + : isSet(object.error_code) + ? errorCodeFromJSON(object.error_code) + : 0, }; }, @@ -645,6 +703,9 @@ export const ProduceResponse: MessageFns = { if (message.errorMessage !== "") { obj.errorMessage = message.errorMessage; } + if (message.errorCode !== 0) { + obj.errorCode = errorCodeToJSON(message.errorCode); + } return obj; }, @@ -656,6 +717,7 @@ export const ProduceResponse: MessageFns = { message.success = object.success ?? false; message.offset = object.offset ?? 0; message.errorMessage = object.errorMessage ?? ""; + message.errorCode = object.errorCode ?? 0; return message; }, }; @@ -739,7 +801,7 @@ export const ProduceBatchRequest: MessageFns = { }; function createBaseProduceBatchRequest_BatchEntry(): ProduceBatchRequest_BatchEntry { - return { payload: Buffer.alloc(0), key: undefined, clientTimestamp: 0 }; + return { payload: new Uint8Array(0), key: undefined, clientTimestamp: 0 }; } export const ProduceBatchRequest_BatchEntry: MessageFns = { @@ -768,7 +830,7 @@ export const ProduceBatchRequest_BatchEntry: MessageFns = { @@ -854,6 +916,9 @@ export const ProduceBatchResponse: MessageFns = { if (message.errorMessage !== "") { writer.uint32(34).string(message.errorMessage); } + if (message.errorCode !== 0) { + writer.uint32(40).int32(message.errorCode); + } return writer; }, @@ -896,6 +961,14 @@ export const ProduceBatchResponse: MessageFns = { message.errorMessage = reader.string(); continue; } + case 5: { + if (tag !== 40) { + break; + } + + message.errorCode = reader.int32() as any; + continue; + } } if ((tag & 7) === 4 || tag === 0) { break; @@ -919,6 +992,11 @@ export const ProduceBatchResponse: MessageFns = { : isSet(object.error_message) ? globalThis.String(object.error_message) : "", + errorCode: isSet(object.errorCode) + ? errorCodeFromJSON(object.errorCode) + : isSet(object.error_code) + ? errorCodeFromJSON(object.error_code) + : 0, }; }, @@ -936,6 +1014,9 @@ export const ProduceBatchResponse: MessageFns = { if (message.errorMessage !== "") { obj.errorMessage = message.errorMessage; } + if (message.errorCode !== 0) { + obj.errorCode = errorCodeToJSON(message.errorCode); + } return obj; }, @@ -948,12 +1029,13 @@ export const ProduceBatchResponse: MessageFns = { message.baseOffset = object.baseOffset ?? 0; message.count = object.count ?? 0; message.errorMessage = object.errorMessage ?? ""; + message.errorCode = object.errorCode ?? 0; return message; }, }; function createBaseStoredMessage(): StoredMessage { - return { offset: 0, topic: "", payload: Buffer.alloc(0), key: undefined, timestamp: 0, storedAt: 0 }; + return { offset: 0, topic: "", payload: new Uint8Array(0), key: undefined, timestamp: 0, storedAt: 0 }; } export const StoredMessage: MessageFns = { @@ -1007,7 +1089,7 @@ export const StoredMessage: MessageFns = { break; } - message.payload = Buffer.from(reader.bytes()); + message.payload = reader.bytes(); continue; } case 4: { @@ -1047,7 +1129,7 @@ export const StoredMessage: MessageFns = { return { offset: isSet(object.offset) ? globalThis.Number(object.offset) : 0, topic: isSet(object.topic) ? globalThis.String(object.topic) : "", - payload: isSet(object.payload) ? Buffer.from(bytesFromBase64(object.payload)) : Buffer.alloc(0), + payload: isSet(object.payload) ? bytesFromBase64(object.payload) : new Uint8Array(0), key: isSet(object.key) ? globalThis.String(object.key) : undefined, timestamp: isSet(object.timestamp) ? globalThis.Number(object.timestamp) : 0, storedAt: isSet(object.storedAt) @@ -1088,7 +1170,7 @@ export const StoredMessage: MessageFns = { const message = createBaseStoredMessage(); message.offset = object.offset ?? 0; message.topic = object.topic ?? ""; - message.payload = object.payload ?? Buffer.alloc(0); + message.payload = object.payload ?? new Uint8Array(0); message.key = object.key ?? undefined; message.timestamp = object.timestamp ?? 0; message.storedAt = object.storedAt ?? 0; @@ -1943,7 +2025,7 @@ export const NackResponse: MessageFns = { }; function createBaseMessageEnvelope(): MessageEnvelope { - return { type: 0, payload: Buffer.alloc(0) }; + return { type: 0, payload: new Uint8Array(0) }; } export const MessageEnvelope: MessageFns = { @@ -1977,7 +2059,7 @@ export const MessageEnvelope: MessageFns = { break; } - message.payload = Buffer.from(reader.bytes()); + message.payload = reader.bytes(); continue; } } @@ -1992,7 +2074,7 @@ export const MessageEnvelope: MessageFns = { fromJSON(object: any): MessageEnvelope { return { type: isSet(object.type) ? messageTypeFromJSON(object.type) : 0, - payload: isSet(object.payload) ? Buffer.from(bytesFromBase64(object.payload)) : Buffer.alloc(0), + payload: isSet(object.payload) ? bytesFromBase64(object.payload) : new Uint8Array(0), }; }, @@ -2013,7 +2095,7 @@ export const MessageEnvelope: MessageFns = { fromPartial, I>>(object: I): MessageEnvelope { const message = createBaseMessageEnvelope(); message.type = object.type ?? 0; - message.payload = object.payload ?? Buffer.alloc(0); + message.payload = object.payload ?? new Uint8Array(0); return message; }, }; @@ -2023,7 +2105,7 @@ function createBaseRaftEntry(): RaftEntry { term: 0, index: 0, topic: "", - payload: Buffer.alloc(0), + payload: new Uint8Array(0), key: undefined, timestamp: 0, commandType: 0, @@ -2100,7 +2182,7 @@ export const RaftEntry: MessageFns = { break; } - message.payload = Buffer.from(reader.bytes()); + message.payload = reader.bytes(); continue; } case 5: { @@ -2157,7 +2239,7 @@ export const RaftEntry: MessageFns = { term: isSet(object.term) ? globalThis.Number(object.term) : 0, index: isSet(object.index) ? globalThis.Number(object.index) : 0, topic: isSet(object.topic) ? globalThis.String(object.topic) : "", - payload: isSet(object.payload) ? Buffer.from(bytesFromBase64(object.payload)) : Buffer.alloc(0), + payload: isSet(object.payload) ? bytesFromBase64(object.payload) : new Uint8Array(0), key: isSet(object.key) ? globalThis.String(object.key) : undefined, timestamp: isSet(object.timestamp) ? globalThis.Number(object.timestamp) : 0, commandType: isSet(object.commandType) @@ -2218,7 +2300,7 @@ export const RaftEntry: MessageFns = { message.term = object.term ?? 0; message.index = object.index ?? 0; message.topic = object.topic ?? ""; - message.payload = object.payload ?? Buffer.alloc(0); + message.payload = object.payload ?? new Uint8Array(0); message.key = object.key ?? undefined; message.timestamp = object.timestamp ?? 0; message.commandType = object.commandType ?? 0; @@ -2889,7 +2971,7 @@ function createBaseInstallSnapshotRequest(): InstallSnapshotRequest { lastIncludedIndex: 0, lastIncludedTerm: 0, offset: 0, - data: Buffer.alloc(0), + data: new Uint8Array(0), done: false, }; } @@ -2972,7 +3054,7 @@ export const InstallSnapshotRequest: MessageFns = { break; } - message.data = Buffer.from(reader.bytes()); + message.data = reader.bytes(); continue; } case 7: { @@ -3011,7 +3093,7 @@ export const InstallSnapshotRequest: MessageFns = { ? globalThis.Number(object.last_included_term) : 0, offset: isSet(object.offset) ? globalThis.Number(object.offset) : 0, - data: isSet(object.data) ? Buffer.from(bytesFromBase64(object.data)) : Buffer.alloc(0), + data: isSet(object.data) ? bytesFromBase64(object.data) : new Uint8Array(0), done: isSet(object.done) ? globalThis.Boolean(object.done) : false, }; }, @@ -3052,7 +3134,7 @@ export const InstallSnapshotRequest: MessageFns = { message.lastIncludedIndex = object.lastIncludedIndex ?? 0; message.lastIncludedTerm = object.lastIncludedTerm ?? 0; message.offset = object.offset ?? 0; - message.data = object.data ?? Buffer.alloc(0); + message.data = object.data ?? new Uint8Array(0); message.done = object.done ?? false; return message; }, @@ -3117,11 +3199,28 @@ export const InstallSnapshotResponse: MessageFns = { }; function bytesFromBase64(b64: string): Uint8Array { - return Uint8Array.from(globalThis.Buffer.from(b64, "base64")); + if ((globalThis as any).Buffer) { + return Uint8Array.from((globalThis as any).Buffer.from(b64, "base64")); + } else { + const bin = globalThis.atob(b64); + const arr = new Uint8Array(bin.length); + for (let i = 0; i < bin.length; ++i) { + arr[i] = bin.charCodeAt(i); + } + return arr; + } } function base64FromBytes(arr: Uint8Array): string { - return globalThis.Buffer.from(arr).toString("base64"); + if ((globalThis as any).Buffer) { + return (globalThis as any).Buffer.from(arr).toString("base64"); + } else { + const bin: string[] = []; + arr.forEach((byte) => { + bin.push(globalThis.String.fromCharCode(byte)); + }); + return globalThis.btoa(bin.join("")); + } } type Builtin = Date | Function | Uint8Array | string | number | boolean | undefined;