Добавить resync блокчейна при рассинхроне

This commit is contained in:
AidarKC
2026-06-26 15:25:11 +04:00
parent 23edad416c
commit be4f76834a
11 changed files with 783 additions and 17 deletions
@@ -0,0 +1,394 @@
package shine.db.dao;
import shine.db.DatabaseInitializer;
import shine.db.SqliteDbController;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
/**
* BlockchainResyncCleanupDAO — подготовительный "жёсткий reset" одной blockchain-цепочки
* перед её полной повторной загрузкой от сервера-партнёра.
*
* Что делает этот DAO:
* 1) в ОДНОЙ SQL-транзакции сначала аккуратно уменьшает чужие агрегаты,
* которые были увеличены блоками удаляемой цепочки:
* - likes_count
* - replies_count
* 2) затем удаляет все локальные записи самой цепочки и её производные состояния.
*
* Почему это вынесено в отдельный DAO-метод, а не в триггеры DELETE:
* - нам нужен один понятный "блок операции", который можно вызвать из resync-flow;
* - эта схема проще и прозрачнее, чем много обратных триггеров по разным таблицам;
* - если любой шаг не удался, делаем rollback и БД остаётся в исходном состоянии;
* - файловые действия (.bch / .tmp_bch) сознательно НЕ входят в эту транзакцию:
* SQLite не может атомарно закоммитить и SQL, и файловую систему сразу;
* поэтому БД-чистка делается здесь, а файловая чистка будет следующим шагом
* отдельным recovery/resync-слоем после успешного commit.
*
* Важный смысл текущей реализации:
* - мы НЕ трогаем identity-слой (`solana_users`) и НЕ трогаем DM-таблицы;
* - мы очищаем только блокчейн пользователя и derived-state, который строится из неё;
* - висячие cross-chain ссылки в чужих blocks допускаются как нормальное поведение системы.
*/
public final class BlockchainResyncCleanupDAO {
private static final int BLOCKCHAIN_LOGIN_SUFFIX_LEN = 4; // "-001"
private static volatile BlockchainResyncCleanupDAO instance;
private final SqliteDbController db = SqliteDbController.getInstance();
private BlockchainResyncCleanupDAO() {}
public static BlockchainResyncCleanupDAO getInstance() {
if (instance == null) {
synchronized (BlockchainResyncCleanupDAO.class) {
if (instance == null) instance = new BlockchainResyncCleanupDAO();
}
}
return instance;
}
/**
* Полностью очищает одну blockchain-цепочку и локальные derived-state, собранные из неё.
*
* Порядок внутри транзакции намеренно такой:
* 1. Сначала уменьшаем чужие likes_count для тех целей, где финальное состояние
* реакции этой цепочки было LIKE.
* 2. Сначала уменьшаем чужие replies_count для reply-блоков этой цепочки.
* 3. После этого удаляем локальные derived-state самой цепочки.
* 4. В конце удаляем blocks и blockchain_state.
*
* Это правильно потому, что агрегаты (`message_stats`) должны видеть исходные blocks
* и reactions_state на момент пересчёта. Если удалить blocks раньше, мы потеряем
* источник правды для корректного уменьшения счётчиков.
*
* Метод идемпотентен по смыслу:
* - если часть данных уже удалена раньше, повторный вызов просто удалит "0 строк";
* - если blockchain_state уже отсутствует, login берём из blockchainName.
*
* Отдельно важно:
* - здесь НЕТ удаления .bch/.tmp_bch;
* - здесь НЕТ повторной загрузки цепочки;
* - это только атомарная SQL-очистка БД, на которую потом будет опираться resync-flow.
*/
public CleanupResult cleanupBlockchainForFullResync(String blockchainName) throws SQLException {
if (blockchainName == null || blockchainName.isBlank()) {
throw new IllegalArgumentException("blockchainName is blank");
}
try (Connection c = db.getConnection()) {
boolean oldAutoCommit = c.getAutoCommit();
c.setAutoCommit(false);
try {
String login = resolveLoginForCleanup(c, blockchainName);
int likesAdjusted = decreaseForeignLikesCount(c, blockchainName);
int repliesAdjusted = decreaseForeignRepliesCount(c, blockchainName);
int deletedMessageStats = deleteMessageStatsForOwnTargets(c, blockchainName);
int deletedReactionsState = deleteReactionsStateForActorChain(c, blockchainName);
int deletedConnectionsState = deleteConnectionsStateForLogin(c, login);
int deletedUsersParams = deleteUsersParamsForLogin(c, login);
int deletedChannelNames = deleteChannelNamesForOwnerChain(c, blockchainName);
int deletedChat200State = deleteChat200StateForOwnerChain(c, blockchainName);
int deletedChat200Members = deleteChat200MembersForOwnerChain(c, blockchainName);
int deletedBlocks = deleteBlocksForChain(c, blockchainName);
int deletedBlockchainState = deleteBlockchainStateForChain(c, blockchainName);
c.commit();
return new CleanupResult(
login,
likesAdjusted,
repliesAdjusted,
deletedMessageStats,
deletedReactionsState,
deletedConnectionsState,
deletedUsersParams,
deletedChannelNames,
deletedChat200State,
deletedChat200Members,
deletedBlocks,
deletedBlockchainState
);
} catch (Exception e) {
try { c.rollback(); } catch (Exception ignored) {}
if (e instanceof SQLException sqlEx) throw sqlEx;
throw new SQLException("Не удалось очистить blockchain для полного resync: " + blockchainName, e);
} finally {
try { c.setAutoCommit(oldAutoCommit); } catch (Exception ignored) {}
}
}
}
private String resolveLoginForCleanup(Connection c, String blockchainName) throws SQLException {
String sql = """
SELECT login
FROM blockchain_state
WHERE blockchain_name = ?
LIMIT 1
""";
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, blockchainName);
try (ResultSet rs = ps.executeQuery()) {
if (rs.next()) {
String login = rs.getString("login");
if (login != null && !login.isBlank()) {
return login;
}
}
}
}
String loginFromName = loginFromBlockchainName(blockchainName);
if (loginFromName == null || loginFromName.isBlank()) {
throw new IllegalArgumentException("Cannot derive login from blockchainName: " + blockchainName);
}
return loginFromName;
}
/**
* DAO остаётся в модуле БД и не тянет зависимость на blockchain-utils модуль.
* Поэтому здесь локально повторяем минимальное правило имени chain:
* login + "-NNN".
*/
private String loginFromBlockchainName(String blockchainName) {
if (blockchainName == null) return null;
String s = blockchainName.trim();
if (s.length() <= BLOCKCHAIN_LOGIN_SUFFIX_LEN) return null;
int dashPos = s.length() - BLOCKCHAIN_LOGIN_SUFFIX_LEN;
if (s.charAt(dashPos) != '-') return null;
for (int i = dashPos + 1; i < s.length(); i++) {
char ch = s.charAt(i);
if (ch < '0' || ch > '9') return null;
}
return s.substring(0, dashPos);
}
/**
* Уменьшаем likes_count только для ЧУЖИХ целей.
*
* Логика:
* - если у удаляемой цепочки финальное состояние реакции на цель = LIKE,
* значит при полном удалении цепочки этот активный лайк исчезает;
* - значит у message_stats этой чужой цели нужно сделать -1;
* - для целей внутри этой же chain этого делать не нужно, потому что сами цели
* тоже будут удалены вместе с цепочкой.
*/
private int decreaseForeignLikesCount(Connection c, String blockchainName) throws SQLException {
String sql = """
UPDATE message_stats
SET likes_count = MAX(
0,
likes_count - (
SELECT COUNT(*)
FROM reactions_state rs
WHERE rs.from_bch_name = ?
AND rs.reaction_type = ?
AND rs.last_sub_type = ?
AND rs.to_login = message_stats.to_login
AND rs.to_bch_name = message_stats.to_bch_name
AND rs.to_block_number = message_stats.to_block_number
AND rs.to_block_hash = message_stats.to_block_hash
AND rs.to_bch_name <> ?
)
)
WHERE EXISTS (
SELECT 1
FROM reactions_state rs
WHERE rs.from_bch_name = ?
AND rs.reaction_type = ?
AND rs.last_sub_type = ?
AND rs.to_login = message_stats.to_login
AND rs.to_bch_name = message_stats.to_bch_name
AND rs.to_block_number = message_stats.to_block_number
AND rs.to_block_hash = message_stats.to_block_hash
AND rs.to_bch_name <> ?
)
""";
try (PreparedStatement ps = c.prepareStatement(sql)) {
int i = 1;
ps.setString(i++, blockchainName);
ps.setInt(i++, DatabaseInitializer.REACTION_LIKE);
ps.setInt(i++, DatabaseInitializer.REACTION_LIKE);
ps.setString(i++, blockchainName);
ps.setString(i++, blockchainName);
ps.setInt(i++, DatabaseInitializer.REACTION_LIKE);
ps.setInt(i++, DatabaseInitializer.REACTION_LIKE);
ps.setString(i++, blockchainName);
return ps.executeUpdate();
}
}
/**
* Уменьшаем replies_count только для ЧУЖИХ целей.
*
* Если reply этой цепочки ссылался на сообщение из другой цепочки,
* значит после удаления blocks этой цепочки чужой replies_count должен уменьшиться.
* Reply на собственные сообщения здесь игнорируем: целевая цепочка тоже будет удалена.
*/
private int decreaseForeignRepliesCount(Connection c, String blockchainName) throws SQLException {
String sql = """
UPDATE message_stats
SET replies_count = MAX(
0,
replies_count - COALESCE((
SELECT COUNT(*)
FROM blocks b
WHERE b.bch_name = ?
AND b.msg_type = 1
AND b.msg_sub_type = ?
AND b.to_login = message_stats.to_login
AND b.to_bch_name = message_stats.to_bch_name
AND b.to_block_number = message_stats.to_block_number
AND b.to_block_hash = message_stats.to_block_hash
AND b.to_bch_name <> ?
), 0)
)
WHERE EXISTS (
SELECT 1
FROM blocks b
WHERE b.bch_name = ?
AND b.msg_type = 1
AND b.msg_sub_type = ?
AND b.to_login = message_stats.to_login
AND b.to_bch_name = message_stats.to_bch_name
AND b.to_block_number = message_stats.to_block_number
AND b.to_block_hash = message_stats.to_block_hash
AND b.to_bch_name <> ?
)
""";
try (PreparedStatement ps = c.prepareStatement(sql)) {
int i = 1;
ps.setString(i++, blockchainName);
ps.setInt(i++, DatabaseInitializer.TEXT_REPLY);
ps.setString(i++, blockchainName);
ps.setString(i++, blockchainName);
ps.setInt(i++, DatabaseInitializer.TEXT_REPLY);
ps.setString(i++, blockchainName);
return ps.executeUpdate();
}
}
/**
* Статистика сообщений самой удаляемой цепочки после reset не нужна,
* потому что её цели исчезают вместе с chain source data.
*/
private int deleteMessageStatsForOwnTargets(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM message_stats
WHERE to_bch_name = ?
""", blockchainName);
}
/**
* reactions_state хранит финальное состояние реакций АКТОРА.
* После удаления всей цепочки актор этой цепочки исчезает, поэтому
* достаточно удалить все строки по from_bch_name.
*/
private int deleteReactionsStateForActorChain(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM reactions_state
WHERE from_bch_name = ?
""", blockchainName);
}
/**
* connections_state — текущее состояние связей, выставленных этим login.
* Чистим по владельцу состояния.
*/
private int deleteConnectionsStateForLogin(Connection c, String login) throws SQLException {
return executeDelete(c, """
DELETE FROM connections_state
WHERE login = ? COLLATE NOCASE
""", login);
}
/**
* users_params — актуальные параметры, собранные из блоков пользователя.
*/
private int deleteUsersParamsForLogin(Connection c, String login) throws SQLException {
return executeDelete(c, """
DELETE FROM users_params
WHERE login = ? COLLATE NOCASE
""", login);
}
/**
* Каналы принадлежат owner_bch_name.
*/
private int deleteChannelNamesForOwnerChain(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM channel_names_state
WHERE owner_bch_name = ?
""", blockchainName);
}
private int deleteChat200StateForOwnerChain(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM chat200_state
WHERE owner_bch_name = ?
""", blockchainName);
}
private int deleteChat200MembersForOwnerChain(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM chat200_members_state
WHERE owner_bch_name = ?
""", blockchainName);
}
/**
* blocks удаляем в конце, потому что до этого шага они нужны как источник правды
* для уменьшения replies_count.
*/
private int deleteBlocksForChain(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM blocks
WHERE bch_name = ?
""", blockchainName);
}
/**
* blockchain_state удаляем после blocks, чтобы не нарушать FK-связь blocks -> blockchain_state.
*/
private int deleteBlockchainStateForChain(Connection c, String blockchainName) throws SQLException {
return executeDelete(c, """
DELETE FROM blockchain_state
WHERE blockchain_name = ?
""", blockchainName);
}
private int executeDelete(Connection c, String sql, String value) throws SQLException {
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, value);
return ps.executeUpdate();
}
}
/**
* Технический результат cleanup-операции.
* Нужен для будущего логирования и ручной диагностики resync-flow.
*/
public record CleanupResult(
String login,
int likesAdjustedRows,
int repliesAdjustedRows,
int deletedMessageStatsRows,
int deletedReactionsStateRows,
int deletedConnectionsStateRows,
int deletedUsersParamsRows,
int deletedChannelNamesRows,
int deletedChat200StateRows,
int deletedChat200MembersRows,
int deletedBlocksRows,
int deletedBlockchainStateRows
) {}
}