From ce2381119e714643d5393035c3e30ad0bcaa5bd2 Mon Sep 17 00:00:00 2001
From: KKSU <15274802129@163.com>
Date: Mon, 17 Jun 2024 15:11:05 +0800
Subject: [PATCH] 后台
---
src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java | 254 +++++++++++++++++++++++++++++++++++---------------
1 files changed, 176 insertions(+), 78 deletions(-)
diff --git a/src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java b/src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java
index 135f689..f21a835 100644
--- a/src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java
+++ b/src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java
@@ -1,27 +1,27 @@
package cc.mrbird.febs.dapp.chain;
import cc.mrbird.febs.common.exception.FebsException;
-import cn.hutool.core.util.StrUtil;
-import cn.hutool.http.HttpUtil;
-import com.alibaba.fastjson.JSONObject;
+import io.reactivex.Flowable;
+import io.reactivex.disposables.Disposable;
+import io.reactivex.schedulers.Schedulers;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.data.repository.query.ParameterOutOfBoundsException;
import org.web3j.crypto.Credentials;
import org.web3j.protocol.Web3j;
import org.web3j.protocol.core.DefaultBlockParameter;
import org.web3j.protocol.core.DefaultBlockParameterName;
import org.web3j.protocol.core.DefaultBlockParameterNumber;
import org.web3j.protocol.core.methods.request.EthFilter;
-import org.web3j.protocol.core.methods.response.TransactionReceipt;
import org.web3j.protocol.http.HttpService;
+import org.web3j.protocol.websocket.WebSocketClient;
+import org.web3j.protocol.websocket.WebSocketService;
import org.web3j.tx.gas.StaticGasProvider;
-import java.math.BigDecimal;
import java.math.BigInteger;
-import java.rmi.activation.UnknownObjectException;
+import java.net.URI;
import java.util.HashMap;
-import java.util.List;
import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
/**
* @author
@@ -29,32 +29,17 @@
**/
@Slf4j
public class ChainService {
+ private final static Map<String, ContractChainService> contractMap = new HashMap<>();
- private final static String TRX_ADDRESS = "TUFzqZRpLwLWJU4jcdf77RKS3Ts2uEhmWL";
- private final static String TRX_PRIVATE = "e08dce7a4626f97b790e791bcdec31cffab46233744bb1aa133f69f98623d3fb";
- private final static String TRX_CONTRACT_ADDRESS = "TR7NHqjeKQxGTCi8q8ZY4pL8otSzgjLj6t";
- private final static String API_KEY = "9d461be6-9796-47b9-85d8-b150cbabbb54";
-
- private final static String ETH_URL = "https://mainnet.infura.io/v3/f54a5887a3894ebb9425920701a97fe0";
- private final static String ETH_ADDRESS = "0x6c5640c572504a75121e57760909a9dd0E672f2D";
- private final static String ETH_PRIVATE = "77f650768ff50a4243c008fbae1be9ffe74c52908ee9081e2e15f3d3411690bb";
- private final static String ETH_CONTRACT_ADDRESS = "0xdac17f958d2ee523a2206206994597c13d831ec7";
-
- private final static String BSC_URL = "https://bsc-dataseed1.ninicoin.io";
-// private final static String BSC_ADDRESS = "0x971c09aA9735EB98459B17EC8b48932D24CbB931";
-// private final static String BSC_PRIVATE = "0x5f38d0e63157f535fc21f89ea13ec3cd245691c20795c1d2cb60233b3ba7bb47";
-// private final static String BSC_CONTRACT_ADDRESS = "0x55d398326f99059fF775485246999027B3197955";
-
- private final static String BSC_ADDRESS = "0x977a9ddfb965a9a3416fa72ca7f91c4949c18f25";
- private final static String BSC_PRIVATE = "0xefe98e00cd227b6322e892c82fcbd8eadf119c3188b7e574bc624f65405d61bf";
- private final static String BSC_CONTRACT_ADDRESS = "0x6c6835e60e7dbad7a60112a6371271e8eb79ee68";
-
- private final static ContractChainService ETH = new EthService(ETH_URL, ETH_ADDRESS, ETH_PRIVATE, ETH_CONTRACT_ADDRESS);
- private final static ContractChainService BSC = new EthService(BSC_URL, BSC_ADDRESS, BSC_PRIVATE, BSC_CONTRACT_ADDRESS);
- private final static ContractChainService TRX = new TrxService(TRX_ADDRESS, TRX_PRIVATE, TRX_CONTRACT_ADDRESS, API_KEY);
- private final static ContractChainService BSC_TFC = new EthService(ChainEnum.BSC_TFC.getUrl(), ChainEnum.BSC_TFC.getAddress(), ChainEnum.BSC_TFC.getPrivateKey(), ChainEnum.BSC_TFC.getContractAddress());
-
- private final String ETH_PREFIX = "0x";
+ static {
+ for (ChainEnum chain : ChainEnum.values()) {
+ if ("TRX".equals(chain.getChain())) {
+ contractMap.put(chain.name(), new TrxService(chain.getAddress(), chain.getPrivateKey(), chain.getContractAddress(), chain.getApiKey()));
+ } else {
+ contractMap.put(chain.name(), new EthService(chain.getUrl(), chain.getAddress(), chain.getPrivateKey(), chain.getContractAddress()));
+ }
+ }
+ }
private ChainService() {
}
@@ -62,81 +47,194 @@
public final static ChainService INSTANCE = new ChainService();
public static ContractChainService getInstance(String chainType) {
- switch (chainType) {
- case "ETH" :
- return ETH;
- case "BSC" :
- return BSC;
- case "TRX" :
- return TRX;
- case "BSC_TFC":
- return BSC_TFC;
- default:
- break;
+ ContractChainService contract = contractMap.get(chainType);
+ if (contract == null) {
+ throw new FebsException("参数错误");
}
- throw new FebsException("参数错误");
+ return contract;
}
/**
* 监听合约事件
+ *
* @param startBlock 开始区块
*/
public static void contractEventListener(BigInteger startBlock, ContractEventService event, String type) {
+ contractEventListener(startBlock, null, event, type);
+ }
+
+ public static void contractEventListener(BigInteger startBlock, BigInteger endBlock, ContractEventService event, String type) {
ChainEnum chain = ChainEnum.getValueByName(type);
assert chain != null;
EthUsdtContract contract = contract(chain.getPrivateKey(), chain.getContractAddress(), chain.getUrl());
- EthFilter filter = getFilter(startBlock, chain.getContractAddress());
+ EthFilter filter = getFilter(startBlock, endBlock, chain.getContractAddress());
- contract.transferEventFlowable(filter).subscribe(e -> {
+ Flowable<EthUsdtContract.TransferEventResponse> eventFlowable = contract.transferEventFlowable(filter);
+ eventFlowable.subscribe(e -> {
event.compile(e);
}, error -> {
- log.error("--->", error);
+ log.error("合约监听启动报错", error);
});
}
- private static EthUsdtContract contract(String privateKey, String contractAddress, String url) {
- Credentials credentials = Credentials.create(privateKey);
- return EthUsdtContract.load(contractAddress, Web3j.build(new HttpService(url)), credentials, new StaticGasProvider(BigInteger.valueOf(4500000L), BigInteger.valueOf(200000L)));
- }
-
- // 18097238 18098663
- private static EthFilter getFilter(BigInteger startBlock, String contractAddress) {
- DefaultBlockParameter parameterName = null;
- if (startBlock != null) {
- parameterName = new DefaultBlockParameterNumber(startBlock);
- } else {
- parameterName = DefaultBlockParameterName.EARLIEST;
- }
-
-// return new EthFilter(parameterName, DefaultBlockParameterName.LATEST, contractAddress);
- return new EthFilter(parameterName, new DefaultBlockParameterNumber(new BigInteger("18098663")), contractAddress);
- }
-
- public static void main(String[] args) {
- ChainEnum chain = ChainEnum.getValueByName(ChainEnum.BSC_TFC.name());
+ public static void sdmUSDTEventListener(BigInteger startBlock, BigInteger endBlock, ContractEventService event, String type) {
+ ChainEnum chain = ChainEnum.getValueByName(type);
assert chain != null;
EthUsdtContract contract = contract(chain.getPrivateKey(), chain.getContractAddress(), chain.getUrl());
- EthFilter filter = getFilter(new BigInteger("18097238"), new BigInteger("18098663"), chain.getContractAddress());
+ EthFilter filter = getFilter(startBlock, endBlock, chain.getContractAddress());
- contract.transferEventFlowable(filter).subscribe(e -> {
- System.out.println(1);
- }, error -> {
- log.error("--->", error);
- });
+ Flowable<EthUsdtContract.TransferEventResponse> eventFlowable = contract.transferEventFlowable(filter)
+ .doOnError(throwable ->
+ log.error("合约事件监听发生错误: " + throwable.getMessage(), throwable)) // 更具体的错误日志记录
+ .retryWhen(errors -> {
+ AtomicInteger counter = new AtomicInteger();
+ return errors.takeWhile(e -> counter.getAndIncrement() != 3)
+ .flatMap(e -> {
+ System.out.println("delay retry by " + counter.get() + " second(s)");
+ return Flowable.timer(counter.get(), TimeUnit.SECONDS);
+ });
+ })
+ .subscribeOn(Schedulers.io()); // 指定subscribe操作在IO线程中执行,避免阻塞主线程
+
+ eventFlowable.subscribe(
+ e -> {
+ try {
+ event.sdmUSDT(e); // 处理事件
+ } catch (Exception ex) {
+ // 处理事件时可能出现的异常
+ log.error("处理合约事件时出错", ex);
+ }
+ },
+ Throwable::printStackTrace, // 打印错误堆栈,或者可以替换为更具体的错误处理逻辑
+ () -> log.info("合约事件监听已完成") // 在Flowable完成时执行的逻辑,如记录日志等
+ );
+
+ }
+
+
+ public static void sdmChargeEventListener(BigInteger startBlock, BigInteger endBlock, ContractEventService event, String type) {
+ ChainEnum chain = ChainEnum.getValueByName(type);
+ assert chain != null;
+
+ EthUsdtContract contract = contract(chain.getPrivateKey(), chain.getContractAddress(), chain.getUrl());
+ EthFilter filter = getFilter(startBlock, endBlock, chain.getContractAddress());
+
+ Flowable<EthUsdtContract.TransferEventResponse> eventFlowable = contract.transferEventFlowable(filter)
+ .doOnError(throwable ->
+ log.error("合约事件监听发生错误: " + throwable.getMessage(), throwable)) // 更具体的错误日志记录
+ .retryWhen(errors -> {
+ AtomicInteger counter = new AtomicInteger();
+ return errors.takeWhile(e -> counter.getAndIncrement() != 3)
+ .flatMap(e -> {
+ System.out.println("delay retry by " + counter.get() + " second(s)");
+ return Flowable.timer(counter.get(), TimeUnit.SECONDS);
+ });
+ })
+ .subscribeOn(Schedulers.io()); // 指定subscribe操作在IO线程中执行,避免阻塞主线程
+
+ eventFlowable.subscribe(
+ e -> {
+ try {
+ event.compile(e); // 处理事件
+ } catch (Exception ex) {
+ // 处理事件时可能出现的异常
+ log.error("处理合约事件时出错", ex);
+ }
+ },
+ Throwable::printStackTrace, // 打印错误堆栈,或者可以替换为更具体的错误处理逻辑
+ () -> log.info("合约事件监听已完成") // 在Flowable完成时执行的逻辑,如记录日志等
+ );
+
+ }
+
+ public static void wssContractEventListener(BigInteger startBlock, ContractEventService event, String type) {
+ WebSocketService ws = null;
+ WebSocketClient webSocketClient = null;
+ Web3j web3j = null;
+
+ try {
+ webSocketClient = new WebSocketClient(new URI("wss://bsc-mainnet.nodereal.io/ws/v1/78074065950e4915aef4f12b6f357d16"));
+ ws = new WebSocketService(webSocketClient, true);
+ ws.connect();
+ web3j = Web3j.build(ws);
+ ChainEnum chain = ChainEnum.getValueByName(type);
+ assert chain != null;
+
+ EthUsdtContract ethUsdtContract = wssContract(chain.getPrivateKey(), chain.getContractAddress(), web3j);
+ EthFilter filter = getFilter(startBlock, null, chain.getContractAddress());
+
+
+ Flowable<EthUsdtContract.TransferEventResponse> eventFlowable = ethUsdtContract.transferEventFlowable(filter);
+ Disposable subscribe = eventFlowable.subscribe(event::compile, error -> {
+ log.error("币安监听异常", error);
+ });
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+
+ }
+
+
+ private static EthUsdtContract contract(String privateKey, String contractAddress, String url) {
+ Credentials credentials = Credentials.create(privateKey);
+ HttpService httpService = new HttpService(url);
+// httpService.addHeader("Authorization", "Bearer " + Base64.encode("tfc:tfc123".getBytes()));
+// httpService.addHeader("Authorization", "Bearer eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJhdWQiOiJwdWJsaWMiLCJleHAiOjE2NTk5MzcxOTAsImp0aSI6IjRiMjNkYTVjLWRlZWEtNDYzNi04YjMwLWNmMmZmMjVkM2NlYyIsImlhdCI6MTY1OTkzMzU5MCwiaXNzIjoiQW5rciIsIm5iZiI6MTY1OTkzMzU5MCwic3ViIjoiZmNiNjY0YjItOGEwNC00N2E5LTg3ZjMtNTJhMjE2ODVlMzEzIn0.YfEwvDByU2MGHywsblZpEmKMIbjv4cWYkn5CaFglXY0TSANzd2pCSbIe40yU_R9_nV6xZeE8Uk74jJOdd_QvMpFyUgo-MMNWZP6uiEaYvK_K3tlpk5yzeZq9D4ruWaq8rFKggr-iaRGzu6coRSAOFv2prWll3a7NdEbmkM-y5Y85xYD6g1N-TPIpE_Y-_-WPf3JUavk744kG8YyHhGvAmk2IL0N2xePfC6CHesdJhwvmJJXzr_53dbPwit1y5KljS0iTZz3mGTML2bq4hGaEHbQxeY2fBpZOSm8sPMz-zB9IVJQKzH5-DXlPKz01mJ9XiBJlubfHsN72RdqFD-O2Tw");
+ return EthUsdtContract.load(contractAddress,
+ Web3j.build(httpService),
+ credentials,
+ new StaticGasProvider(BigInteger.valueOf(4500000L), BigInteger.valueOf(200000L)));
+ }
+
+ private static EthUsdtContract wssContract(String privateKey, String contractAddress, Web3j web3j) {
+ Credentials credentials = Credentials.create(privateKey);
+// httpService.addHeader("Authorization", "Bearer " + Base64.encode("tfc:tfc123".getBytes()));
+// httpService.addHeader("Authorization", "Bearer eyJhbGciOiJSUzI1NiIsInR5cCI6IkpXVCJ9.eyJhdWQiOiJwdWJsaWMiLCJleHAiOjE2NTk5MzcxOTAsImp0aSI6IjRiMjNkYTVjLWRlZWEtNDYzNi04YjMwLWNmMmZmMjVkM2NlYyIsImlhdCI6MTY1OTkzMzU5MCwiaXNzIjoiQW5rciIsIm5iZiI6MTY1OTkzMzU5MCwic3ViIjoiZmNiNjY0YjItOGEwNC00N2E5LTg3ZjMtNTJhMjE2ODVlMzEzIn0.YfEwvDByU2MGHywsblZpEmKMIbjv4cWYkn5CaFglXY0TSANzd2pCSbIe40yU_R9_nV6xZeE8Uk74jJOdd_QvMpFyUgo-MMNWZP6uiEaYvK_K3tlpk5yzeZq9D4ruWaq8rFKggr-iaRGzu6coRSAOFv2prWll3a7NdEbmkM-y5Y85xYD6g1N-TPIpE_Y-_-WPf3JUavk744kG8YyHhGvAmk2IL0N2xePfC6CHesdJhwvmJJXzr_53dbPwit1y5KljS0iTZz3mGTML2bq4hGaEHbQxeY2fBpZOSm8sPMz-zB9IVJQKzH5-DXlPKz01mJ9XiBJlubfHsN72RdqFD-O2Tw");
+ return EthUsdtContract.load(contractAddress,
+ web3j,
+ credentials,
+ new StaticGasProvider(BigInteger.valueOf(4500000L), BigInteger.valueOf(200000L)));
+ }
+
+ private static EthFilter getFilter(BigInteger startBlock, String contractAddress) {
+ return getFilter(startBlock, null, contractAddress);
}
private static EthFilter getFilter(BigInteger startBlock, BigInteger endBlock, String contractAddress) {
- DefaultBlockParameter parameterName = null;
+ DefaultBlockParameter startParameterName = null;
+ DefaultBlockParameter endParameterName = null;
if (startBlock != null) {
- parameterName = new DefaultBlockParameterNumber(startBlock);
+ startParameterName = new DefaultBlockParameterNumber(startBlock);
} else {
- parameterName = DefaultBlockParameterName.EARLIEST;
+ startParameterName = DefaultBlockParameterName.EARLIEST;
}
- return new EthFilter(parameterName, new DefaultBlockParameterNumber(endBlock), contractAddress);
+ if (endBlock != null) {
+ endParameterName = new DefaultBlockParameterNumber(endBlock);
+ } else {
+ endParameterName = DefaultBlockParameterName.LATEST;
+ }
+
+ return new EthFilter(startParameterName, endParameterName, contractAddress);
+ }
+
+ public static void main(String[] args) {
+// ChainEnum chain = ChainEnum.getValueByName(ChainEnum.BSC_TFC.name());
+// assert chain != null;
+//
+// EthUsdtContract contract = contract(chain.getPrivateKey(), chain.getContractAddress(), chain.getUrl());
+// EthFilter filter = getFilter(new BigInteger("18097238"), chain.getContractAddress());
+//
+// contract.transferEventFlowable(filter).subscribe(e -> {
+// System.out.println(1);
+// }, error -> {
+// log.error("--->", error);
+// });
+
+ System.out.println(ChainService.getInstance(ChainEnum.BSC_TFC.name()).totalSupply());
}
}
--
Gitblit v1.9.1