SHA256
510 lines
23 KiB
Java
510 lines
23 KiB
Java
package shine.db.dao;
|
||
|
||
import shine.db.DatabaseInitializer;
|
||
import shine.db.DbController;
|
||
|
||
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 и БД остаётся в исходном состоянии;
|
||
* - пользовательские blockchain-файлы больше не существуют; cleanup полностью DB-only.
|
||
* SQL-транзакция и файловая система не коммитятся атомарно вместе;
|
||
* поэтому БД-чистка делается здесь, а файловая чистка будет следующим шагом
|
||
* отдельным recovery/resync-слоем после успешного commit.
|
||
*
|
||
* Важный смысл текущей реализации:
|
||
* - мы НЕ трогаем current users слой (`solana_user_pda_current`) и НЕ трогаем 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 DbController db = DbController.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.
|
||
*
|
||
* Отдельно важно:
|
||
* - filesystem blockchain storage отсутствует;
|
||
* - здесь НЕТ повторной загрузки цепочки;
|
||
* - это только атомарная 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);
|
||
|
||
rebuildStatsState(c);
|
||
|
||
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 = GREATEST(
|
||
0,
|
||
likes_count - COALESCE((
|
||
SELECT COUNT(*)::int
|
||
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 <> ?
|
||
), 0)
|
||
)
|
||
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 = GREATEST(
|
||
0,
|
||
replies_count - COALESCE((
|
||
SELECT COUNT(*)::int
|
||
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 LOWER(login) = LOWER(?)
|
||
""", login);
|
||
}
|
||
|
||
/**
|
||
* users_params — актуальные параметры, собранные из блоков пользователя.
|
||
*/
|
||
private int deleteUsersParamsForLogin(Connection c, String login) throws SQLException {
|
||
return executeDelete(c, """
|
||
DELETE FROM users_params
|
||
WHERE LOWER(login) = LOWER(?)
|
||
""", 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 void rebuildStatsState(Connection c) throws SQLException {
|
||
try (PreparedStatement truncate = c.prepareStatement("""
|
||
TRUNCATE TABLE user_stats_state, channel_stats_state
|
||
""")) {
|
||
truncate.executeUpdate();
|
||
}
|
||
|
||
try (PreparedStatement ps = c.prepareStatement("""
|
||
INSERT INTO user_stats_state (
|
||
login,
|
||
owned_public_channels_count,
|
||
following_users_count,
|
||
following_channels_count,
|
||
close_friends_count,
|
||
updated_at_ms
|
||
)
|
||
SELECT
|
||
u.login,
|
||
COALESCE(own.owned_public_channels_count, 0),
|
||
COALESCE(fu.following_users_count, 0),
|
||
COALESCE(fc.following_channels_count, 0),
|
||
COALESCE(cf.close_friends_count, 0),
|
||
CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
|
||
FROM solana_user_pda_current u
|
||
LEFT JOIN (
|
||
SELECT owner_login, COUNT(*)::INTEGER AS owned_public_channels_count
|
||
FROM channel_names_state
|
||
WHERE channel_type_code = 1
|
||
GROUP BY owner_login
|
||
) own ON LOWER(own.owner_login) = LOWER(u.login)
|
||
LEFT JOIN (
|
||
SELECT login, COUNT(*)::INTEGER AS following_users_count
|
||
FROM connections_state
|
||
WHERE rel_type = 30
|
||
AND to_block_number = 0
|
||
GROUP BY login
|
||
) fu ON LOWER(fu.login) = LOWER(u.login)
|
||
LEFT JOIN (
|
||
SELECT cs.login, COUNT(*)::INTEGER AS following_channels_count
|
||
FROM connections_state cs
|
||
JOIN channel_names_state cn
|
||
ON cn.owner_bch_name = cs.to_bch_name
|
||
AND cn.channel_root_block_number = cs.to_block_number
|
||
AND cn.channel_root_block_hash = cs.to_block_hash
|
||
WHERE cs.rel_type = 30
|
||
AND cn.channel_type_code = 1
|
||
GROUP BY cs.login
|
||
) fc ON LOWER(fc.login) = LOWER(u.login)
|
||
LEFT JOIN (
|
||
SELECT login, COUNT(*)::INTEGER AS close_friends_count
|
||
FROM connections_state
|
||
WHERE rel_type = 10
|
||
GROUP BY login
|
||
) cf ON LOWER(cf.login) = LOWER(u.login)
|
||
""")) {
|
||
ps.executeUpdate();
|
||
}
|
||
|
||
|
||
try (PreparedStatement ps = c.prepareStatement("SELECT shine_refresh_user_profile(login) FROM solana_user_pda_current")) {
|
||
ps.execute();
|
||
}
|
||
try (PreparedStatement ps = c.prepareStatement("SELECT shine_refresh_user_stats(login) FROM solana_user_pda_current")) {
|
||
ps.execute();
|
||
}
|
||
try (PreparedStatement ps = c.prepareStatement("""
|
||
UPDATE message_stats ms SET
|
||
primary_likes_count=(SELECT COUNT(*)::INTEGER FROM reactions_state rs WHERE rs.reaction_type=1 AND rs.last_sub_type=1
|
||
AND rs.to_login=ms.to_login AND rs.to_bch_name=ms.to_bch_name AND rs.to_block_number=ms.to_block_number AND rs.to_block_hash=ms.to_block_hash
|
||
AND shine_is_primary(rs.from_login)),
|
||
shining_likes_count=(SELECT COUNT(*)::INTEGER FROM reactions_state rs WHERE rs.reaction_type=1 AND rs.last_sub_type=1
|
||
AND rs.to_login=ms.to_login AND rs.to_bch_name=ms.to_bch_name AND rs.to_block_number=ms.to_block_number AND rs.to_block_hash=ms.to_block_hash
|
||
AND shine_is_primary(rs.from_login) AND shine_is_shining(rs.from_login))
|
||
""")) {
|
||
ps.executeUpdate();
|
||
}
|
||
|
||
try (PreparedStatement ps = c.prepareStatement("""
|
||
INSERT INTO channel_stats_state (
|
||
owner_bch_name,
|
||
channel_root_block_number,
|
||
channel_root_block_hash,
|
||
owner_login,
|
||
channel_type_code,
|
||
subscribers_count,
|
||
updated_at_ms
|
||
)
|
||
SELECT
|
||
cn.owner_bch_name,
|
||
cn.channel_root_block_number,
|
||
cn.channel_root_block_hash,
|
||
cn.owner_login,
|
||
cn.channel_type_code,
|
||
COUNT(DISTINCT cs.login)::INTEGER AS subscribers_count,
|
||
CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
|
||
FROM channel_names_state cn
|
||
LEFT JOIN connections_state cs
|
||
ON cs.rel_type = 30
|
||
AND cs.to_bch_name = cn.owner_bch_name
|
||
AND cs.to_block_number = cn.channel_root_block_number
|
||
AND cs.to_block_hash = cn.channel_root_block_hash
|
||
WHERE cn.channel_type_code = 1
|
||
GROUP BY
|
||
cn.owner_bch_name,
|
||
cn.channel_root_block_number,
|
||
cn.channel_root_block_hash,
|
||
cn.owner_login,
|
||
cn.channel_type_code
|
||
""")) {
|
||
ps.executeUpdate();
|
||
}
|
||
}
|
||
|
||
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
|
||
) {}
|
||
}
|