Files
SHiNE-server/SHiNE-server/shine-server-net-protocol/src/main/java/server/sync/AddBlockSyncService.java
T
AidarKC 3a5851939e Починить межсерверную репликацию личных сообщений
Репликация личных сообщений между серверами теперь работает корректно.
2026-08-25 19:31:38 +04:00

251 lines
10 KiB
Java

package server.sync;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import server.logic.ws_protocol.Base64Ws;
import shine.db.dao.BlocksDAO;
import shine.db.dao.SyncServersDAO;
import shine.db.entities.BlockEntry;
import shine.db.entities.SyncServerEntry;
import utils.blockchain.BlockchainNameUtil;
import java.util.List;
import java.util.Locale;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;
/**
* Фоновая репликация AddBlock через общий постоянный WSS-пул на серверы
* из локальной таблицы sync_servers.
*/
public final class AddBlockSyncService {
private static final Logger log = LoggerFactory.getLogger(AddBlockSyncService.class);
private static final ObjectMapper MAPPER = new ObjectMapper();
private static final ExecutorService EXECUTOR = new ThreadPoolExecutor(
1,
Math.max(2, Runtime.getRuntime().availableProcessors()),
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(10_000),
new ThreadFactory() {
private final AtomicLong n = new AtomicLong(1);
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, "sync-addblock-" + n.getAndIncrement());
t.setDaemon(true);
return t;
}
},
new ThreadPoolExecutor.DiscardPolicy()
);
private final BlocksDAO blocksDAO = BlocksDAO.getInstance();
private final SyncServersDAO syncServersDAO = SyncServersDAO.getInstance();
public void replicateAsync(String blockchainName, int blockNumber) {
EXECUTOR.execute(() -> {
try {
replicate(blockchainName, blockNumber);
} catch (Exception e) {
log.error("AddBlock sync failed unexpectedly (blockchainName={}, blockNumber={})",
blockchainName, blockNumber, e);
}
});
}
private void replicate(String blockchainName, int blockNumber) throws Exception {
String ownerLogin = normalize(BlockchainNameUtil.loginFromBlockchainName(blockchainName));
if (ownerLogin == null) {
log.warn("AddBlock sync skipped: cannot derive owner login from blockchainName={}", blockchainName);
return;
}
List<SyncServerEntry> partners = syncServersDAO.listAll();
if (partners.isEmpty()) {
return;
}
BlockEntry currentBlock = blocksDAO.getByNumber(blockchainName, blockNumber);
if (currentBlock == null || currentBlock.getBlockBytes() == null) {
log.warn("AddBlock sync skipped: block not found in DB (blockchainName={}, blockNumber={})",
blockchainName, blockNumber);
return;
}
for (SyncServerEntry partner : partners) {
if (partner == null) continue;
String partnerLogin = normalize(partner.getLogin());
if (partnerLogin == null) continue;
if (partnerLogin.equals(ownerLogin)) {
continue;
}
try {
replicateToPartner(partner, blockchainName, blockNumber, currentBlock);
} catch (Exception e) {
log.warn("AddBlock sync aborted for partner login={} blockchainName={} blockNumber={} reason={}",
partnerLogin, blockchainName, blockNumber, e.toString());
}
}
}
private void replicateToPartner(SyncServerEntry partner, String blockchainName, int blockNumber, BlockEntry currentBlock) throws Exception {
String wsUrl = ServerConnectionPool.buildWsUrl(partner.getServerAddress());
if (wsUrl == null) {
log.warn("AddBlock sync skipped: invalid server_address for partner login={} address={}",
partner.getLogin(), partner.getServerAddress());
return;
}
AddBlockPushResult firstTry = pushBlock(partner, blockchainName, currentBlock);
if (firstTry.ok()) {
log.info("AddBlock sync ok: partner={} blockchainName={} blockNumber={}",
partner.getLogin(), blockchainName, blockNumber);
return;
}
if (firstTry.serverAlreadyHasBlock(blockNumber)) {
log.info("AddBlock sync skipped: partner already has block. partner={} blockchainName={} blockNumber={} remoteLast={}",
partner.getLogin(), blockchainName, blockNumber, firstTry.serverLastGlobalNumber());
return;
}
if (!firstTry.needsBackfill()) {
log.warn("AddBlock sync failed without backfill: partner={} blockchainName={} blockNumber={} code={}",
partner.getLogin(), blockchainName, blockNumber, firstTry.code());
return;
}
int remoteLast = firstTry.serverLastGlobalNumber();
int fromBlockNumber = remoteLast + 1;
if (fromBlockNumber > blockNumber) {
log.info("AddBlock sync skipped: partner already caught up during backfill window. partner={} blockchainName={} remoteLast={} target={}",
partner.getLogin(), blockchainName, remoteLast, blockNumber);
return;
}
List<BlockEntry> missingBlocks = blocksDAO.listRangeByNumber(blockchainName, fromBlockNumber, blockNumber);
if (missingBlocks.isEmpty()) {
log.warn("AddBlock sync backfill failed: local range empty partner={} blockchainName={} from={} to={}",
partner.getLogin(), blockchainName, fromBlockNumber, blockNumber);
return;
}
for (BlockEntry blockEntry : missingBlocks) {
AddBlockPushResult backfillResult = pushBlock(partner, blockchainName, blockEntry);
if (!backfillResult.ok()) {
log.warn("AddBlock sync backfill failed: partner={} blockchainName={} blockNumber={} code={}",
partner.getLogin(), blockchainName, blockEntry.getBlockNumber(), backfillResult.code());
return;
}
}
log.info("AddBlock sync backfill ok: partner={} blockchainName={} from={} to={}",
partner.getLogin(), blockchainName, fromBlockNumber, blockNumber);
}
private AddBlockPushResult pushBlock(SyncServerEntry partner, String blockchainName, BlockEntry blockEntry) throws Exception {
JsonNode response = sendAddBlock(partner, blockchainName, blockEntry);
int status = response.path("status").asInt(500);
if (status >= 200 && status < 300) {
return AddBlockPushResult.success();
}
String code = textOrEmpty(response, "code");
if (code.isBlank()) {
code = textOrEmpty(response, "error");
}
JsonNode payload = response.path("payload");
int serverLastGlobalNumber = payload.path("serverLastGlobalNumber").asInt(Integer.MIN_VALUE);
String serverLastGlobalHash = payload.path("serverLastGlobalHash").asText("");
return new AddBlockPushResult(false, status, code, serverLastGlobalNumber, serverLastGlobalHash);
}
private JsonNode sendAddBlock(SyncServerEntry partner, String blockchainName, BlockEntry blockEntry) throws Exception {
String jsonTemplate = buildAddBlockJsonTemplate(blockchainName, blockEntry);
return ServerConnectionPool.getInstance().request(
partner.getLogin(),
partner.getServerAddress(),
jsonTemplate,
ServerConnectionPool.Priority.BULK);
}
private String buildAddBlockJsonTemplate(String blockchainName, BlockEntry blockEntry) throws Exception {
String prevHashHex = blockEntry.getBlockNumber() <= 0
? ""
: toHex(extractPrevHash32(blockEntry.getBlockBytes()));
String blockBytesB64 = Base64Ws.encode(blockEntry.getBlockBytes());
String safeBlockchainName = MAPPER.writeValueAsString(blockchainName);
String safePrevHashHex = MAPPER.writeValueAsString(prevHashHex);
String safeBlockBytes = MAPPER.writeValueAsString(blockBytesB64);
return """
{
"op":"AddBlock",
"requestId":%s,
"payload":{
"blockchainName":%s,
"blockNumber":%d,
"prevBlockHash":%s,
"blockBytesB64":%s
}
}
""".formatted("%s", safeBlockchainName, blockEntry.getBlockNumber(), safePrevHashHex, safeBlockBytes);
}
private static byte[] extractPrevHash32(byte[] blockBytes) {
if (blockBytes == null || blockBytes.length < 44) {
return new byte[32];
}
byte[] out = new byte[32];
System.arraycopy(blockBytes, 12, out, 0, 32);
return out;
}
private static String textOrEmpty(JsonNode node, String field) {
return node == null ? "" : node.path(field).asText("");
}
private static String normalize(String value) {
if (value == null) return null;
String s = value.trim().toLowerCase(Locale.ROOT);
return s.isEmpty() ? null : s;
}
private static String toHex(byte[] bytes) {
if (bytes == null) return "";
StringBuilder sb = new StringBuilder(bytes.length * 2);
for (byte b : bytes) {
sb.append(Character.forDigit((b >>> 4) & 0xF, 16));
sb.append(Character.forDigit(b & 0xF, 16));
}
return sb.toString();
}
private record AddBlockPushResult(
boolean ok,
int status,
String code,
int serverLastGlobalNumber,
String serverLastGlobalHash
) {
static AddBlockPushResult success() {
return new AddBlockPushResult(true, 200, "", Integer.MIN_VALUE, "");
}
boolean needsBackfill() {
return !ok && ("bad_prev_hash".equalsIgnoreCase(code) || "bad_block_number".equalsIgnoreCase(code));
}
boolean serverAlreadyHasBlock(int targetBlockNumber) {
return !ok
&& serverLastGlobalNumber != Integer.MIN_VALUE
&& serverLastGlobalNumber >= targetBlockNumber;
}
}
}