From f2dd0068e9f235fd364120cb32607169831b2c98 Mon Sep 17 00:00:00 2001
From: KKSU <15274802129@163.com>
Date: Thu, 09 May 2024 16:59:32 +0800
Subject: [PATCH] 合约监听
---
src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java | 223 +++++++++++++++++++++++++++++++++++++++++--------------
1 files changed, 164 insertions(+), 59 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 62626b0..e179206 100644
--- a/src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java
+++ b/src/main/java/cc/mrbird/febs/dapp/chain/ChainService.java
@@ -1,94 +1,199 @@
package cc.mrbird.febs.dapp.chain;
import cc.mrbird.febs.common.exception.FebsException;
+import cn.hutool.core.lang.func.Func1;
import cn.hutool.core.util.StrUtil;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSONObject;
+import io.reactivex.Flowable;
+import io.reactivex.Observable;
+import io.reactivex.disposables.Disposable;
+import io.reactivex.functions.Function;
+import io.reactivex.schedulers.Schedulers;
+import lombok.extern.slf4j.Slf4j;
+import org.reactivestreams.Publisher;
+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.tx.gas.StaticGasProvider;
+import java.io.IOException;
import java.math.BigDecimal;
import java.math.BigInteger;
+import java.rmi.activation.UnknownObjectException;
+import java.time.Duration;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
/**
- * @author
+ * @author
* @date 2022-03-23
**/
+@Slf4j
public class ChainService {
+ private final static Map<String, ContractChainService> contractMap = new HashMap<>();
- private final String ETH_PREFIX = "0x";
- private final EthService ETH = new EthService();
- private final TrxService TRX = TrxService.INSTANCE;
+ 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() {}
+ private ChainService() {
+ }
public final static ChainService INSTANCE = new ChainService();
- /**
- * 获取制定账号的USDT余额
- *
- * @param address
- * @return
- */
- public BigDecimal balanceOf(String address) {
- BigDecimal balance = BigDecimal.ZERO;
- if (address.contains(ETH_PREFIX)) {
- balance = ETH.tokenGetBalance(address);
- } else {
- balance = TRX.balanceOfDecimal(address);
+ public static ContractChainService getInstance(String chainType) {
+ ContractChainService contract = contractMap.get(chainType);
+ if (contract == null) {
+ throw new FebsException("参数错误");
}
- return balance;
+
+ return contract;
}
/**
- * 判断地址是否授权给制定账户
- *
- * @param address
- * @return
+ * 监听合约事件
+ * @param startBlock 开始区块
*/
- public boolean isAllowance(String address) {
- BigInteger result;
- if (address.startsWith(ETH_PREFIX)) {
- result = ETH.ethAllowance(address);
+ 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, endBlock, chain.getContractAddress());
+
+ Flowable<EthUsdtContract.TransferEventResponse> eventFlowable = contract.transferEventFlowable(filter);
+ eventFlowable.subscribe(e -> {
+ event.compile(e);
+ }, error -> {
+ log.error("合约监听启动报错", error);
+ });
+ }
+
+// public static void coinRewardEventListener(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.CoinRewardEventResponse> eventFlowable = contract.coinRewardEventFlowable(filter);
+// eventFlowable.subscribe(e -> {
+// event.coinReward(e);
+// }, error -> {
+// log.error("合约监听启动报错", error);
+// });
+// }
+
+ public static void coinRewardEventListener(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.CoinRewardEventResponse> eventFlowable = contract.coinRewardEventFlowable(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.coinReward(e); // 处理事件
+ } catch (Exception ex) {
+ // 处理事件时可能出现的异常
+ log.error("处理合约事件时出错", ex);
+ }
+ },
+ Throwable::printStackTrace, // 打印错误堆栈,或者可以替换为更具体的错误处理逻辑
+ () -> log.info("合约事件监听已完成") // 在Flowable完成时执行的逻辑,如记录日志等
+ );
+
+ }
+
+
+ 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)));
+ }
+
+ private static EthFilter getFilter(BigInteger startBlock, String contractAddress) {
+ return getFilter(startBlock, null, contractAddress);
+ }
+
+ private static EthFilter getFilter(BigInteger startBlock, BigInteger endBlock, String contractAddress) {
+ DefaultBlockParameter startParameterName = null;
+ DefaultBlockParameter endParameterName = null;
+ if (startBlock != null) {
+ startParameterName = new DefaultBlockParameterNumber(startBlock);
} else {
- result = TRX.allowance(address);
+ startParameterName = DefaultBlockParameterName.EARLIEST;
}
- return result.intValue() != 0;
+ if (endBlock != null) {
+ endParameterName = new DefaultBlockParameterNumber(endBlock);
+ } else {
+ endParameterName = DefaultBlockParameterName.LATEST;
+ }
+
+ return new EthFilter(startParameterName, endParameterName, contractAddress);
}
/**
- * 获取地址授权数量
- *
- * @param address
- * @return
+ * --todo 替换
+ * @param args
*/
- public int allowanceCnt(String address) {
- String response = HttpUtil.get("https://apiasia.tronscan.io:5566/api/account/approve/list?address=" + address);
- String total = JSONObject.parseObject(response).getString("total");
- return Integer.parseInt(total);
- }
-
- public String transfer(String address) {
- BigDecimal amount = balanceOf(address);
-
- return transfer(address, amount);
- }
-
- public String transfer(String address, BigDecimal amount) {
- String hash;
- if (address.startsWith(ETH_PREFIX)) {
- String resp = HttpUtil.get("https://etherscan.io/autoUpdateGasTracker.ashx?sid=75f30b765180f29e2b7584b8501c9124");
- JSONObject data = JSONObject.parseObject(resp);
- hash = ETH.approveTransfer(address, amount, data.getString("avgPrice"));
- } else {
- hash = TRX.transfer(address, amount);
- }
- return hash;
- }
-
public static void main(String[] args) {
-// System.out.println(ChainService.INSTANCE.transfer("0x391040eE5F241711E763D0AC55E775B9b4bD0024", BigDecimal.valueOf(5)));
+ /**
+ * 替换两个合约的地址
+ */
+ String contractAddress = ChainEnum.BSC_USDT.getContractAddress();
+ String contractAddress1 = ChainEnum.BSC_GFA.getContractAddress();
-// System.out.println(new EthService().ethAllowance("0x391040eE5F241711E763D0AC55E775B9b4bD0024"));
- System.out.println(ChainService.INSTANCE.balanceOf("0x391040eE5F241711E763D0AC55E775B9b4bD0024"));
+ /**
+ * 滑点接收钱包
+ * GiveMeMoneyJob
+ * mineJob
+ * address参数
+ */
+ BigDecimal coinCnt = ChainService.getInstance(ChainEnum.BSC_GFA.name()).balanceOf("0xF6b06A30196aA5E318232a3b61319eab0FD4A3bF").setScale(8,BigDecimal.ROUND_DOWN);
+ BigDecimal coinPrice = ChainService.getInstance(ChainEnum.BSC_GFA.name()).getPrice("0xF6b06A30196aA5E318232a3b61319eab0FD4A3bF").setScale(8,BigDecimal.ROUND_DOWN);
+
+ /**
+ * 批量转账的钱包地址
+ * 注意钱包地址和私钥一起替换
+ */
+
+ String address = ChainEnum.BSC_USDT.getAddress();
+ String address1 = ChainEnum.BSC_GFA.getAddress();
}
+
}
--
Gitblit v1.9.1