Сервер: дочистить legacy-слои и доки

This commit is contained in:
AidarKC
2026-07-27 19:29:10 +04:00
parent 0db3c3af5a
commit 6ed91a105e
33 changed files with 127 additions and 165 deletions
@@ -29,7 +29,7 @@ import java.sql.SQLException;
* отдельным recovery/resync-слоем после успешного commit.
*
* Важный смысл текущей реализации:
* - мы НЕ трогаем identity-слой (`solana_users`) и НЕ трогаем DM-таблицы;
* - мы НЕ трогаем current users слой (`solana_user_pda_current`) и НЕ трогаем DM-таблицы;
* - мы очищаем только блокчейн пользователя и derived-state, который строится из неё;
* - висячие cross-chain ссылки в чужих blocks допускаются как нормальное поведение системы.
*/
@@ -142,7 +142,7 @@ public final class BlockchainStateDAO {
* Строгая вставка state только если записи ещё нет.
*
* Нужна для recovery / resync:
* - identity пользователя уже может существовать в solana_users;
* - runtime-проекция пользователя уже может существовать в current users слое;
* - в таком случае нам надо восстановить только blockchain_state;
* - если запись уже есть, метод просто ничего не меняет.
*/
@@ -1,7 +1,7 @@
package shine.db.dao;
import shine.db.DbController;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.entities.SignedMessageEntry;
import java.sql.Connection;
import java.sql.PreparedStatement;
@@ -11,7 +11,7 @@ import java.sql.Statement;
import java.util.ArrayList;
import java.util.List;
public final class SignedMessagesV2DAO {
public final class SignedMessagesDAO {
public enum ApplyStatus {
APPLIED,
@@ -24,21 +24,21 @@ public final class SignedMessagesV2DAO {
}
}
private static volatile SignedMessagesV2DAO instance;
private static volatile SignedMessagesDAO instance;
private final DbController db = DbController.getInstance();
private SignedMessagesV2DAO() {}
private SignedMessagesDAO() {}
public static SignedMessagesV2DAO getInstance() {
public static SignedMessagesDAO getInstance() {
if (instance == null) {
synchronized (SignedMessagesV2DAO.class) {
if (instance == null) instance = new SignedMessagesV2DAO();
synchronized (SignedMessagesDAO.class) {
if (instance == null) instance = new SignedMessagesDAO();
}
}
return instance;
}
public ApplyStatus insertIfAbsent(SignedMessageV2Entry e) throws Exception {
public ApplyStatus insertIfAbsent(SignedMessageEntry e) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
if (isBlockedByConversationDelete(c, e.getFromLogin(), e.getToLogin(), e.getTimeMs())) {
@@ -65,7 +65,7 @@ public final class SignedMessagesV2DAO {
});
}
public boolean insertPairBothOrNothing(SignedMessageV2Entry first, SignedMessageV2Entry second) throws Exception {
public boolean insertPairBothOrNothing(SignedMessageEntry first, SignedMessageEntry second) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
boolean prevAutoCommit = c.getAutoCommit();
@@ -94,7 +94,7 @@ public final class SignedMessagesV2DAO {
});
}
public ApplyStatus upsertContentPair(SignedMessageV2Entry incoming, SignedMessageV2Entry outgoing) throws Exception {
public ApplyStatus upsertContentPair(SignedMessageEntry incoming, SignedMessageEntry outgoing) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
boolean prevAutoCommit = c.getAutoCommit();
@@ -135,7 +135,7 @@ public final class SignedMessagesV2DAO {
});
}
public ApplyStatus upsertIncomingCopy(SignedMessageV2Entry incoming) throws Exception {
public ApplyStatus upsertIncomingCopy(SignedMessageEntry incoming) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
boolean prevAutoCommit = c.getAutoCommit();
@@ -172,7 +172,7 @@ public final class SignedMessagesV2DAO {
});
}
public ApplyStatus applyDeleteMessage(SignedMessageV2Entry tombstone) throws Exception {
public ApplyStatus applyDeleteMessage(SignedMessageEntry tombstone) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
boolean prevAutoCommit = c.getAutoCommit();
@@ -203,7 +203,7 @@ public final class SignedMessagesV2DAO {
});
}
public ApplyStatus applyDeleteConversation(SignedMessageV2Entry tombstone) throws Exception {
public ApplyStatus applyDeleteConversation(SignedMessageEntry tombstone) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
boolean prevAutoCommit = c.getAutoCommit();
@@ -231,7 +231,7 @@ public final class SignedMessagesV2DAO {
});
}
public SignedMessageV2Entry getByMessageKey(String messageKey) throws Exception {
public SignedMessageEntry getByMessageKey(String messageKey) throws Exception {
try (Connection c = db.getConnection()) {
String sql = """
SELECT
@@ -252,7 +252,7 @@ public final class SignedMessagesV2DAO {
}
}
public SignedMessageV2Entry getLatestConversationDelete(String fromLogin, String toLogin) throws Exception {
public SignedMessageEntry getLatestConversationDelete(String fromLogin, String toLogin) throws Exception {
try (Connection c = db.getConnection()) {
return getLatestConversationDelete(c, fromLogin, toLogin);
}
@@ -314,7 +314,7 @@ public final class SignedMessagesV2DAO {
});
}
public List<SignedMessageV2Entry> listPendingForSession(String login, String sessionId) throws Exception {
public List<SignedMessageEntry> listPendingForSession(String login, String sessionId) throws Exception {
return withBusyRetry(() -> {
try (Connection c = db.getConnection()) {
String fillSql = """
@@ -354,7 +354,7 @@ public final class SignedMessagesV2DAO {
WHERE d.session_id = ? AND d.delivered = 0
ORDER BY m.time_ms ASC, m.revision_time_ms ASC, m.reencrypted_at_ms ASC, m.created_at_ms ASC
""".formatted(messagesTable());
List<SignedMessageV2Entry> out = new ArrayList<>();
List<SignedMessageEntry> out = new ArrayList<>();
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, sessionId);
try (ResultSet rs = ps.executeQuery()) {
@@ -366,7 +366,7 @@ public final class SignedMessagesV2DAO {
});
}
public List<SignedMessageV2Entry> listConversationPage(
public List<SignedMessageEntry> listConversationPage(
String login,
String peerLogin,
long beforeTimeMs,
@@ -395,7 +395,7 @@ public final class SignedMessagesV2DAO {
ORDER BY time_ms DESC, message_key DESC
LIMIT ?
""".formatted(messagesTable());
List<SignedMessageV2Entry> out = new ArrayList<>();
List<SignedMessageEntry> out = new ArrayList<>();
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, login);
ps.setString(2, login);
@@ -416,7 +416,7 @@ public final class SignedMessagesV2DAO {
}
}
private void upsertMessage(Connection c, SignedMessageV2Entry e) throws SQLException {
private void upsertMessage(Connection c, SignedMessageEntry e) throws SQLException {
String sql = """
INSERT INTO %s (
message_key, base_key, target_login, from_login, to_login,
@@ -448,7 +448,7 @@ public final class SignedMessagesV2DAO {
}
}
private void markMessageReadByReceipt(Connection c, SignedMessageV2Entry entry) throws SQLException {
private void markMessageReadByReceipt(Connection c, SignedMessageEntry entry) throws SQLException {
if (entry == null) return;
int messageType = entry.getMessageType();
if (messageType != 3 && messageType != 4) return;
@@ -552,7 +552,7 @@ public final class SignedMessagesV2DAO {
}
}
private SignedMessageV2Entry getLatestConversationDelete(Connection c, String fromLogin, String toLogin) throws Exception {
private SignedMessageEntry getLatestConversationDelete(Connection c, String fromLogin, String toLogin) throws Exception {
String sql = """
SELECT
message_key, base_key, target_login, from_login, to_login,
@@ -646,7 +646,7 @@ public final class SignedMessagesV2DAO {
}
}
private int insertStrict(Connection c, SignedMessageV2Entry e) throws SQLException {
private int insertStrict(Connection c, SignedMessageEntry e) throws SQLException {
String sql = """
INSERT INTO %s (
message_key, base_key, target_login, from_login, to_login,
@@ -661,7 +661,7 @@ public final class SignedMessagesV2DAO {
}
}
private void bindSignedMessage(PreparedStatement ps, SignedMessageV2Entry e) throws SQLException {
private void bindSignedMessage(PreparedStatement ps, SignedMessageEntry e) throws SQLException {
ps.setString(1, e.getMessageKey());
ps.setString(2, e.getBaseKey());
ps.setString(3, e.getTargetLogin());
@@ -718,8 +718,8 @@ public final class SignedMessagesV2DAO {
return "signed_messages";
}
private SignedMessageV2Entry mapRow(ResultSet rs) throws Exception {
SignedMessageV2Entry e = new SignedMessageV2Entry();
private SignedMessageEntry mapRow(ResultSet rs) throws Exception {
SignedMessageEntry e = new SignedMessageEntry();
e.setMessageKey(rs.getString("message_key"));
e.setBaseKey(rs.getString("base_key"));
e.setTargetLogin(rs.getString("target_login"));
@@ -743,7 +743,7 @@ public final class SignedMessagesV2DAO {
}
private record RevisionMarker(long revisionTimeMs, long reencryptedAtMs) {
private static RevisionMarker of(SignedMessageV2Entry entry) {
private static RevisionMarker of(SignedMessageEntry entry) {
return new RevisionMarker(entry.getRevisionTimeMs(), entry.getReencryptedAtMs());
}
}
@@ -1,6 +1,6 @@
package shine.db.entities;
public class SignedMessageV2Entry {
public class SignedMessageEntry {
private String messageKey;
private String baseKey;
private String targetLogin;
@@ -25,8 +25,8 @@ public class ConnectionContext {
public static final int AUTH_STATUS_AUTH_IN_PROGRESS = 1; // выполнен challenge (AuthChallenge или SessionChallenge)
public static final int AUTH_STATUS_USER = 2; // авторизованный пользователь
// Полный пользователь из БД (solana_users)
private CurrentUserEntry solanaUserEntry;
// Полный пользователь из runtime БД (current users / solana_user_pda_current)
private CurrentUserEntry currentUserEntry;
// Активная сессия из БД (active_sessions)
private ActiveSessionEntry activeSessionEntry;
@@ -89,12 +89,12 @@ public class ConnectionContext {
// --- SolanaUser / ActiveSession ---
public CurrentUserEntry getSolanaUser() {
return solanaUserEntry;
public CurrentUserEntry getCurrentUser() {
return currentUserEntry;
}
public void setSolanaUser(CurrentUserEntry solanaUserEntry) {
this.solanaUserEntry = solanaUserEntry;
public void setCurrentUser(CurrentUserEntry currentUserEntry) {
this.currentUserEntry = currentUserEntry;
}
public ActiveSessionEntry getActiveSession() {
@@ -108,7 +108,7 @@ public class ConnectionContext {
// --- Удобный геттер для логина ---
public String getLogin() {
return solanaUserEntry != null ? solanaUserEntry.getLogin() : null;
return currentUserEntry != null ? currentUserEntry.getLogin() : null;
}
// --- sessionId ---
@@ -176,7 +176,7 @@ public class ConnectionContext {
}
public void reset() {
solanaUserEntry = null;
currentUserEntry = null;
activeSessionEntry = null;
sessionId = null;
@@ -198,4 +198,4 @@ public class ConnectionContext {
", authenticationStatus=" + authenticationStatus +
'}';
}
}
}
@@ -70,7 +70,7 @@ public class Net_AuthChallenge_Handler implements JsonMessageHandler {
);
}
ctx.setSolanaUser(solanaUserEntry);
ctx.setCurrentUser(solanaUserEntry);
ctx.setAuthenticationStatus(ConnectionContext.AUTH_STATUS_AUTH_IN_PROGRESS);
byte[] buf = new byte[32];
@@ -42,7 +42,7 @@ public class Net_CloseActiveSession_Handler implements JsonMessageHandler {
public Net_Response handle(Net_Request baseReq, ConnectionContext ctx) throws Exception {
Net_CloseActiveSession_Request req = (Net_CloseActiveSession_Request) baseReq;
if (ctx == null || ctx.getSolanaUser() == null || ctx.getAuthenticationStatus() != ConnectionContext.AUTH_STATUS_USER) {
if (ctx == null || ctx.getCurrentUser() == null || ctx.getAuthenticationStatus() != ConnectionContext.AUTH_STATUS_USER) {
return NetExceptionResponseFactory.error(
req,
WireCodes.Status.UNVERIFIED,
@@ -51,7 +51,7 @@ public class Net_CloseActiveSession_Handler implements JsonMessageHandler {
);
}
CurrentUserEntry user = ctx.getSolanaUser();
CurrentUserEntry user = ctx.getCurrentUser();
String currentLogin = user.getLogin();
String targetSessionId = req.getSessionId();
@@ -152,4 +152,4 @@ public class Net_CloseActiveSession_Handler implements JsonMessageHandler {
);
}
}
}
}
@@ -58,7 +58,7 @@ public class Net_CreateAuthSession__Handler implements JsonMessageHandler {
Net_CreateAuthSession_Request req = (Net_CreateAuthSession_Request) baseReq;
if (ctx == null
|| ctx.getSolanaUser() == null
|| ctx.getCurrentUser() == null
|| ctx.getAuthNonce() == null
|| ctx.getAuthenticationStatus() != ConnectionContext.AUTH_STATUS_AUTH_IN_PROGRESS) {
@@ -72,7 +72,7 @@ public class Net_CreateAuthSession__Handler implements JsonMessageHandler {
return err;
}
CurrentUserEntry userFromContext = ctx.getSolanaUser();
CurrentUserEntry userFromContext = ctx.getCurrentUser();
String loginFromContext = userFromContext.getLogin();
String loginFromReq = req.getLogin();
if (loginFromReq == null || loginFromReq.isBlank()) {
@@ -36,7 +36,7 @@ public class Net_ListSessions_Handler implements JsonMessageHandler {
public Net_Response handle(Net_Request baseReq, ConnectionContext ctx) throws Exception {
Net_ListSessions_Request req = (Net_ListSessions_Request) baseReq;
if (ctx == null || ctx.getSolanaUser() == null || ctx.getAuthenticationStatus() != ConnectionContext.AUTH_STATUS_USER) {
if (ctx == null || ctx.getCurrentUser() == null || ctx.getAuthenticationStatus() != ConnectionContext.AUTH_STATUS_USER) {
return NetExceptionResponseFactory.error(
req,
WireCodes.Status.UNVERIFIED,
@@ -45,7 +45,7 @@ public class Net_ListSessions_Handler implements JsonMessageHandler {
);
}
CurrentUserEntry user = ctx.getSolanaUser();
CurrentUserEntry user = ctx.getCurrentUser();
String currentLogin = user.getLogin();
List<ActiveSessionEntry> sessions;
@@ -294,7 +294,7 @@ public class Net_SessionLogin_Handler implements JsonMessageHandler {
// ctx
ctx.setActiveSession(session);
ctx.setSolanaUser(user);
ctx.setCurrentUser(user);
ctx.setSessionId(sessionId);
ctx.setAuthenticationStatus(ConnectionContext.AUTH_STATUS_USER);
@@ -54,7 +54,7 @@ public class Net_GetFriendsLists_Handler implements JsonMessageHandler {
try (Connection c = db.getConnection()) {
// 1) Канонизируем login через solana_users (NOCASE)
// 1) Канонизируем login через current users слой (NOCASE)
String canonicalLogin = findCanonicalLogin(c, loginAnyCase);
if (canonicalLogin == null) {
return NetExceptionResponseFactory.error(
@@ -17,8 +17,8 @@ import shine.db.entities.CurrentUserEntry;
/**
* GetSyncUserProfile server-to-server профиль пользователя для межсерверной синхронизации.
* Нужен, чтобы принимающий сервер мог создать локальные solana_users + blockchain_state
* без прямого запроса в Solana RPC.
* Нужен, чтобы принимающий сервер мог создать локальную runtime-проекцию пользователя
* и blockchain_state без прямого запроса в Solana RPC.
*/
public final class Net_GetSyncUserProfile_Handler implements JsonMessageHandler {
@@ -8,7 +8,7 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_AckSessionDelivery_Re
import server.logic.ws_protocol.JSON.messages.entyties.Net_AckSessionDelivery_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.dao.SignedMessagesDAO;
public class Net_AckSessionDelivery_Handler implements JsonMessageHandler {
@Override
@@ -22,7 +22,7 @@ public class Net_AckSessionDelivery_Handler implements JsonMessageHandler {
}
String messageKey = req.getMessageKey().trim();
SignedMessagesV2DAO.getInstance().markDelivered(messageKey, ctx.getSessionId(), System.currentTimeMillis());
SignedMessagesDAO.getInstance().markDelivered(messageKey, ctx.getSessionId(), System.currentTimeMillis());
Net_AckSessionDelivery_Response resp = new Net_AckSessionDelivery_Response();
resp.setOp(req.getOp());
@@ -9,8 +9,8 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteConversation_Re
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import server.sync.DmFederationService;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.SignedMessageEntry;
public class Net_DeleteConversation_Handler implements JsonMessageHandler {
@Override
@@ -38,8 +38,8 @@ public class Net_DeleteConversation_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, status, code, "Сообщение не прошло проверку");
}
SignedMessageV2Entry entry = SignedMessagesCore.toEntry(block, "DeleteConversation", null);
SignedMessagesV2DAO.ApplyStatus status = SignedMessagesV2DAO.getInstance().applyDeleteConversation(entry);
SignedMessageEntry entry = SignedMessagesCore.toEntry(block, "DeleteConversation", null);
SignedMessagesDAO.ApplyStatus status = SignedMessagesDAO.getInstance().applyDeleteConversation(entry);
SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters();
if (status.applied()) {
@@ -9,8 +9,8 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteMessage_Respons
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import server.sync.DmFederationService;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.SignedMessageEntry;
public class Net_DeleteMessage_Handler implements JsonMessageHandler {
@Override
@@ -38,8 +38,8 @@ public class Net_DeleteMessage_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, status, code, "Сообщение не прошло проверку");
}
SignedMessageV2Entry entry = SignedMessagesCore.toEntry(block, "DeleteMessage", null);
SignedMessagesV2DAO.ApplyStatus status = SignedMessagesV2DAO.getInstance().applyDeleteMessage(entry);
SignedMessageEntry entry = SignedMessagesCore.toEntry(block, "DeleteMessage", null);
SignedMessagesDAO.ApplyStatus status = SignedMessagesDAO.getInstance().applyDeleteMessage(entry);
SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters();
if (status.applied()) {
@@ -10,8 +10,8 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Req
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.SignedMessageEntry;
import java.util.ArrayList;
import java.util.Base64;
@@ -43,7 +43,7 @@ public class Net_GetDirectMessages_Handler implements JsonMessageHandler {
String beforeMessageKey = req.getBeforeMessageKey() == null ? "" : req.getBeforeMessageKey().trim();
try {
List<SignedMessageV2Entry> page = SignedMessagesV2DAO.getInstance().listConversationPage(
List<SignedMessageEntry> page = SignedMessagesDAO.getInstance().listConversationPage(
login,
peerLogin,
beforeTimeMs,
@@ -66,7 +66,7 @@ public class Net_GetDirectMessages_Handler implements JsonMessageHandler {
resp.setHasMore(hasMore);
List<Net_GetDirectMessages_Response.MessageItem> items = new ArrayList<>();
for (SignedMessageV2Entry entry : page) {
for (SignedMessageEntry entry : page) {
Net_GetDirectMessages_Response.MessageItem item = new Net_GetDirectMessages_Response.MessageItem();
item.setMessageKey(entry.getMessageKey());
item.setBaseKey(entry.getBaseKey());
@@ -8,8 +8,8 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_ReceiveIncomingMessag
import server.logic.ws_protocol.JSON.messages.entyties.Net_ReceiveIncomingMessage_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.SignedMessageEntry;
import java.util.Base64;
@@ -40,21 +40,21 @@ public class Net_ReceiveIncomingMessage_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, status, code, "Сообщение не прошло проверку");
}
final SignedMessageV2Entry entry;
final SignedMessageEntry entry;
try {
entry = SignedMessagesCore.toEntry(incoming, "ReceiveIncomingMessage", null);
} catch (IllegalArgumentException ex) {
return NetExceptionResponseFactory.error(req, WireCodes.Status.BAD_REQUEST, ex.getMessage(), "Некорректный payload подтверждения");
}
SignedMessagesV2DAO.ApplyStatus status = incoming.isContentType()
? SignedMessagesV2DAO.getInstance().upsertIncomingCopy(entry)
: SignedMessagesV2DAO.getInstance().insertIfAbsent(entry);
SignedMessagesDAO.ApplyStatus status = incoming.isContentType()
? SignedMessagesDAO.getInstance().upsertIncomingCopy(entry)
: SignedMessagesDAO.getInstance().insertIfAbsent(entry);
SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters();
if (status.applied()) {
counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, incoming);
}
if (status == SignedMessagesV2DAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) {
if (status == SignedMessagesDAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) {
bounceConversationDeleteIfKnown(incoming.fromLogin, incoming.toLogin);
}
@@ -74,7 +74,7 @@ public class Net_ReceiveIncomingMessage_Handler implements JsonMessageHandler {
}
private void bounceConversationDeleteIfKnown(String fromLogin, String toLogin) throws Exception {
SignedMessageV2Entry tombstone = SignedMessagesV2DAO.getInstance().getLatestConversationDelete(fromLogin, toLogin);
SignedMessageEntry tombstone = SignedMessagesDAO.getInstance().getLatestConversationDelete(fromLogin, toLogin);
if (tombstone == null || tombstone.getRawBlock() == null || tombstone.getRawBlock().length == 0) return;
server.sync.DmFederationService.fanOutDeleteConversation(
tombstone.getFromLogin(),
@@ -9,8 +9,8 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_SendMessagePair_Respo
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import server.sync.DmFederationService;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.SignedMessageEntry;
import java.util.Base64;
@@ -41,8 +41,8 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, status, code, "Сообщение не прошло проверку");
}
SignedMessageV2Entry incomingEntry;
SignedMessageV2Entry outgoingEntry;
SignedMessageEntry incomingEntry;
SignedMessageEntry outgoingEntry;
try {
String sourceApi = "SendMessagePair";
String originSessionId = (ctx != null && !isBlank(ctx.getSessionId())) ? ctx.getSessionId() : null;
@@ -52,15 +52,15 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, WireCodes.Status.BAD_REQUEST, ex.getMessage(), "Некорректный payload подтверждения");
}
SignedMessagesV2DAO.ApplyStatus pairStatus;
SignedMessagesDAO.ApplyStatus pairStatus;
if (incoming.isContentType()) {
pairStatus = SignedMessagesV2DAO.getInstance().upsertContentPair(
pairStatus = SignedMessagesDAO.getInstance().upsertContentPair(
incomingEntry, outgoingEntry
);
} else {
pairStatus = SignedMessagesV2DAO.getInstance().insertPairBothOrNothing(incomingEntry, outgoingEntry)
? SignedMessagesV2DAO.ApplyStatus.APPLIED
: SignedMessagesV2DAO.ApplyStatus.DUPLICATE_OR_OLDER;
pairStatus = SignedMessagesDAO.getInstance().insertPairBothOrNothing(incomingEntry, outgoingEntry)
? SignedMessagesDAO.ApplyStatus.APPLIED
: SignedMessagesDAO.ApplyStatus.DUPLICATE_OR_OLDER;
}
SignedMessagesRealtime.DeliveryCounters inCounters = new SignedMessagesRealtime.DeliveryCounters();
@@ -77,7 +77,7 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
outCounters = SignedMessagesRealtime.deliverToRelevantSessions(outgoingEntry, outgoing, excludeSessionId);
}
if (pairStatus == SignedMessagesV2DAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) {
if (pairStatus == SignedMessagesDAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) {
bounceConversationDeleteIfKnown(incoming.fromLogin, incoming.toLogin);
}
@@ -107,7 +107,7 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
}
private void bounceConversationDeleteIfKnown(String fromLogin, String toLogin) throws Exception {
SignedMessageV2Entry tombstone = SignedMessagesV2DAO.getInstance().getLatestConversationDelete(fromLogin, toLogin);
SignedMessageEntry tombstone = SignedMessagesDAO.getInstance().getLatestConversationDelete(fromLogin, toLogin);
if (tombstone == null || tombstone.getRawBlock() == null || tombstone.getRawBlock().length == 0) return;
DmFederationService.fanOutDeleteConversation(
tombstone.getFromLogin(),
@@ -74,7 +74,7 @@ public class Net_SendSignal_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, WireCodes.Status.BAD_REQUEST, "TIME_SKEW", "Время клиента отличается от сервера более чем на 30 секунд");
}
CurrentUserEntry senderUser = ctx.getSolanaUser();
CurrentUserEntry senderUser = ctx.getCurrentUser();
if (senderUser == null || senderUser.getClientKey() == null || senderUser.getClientKey().isBlank()) {
return NetExceptionResponseFactory.error(req, WireCodes.Status.SERVER_DATA_ERROR, "NO_CLIENT_KEY", "Для пользователя не найден client key");
}
@@ -1,7 +1,7 @@
package server.logic.ws_protocol.JSON.messages;
import shine.db.dao.CurrentUsersDAO;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.entities.SignedMessageEntry;
import shine.db.entities.CurrentUserEntry;
import utils.crypto.Ed25519Util;
@@ -96,11 +96,11 @@ final class SignedMessagesCore {
}
}
static SignedMessageV2Entry toEntry(SignedMessageBlock block, String sourceApi, String originSessionId) {
static SignedMessageEntry toEntry(SignedMessageBlock block, String sourceApi, String originSessionId) {
String baseKey = SignedMessageKeys.baseKey(block.toLogin, block.fromLogin, block.timeMs, block.nonce);
String messageKey = SignedMessageKeys.messageKey(block.toLogin, block.fromLogin, block.timeMs, block.nonce, block.messageType);
SignedMessageV2Entry entry = new SignedMessageV2Entry();
SignedMessageEntry entry = new SignedMessageEntry();
entry.setMessageKey(messageKey);
entry.setBaseKey(baseKey);
entry.setTargetLogin(primaryTargetLogin(block));
@@ -9,9 +9,9 @@ import server.logic.ws_protocol.JSON.ConnectionContext;
import server.logic.ws_protocol.JSON.push.WebPushSender;
import server.logic.ws_protocol.JSON.push.WsEventSender;
import shine.db.dao.ActiveSessionsDAO;
import shine.db.dao.SignedMessagesV2DAO;
import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.ActiveSessionEntry;
import shine.db.entities.SignedMessageV2Entry;
import shine.db.entities.SignedMessageEntry;
import java.util.ArrayList;
import java.util.Base64;
@@ -37,12 +37,12 @@ public final class SignedMessagesRealtime {
private SignedMessagesRealtime() {}
static DeliveryCounters deliverToRelevantSessions(SignedMessageV2Entry message, SignedMessageBlock block) throws Exception {
static DeliveryCounters deliverToRelevantSessions(SignedMessageEntry message, SignedMessageBlock block) throws Exception {
return deliverToRelevantSessions(message, block, null);
}
static DeliveryCounters deliverToRelevantSessions(
SignedMessageV2Entry message,
SignedMessageEntry message,
SignedMessageBlock block,
String excludeSessionId
) throws Exception {
@@ -58,7 +58,7 @@ public final class SignedMessagesRealtime {
}
sessionIdsToTrack.add(sessionId);
}
SignedMessagesV2DAO.getInstance().ensureDeliveryRows(message.getMessageKey(), sessionIdsToTrack, now);
SignedMessagesDAO.getInstance().ensureDeliveryRows(message.getMessageKey(), sessionIdsToTrack, now);
for (ActiveSessionEntry s : sessions) {
String sessionId = s.getSessionId();
if (excludeSessionId != null && excludeSessionId.equals(sessionId)) {
@@ -97,9 +97,9 @@ public final class SignedMessagesRealtime {
private static void dispatchPendingForSession(String login, String sessionId) {
try {
List<SignedMessageV2Entry> pending = SignedMessagesV2DAO.getInstance()
List<SignedMessageEntry> pending = SignedMessagesDAO.getInstance()
.listPendingForSession(login, sessionId);
for (SignedMessageV2Entry e : pending) {
for (SignedMessageEntry e : pending) {
sendEventToSessionIfOnline(sessionId, login, e, true);
}
} catch (Exception e) {
@@ -110,7 +110,7 @@ public final class SignedMessagesRealtime {
private static boolean sendEventToSessionIfOnline(
String sessionId,
String actualTargetLogin,
SignedMessageV2Entry message,
SignedMessageEntry message,
boolean backlog
) {
ConnectionContext targetCtx = ActiveConnectionsRegistry.getInstance().getBySessionId(sessionId);
@@ -134,14 +134,14 @@ public final class SignedMessagesRealtime {
return WsEventSender.sendEvent(targetCtx, "SignedMessageArrived", message.getMessageKey(), payload);
}
private static boolean shouldPushNewIncomingMessage(String targetLogin, SignedMessageV2Entry message, SignedMessageBlock block) {
private static boolean shouldPushNewIncomingMessage(String targetLogin, SignedMessageEntry message, SignedMessageBlock block) {
if (block == null) return false;
if (message.getMessageType() != SignedMessageBlock.TYPE_INCOMING_TEXT) return false;
if (!targetLogin.equalsIgnoreCase(message.getToLogin())) return false;
return block.revisionTimeMs == 0;
}
private static boolean pushNewMessageNotification(ActiveSessionEntry session, SignedMessageV2Entry message) {
private static boolean pushNewMessageNotification(ActiveSessionEntry session, SignedMessageEntry message) {
try {
if (session == null) return false;
if (isBlank(session.getPushEndpoint()) || isBlank(session.getPushP256dhKey()) || isBlank(session.getPushAuthKey())) {
@@ -160,7 +160,7 @@ public final class SignedMessagesRealtime {
}
}
private static List<String> targetLoginsForMessage(SignedMessageV2Entry message) {
private static List<String> targetLoginsForMessage(SignedMessageEntry message) {
Set<String> out = new LinkedHashSet<>();
int type = message.getMessageType();
if (type == SignedMessageBlock.TYPE_INCOMING_TEXT || type == SignedMessageBlock.TYPE_READ_INCOMING) {
@@ -57,7 +57,7 @@ public final class WsConnectionUtils {
final String sessionId = safeString(ctx.getSessionId());
final int authStatus = safeAuthStatus(ctx);
final CurrentUserEntry user = ctx.getSolanaUser();
final CurrentUserEntry user = ctx.getCurrentUser();
final String login = (user != null ? safeString(user.getLogin()) : "");
final String activeSessionId =
@@ -152,4 +152,4 @@ public final class WsConnectionUtils {
return "remote=" + remote + ", local=" + local;
}
}
}
@@ -20,7 +20,7 @@ solana.users.sync.pollIntervalSeconds=300
# false - брать профиль пользователя напрямую из Solana PDA (обычный режим).
# true - не ходить в Solana RPC, а запрашивать у сервера-партнёра специальный
# sync-профиль пользователя и по нему локально создавать
# solana_users + blockchain_state.
# current user runtime-проекцию + blockchain_state.
# Эта настройка нужна как временный обход лимитов Solana RPC (например 429),
# чтобы чистый сервер мог восстановить цепочки от партнёра без зависимости
# от внешнего Solana endpoint.