diff --git a/command/src/main/java/io/openmessaging/storage/dledger/command/AppendCommand.java b/command/src/main/java/io/openmessaging/storage/dledger/command/AppendCommand.java index 7296f579..a01381d8 100644 --- a/command/src/main/java/io/openmessaging/storage/dledger/command/AppendCommand.java +++ b/command/src/main/java/io/openmessaging/storage/dledger/command/AppendCommand.java @@ -16,10 +16,10 @@ package io.openmessaging.storage.dledger.command; -import com.alibaba.fastjson.JSON; import com.beust.jcommander.Parameter; import io.openmessaging.storage.dledger.client.DLedgerClient; import io.openmessaging.storage.dledger.protocol.AppendEntryResponse; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -45,7 +45,7 @@ public void doCommand() { dLedgerClient.startup(); for (int i = 0; i < count; i++) { AppendEntryResponse response = dLedgerClient.append(data.getBytes()); - logger.info("Append Result:{}", JSON.toJSONString(response)); + logger.info("Append Result:{}", DLedgerJsonUtils.toJsonString(response)); } dLedgerClient.shutdown(); } diff --git a/command/src/main/java/io/openmessaging/storage/dledger/command/DLedger.java b/command/src/main/java/io/openmessaging/storage/dledger/command/DLedger.java index 9c70ef4b..1a829553 100644 --- a/command/src/main/java/io/openmessaging/storage/dledger/command/DLedger.java +++ b/command/src/main/java/io/openmessaging/storage/dledger/command/DLedger.java @@ -16,12 +16,12 @@ package io.openmessaging.storage.dledger.command; -import com.alibaba.fastjson.JSON; import com.beust.jcommander.JCommander; import io.openmessaging.storage.dledger.DLedgerConfig; import io.openmessaging.storage.dledger.proxy.DLedgerProxy; import io.openmessaging.storage.dledger.proxy.DLedgerProxyConfig; import io.openmessaging.storage.dledger.proxy.util.ConfigUtils; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import java.util.Collections; import java.util.LinkedList; import java.util.List; @@ -58,7 +58,7 @@ public static void bootstrapDLedger(List dLedgerConfigs) { } DLedgerProxy dLedgerProxy = new DLedgerProxy(dLedgerConfigs); dLedgerProxy.startup(); - logger.info("DLedgers start ok with config {}", JSON.toJSONString(dLedgerConfigs)); + logger.info("DLedgers start ok with config {}", DLedgerJsonUtils.toJsonString(dLedgerConfigs)); Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() { private volatile boolean hasShutdown = false; diff --git a/command/src/main/java/io/openmessaging/storage/dledger/command/GetCommand.java b/command/src/main/java/io/openmessaging/storage/dledger/command/GetCommand.java index 46c625ce..08ff5117 100644 --- a/command/src/main/java/io/openmessaging/storage/dledger/command/GetCommand.java +++ b/command/src/main/java/io/openmessaging/storage/dledger/command/GetCommand.java @@ -16,12 +16,12 @@ package io.openmessaging.storage.dledger.command; -import com.alibaba.fastjson.JSON; import com.beust.jcommander.Parameter; import com.beust.jcommander.Parameters; import io.openmessaging.storage.dledger.client.DLedgerClient; import io.openmessaging.storage.dledger.entry.DLedgerEntry; import io.openmessaging.storage.dledger.protocol.GetEntriesResponse; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -44,7 +44,7 @@ public void doCommand() { DLedgerClient dLedgerClient = new DLedgerClient(group, peers); dLedgerClient.startup(); GetEntriesResponse response = dLedgerClient.get(index); - logger.info("Get Result:{}", JSON.toJSONString(response)); + logger.info("Get Result:{}", DLedgerJsonUtils.toJsonString(response)); if (response.getEntries() != null && response.getEntries().size() > 0) { for (DLedgerEntry entry : response.getEntries()) { logger.info("Get Result index:{} {}", entry.getIndex(), new String(entry.getBody())); diff --git a/command/src/main/java/io/openmessaging/storage/dledger/command/LeadershipTransferCommand.java b/command/src/main/java/io/openmessaging/storage/dledger/command/LeadershipTransferCommand.java index 4f7861ac..559bff2b 100644 --- a/command/src/main/java/io/openmessaging/storage/dledger/command/LeadershipTransferCommand.java +++ b/command/src/main/java/io/openmessaging/storage/dledger/command/LeadershipTransferCommand.java @@ -16,12 +16,12 @@ package io.openmessaging.storage.dledger.command; -import com.alibaba.fastjson.JSON; import com.beust.jcommander.Parameter; import com.beust.jcommander.Parameters; import io.openmessaging.storage.dledger.client.DLedgerClient; import io.openmessaging.storage.dledger.protocol.DLedgerResponseCode; import io.openmessaging.storage.dledger.protocol.LeadershipTransferResponse; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -51,7 +51,7 @@ public void doCommand() { dLedgerClient.startup(); LeadershipTransferResponse response = dLedgerClient.leadershipTransfer(leaderId, transfereeId, term); LOGGER.info("LeadershipTransfer code={}, Result:{}", DLedgerResponseCode.valueOf(response.getCode()), - JSON.toJSONString(response)); + DLedgerJsonUtils.toJsonString(response)); dLedgerClient.shutdown(); } } diff --git a/command/src/main/java/io/openmessaging/storage/dledger/command/ReadFileCommand.java b/command/src/main/java/io/openmessaging/storage/dledger/command/ReadFileCommand.java index adfd186b..2a1a9bb6 100644 --- a/command/src/main/java/io/openmessaging/storage/dledger/command/ReadFileCommand.java +++ b/command/src/main/java/io/openmessaging/storage/dledger/command/ReadFileCommand.java @@ -16,7 +16,6 @@ package io.openmessaging.storage.dledger.command; -import com.alibaba.fastjson.JSON; import com.beust.jcommander.Parameter; import com.beust.jcommander.Parameters; import io.openmessaging.storage.dledger.entry.DLedgerEntry; @@ -25,6 +24,7 @@ import io.openmessaging.storage.dledger.store.file.MmapFile; import io.openmessaging.storage.dledger.store.file.MmapFileList; import io.openmessaging.storage.dledger.store.file.SelectMmapBufferResult; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import java.nio.ByteBuffer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -74,7 +74,7 @@ public void doCommand() { logger.info("magic={} pos={} size={} index={} term={}", buffer.getInt(), buffer.getLong(), buffer.getInt(), buffer.getLong(), buffer.getLong()); } else { DLedgerEntry entry = DLedgerEntryCoder.decode(buffer, readBody); - logger.info(JSON.toJSONString(entry)); + logger.info(DLedgerJsonUtils.toJsonString(entry)); } } } diff --git a/dledger/pom.xml b/dledger/pom.xml index 2c378362..5ad00bd3 100644 --- a/dledger/pom.xml +++ b/dledger/pom.xml @@ -30,6 +30,16 @@ org.apache.rocketmq rocketmq-remoting + + + com.alibaba + fastjson + + + + + com.alibaba.fastjson2 + fastjson2 org.slf4j @@ -56,4 +66,4 @@ - \ No newline at end of file + diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerEntryPusher.java b/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerEntryPusher.java index c2143a66..c0b23b6a 100644 --- a/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerEntryPusher.java +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerEntryPusher.java @@ -16,7 +16,6 @@ package io.openmessaging.storage.dledger; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.entry.DLedgerEntry; import io.openmessaging.storage.dledger.exception.DLedgerException; import io.openmessaging.storage.dledger.protocol.AppendEntryResponse; @@ -27,6 +26,7 @@ import io.openmessaging.storage.dledger.store.DLedgerMemoryStore; import io.openmessaging.storage.dledger.store.DLedgerStore; import io.openmessaging.storage.dledger.store.file.DLedgerMmapFileStore; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.DLedgerUtils; import io.openmessaging.storage.dledger.utils.Pair; import io.openmessaging.storage.dledger.utils.PreConditions; @@ -270,10 +270,10 @@ public void doWork() { if (DLedgerEntryPusher.this.fsmCaller.isPresent()) { final long lastAppliedIndex = DLedgerEntryPusher.this.fsmCaller.get().getLastAppliedIndex(); logger.info("[{}][{}] term={} ledgerBegin={} ledgerEnd={} committed={} watermarks={} appliedIndex={}", - memberState.getSelfId(), memberState.getRole(), memberState.currTerm(), dLedgerStore.getLedgerBeginIndex(), dLedgerStore.getLedgerEndIndex(), dLedgerStore.getCommittedIndex(), JSON.toJSONString(peerWaterMarksByTerm), lastAppliedIndex); + memberState.getSelfId(), memberState.getRole(), memberState.currTerm(), dLedgerStore.getLedgerBeginIndex(), dLedgerStore.getLedgerEndIndex(), dLedgerStore.getCommittedIndex(), DLedgerJsonUtils.toJsonString(peerWaterMarksByTerm), lastAppliedIndex); } else { logger.info("[{}][{}] term={} ledgerBegin={} ledgerEnd={} committed={} watermarks={}", - memberState.getSelfId(), memberState.getRole(), memberState.currTerm(), dLedgerStore.getLedgerBeginIndex(), dLedgerStore.getLedgerEndIndex(), dLedgerStore.getCommittedIndex(), JSON.toJSONString(peerWaterMarksByTerm)); + memberState.getSelfId(), memberState.getRole(), memberState.currTerm(), dLedgerStore.getLedgerBeginIndex(), dLedgerStore.getLedgerEndIndex(), dLedgerStore.getCommittedIndex(), DLedgerJsonUtils.toJsonString(peerWaterMarksByTerm)); } lastPrintWatermarkTimeMs = System.currentTimeMillis(); } diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerLeaderElector.java b/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerLeaderElector.java index 6c90ed28..4e990506 100644 --- a/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerLeaderElector.java +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerLeaderElector.java @@ -16,7 +16,6 @@ package io.openmessaging.storage.dledger; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.protocol.DLedgerResponseCode; import io.openmessaging.storage.dledger.protocol.HeartBeatRequest; import io.openmessaging.storage.dledger.protocol.HeartBeatResponse; @@ -24,6 +23,7 @@ import io.openmessaging.storage.dledger.protocol.LeadershipTransferResponse; import io.openmessaging.storage.dledger.protocol.VoteRequest; import io.openmessaging.storage.dledger.protocol.VoteResponse; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.DLedgerUtils; import java.sql.Timestamp; @@ -448,7 +448,7 @@ private void maintainAsCandidate() throws Exception { if (ex != null) { throw ex; } - LOGGER.info("[{}][GetVoteResponse] {}", memberState.getSelfId(), JSON.toJSONString(x)); + LOGGER.info("[{}][GetVoteResponse] {}", memberState.getSelfId(), DLedgerJsonUtils.toJsonString(x)); if (x.getVoteResult() != VoteResponse.RESULT.UNKNOWN) { validNum.incrementAndGet(); } diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerRpcNettyService.java b/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerRpcNettyService.java index d3be885c..74ebec81 100644 --- a/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerRpcNettyService.java +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/DLedgerRpcNettyService.java @@ -16,7 +16,6 @@ package io.openmessaging.storage.dledger; -import com.alibaba.fastjson.JSON; import io.netty.channel.ChannelHandlerContext; import io.openmessaging.storage.dledger.protocol.AppendEntryRequest; import io.openmessaging.storage.dledger.protocol.AppendEntryResponse; @@ -37,6 +36,7 @@ import io.openmessaging.storage.dledger.protocol.RequestOrResponse; import io.openmessaging.storage.dledger.protocol.VoteRequest; import io.openmessaging.storage.dledger.protocol.VoteResponse; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.DLedgerUtils; import java.util.concurrent.CompletableFuture; @@ -44,11 +44,13 @@ import java.util.concurrent.Executors; import org.apache.rocketmq.remoting.ChannelEventListener; +import org.apache.rocketmq.remoting.InvokeCallback; import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.apache.rocketmq.remoting.netty.NettyRemotingClient; import org.apache.rocketmq.remoting.netty.NettyRemotingServer; import org.apache.rocketmq.remoting.netty.NettyRequestProcessor; import org.apache.rocketmq.remoting.netty.NettyServerConfig; +import org.apache.rocketmq.remoting.netty.ResponseFuture; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -138,14 +140,23 @@ public CompletableFuture heartBeat(HeartBeatRequest request) heartBeatInvokeExecutor.execute(() -> { try { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.HEART_BEAT.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); - remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, responseFuture -> { - RemotingCommand responseCommand = responseFuture.getResponseCommand(); - if (responseCommand != null) { - HeartBeatResponse response = JSON.parseObject(responseCommand.getBody(), HeartBeatResponse.class); - future.complete(response); - } else { - LOGGER.error("HeartBeat request time out, {}", request.baseInfo()); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); + remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, new InvokeCallback() { + @Override + public void operationComplete(ResponseFuture responseFuture) { + RemotingCommand responseCommand = responseFuture.getResponseCommand(); + if (responseCommand != null) { + HeartBeatResponse response = DLedgerJsonUtils.parseObject(responseCommand.getBody(), HeartBeatResponse.class); + future.complete(response); + } else { + LOGGER.error("HeartBeat request time out, {}", request.baseInfo()); + future.complete(new HeartBeatResponse().code(DLedgerResponseCode.NETWORK_ERROR.getCode())); + } + } + + @Override + public void operationFail(Throwable throwable) { + LOGGER.error("HeartBeat request failed due to network error, {}", request.baseInfo(), throwable); future.complete(new HeartBeatResponse().code(DLedgerResponseCode.NETWORK_ERROR.getCode())); } }); @@ -163,14 +174,23 @@ public CompletableFuture vote(VoteRequest request) { voteInvokeExecutor.execute(() -> { try { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.VOTE.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); - this.remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, responseFuture -> { - RemotingCommand responseCommand = responseFuture.getResponseCommand(); - if (responseCommand != null) { - VoteResponse response = JSON.parseObject(responseCommand.getBody(), VoteResponse.class); - future.complete(response); - } else { - LOGGER.error("Vote request time out, {}", request.baseInfo()); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); + this.remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, new InvokeCallback() { + @Override + public void operationComplete(ResponseFuture responseFuture) { + RemotingCommand responseCommand = responseFuture.getResponseCommand(); + if (responseCommand != null) { + VoteResponse response = DLedgerJsonUtils.parseObject(responseCommand.getBody(), VoteResponse.class); + future.complete(response); + } else { + LOGGER.error("Vote request time out, {}", request.baseInfo()); + future.complete(new VoteResponse()); + } + } + + @Override + public void operationFail(Throwable throwable) { + LOGGER.error("Vote request failed due to network error, {}", request.baseInfo(), throwable); future.complete(new VoteResponse()); } }); @@ -194,19 +214,29 @@ public CompletableFuture append(AppendEntryRequest request) CompletableFuture future = new CompletableFuture<>(); try { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.APPEND.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); - remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, responseFuture -> { - RemotingCommand responseCommand = responseFuture.getResponseCommand(); - - AppendEntryResponse response; - if (responseCommand != null) { - response = JSON.parseObject(responseCommand.getBody(), AppendEntryResponse.class); - } else { - response = new AppendEntryResponse(); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); + remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, new InvokeCallback() { + @Override + public void operationComplete(ResponseFuture responseFuture) { + RemotingCommand responseCommand = responseFuture.getResponseCommand(); + AppendEntryResponse response; + if (responseCommand != null) { + response = DLedgerJsonUtils.parseObject(responseCommand.getBody(), AppendEntryResponse.class); + } else { + response = new AppendEntryResponse(); + response.copyBaseInfo(request); + response.setCode(DLedgerResponseCode.NETWORK_ERROR.getCode()); + } + future.complete(response); + } + + @Override + public void operationFail(Throwable throwable) { + AppendEntryResponse response = new AppendEntryResponse(); response.copyBaseInfo(request); response.setCode(DLedgerResponseCode.NETWORK_ERROR.getCode()); + future.complete(response); } - future.complete(response); }); } catch (Throwable t) { LOGGER.error("Send append request failed, {}", request.baseInfo(), t); @@ -228,9 +258,9 @@ public CompletableFuture metadata(MetadataRequest request) { @Override public CompletableFuture pull(PullEntriesRequest request) throws Exception { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.PULL.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); RemotingCommand wrapperResponse = remotingClient.invokeSync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000); - PullEntriesResponse response = JSON.parseObject(wrapperResponse.getBody(), PullEntriesResponse.class); + PullEntriesResponse response = DLedgerJsonUtils.parseObject(wrapperResponse.getBody(), PullEntriesResponse.class); return CompletableFuture.completedFuture(response); } @@ -239,19 +269,29 @@ public CompletableFuture push(PushEntryRequest request) { CompletableFuture future = new CompletableFuture<>(); try { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.PUSH.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); - remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, responseFuture -> { - RemotingCommand responseCommand = responseFuture.getResponseCommand(); - - PushEntryResponse response; - if (responseCommand != null) { - response = JSON.parseObject(responseCommand.getBody(), PushEntryResponse.class); - } else { - response = new PushEntryResponse(); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); + remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, new InvokeCallback() { + @Override + public void operationComplete(ResponseFuture responseFuture) { + RemotingCommand responseCommand = responseFuture.getResponseCommand(); + PushEntryResponse response; + if (responseCommand != null) { + response = DLedgerJsonUtils.parseObject(responseCommand.getBody(), PushEntryResponse.class); + } else { + response = new PushEntryResponse(); + response.copyBaseInfo(request); + response.setCode(DLedgerResponseCode.NETWORK_ERROR.getCode()); + } + future.complete(response); + } + + @Override + public void operationFail(Throwable throwable) { + PushEntryResponse response = new PushEntryResponse(); response.copyBaseInfo(request); response.setCode(DLedgerResponseCode.NETWORK_ERROR.getCode()); + future.complete(response); } - future.complete(response); }); } catch (Throwable t) { LOGGER.error("Send push request failed, {}", request.baseInfo(), t); @@ -270,19 +310,29 @@ public CompletableFuture leadershipTransfer( CompletableFuture future = new CompletableFuture<>(); try { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.LEADERSHIP_TRANSFER.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); - remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, responseFuture -> { - RemotingCommand responseCommand = responseFuture.getResponseCommand(); - - LeadershipTransferResponse response; - if (responseCommand != null) { - response = JSON.parseObject(responseFuture.getResponseCommand().getBody(), LeadershipTransferResponse.class); - } else { - response = new LeadershipTransferResponse(); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); + remotingClient.invokeAsync(getPeerAddr(request.getGroup(), request.getRemoteId()), wrapperRequest, 3000, new InvokeCallback() { + @Override + public void operationComplete(ResponseFuture responseFuture) { + RemotingCommand responseCommand = responseFuture.getResponseCommand(); + LeadershipTransferResponse response; + if (responseCommand != null) { + response = DLedgerJsonUtils.parseObject(responseFuture.getResponseCommand().getBody(), LeadershipTransferResponse.class); + } else { + response = new LeadershipTransferResponse(); + response.copyBaseInfo(request); + response.setCode(DLedgerResponseCode.NETWORK_ERROR.getCode()); + } + future.complete(response); + } + + @Override + public void operationFail(Throwable throwable) { + LeadershipTransferResponse response = new LeadershipTransferResponse(); response.copyBaseInfo(request); response.setCode(DLedgerResponseCode.NETWORK_ERROR.getCode()); + future.complete(response); } - future.complete(response); }); } catch (Throwable t) { LOGGER.error("Send leadershipTransfer request failed, {}", request.baseInfo(), t); @@ -330,50 +380,50 @@ public RemotingCommand processRequest(ChannelHandlerContext ctx, RemotingCommand DLedgerRequestCode requestCode = DLedgerRequestCode.valueOf(request.getCode()); switch (requestCode) { case METADATA: { - MetadataRequest metadataRequest = JSON.parseObject(request.getBody(), MetadataRequest.class); + MetadataRequest metadataRequest = DLedgerJsonUtils.parseObject(request.getBody(), MetadataRequest.class); CompletableFuture future = handleMetadata(metadataRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case APPEND: { - AppendEntryRequest appendEntryRequest = JSON.parseObject(request.getBody(), AppendEntryRequest.class); + AppendEntryRequest appendEntryRequest = DLedgerJsonUtils.parseObject(request.getBody(), AppendEntryRequest.class); CompletableFuture future = handleAppend(appendEntryRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case GET: { - GetEntriesRequest getEntriesRequest = JSON.parseObject(request.getBody(), GetEntriesRequest.class); + GetEntriesRequest getEntriesRequest = DLedgerJsonUtils.parseObject(request.getBody(), GetEntriesRequest.class); CompletableFuture future = handleGet(getEntriesRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case PULL: { - PullEntriesRequest pullEntriesRequest = JSON.parseObject(request.getBody(), PullEntriesRequest.class); + PullEntriesRequest pullEntriesRequest = DLedgerJsonUtils.parseObject(request.getBody(), PullEntriesRequest.class); CompletableFuture future = handlePull(pullEntriesRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case PUSH: { - PushEntryRequest pushEntryRequest = JSON.parseObject(request.getBody(), PushEntryRequest.class); + PushEntryRequest pushEntryRequest = DLedgerJsonUtils.parseObject(request.getBody(), PushEntryRequest.class); CompletableFuture future = handlePush(pushEntryRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case VOTE: { - VoteRequest voteRequest = JSON.parseObject(request.getBody(), VoteRequest.class); + VoteRequest voteRequest = DLedgerJsonUtils.parseObject(request.getBody(), VoteRequest.class); CompletableFuture future = handleVote(voteRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case HEART_BEAT: { - HeartBeatRequest heartBeatRequest = JSON.parseObject(request.getBody(), HeartBeatRequest.class); + HeartBeatRequest heartBeatRequest = DLedgerJsonUtils.parseObject(request.getBody(), HeartBeatRequest.class); CompletableFuture future = handleHeartBeat(heartBeatRequest); future.whenCompleteAsync((x, y) -> writeResponse(x, y, request, ctx), futureExecutor); break; } case LEADERSHIP_TRANSFER: { long start = System.currentTimeMillis(); - LeadershipTransferRequest leadershipTransferRequest = JSON.parseObject(request.getBody(), LeadershipTransferRequest.class); + LeadershipTransferRequest leadershipTransferRequest = DLedgerJsonUtils.parseObject(request.getBody(), LeadershipTransferRequest.class); CompletableFuture future = handleLeadershipTransfer(leadershipTransferRequest); future.whenCompleteAsync((x, y) -> { writeResponse(x, y, request, ctx); @@ -432,7 +482,7 @@ public CompletableFuture handlePush(PushEntryRequest request) public RemotingCommand handleResponse(RequestOrResponse response, RemotingCommand request) { RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(DLedgerResponseCode.SUCCESS.getCode(), null); - remotingCommand.setBody(JSON.toJSONBytes(response)); + remotingCommand.setBody(DLedgerJsonUtils.toJsonBytes(response)); remotingCommand.setOpaque(request.getOpaque()); return remotingCommand; } diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/client/DLedgerClientRpcNettyService.java b/dledger/src/main/java/io/openmessaging/storage/dledger/client/DLedgerClientRpcNettyService.java index 5038b464..ab0cb166 100644 --- a/dledger/src/main/java/io/openmessaging/storage/dledger/client/DLedgerClientRpcNettyService.java +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/client/DLedgerClientRpcNettyService.java @@ -16,7 +16,6 @@ package io.openmessaging.storage.dledger.client; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.protocol.AppendEntryRequest; import io.openmessaging.storage.dledger.protocol.AppendEntryResponse; import io.openmessaging.storage.dledger.protocol.DLedgerRequestCode; @@ -26,6 +25,7 @@ import io.openmessaging.storage.dledger.protocol.MetadataResponse; import io.openmessaging.storage.dledger.protocol.LeadershipTransferResponse; import io.openmessaging.storage.dledger.protocol.LeadershipTransferRequest; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.apache.rocketmq.remoting.netty.NettyRemotingClient; @@ -42,36 +42,36 @@ public DLedgerClientRpcNettyService() { @Override public CompletableFuture append(AppendEntryRequest request) throws Exception { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.APPEND.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); RemotingCommand wrapperResponse = this.remotingClient.invokeSync(getPeerAddr(request.getRemoteId()), wrapperRequest, 3000); - AppendEntryResponse response = JSON.parseObject(wrapperResponse.getBody(), AppendEntryResponse.class); + AppendEntryResponse response = DLedgerJsonUtils.parseObject(wrapperResponse.getBody(), AppendEntryResponse.class); return CompletableFuture.completedFuture(response); } @Override public CompletableFuture metadata(MetadataRequest request) throws Exception { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.METADATA.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); RemotingCommand wrapperResponse = this.remotingClient.invokeSync(getPeerAddr(request.getRemoteId()), wrapperRequest, 3000); - MetadataResponse response = JSON.parseObject(wrapperResponse.getBody(), MetadataResponse.class); + MetadataResponse response = DLedgerJsonUtils.parseObject(wrapperResponse.getBody(), MetadataResponse.class); return CompletableFuture.completedFuture(response); } @Override public CompletableFuture leadershipTransfer(LeadershipTransferRequest request) throws Exception { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.LEADERSHIP_TRANSFER.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); RemotingCommand wrapperResponse = this.remotingClient.invokeSync(getPeerAddr(request.getRemoteId()), wrapperRequest, 10000); - LeadershipTransferResponse response = JSON.parseObject(wrapperResponse.getBody(), LeadershipTransferResponse.class); + LeadershipTransferResponse response = DLedgerJsonUtils.parseObject(wrapperResponse.getBody(), LeadershipTransferResponse.class); return CompletableFuture.completedFuture(response); } @Override public CompletableFuture get(GetEntriesRequest request) throws Exception { RemotingCommand wrapperRequest = RemotingCommand.createRequestCommand(DLedgerRequestCode.GET.getCode(), null); - wrapperRequest.setBody(JSON.toJSONBytes(request)); + wrapperRequest.setBody(DLedgerJsonUtils.toJsonBytes(request)); RemotingCommand wrapperResponse = this.remotingClient.invokeSync(getPeerAddr(request.getRemoteId()), wrapperRequest, 3000); - GetEntriesResponse response = JSON.parseObject(wrapperResponse.getBody(), GetEntriesResponse.class); + GetEntriesResponse response = DLedgerJsonUtils.parseObject(wrapperResponse.getBody(), GetEntriesResponse.class); return CompletableFuture.completedFuture(response); } diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotReader.java b/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotReader.java index 7afc6738..9a9484f3 100644 --- a/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotReader.java +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotReader.java @@ -16,10 +16,10 @@ package io.openmessaging.storage.dledger.snapshot.file; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.snapshot.SnapshotManager; import io.openmessaging.storage.dledger.snapshot.SnapshotMeta; import io.openmessaging.storage.dledger.snapshot.SnapshotReader; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.IOUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,7 +40,7 @@ public FileSnapshotReader(String snapshotStorePath) { @Override public SnapshotMeta load() throws IOException { - SnapshotMeta snapshotMetaFromJSON = JSON.parseObject(IOUtils.file2String(this.snapshotStorePath + + SnapshotMeta snapshotMetaFromJSON = DLedgerJsonUtils.parseObject(IOUtils.file2String(this.snapshotStorePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE), SnapshotMeta.class); if (snapshotMetaFromJSON == null) { return null; diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotWriter.java b/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotWriter.java index 11ef354d..0205071a 100644 --- a/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotWriter.java +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/snapshot/file/FileSnapshotWriter.java @@ -16,11 +16,11 @@ package io.openmessaging.storage.dledger.snapshot.file; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.snapshot.SnapshotManager; import io.openmessaging.storage.dledger.snapshot.SnapshotMeta; import io.openmessaging.storage.dledger.snapshot.SnapshotStatus; import io.openmessaging.storage.dledger.snapshot.SnapshotWriter; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.IOUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -85,7 +85,7 @@ public void save(SnapshotStatus status) throws IOException { } private void sync() throws IOException { - IOUtils.string2File(JSON.toJSONString(this.snapshotMeta), this.snapshotStorePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE); + IOUtils.string2File(DLedgerJsonUtils.toJsonString(this.snapshotMeta), this.snapshotStorePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE); } @Override diff --git a/dledger/src/main/java/io/openmessaging/storage/dledger/utils/DLedgerJsonUtils.java b/dledger/src/main/java/io/openmessaging/storage/dledger/utils/DLedgerJsonUtils.java new file mode 100644 index 00000000..4a3f57d1 --- /dev/null +++ b/dledger/src/main/java/io/openmessaging/storage/dledger/utils/DLedgerJsonUtils.java @@ -0,0 +1,43 @@ +/* + * Copyright 2017-2022 The DLedger Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.openmessaging.storage.dledger.utils; + +import com.alibaba.fastjson2.JSON; +import com.alibaba.fastjson2.JSONReader; +import com.alibaba.fastjson2.JSONWriter; + +public final class DLedgerJsonUtils { + + private DLedgerJsonUtils() { + } + + public static byte[] toJsonBytes(Object object) { + return JSON.toJSONBytes(object, JSONWriter.Feature.WriteByteArrayAsBase64); + } + + public static String toJsonString(Object object) { + return JSON.toJSONString(object, JSONWriter.Feature.WriteByteArrayAsBase64); + } + + public static T parseObject(byte[] bytes, Class clazz) { + return JSON.parseObject(bytes, clazz, JSONReader.Feature.Base64StringAsByteArray); + } + + public static T parseObject(String json, Class clazz) { + return JSON.parseObject(json, clazz, JSONReader.Feature.Base64StringAsByteArray); + } +} diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/DLedgerRpcNettyServiceTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/DLedgerRpcNettyServiceTest.java new file mode 100644 index 00000000..7dc3d34e --- /dev/null +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/DLedgerRpcNettyServiceTest.java @@ -0,0 +1,133 @@ +/* + * Copyright 2017-2022 The DLedger Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.openmessaging.storage.dledger; + +import io.openmessaging.storage.dledger.client.DLedgerClient; +import io.openmessaging.storage.dledger.protocol.AppendEntryRequest; +import io.openmessaging.storage.dledger.protocol.AppendEntryResponse; +import io.openmessaging.storage.dledger.protocol.DLedgerResponseCode; +import io.openmessaging.storage.dledger.protocol.GetEntriesResponse; +import java.io.BufferedReader; +import java.io.File; +import java.io.InputStreamReader; +import java.nio.charset.StandardCharsets; +import java.util.UUID; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.remoting.protocol.RemotingSerializable; +import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +public class DLedgerRpcNettyServiceTest extends ServerTestHarness { + + @Test + public void testRemotingCodecColdStart() throws Exception { + runColdStartProbe(); + } + + private void runColdStartProbe() throws Exception { + String javaExecutable = System.getProperty("java.home") + File.separator + "bin" + File.separator + "java"; + String classPath = System.getProperty("surefire.test.class.path", System.getProperty("java.class.path")); + ProcessBuilder processBuilder = new ProcessBuilder(javaExecutable, "-cp", classPath, ColdStartProbe.class.getName()); + processBuilder.environment().remove("JAVA_TOOL_OPTIONS"); + processBuilder.environment().remove("_JAVA_OPTIONS"); + processBuilder.environment().remove("JDK_JAVA_OPTIONS"); + Process process = processBuilder.redirectErrorStream(true).start(); + + boolean finished = process.waitFor(10, TimeUnit.SECONDS); + if (!finished) { + process.destroyForcibly(); + Assertions.fail("Cold-start probe did not finish"); + } + StringBuilder output = new StringBuilder(); + try (BufferedReader reader = new BufferedReader( + new InputStreamReader(process.getInputStream(), StandardCharsets.UTF_8))) { + String line; + while ((line = reader.readLine()) != null) { + output.append(line).append(System.lineSeparator()); + } + } + Assertions.assertEquals(0, process.exitValue(), output.toString()); + } + + @Test + public void testAppendAndGetBinaryBodyOverNetwork() { + String group = UUID.randomUUID().toString(); + String peers = "n0-localhost:" + nextPort(); + DLedgerServer server = launchServer(group, peers, "n0", "n0", DLedgerConfig.MEMORY); + DLedgerClient client = launchClient(group, peers); + byte[] body = new byte[] {0, 1, (byte) 0xFF, 65}; + try { + AppendEntryResponse appendResponse = client.append(body); + Assertions.assertEquals(DLedgerResponseCode.SUCCESS.getCode(), appendResponse.getCode()); + + GetEntriesResponse getResponse = client.get(appendResponse.getIndex()); + Assertions.assertEquals(1, getResponse.getEntries().size()); + Assertions.assertArrayEquals(body, getResponse.getEntries().get(0).getBody()); + } finally { + client.shutdown(); + server.shutdown(); + } + } + + @Test + public void testAppendCompletesWhenConnectionFails() throws Exception { + AbstractDLedgerServer server = Mockito.mock(AbstractDLedgerServer.class); + Mockito.when(server.getListenAddress()).thenReturn("localhost:" + nextPort()); + Mockito.when(server.getPeerAddr(Mockito.anyString(), Mockito.anyString())) + .thenReturn("localhost:1"); + + DLedgerRpcNettyService rpcService = new DLedgerRpcNettyService(server); + rpcService.startup(); + try { + AppendEntryRequest request = new AppendEntryRequest(); + request.setGroup("group"); + request.setRemoteId("n0"); + request.setBody(new byte[] {1}); + + AppendEntryResponse response = rpcService.append(request).get(5, TimeUnit.SECONDS); + + Assertions.assertEquals(DLedgerResponseCode.NETWORK_ERROR.getCode(), response.getCode()); + } finally { + rpcService.shutdown(); + } + } + + public static final class ColdStartProbe { + + private ColdStartProbe() { + } + + public static void main(String[] args) { + try { + ConsumerConnection connection = new ConsumerConnection(); + String json = RemotingSerializable.toJson(connection, false); + ConsumerConnection decoded = RemotingSerializable.fromJson(json, ConsumerConnection.class); + if (decoded == null || decoded.getConnectionSet() == null) { + throw new AssertionError("Remoting codec returned an incomplete ConsumerConnection: " + json); + } + Runtime.getRuntime().halt(0); + } catch (Throwable t) { + t.printStackTrace(System.err); + System.err.flush(); + Runtime.getRuntime().halt(1); + } + } + } + +} diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotManagerTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotManagerTest.java index a71334b7..e05aa27e 100644 --- a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotManagerTest.java +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotManagerTest.java @@ -1,6 +1,5 @@ package io.openmessaging.storage.dledger.snapshot; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.DLedgerConfig; import io.openmessaging.storage.dledger.DLedgerServer; import io.openmessaging.storage.dledger.ServerTestHarness; @@ -10,6 +9,7 @@ import io.openmessaging.storage.dledger.statemachine.MockStateMachine; import io.openmessaging.storage.dledger.statemachine.StateMachineCaller; import io.openmessaging.storage.dledger.util.FileTestUtil; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.IOUtils; import org.junit.jupiter.api.Test; @@ -109,7 +109,7 @@ public void testLoadErrorSnapshot() throws Exception { // Build error snapshot without state machine data long errorSnapshotIdx1 = 10; String errorSnapshotStoreBasePath1 = snapshotBaseDirPrefix + errorSnapshotIdx1; - IOUtils.string2File(JSON.toJSONString(new SnapshotMeta(errorSnapshotIdx1, 1)), + IOUtils.string2File(DLedgerJsonUtils.toJsonString(new SnapshotMeta(errorSnapshotIdx1, 1)), errorSnapshotStoreBasePath1 + File.separator + SnapshotManager.SNAPSHOT_META_FILE); // Build error snapshot without state machine meta @@ -120,7 +120,7 @@ public void testLoadErrorSnapshot() throws Exception { long snapshotIdx = 8; String snapshotStoreBasePath = snapshotBaseDirPrefix + snapshotIdx; SnapshotMeta snapshotMeta = new SnapshotMeta(snapshotIdx, 1); - IOUtils.string2File(JSON.toJSONString(snapshotMeta), snapshotStoreBasePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE); + IOUtils.string2File(DLedgerJsonUtils.toJsonString(snapshotMeta), snapshotStoreBasePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE); IOUtils.string2File("80", snapshotStoreBasePath + File.separator + SnapshotManager.SNAPSHOT_DATA_FILE); DLedgerServer server = launchServerWithStateMachine(group, peers, "n0", "n0", DLedgerConfig.FILE, 10, 10 * 1024 * 1024); diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotReaderTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotReaderTest.java index f538351b..0f2225c4 100644 --- a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotReaderTest.java +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotReaderTest.java @@ -1,8 +1,8 @@ package io.openmessaging.storage.dledger.snapshot; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.snapshot.file.FileSnapshotReader; import io.openmessaging.storage.dledger.util.FileTestUtil; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.IOUtils; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -18,7 +18,7 @@ public void testReaderLoad() throws IOException { String metaFilePath = FileTestUtil.TEST_BASE + File.separator + SnapshotManager.SNAPSHOT_META_FILE; try { SnapshotMeta snapshotMeta = new SnapshotMeta(10, 0); - IOUtils.string2File(JSON.toJSONString(snapshotMeta), metaFilePath); + IOUtils.string2File(DLedgerJsonUtils.toJsonString(snapshotMeta), metaFilePath); SnapshotReader reader = new FileSnapshotReader(FileTestUtil.TEST_BASE); Assertions.assertNull(reader.getSnapshotMeta()); diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotStoreTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotStoreTest.java index b3208416..0c8382b4 100644 --- a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotStoreTest.java +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotStoreTest.java @@ -1,6 +1,5 @@ package io.openmessaging.storage.dledger.snapshot; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.snapshot.file.FileSnapshotStore; import io.openmessaging.storage.dledger.util.FileTestUtil; import io.openmessaging.storage.dledger.utils.IOUtils; diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotWriterTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotWriterTest.java index 21df4363..6093b9d7 100644 --- a/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotWriterTest.java +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/snapshot/SnapshotWriterTest.java @@ -1,9 +1,9 @@ package io.openmessaging.storage.dledger.snapshot; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.snapshot.file.FileSnapshotStore; import io.openmessaging.storage.dledger.snapshot.file.FileSnapshotWriter; import io.openmessaging.storage.dledger.util.FileTestUtil; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.IOUtils; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -25,7 +25,7 @@ public void testWriterSave() throws IOException { String testDir = FileTestUtil.TEST_BASE + File.separator + SnapshotManager.SNAPSHOT_DIR_PREFIX + lastSnapshotIndex; try { Assertions.assertEquals(IOUtils.file2String(testDir + File.separator + - SnapshotManager.SNAPSHOT_META_FILE), JSON.toJSONString(snapshotMeta)); + SnapshotManager.SNAPSHOT_META_FILE), DLedgerJsonUtils.toJsonString(snapshotMeta)); } finally { IOUtils.deleteFile(new File(testDir)); } diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/statemachine/StateMachineCallerTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/statemachine/StateMachineCallerTest.java index ff4d0fe9..b31663b4 100644 --- a/dledger/src/test/java/io/openmessaging/storage/dledger/statemachine/StateMachineCallerTest.java +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/statemachine/StateMachineCallerTest.java @@ -22,7 +22,6 @@ import java.util.UUID; import java.util.concurrent.CountDownLatch; -import com.alibaba.fastjson.JSON; import io.openmessaging.storage.dledger.DLedgerConfig; import io.openmessaging.storage.dledger.DLedgerServer; import io.openmessaging.storage.dledger.MemberState; @@ -35,6 +34,7 @@ import io.openmessaging.storage.dledger.snapshot.hook.LoadSnapshotHook; import io.openmessaging.storage.dledger.store.file.DLedgerMmapFileStore; import io.openmessaging.storage.dledger.util.FileTestUtil; +import io.openmessaging.storage.dledger.utils.DLedgerJsonUtils; import io.openmessaging.storage.dledger.utils.IOUtils; import org.junit.jupiter.api.Test; @@ -72,7 +72,7 @@ public void testOnCommittedAndOnSnapshotSave() throws Exception { String snapshotMetaJSON = IOUtils.file2String(this.config.getSnapshotStoreBaseDir() + File.separator + SnapshotManager.SNAPSHOT_DIR_PREFIX + fsm.getAppliedIndex() + File.separator + SnapshotManager.SNAPSHOT_META_FILE); - SnapshotMeta snapshotMetaFromJSON = JSON.parseObject(snapshotMetaJSON, SnapshotMeta.class); + SnapshotMeta snapshotMetaFromJSON = DLedgerJsonUtils.parseObject(snapshotMetaJSON, SnapshotMeta.class); assertEquals(snapshotMetaFromJSON.getLastIncludedIndex(), 9); assertEquals(snapshotMetaFromJSON.getLastIncludedTerm(), 0); String snapshotData = IOUtils.file2String(this.config.getSnapshotStoreBaseDir() + File.separator + @@ -96,7 +96,7 @@ public void testOnSnapshotLoad() throws Exception { final long lastIncludedIndex = 10; String snapshotStoreBasePath = this.config.getSnapshotStoreBaseDir() + File.separator + SnapshotManager.SNAPSHOT_DIR_PREFIX + lastIncludedIndex; SnapshotMeta snapshotMeta = new SnapshotMeta(lastIncludedIndex, 1); - IOUtils.string2File(JSON.toJSONString(snapshotMeta), snapshotStoreBasePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE); + IOUtils.string2File(DLedgerJsonUtils.toJsonString(snapshotMeta), snapshotStoreBasePath + File.separator + SnapshotManager.SNAPSHOT_META_FILE); IOUtils.string2File("90", snapshotStoreBasePath + File.separator + SnapshotManager.SNAPSHOT_DATA_FILE); SnapshotReader reader = new FileSnapshotReader(snapshotStoreBasePath); @@ -204,4 +204,4 @@ public void testServerWithStateMachine() throws InterruptedException { assertEquals(10, fsm.getTotalEntries()); } } -} \ No newline at end of file +} diff --git a/dledger/src/test/java/io/openmessaging/storage/dledger/utils/DLedgerJsonUtilsTest.java b/dledger/src/test/java/io/openmessaging/storage/dledger/utils/DLedgerJsonUtilsTest.java new file mode 100644 index 00000000..ce2252bc --- /dev/null +++ b/dledger/src/test/java/io/openmessaging/storage/dledger/utils/DLedgerJsonUtilsTest.java @@ -0,0 +1,96 @@ +/* + * Copyright 2017-2022 The DLedger Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.openmessaging.storage.dledger.utils; + +import io.openmessaging.storage.dledger.entry.DLedgerEntry; +import io.openmessaging.storage.dledger.protocol.AppendEntryRequest; +import io.openmessaging.storage.dledger.protocol.BatchAppendEntryRequest; +import io.openmessaging.storage.dledger.protocol.PushEntryRequest; +import io.openmessaging.storage.dledger.snapshot.SnapshotMeta; +import java.nio.charset.StandardCharsets; +import java.util.Arrays; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +public class DLedgerJsonUtilsTest { + + private static final byte[] BODY = new byte[] {0, 1, (byte) 0xFF, 65}; + + @Test + public void testParseFastjson1Base64Fixture() { + AppendEntryRequest request = DLedgerJsonUtils.parseObject( + "{\"body\":\"AAH/QQ==\"}", AppendEntryRequest.class); + + Assertions.assertArrayEquals(BODY, request.getBody()); + } + + @Test + public void testSerializeAppendEntryRequestBodyAsBase64() { + AppendEntryRequest request = new AppendEntryRequest(); + request.setBody(BODY); + + String json = DLedgerJsonUtils.toJsonString(request); + + Assertions.assertTrue(json.contains("\"body\":\"AAH/QQ==\""), json); + Assertions.assertFalse(json.contains("\"body\":["), json); + } + + @Test + public void testRoundTripBatchAppendEntryRequestBodiesAsBase64() { + BatchAppendEntryRequest request = new BatchAppendEntryRequest(); + request.setBatchMsgs(Arrays.asList(BODY, new byte[] {2, 3})); + + String json = DLedgerJsonUtils.toJsonString(request); + BatchAppendEntryRequest decoded = DLedgerJsonUtils.parseObject(json, BatchAppendEntryRequest.class); + + Assertions.assertTrue(json.contains("\"batchMsgs\":[\"AAH/QQ==\",\"AgM=\"]"), json); + Assertions.assertEquals(2, decoded.getBatchMsgs().size()); + Assertions.assertArrayEquals(BODY, decoded.getBatchMsgs().get(0)); + Assertions.assertArrayEquals(new byte[] {2, 3}, decoded.getBatchMsgs().get(1)); + } + + @Test + public void testRoundTripPushEntryRequestBodyWithByteArrayApi() { + DLedgerEntry entry = new DLedgerEntry(); + entry.setBody(BODY); + PushEntryRequest request = new PushEntryRequest(); + request.setEntry(entry); + + byte[] jsonBytes = DLedgerJsonUtils.toJsonBytes(request); + String json = new String(jsonBytes, StandardCharsets.UTF_8); + PushEntryRequest decoded = DLedgerJsonUtils.parseObject(jsonBytes, PushEntryRequest.class); + + Assertions.assertTrue(json.contains("\"body\":\"AAH/QQ==\""), json); + Assertions.assertFalse(json.contains("\"body\":["), json); + Assertions.assertArrayEquals(BODY, decoded.getEntry().getBody()); + } + + @Test + public void testRoundTripSnapshotMetaFixture() { + SnapshotMeta snapshotMeta = DLedgerJsonUtils.parseObject( + "{\"lastIncludedIndex\":10,\"lastIncludedTerm\":2}", SnapshotMeta.class); + + Assertions.assertEquals(10, snapshotMeta.getLastIncludedIndex()); + Assertions.assertEquals(2, snapshotMeta.getLastIncludedTerm()); + + String json = DLedgerJsonUtils.toJsonString(snapshotMeta); + SnapshotMeta decoded = DLedgerJsonUtils.parseObject(json, SnapshotMeta.class); + + Assertions.assertEquals(10, decoded.getLastIncludedIndex()); + Assertions.assertEquals(2, decoded.getLastIncludedTerm()); + } +} diff --git a/pom.xml b/pom.xml index 2b5a7ff7..46596ea9 100644 --- a/pom.xml +++ b/pom.xml @@ -41,7 +41,8 @@ 1.7.36 1.30 1.72 - 4.9.4 + 5.5.0 + 2.0.64 jacoco 3.12.0 @@ -95,6 +96,11 @@ rocketmq-remoting ${rocketmq-remoting.version} + + com.alibaba.fastjson2 + fastjson2 + ${fastjson2.version} + org.slf4j slf4j-api @@ -335,4 +341,4 @@ - \ No newline at end of file +