diff --git a/SHiNE-server/shine-server-db/src/main/java/shine/db/DatabaseInitializer.java b/SHiNE-server/shine-server-db/src/main/java/shine/db/DatabaseInitializer.java index 73975a76..4586fe81 100644 --- a/SHiNE-server/shine-server-db/src/main/java/shine/db/DatabaseInitializer.java +++ b/SHiNE-server/shine-server-db/src/main/java/shine/db/DatabaseInitializer.java @@ -26,6 +26,7 @@ public final class DatabaseInitializer { public static final int SCHEMA_VERSION_7 = 7; public static final int SCHEMA_VERSION_8 = 8; public static final int SCHEMA_VERSION_9 = 9; + public static final int SCHEMA_VERSION_10 = 10; public static final String POSTGRES_SCHEMA_RESOURCE = "postgres/schema_v1.sql"; public static final String POSTGRES_MIGRATION_V2_RESOURCE = "postgres/migration_v2.sql"; public static final String POSTGRES_MIGRATION_V3_RESOURCE = "postgres/migration_v3.sql"; @@ -35,6 +36,7 @@ public final class DatabaseInitializer { public static final String POSTGRES_MIGRATION_V7_RESOURCE = "postgres/migration_v7.sql"; public static final String POSTGRES_MIGRATION_V8_RESOURCE = "postgres/migration_v8.sql"; public static final String POSTGRES_MIGRATION_V9_RESOURCE = "postgres/migration_v9.sql"; + public static final String POSTGRES_MIGRATION_V10_RESOURCE = "postgres/migration_v10.sql"; private DatabaseInitializer() {} @@ -130,6 +132,10 @@ public final class DatabaseInitializer { } if (currentVersion < SCHEMA_VERSION_9) { runSqlScript(conn, POSTGRES_MIGRATION_V9_RESOURCE); + currentVersion = SCHEMA_VERSION_9; + } + if (currentVersion < SCHEMA_VERSION_10) { + runSqlScript(conn, POSTGRES_MIGRATION_V10_RESOURCE); } } } diff --git a/SHiNE-server/shine-server-db/src/main/java/shine/db/PostgresDbController.java b/SHiNE-server/shine-server-db/src/main/java/shine/db/PostgresDbController.java index a3b77fcc..680af9ad 100644 --- a/SHiNE-server/shine-server-db/src/main/java/shine/db/PostgresDbController.java +++ b/SHiNE-server/shine-server-db/src/main/java/shine/db/PostgresDbController.java @@ -1,6 +1,7 @@ package shine.db; import shine.db.connection.DriverManagerDbProvider; +import shine.db.dao.DmDialogStateDAO; import utils.config.AppConfig; import java.sql.Connection; @@ -34,6 +35,12 @@ public final class PostgresDbController { this.delegate = new DriverManagerDbProvider(jdbcUrl, dbUser, dbPassword, connection -> { connection.setAutoCommit(true); }); + + try (Connection connection = this.delegate.getConnection()) { + DmDialogStateDAO.getInstance().bootstrapIfEmpty(connection); + } catch (SQLException e) { + throw new RuntimeException("DM dialog state bootstrap failed", e); + } } public static PostgresDbController getInstance() { diff --git a/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/ConnectionsStateDAO.java b/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/ConnectionsStateDAO.java index e01fdf52..dc6c01a6 100644 --- a/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/ConnectionsStateDAO.java +++ b/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/ConnectionsStateDAO.java @@ -67,6 +67,36 @@ public final class ConnectionsStateDAO { return out; } + public boolean hasOutgoingByRelTypeCanonical(Connection c, String loginAnyCase, String peerLoginAnyCase, int relType) throws SQLException { + String sql = """ + SELECT 1 + FROM connections_state cs + LEFT JOIN %s + ON LOWER(u_login.login) = LOWER(cs.to_login) + LEFT JOIN %s + ON LOWER(u_bch.blockchain_name) = LOWER(cs.to_bch_name) + WHERE LOWER(cs.login) = LOWER(?) + AND cs.rel_type = ? + AND ( + LOWER(cs.to_login) = LOWER(?) + OR LOWER(COALESCE(u_login.login, u_bch.login, cs.to_login)) = LOWER(?) + ) + LIMIT 1 + """.formatted( + CurrentUsersSql.usersSubquery("u_login"), + CurrentUsersSql.usersSubquery("u_bch") + ); + try (PreparedStatement ps = c.prepareStatement(sql)) { + ps.setString(1, loginAnyCase); + ps.setInt(2, relType); + ps.setString(3, peerLoginAnyCase); + ps.setString(4, peerLoginAnyCase); + try (ResultSet rs = ps.executeQuery()) { + return rs.next(); + } + } + } + /** * Incoming: список логинов (канонических), кто поставил relType пользователю login. */ diff --git a/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/DmDialogStateDAO.java b/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/DmDialogStateDAO.java new file mode 100644 index 00000000..fe417830 --- /dev/null +++ b/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/DmDialogStateDAO.java @@ -0,0 +1,453 @@ +package shine.db.dao; + +import shine.db.MsgSubType; +import shine.db.entities.SignedMessageEntry; + +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Statement; +import java.util.Base64; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; + +public final class DmDialogStateDAO { + private static volatile DmDialogStateDAO instance; + + private DmDialogStateDAO() {} + + public static DmDialogStateDAO getInstance() { + if (instance == null) { + synchronized (DmDialogStateDAO.class) { + if (instance == null) instance = new DmDialogStateDAO(); + } + } + return instance; + } + + public void bootstrapIfEmpty(Connection c) throws SQLException { + if (c == null) return; + if (hasAnyRow(c)) return; + + boolean prevAutoCommit = c.getAutoCommit(); + c.setAutoCommit(false); + try { + for (PairItem pair : listConversationPairs(c)) { + refreshConversationPair(c, pair.ownerLogin(), pair.peerLogin()); + } + c.commit(); + } catch (Exception e) { + try { c.rollback(); } catch (Exception ignored) {} + if (e instanceof SQLException sqlEx) throw sqlEx; + throw new SQLException("Failed to bootstrap dm_dialog_state", e); + } finally { + c.setAutoCommit(prevAutoCommit); + } + } + + public void refreshConversationPair(Connection c, String ownerLogin, String peerLogin) throws SQLException { + String cleanOwner = normalize(ownerLogin); + String cleanPeer = normalize(peerLogin); + if (cleanOwner.isEmpty() || cleanPeer.isEmpty()) return; + if (cleanOwner.equalsIgnoreCase(cleanPeer)) return; + refreshConversation(c, cleanOwner, cleanPeer); + refreshConversation(c, cleanPeer, cleanOwner); + } + + public List listInboxDialogs(Connection c, String ownerLogin) throws SQLException { + String cleanOwner = normalize(ownerLogin); + if (cleanOwner.isEmpty()) return List.of(); + + Map byPeer = new LinkedHashMap<>(); + List contacts = ConnectionsStateDAO.getInstance().listOutgoingByRelTypeCanonical(c, cleanOwner, MsgSubType.CONNECTION_CONTACT); + List closeFriends = ConnectionsStateDAO.getInstance().listOutgoingByRelTypeCanonical(c, cleanOwner, MsgSubType.CONNECTION_CLOSE_FRIEND); + + try (PreparedStatement ps = c.prepareStatement(""" + SELECT owner_login, peer_login, relation_flag, last_message_blob_b64, + last_message_time_ms, unread_count, last_read_receipt_time_ms, + updated_at_ms + FROM dm_dialog_state + WHERE LOWER(owner_login) = LOWER(?) + """)) { + ps.setString(1, cleanOwner); + try (ResultSet rs = ps.executeQuery()) { + while (rs.next()) { + DialogSummary row = new DialogSummary( + rs.getString("owner_login"), + rs.getString("peer_login"), + normalizeRelationFlag(rs.getString("relation_flag")), + rs.getString("last_message_blob_b64"), + rs.getLong("last_message_time_ms"), + rs.getInt("unread_count"), + rs.getLong("last_read_receipt_time_ms"), + rs.getLong("updated_at_ms"), + true + ); + byPeer.put(normKey(row.peerLogin()), row); + } + } + } + + for (String peer : contacts) { + addOrUpdateSummary(byPeer, cleanOwner, peer, "contact", false); + } + for (String peer : closeFriends) { + addOrUpdateSummary(byPeer, cleanOwner, peer, "close_friend", false); + } + + for (DialogSummary row : new ArrayList<>(byPeer.values())) { + String relationFlag = resolveRelationFlag(c, cleanOwner, row.peerLogin()); + byPeer.put(normKey(row.peerLogin()), new DialogSummary( + row.ownerLogin(), + row.peerLogin(), + relationFlag, + row.lastMessageBlobB64(), + row.lastMessageTimeMs(), + row.unreadCount(), + row.lastReadReceiptTimeMs(), + row.updatedAtMs(), + row.hasDialog() + )); + } + + List out = new ArrayList<>(byPeer.values()); + out.sort((a, b) -> { + int cmp = Long.compare(b.lastMessageTimeMs(), a.lastMessageTimeMs()); + if (cmp != 0) return cmp; + return normalize(a.peerLogin()).compareToIgnoreCase(normalize(b.peerLogin())); + }); + return out; + } + + public void refreshFromEntry(Connection c, SignedMessageEntry entry) throws SQLException { + if (entry == null) return; + refreshConversationPair(c, entry.getFromLogin(), entry.getToLogin()); + } + + private void refreshConversation(Connection c, String ownerLogin, String peerLogin) throws SQLException { + DialogSummary summary = loadConversationSummary(c, ownerLogin, peerLogin); + upsert(c, summary); + if (summary.lastReadReceiptTimeMs() > 0) { + syncMessagesToWatermark(c, summary.ownerLogin(), summary.peerLogin(), summary.lastReadReceiptTimeMs()); + } + } + + private DialogSummary loadConversationSummary(Connection c, String ownerLogin, String peerLogin) throws SQLException { + String cleanOwner = normalize(ownerLogin); + String cleanPeer = normalize(peerLogin); + if (cleanOwner.isEmpty() || cleanPeer.isEmpty() || cleanOwner.equalsIgnoreCase(cleanPeer)) { + return new DialogSummary(cleanOwner, cleanPeer, "none", "", 0L, 0, 0L, System.currentTimeMillis(), false); + } + + String relationFlag = resolveRelationFlag(c, cleanOwner, cleanPeer); + long existingWatermark = loadExistingWatermark(c, cleanOwner, cleanPeer); + + String latestSql = """ + SELECT raw_block, time_ms + FROM signed_messages + WHERE ( + (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + OR (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + ) + AND message_type IN (1, 2) + ORDER BY time_ms DESC, revision_time_ms DESC, reencrypted_at_ms DESC, created_at_ms DESC, message_key DESC + LIMIT 1 + """; + String unreadSql = """ + SELECT COUNT(*) + FROM signed_messages + WHERE LOWER(from_login) = LOWER(?) + AND LOWER(to_login) = LOWER(?) + AND message_type = 1 + AND time_ms > ? + AND (read_at_ms IS NULL OR read_at_ms <= 0) + """; + String contentReadSql = """ + SELECT MAX(read_at_ms) + FROM signed_messages + WHERE ( + (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + OR (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + ) + AND message_type IN (1, 2) + AND read_at_ms IS NOT NULL + AND read_at_ms > 0 + """; + String receiptWatermarkSql = """ + SELECT MAX(time_ms) + FROM signed_messages + WHERE ( + (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + OR (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + ) + AND message_type IN (3, 4) + """; + + String lastMessageBlobB64 = ""; + long lastMessageTimeMs = 0L; + int unreadCount = 0; + long lastReadReceiptTimeMs = existingWatermark; + + try (PreparedStatement ps = c.prepareStatement(latestSql)) { + ps.setString(1, cleanOwner); + ps.setString(2, cleanPeer); + ps.setString(3, cleanPeer); + ps.setString(4, cleanOwner); + try (ResultSet rs = ps.executeQuery()) { + if (rs.next()) { + lastMessageTimeMs = rs.getLong("time_ms"); + byte[] rawBlock = rs.getBytes("raw_block"); + if (rawBlock != null && rawBlock.length > 0) { + lastMessageBlobB64 = Base64.getEncoder().encodeToString(rawBlock); + } + } + } + } + + try (PreparedStatement ps = c.prepareStatement(contentReadSql)) { + ps.setString(1, cleanOwner); + ps.setString(2, cleanPeer); + ps.setString(3, cleanPeer); + ps.setString(4, cleanOwner); + try (ResultSet rs = ps.executeQuery()) { + if (rs.next()) { + long value = rs.getLong(1); + if (!rs.wasNull()) lastReadReceiptTimeMs = Math.max(lastReadReceiptTimeMs, value); + } + } + } + + try (PreparedStatement ps = c.prepareStatement(receiptWatermarkSql)) { + ps.setString(1, cleanOwner); + ps.setString(2, cleanPeer); + ps.setString(3, cleanPeer); + ps.setString(4, cleanOwner); + try (ResultSet rs = ps.executeQuery()) { + if (rs.next()) { + long value = rs.getLong(1); + if (!rs.wasNull()) lastReadReceiptTimeMs = Math.max(lastReadReceiptTimeMs, value); + } + } + } + + try (PreparedStatement ps = c.prepareStatement(unreadSql)) { + ps.setString(1, cleanPeer); + ps.setString(2, cleanOwner); + ps.setLong(3, lastReadReceiptTimeMs); + try (ResultSet rs = ps.executeQuery()) { + if (rs.next()) unreadCount = rs.getInt(1); + } + } + + return new DialogSummary( + cleanOwner, + cleanPeer, + relationFlag, + lastMessageBlobB64, + lastMessageTimeMs, + unreadCount, + lastReadReceiptTimeMs, + System.currentTimeMillis(), + lastMessageTimeMs > 0 || unreadCount > 0 || lastReadReceiptTimeMs > 0 + ); + } + + private void upsert(Connection c, DialogSummary summary) throws SQLException { + try (PreparedStatement ps = c.prepareStatement(""" + INSERT INTO dm_dialog_state ( + owner_login, peer_login, relation_flag, last_message_blob_b64, + last_message_time_ms, unread_count, last_read_receipt_time_ms, updated_at_ms + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (owner_login, peer_login) DO UPDATE SET + relation_flag = EXCLUDED.relation_flag, + last_message_blob_b64 = EXCLUDED.last_message_blob_b64, + last_message_time_ms = EXCLUDED.last_message_time_ms, + unread_count = EXCLUDED.unread_count, + last_read_receipt_time_ms = EXCLUDED.last_read_receipt_time_ms, + updated_at_ms = EXCLUDED.updated_at_ms + """)) { + ps.setString(1, summary.ownerLogin()); + ps.setString(2, summary.peerLogin()); + ps.setString(3, normalizeRelationFlag(summary.relationFlag())); + ps.setString(4, summary.lastMessageBlobB64()); + ps.setLong(5, summary.lastMessageTimeMs()); + ps.setInt(6, summary.unreadCount()); + ps.setLong(7, summary.lastReadReceiptTimeMs()); + ps.setLong(8, summary.updatedAtMs()); + ps.executeUpdate(); + } + } + + private String resolveRelationFlag(Connection c, String ownerLogin, String peerLogin) throws SQLException { + if (ConnectionsStateDAO.getInstance().hasOutgoingByRelTypeCanonical(c, ownerLogin, peerLogin, MsgSubType.CONNECTION_CLOSE_FRIEND)) { + return "close_friend"; + } + if (ConnectionsStateDAO.getInstance().hasOutgoingByRelTypeCanonical(c, ownerLogin, peerLogin, MsgSubType.CONNECTION_CONTACT)) { + return "contact"; + } + return "none"; + } + + private boolean hasAnyRow(Connection c) throws SQLException { + try (Statement st = c.createStatement(); + ResultSet rs = st.executeQuery("SELECT 1 FROM dm_dialog_state LIMIT 1")) { + return rs.next(); + } + } + + private List listConversationPairs(Connection c) throws SQLException { + String sql = """ + SELECT DISTINCT owner_login, peer_login + FROM ( + SELECT + m.target_login AS owner_login, + CASE + WHEN LOWER(m.target_login) = LOWER(m.to_login) THEN m.from_login + ELSE m.to_login + END AS peer_login + FROM signed_messages m + WHERE m.message_type IN (1, 2, 3, 4, 5, 6, 7, 8) + + UNION ALL + + SELECT m.from_login AS owner_login, m.to_login AS peer_login + FROM signed_messages m + WHERE m.message_type IN (5, 6, 7, 8) + + UNION ALL + + SELECT m.to_login AS owner_login, m.from_login AS peer_login + FROM signed_messages m + WHERE m.message_type IN (5, 6, 7, 8) + ) pairs + WHERE owner_login IS NOT NULL + AND peer_login IS NOT NULL + AND BTRIM(owner_login) <> '' + AND BTRIM(peer_login) <> '' + AND LOWER(owner_login) <> LOWER(peer_login) + ORDER BY owner_login, peer_login + """; + List out = new ArrayList<>(); + try (Statement st = c.createStatement(); + ResultSet rs = st.executeQuery(sql)) { + while (rs.next()) { + String ownerLogin = rs.getString("owner_login"); + String peerLogin = rs.getString("peer_login"); + if (normalize(ownerLogin).isEmpty() || normalize(peerLogin).isEmpty()) continue; + out.add(new PairItem(ownerLogin, peerLogin)); + } + } + return out; + } + + private void addOrUpdateSummary(Map map, String ownerLogin, String peerLogin, String relationFlag, boolean hasDialog) { + String owner = normalize(ownerLogin); + String peer = normalize(peerLogin); + if (owner.isEmpty() || peer.isEmpty() || owner.equalsIgnoreCase(peer)) return; + String key = normKey(peer); + DialogSummary current = map.get(key); + if (current == null) { + map.put(key, new DialogSummary(owner, peer, relationFlag, "", 0L, 0, 0L, System.currentTimeMillis(), hasDialog)); + return; + } + String nextRelation = current.relationFlag(); + if ("close_friend".equalsIgnoreCase(relationFlag) || "close_friend".equalsIgnoreCase(nextRelation)) { + nextRelation = "close_friend"; + } else if ("contact".equalsIgnoreCase(relationFlag) || "contact".equalsIgnoreCase(nextRelation)) { + nextRelation = "contact"; + } else { + nextRelation = normalizeRelationFlag(nextRelation); + } + map.put(key, new DialogSummary( + current.ownerLogin().isEmpty() ? owner : current.ownerLogin(), + peer, + nextRelation, + current.lastMessageBlobB64(), + current.lastMessageTimeMs(), + current.unreadCount(), + current.lastReadReceiptTimeMs(), + current.updatedAtMs(), + current.hasDialog() || hasDialog + )); + } + + private void syncMessagesToWatermark(Connection c, String ownerLogin, String peerLogin, long watermark) throws SQLException { + if (watermark <= 0) return; + String cleanOwner = normalize(ownerLogin); + String cleanPeer = normalize(peerLogin); + if (cleanOwner.isEmpty() || cleanPeer.isEmpty() || cleanOwner.equalsIgnoreCase(cleanPeer)) return; + try (PreparedStatement ps = c.prepareStatement(""" + UPDATE signed_messages + SET read_at_ms = CASE + WHEN read_at_ms IS NULL OR read_at_ms <= 0 THEN ? + WHEN read_at_ms > ? THEN read_at_ms + ELSE read_at_ms + END + WHERE ( + (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + OR (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)) + ) + AND message_type IN (1, 2) + AND time_ms <= ? + """)) { + ps.setLong(1, watermark); + ps.setLong(2, watermark); + ps.setString(3, cleanOwner); + ps.setString(4, cleanPeer); + ps.setString(5, cleanPeer); + ps.setString(6, cleanOwner); + ps.setLong(7, watermark); + ps.executeUpdate(); + } + } + + private long loadExistingWatermark(Connection c, String ownerLogin, String peerLogin) throws SQLException { + try (PreparedStatement ps = c.prepareStatement(""" + SELECT last_read_receipt_time_ms + FROM dm_dialog_state + WHERE LOWER(owner_login) = LOWER(?) + AND LOWER(peer_login) = LOWER(?) + LIMIT 1 + """)) { + ps.setString(1, ownerLogin); + ps.setString(2, peerLogin); + try (ResultSet rs = ps.executeQuery()) { + if (!rs.next()) return 0L; + long value = rs.getLong(1); + return rs.wasNull() ? 0L : value; + } + } + } + + private String normalize(String value) { + return value == null ? "" : value.trim(); + } + + private String normKey(String value) { + return normalize(value).toLowerCase(Locale.ROOT); + } + + private String normalizeRelationFlag(String value) { + String clean = normalize(value).toLowerCase(Locale.ROOT); + if ("close_friend".equals(clean) || "contact".equals(clean)) return clean; + return "none"; + } + + public record DialogSummary( + String ownerLogin, + String peerLogin, + String relationFlag, + String lastMessageBlobB64, + long lastMessageTimeMs, + int unreadCount, + long lastReadReceiptTimeMs, + long updatedAtMs, + boolean hasDialog + ) {} + + private record PairItem(String ownerLogin, String peerLogin) {} +} diff --git a/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/SignedMessagesDAO.java b/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/SignedMessagesDAO.java index e8d11c6f..988f6c72 100644 --- a/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/SignedMessagesDAO.java +++ b/SHiNE-server/shine-server-db/src/main/java/shine/db/dao/SignedMessagesDAO.java @@ -58,6 +58,7 @@ public final class SignedMessagesDAO { ApplyStatus status = ps.executeUpdate() > 0 ? ApplyStatus.APPLIED : ApplyStatus.DUPLICATE_OR_OLDER; if (status.applied()) { markMessageReadByReceipt(c, e); + DmDialogStateDAO.getInstance().refreshFromEntry(c, e); } return status; } @@ -76,6 +77,7 @@ public final class SignedMessagesDAO { if (insertedFirst == 1 && insertedSecond == 1) { markMessageReadByReceipt(c, first); markMessageReadByReceipt(c, second); + DmDialogStateDAO.getInstance().refreshConversationPair(c, first.getFromLogin(), first.getToLogin()); c.commit(); return true; } @@ -122,6 +124,7 @@ public final class SignedMessagesDAO { markMessageReadByReceipt(c, outgoing); resetDeliveryRows(c, incoming.getMessageKey()); resetDeliveryRows(c, outgoing.getMessageKey()); + DmDialogStateDAO.getInstance().refreshConversationPair(c, incoming.getFromLogin(), incoming.getToLogin()); c.commit(); return ApplyStatus.APPLIED; @@ -160,6 +163,7 @@ public final class SignedMessagesDAO { upsertMessage(c, incoming); markMessageReadByReceipt(c, incoming); resetDeliveryRows(c, incoming.getMessageKey()); + DmDialogStateDAO.getInstance().refreshConversationPair(c, incoming.getFromLogin(), incoming.getToLogin()); c.commit(); return ApplyStatus.APPLIED; } catch (Exception ex) { @@ -190,6 +194,7 @@ public final class SignedMessagesDAO { deleteMessageContentAndReceipts(c, tombstone.getBaseKey()); upsertMessage(c, tombstone); resetDeliveryRows(c, tombstone.getMessageKey()); + DmDialogStateDAO.getInstance().refreshConversationPair(c, tombstone.getFromLogin(), tombstone.getToLogin()); c.commit(); return ApplyStatus.APPLIED; @@ -218,6 +223,7 @@ public final class SignedMessagesDAO { deleteConversationHistoryBefore(c, tombstone.getFromLogin(), tombstone.getToLogin(), tombstone.getTimeMs()); upsertMessage(c, tombstone); resetDeliveryRows(c, tombstone.getMessageKey()); + DmDialogStateDAO.getInstance().refreshConversationPair(c, tombstone.getFromLogin(), tombstone.getToLogin()); c.commit(); return ApplyStatus.APPLIED; diff --git a/SHiNE-server/shine-server-db/src/main/resources/postgres/migration_v10.sql b/SHiNE-server/shine-server-db/src/main/resources/postgres/migration_v10.sql new file mode 100644 index 00000000..42bf9750 --- /dev/null +++ b/SHiNE-server/shine-server-db/src/main/resources/postgres/migration_v10.sql @@ -0,0 +1,30 @@ +BEGIN; + +CREATE TABLE IF NOT EXISTS dm_dialog_state ( + owner_login TEXT NOT NULL, + peer_login TEXT NOT NULL, + relation_flag TEXT NOT NULL DEFAULT 'none', + last_message_blob_b64 TEXT NOT NULL DEFAULT '', + last_message_key TEXT NOT NULL DEFAULT '', + last_message_type INTEGER NOT NULL DEFAULT 0, + last_message_time_ms BIGINT NOT NULL DEFAULT 0, + last_message_revision_ms BIGINT NOT NULL DEFAULT 0, + unread_count INTEGER NOT NULL DEFAULT 0, + last_read_receipt_time_ms BIGINT NOT NULL DEFAULT 0, + updated_at_ms BIGINT NOT NULL, + PRIMARY KEY (owner_login, peer_login) +); + +ALTER TABLE IF EXISTS dm_dialog_state + ADD COLUMN IF NOT EXISTS last_message_blob_b64 TEXT NOT NULL DEFAULT ''; + +CREATE INDEX IF NOT EXISTS idx_dm_dialog_state_owner_last_time + ON dm_dialog_state(owner_login, last_message_time_ms DESC, peer_login); + +INSERT INTO db_schema_version (id, schema_version, updated_at_ms) +VALUES (1, 10, CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)) +ON CONFLICT (id) DO UPDATE SET + schema_version = EXCLUDED.schema_version, + updated_at_ms = EXCLUDED.updated_at_ms; + +COMMIT; diff --git a/SHiNE-server/shine-server-db/src/main/resources/postgres/schema_v1.sql b/SHiNE-server/shine-server-db/src/main/resources/postgres/schema_v1.sql index ff1c91db..4a8191f5 100644 --- a/SHiNE-server/shine-server-db/src/main/resources/postgres/schema_v1.sql +++ b/SHiNE-server/shine-server-db/src/main/resources/postgres/schema_v1.sql @@ -751,6 +751,24 @@ CREATE TABLE IF NOT EXISTS signed_message_session_delivery ( CREATE INDEX IF NOT EXISTS idx_signed_message_delivery_session ON signed_message_session_delivery(session_id, delivered); +CREATE TABLE IF NOT EXISTS dm_dialog_state ( + owner_login TEXT NOT NULL, + peer_login TEXT NOT NULL, + relation_flag TEXT NOT NULL DEFAULT 'none', + last_message_blob_b64 TEXT NOT NULL DEFAULT '', + last_message_key TEXT NOT NULL DEFAULT '', + last_message_type INTEGER NOT NULL DEFAULT 0, + last_message_time_ms BIGINT NOT NULL DEFAULT 0, + last_message_revision_ms BIGINT NOT NULL DEFAULT 0, + unread_count INTEGER NOT NULL DEFAULT 0, + last_read_receipt_time_ms BIGINT NOT NULL DEFAULT 0, + updated_at_ms BIGINT NOT NULL, + PRIMARY KEY (owner_login, peer_login) +); + +CREATE INDEX IF NOT EXISTS idx_dm_dialog_state_owner_last_time + ON dm_dialog_state(owner_login, last_message_time_ms DESC, peer_login); + CREATE TABLE IF NOT EXISTS dm_sync_peer_state ( owner_login TEXT NOT NULL, remote_server_login TEXT NOT NULL, diff --git a/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/Net_ListContacts_Handler.java b/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/Net_ListContacts_Handler.java index 92dcd2df..52888466 100644 --- a/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/Net_ListContacts_Handler.java +++ b/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/Net_ListContacts_Handler.java @@ -8,10 +8,10 @@ import server.logic.ws_protocol.JSON.handlers.connections.entyties.Net_ListConta import server.logic.ws_protocol.JSON.handlers.connections.entyties.Net_ListContacts_Response; import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory; import server.logic.ws_protocol.WireCodes; -import shine.db.MsgSubType; -import shine.db.dao.ConnectionsStateDAO; +import shine.db.dao.DmDialogStateDAO; import java.sql.Connection; +import java.util.ArrayList; import java.util.List; public class Net_ListContacts_Handler implements JsonMessageHandler { @@ -23,14 +23,30 @@ public class Net_ListContacts_Handler implements JsonMessageHandler { } try (Connection c = shine.db.DbController.getInstance().getConnection()) { - List contacts = ConnectionsStateDAO.getInstance().listOutgoingByRelTypeCanonical(c, ctx.getLogin(), MsgSubType.CONNECTION_CONTACT); + List dialogs = DmDialogStateDAO.getInstance().listInboxDialogs(c, ctx.getLogin()); Net_ListContacts_Response resp = new Net_ListContacts_Response(); resp.setOp(req.getOp()); resp.setRequestId(req.getRequestId()); resp.setStatus(WireCodes.Status.OK); resp.setLogin(ctx.getLogin()); - resp.setContacts(contacts); + resp.setDialogs(toDialogItems(dialogs)); return resp; } } + + private List toDialogItems(List dialogs) { + List items = new ArrayList<>(); + if (dialogs == null) return items; + for (DmDialogStateDAO.DialogSummary dialog : dialogs) { + Net_ListContacts_Response.DialogItem item = new Net_ListContacts_Response.DialogItem(); + item.setPeerLogin(dialog.peerLogin()); + item.setRelationFlag(dialog.relationFlag()); + item.setLastMessageBlobB64(dialog.lastMessageBlobB64()); + item.setLastMessageTimeMs(dialog.lastMessageTimeMs()); + item.setUnreadCount(dialog.unreadCount()); + item.setHasDialog(dialog.hasDialog()); + items.add(item); + } + return items; + } } diff --git a/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/entyties/Net_ListContacts_Response.java b/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/entyties/Net_ListContacts_Response.java index cdb2f5b0..ff827a23 100644 --- a/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/entyties/Net_ListContacts_Response.java +++ b/SHiNE-server/shine-server-net-protocol/src/main/java/server/logic/ws_protocol/JSON/handlers/connections/entyties/Net_ListContacts_Response.java @@ -7,10 +7,32 @@ import java.util.List; public class Net_ListContacts_Response extends Net_Response { private String login; - private List contacts = new ArrayList<>(); + private List dialogs = new ArrayList<>(); public String getLogin() { return login; } public void setLogin(String login) { this.login = login; } - public List getContacts() { return contacts; } - public void setContacts(List contacts) { this.contacts = contacts; } + public List getDialogs() { return dialogs; } + public void setDialogs(List dialogs) { this.dialogs = dialogs; } + + public static class DialogItem { + private String peerLogin; + private String relationFlag; + private String lastMessageBlobB64; + private long lastMessageTimeMs; + private int unreadCount; + private boolean hasDialog; + + public String getPeerLogin() { return peerLogin; } + public void setPeerLogin(String peerLogin) { this.peerLogin = peerLogin; } + public String getRelationFlag() { return relationFlag; } + public void setRelationFlag(String relationFlag) { this.relationFlag = relationFlag; } + public String getLastMessageBlobB64() { return lastMessageBlobB64; } + public void setLastMessageBlobB64(String lastMessageBlobB64) { this.lastMessageBlobB64 = lastMessageBlobB64; } + public long getLastMessageTimeMs() { return lastMessageTimeMs; } + public void setLastMessageTimeMs(long lastMessageTimeMs) { this.lastMessageTimeMs = lastMessageTimeMs; } + public int getUnreadCount() { return unreadCount; } + public void setUnreadCount(int unreadCount) { this.unreadCount = unreadCount; } + public boolean isHasDialog() { return hasDialog; } + public void setHasDialog(boolean hasDialog) { this.hasDialog = hasDialog; } + } } diff --git a/docs/API/11_Connections_API.md b/docs/API/11_Connections_API.md index f0bd8c42..c7acfa3e 100644 --- a/docs/API/11_Connections_API.md +++ b/docs/API/11_Connections_API.md @@ -66,11 +66,44 @@ "ok": true, "payload": { "login": "Alice", - "contacts": ["Bob", "Kate"] + "dialogs": [ + { + "peerLogin": "Bob", + "relationFlag": "close_friend", + "lastMessageBlobB64": "U0hpTkVfRE0B...", + "lastMessageTimeMs": 1774700000123, + "unreadCount": 2, + "hasDialog": true + }, + { + "peerLogin": "Kate", + "relationFlag": "contact", + "lastMessageBlobB64": "", + "lastMessageTimeMs": 0, + "unreadCount": 0, + "hasDialog": false + }, + { + "peerLogin": "Mira", + "relationFlag": "none", + "lastMessageBlobB64": "U0hpTkVfRE0B...", + "lastMessageTimeMs": 1774700000555, + "unreadCount": 1, + "hasDialog": true + } + ] } } ``` +### Примечание + +- `dialogs` это серверный inbox-проекционный список диалогов; +- `relationFlag` возвращается как `close_friend`, `contact` или `none`; +- если один и тот же человек есть и в `contact`, и в `close_friend`, в `dialogs` он приходит как `close_friend`. +- `lastMessageBlobB64` содержит полный signed DM block последнего контентного сообщения в base64; +- для чатов без сообщений поле `lastMessageBlobB64` пустое. + --- ## 3. `GetUserConnectionsGraph` diff --git a/docs/API/12_Direct_Messages_Push_Calls_API.md b/docs/API/12_Direct_Messages_Push_Calls_API.md index 6bd5d554..2efb1682 100644 --- a/docs/API/12_Direct_Messages_Push_Calls_API.md +++ b/docs/API/12_Direct_Messages_Push_Calls_API.md @@ -11,6 +11,8 @@ - для DM v1 нужно использовать `SendMessagePair`, `ReceiveOutcomingMessage`, `ReceiveIncomingMessage`, `DeleteMessage`, `DeleteConversation`, `GetDirectMessages`; - `DmSyncBatch` предназначен для межсерверной догоняющей синхронизации, не для обычного клиентского UI. +- сервер поддерживает материализованный слой диалогов `dm_dialog_state`; `read receipt` обновляет серверный watermark и `unreadCount`, а не только локальный клиентский флаг. +- в `dm_dialog_state` сервер также хранит `last_message_blob_b64` для последнего контентного DM в base64, чтобы клиент мог отрисовать список чатов без дополнительного запроса. ## 1. `UpsertPushToken` @@ -147,6 +149,12 @@ `sourceServerLogin` необязателен. Если поле есть, сервер использует его как подсказку, чтобы не отправлять событие обратно серверу-источнику. +### Примечание + +- входящий `type=3` не только сохраняется как событие прочтения, но и обновляет серверный watermark диалога; +- если подтверждение прочтения приходит не по порядку, сервер сохраняет наибольший watermark и пересчитывает `unreadCount` по фактическому состоянию сообщений; +- это нужно, чтобы разные устройства не расходились по счётчику непрочитанных. + ## 5. `DeleteMessage` Принимает один signed DM-блок `type=5` или `type=6`. @@ -359,7 +367,7 @@ - все DM-типы `1..8` используют `SHiNE_DM` - `GetUser` может lazy-import пользователя из Solana PDA, поэтому именно через него клиент обычно получает `clientKey` адресата для E2EE -- сервер не расшифровывает DM и не использует ciphertext как preview текста +- сервер не расшифровывает DM; в списке диалогов он отдаёт последний signed block как `lastMessageBlobB64`, а не извлекает plaintext preview - сервер хранит последнюю применённую версию контентного сообщения по правилу `revisionTimeMs`, а при равенстве — по `reencryptedAtMs` - если сервер уже знает tombstone удаления переписки и получает старое сообщение до этой границы, он перерассылает известный `DeleteConversation` на `access_servers` обеих сторон - HTTP endpoints для DM-файлов сейчас отсутствуют diff --git a/docs/Personal_Messages/Протокол_DM_v1.md b/docs/Personal_Messages/Протокол_DM_v1.md index b6843f3f..807b5adf 100644 --- a/docs/Personal_Messages/Протокол_DM_v1.md +++ b/docs/Personal_Messages/Протокол_DM_v1.md @@ -334,6 +334,25 @@ ### 8.2. Новые методы, которые нужны +### 8.3. Серверный слой диалогов + +Помимо хранения самих DM-сообщений сервер поддерживает материализованный слой состояния диалогов: + +- отдельная запись на пару `owner_login` + `peer_login`; +- `relation_flag` со значениями `close_friend`, `contact`, `none`; +- `last_message_blob_b64` как последний контентный signed DM block в base64; +- `last_message_time_ms`; +- `unread_count`; +- `last_read_receipt_time_ms` как watermark последнего подтверждения прочтения. + +Ключевые правила: + +- `close_friend` всегда имеет приоритет над `contact`; +- если `read receipt` приходит не по порядку, сервер хранит наибольший watermark и не откатывает состояние назад; +- `unread_count` пересчитывается сервером по сообщениям диалога с учётом watermark и `read_at_ms`; +- старые исторические данные восстанавливаются из `signed_messages` при инициализации/миграции; +- UI не должен собирать inbox только из локального кеша, когда ему доступен серверный список диалогов. + ## 9. Правила валидации и применения ### 9.1. Общее правило по ревизиям diff --git a/docs/Personal_Messages/Формат_DM_v1.md b/docs/Personal_Messages/Формат_DM_v1.md index 98b75add..b2bdaf10 100644 --- a/docs/Personal_Messages/Формат_DM_v1.md +++ b/docs/Personal_Messages/Формат_DM_v1.md @@ -252,6 +252,15 @@ ReadReceiptBody_v1_0 Различается только `messageType` и формат `body`. +### Серверное примечание + +Внешний байтовый формат `type=3/4` не меняется, но сервер использует такие контейнеры как вход для обновления `dm_dialog_state`: + +- `read receipt` обновляет серверный watermark диалога; +- `unreadCount` пересчитывается на сервере, а не только на клиенте; +- если подтверждение прочтения приходит в другом порядке, сервер сохраняет максимальный watermark и не откатывает счётчик назад. +- в списке диалогов сервер может отдавать последний signed block как `lastMessageBlobB64` без попытки извлечь plaintext preview. + ## 9. Контент типов `5/6` Типы: