Сохранить текущий снимок DM-синхронизации

This commit is contained in:
AidarKC
2026-08-24 17:50:36 +04:00
parent 015fade4c0
commit d8e0c77951
48 changed files with 2498 additions and 383 deletions
@@ -28,6 +28,8 @@ public final class DatabaseInitializer {
public static final int SCHEMA_VERSION_9 = 9; public static final int SCHEMA_VERSION_9 = 9;
public static final int SCHEMA_VERSION_10 = 10; public static final int SCHEMA_VERSION_10 = 10;
public static final int SCHEMA_VERSION_11 = 11; public static final int SCHEMA_VERSION_11 = 11;
public static final int SCHEMA_VERSION_12 = 12;
public static final int SCHEMA_VERSION_13 = 13;
public static final String POSTGRES_SCHEMA_RESOURCE = "postgres/schema_v1.sql"; 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_V2_RESOURCE = "postgres/migration_v2.sql";
public static final String POSTGRES_MIGRATION_V3_RESOURCE = "postgres/migration_v3.sql"; public static final String POSTGRES_MIGRATION_V3_RESOURCE = "postgres/migration_v3.sql";
@@ -39,6 +41,8 @@ public final class DatabaseInitializer {
public static final String POSTGRES_MIGRATION_V9_RESOURCE = "postgres/migration_v9.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"; public static final String POSTGRES_MIGRATION_V10_RESOURCE = "postgres/migration_v10.sql";
public static final String POSTGRES_MIGRATION_V11_RESOURCE = "postgres/migration_v11.sql"; public static final String POSTGRES_MIGRATION_V11_RESOURCE = "postgres/migration_v11.sql";
public static final String POSTGRES_MIGRATION_V12_RESOURCE = "postgres/migration_v12.sql";
public static final String POSTGRES_MIGRATION_V13_RESOURCE = "postgres/migration_v13.sql";
private DatabaseInitializer() {} private DatabaseInitializer() {}
@@ -142,6 +146,14 @@ public final class DatabaseInitializer {
runSqlScript(conn, POSTGRES_MIGRATION_V11_RESOURCE); runSqlScript(conn, POSTGRES_MIGRATION_V11_RESOURCE);
currentVersion = SCHEMA_VERSION_11; currentVersion = SCHEMA_VERSION_11;
} }
if (currentVersion < SCHEMA_VERSION_12) {
runSqlScript(conn, POSTGRES_MIGRATION_V12_RESOURCE);
currentVersion = SCHEMA_VERSION_12;
}
if (currentVersion < SCHEMA_VERSION_13) {
runSqlScript(conn, POSTGRES_MIGRATION_V13_RESOURCE);
currentVersion = SCHEMA_VERSION_13;
}
} }
} }
@@ -0,0 +1,518 @@
package shine.db.dao;
import shine.db.DbController;
import shine.db.entities.DmDeliveryStateEntry;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Types;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/** Хранилище изменяемого состояния межсерверной доставки DM-пары. */
public final class DmDeliveryStateDAO {
private static volatile DmDeliveryStateDAO instance;
private final DbController db = DbController.getInstance();
private DmDeliveryStateDAO() {}
public static DmDeliveryStateDAO getInstance() {
if (instance == null) {
synchronized (DmDeliveryStateDAO.class) {
if (instance == null) instance = new DmDeliveryStateDAO();
}
}
return instance;
}
public DmDeliveryStateEntry upsertPair(
String outgoingMessageKey,
String eventId,
String baseKey,
String fromLogin,
String toLogin,
String incomingMessageKey,
long createdAtMs,
long expiresAtMs,
int initialState,
String deliveredServerLogin,
String routesHash,
boolean assistImmediately
) throws SQLException {
long now = System.currentTimeMillis();
int safeState = normalizeState(initialState, deliveredServerLogin);
if (expiresAtMs <= now && safeState == DmDeliveryStateEntry.PENDING_NONE) {
safeState = DmDeliveryStateEntry.FAILED_FINAL;
}
Long nextAttemptAt = safeState == DmDeliveryStateEntry.ACCEPTED
? now
: null;
String safeDeliveredLogin = safeState == DmDeliveryStateEntry.DELIVERED_ONE
? normalize(deliveredServerLogin)
: null;
try (Connection c = db.getConnection()) {
String sql = """
INSERT INTO dm_delivery_state (
outgoing_message_key, event_id, base_key, from_login, to_login,
incoming_message_key, created_at_ms, delivery_expires_at_ms,
delivery_state, delivered_server_login, recipient_routes_hash,
attempt_index, next_attempt_at_ms, last_attempt_at_ms, last_error, updated_at_ms
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, ?, NULL, NULL, ?)
ON CONFLICT (outgoing_message_key) DO UPDATE SET
event_id = EXCLUDED.event_id,
base_key = EXCLUDED.base_key,
from_login = EXCLUDED.from_login,
to_login = EXCLUDED.to_login,
incoming_message_key = EXCLUDED.incoming_message_key,
created_at_ms = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.created_at_ms
ELSE EXCLUDED.created_at_ms
END,
delivery_expires_at_ms = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.delivery_expires_at_ms
ELSE EXCLUDED.delivery_expires_at_ms
END,
delivery_state = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.delivery_state
ELSE EXCLUDED.delivery_state
END,
delivered_server_login = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.delivered_server_login
ELSE EXCLUDED.delivered_server_login
END,
recipient_routes_hash = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN COALESCE(dm_delivery_state.recipient_routes_hash, EXCLUDED.recipient_routes_hash)
ELSE EXCLUDED.recipient_routes_hash
END,
attempt_index = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.attempt_index
ELSE 0
END,
next_attempt_at_ms = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.next_attempt_at_ms
ELSE EXCLUDED.next_attempt_at_ms
END,
last_attempt_at_ms = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.last_attempt_at_ms
ELSE NULL
END,
last_error = CASE
WHEN dm_delivery_state.event_id = EXCLUDED.event_id THEN dm_delivery_state.last_error
ELSE NULL
END,
updated_at_ms = EXCLUDED.updated_at_ms
""";
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, outgoingMessageKey);
ps.setString(2, eventId);
ps.setString(3, baseKey);
ps.setString(4, normalize(fromLogin));
ps.setString(5, normalize(toLogin));
ps.setString(6, incomingMessageKey);
ps.setLong(7, createdAtMs);
ps.setLong(8, expiresAtMs);
ps.setInt(9, safeState);
setNullableString(ps, 10, safeDeliveredLogin);
setNullableString(ps, 11, normalizeBlank(routesHash));
setNullableLong(ps, 12, nextAttemptAt);
ps.setLong(13, now);
ps.executeUpdate();
}
}
DmDeliveryStateEntry result;
if (initialState != DmDeliveryStateEntry.PENDING_NONE || !isBlank(deliveredServerLogin)) {
result = mergeRemote(eventId, initialState, deliveredServerLogin, routesHash, now);
} else {
result = getByEventId(eventId);
}
if (result != null && result.getDeliveryExpiresAtMs() <= now
&& result.getDeliveryState() == DmDeliveryStateEntry.PENDING_NONE) {
return finishAtExpiry(eventId, now);
}
return result;
}
public DmDeliveryStateEntry getByEventId(String eventId) throws SQLException {
if (isBlank(eventId)) return null;
try (Connection c = db.getConnection();
PreparedStatement ps = c.prepareStatement(selectColumns() + " WHERE event_id = ?")) {
ps.setString(1, eventId.trim());
try (ResultSet rs = ps.executeQuery()) {
return rs.next() ? mapRow(rs) : null;
}
}
}
public DmDeliveryStateEntry getByOutgoingMessageKey(String outgoingMessageKey) throws SQLException {
if (isBlank(outgoingMessageKey)) return null;
try (Connection c = db.getConnection();
PreparedStatement ps = c.prepareStatement(selectColumns() + " WHERE outgoing_message_key = ?")) {
ps.setString(1, outgoingMessageKey.trim());
try (ResultSet rs = ps.executeQuery()) {
return rs.next() ? mapRow(rs) : null;
}
}
}
public Map<String, DmDeliveryStateEntry> listByOutgoingMessageKeys(List<String> keys) throws SQLException {
Map<String, DmDeliveryStateEntry> out = new HashMap<>();
if (keys == null || keys.isEmpty()) return out;
String placeholders = String.join(",", java.util.Collections.nCopies(keys.size(), "?"));
String sql = selectColumns() + " WHERE outgoing_message_key IN (" + placeholders + ")";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
for (int i = 0; i < keys.size(); i++) ps.setString(i + 1, keys.get(i));
try (ResultSet rs = ps.executeQuery()) {
while (rs.next()) {
DmDeliveryStateEntry row = mapRow(rs);
out.put(row.getOutgoingMessageKey(), row);
}
}
}
return out;
}
public List<DmDeliveryStateEntry> listDue(long nowMs, int limit) throws SQLException {
String sql = selectColumns() + """
WHERE delivery_state = 0
AND next_attempt_at_ms IS NOT NULL
AND next_attempt_at_ms <= ?
ORDER BY next_attempt_at_ms ASC, created_at_ms ASC
LIMIT ?
""";
List<DmDeliveryStateEntry> out = new ArrayList<>();
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setLong(1, nowMs);
ps.setInt(2, Math.max(1, limit));
try (ResultSet rs = ps.executeQuery()) {
while (rs.next()) out.add(mapRow(rs));
}
}
return out;
}
public boolean claimAttempt(DmDeliveryStateEntry row, long nowMs, Long nextAttemptAtMs) throws SQLException {
String sql = """
UPDATE dm_delivery_state
SET attempt_index = attempt_index + 1,
next_attempt_at_ms = ?,
last_attempt_at_ms = ?,
last_error = NULL,
updated_at_ms = ?
WHERE outgoing_message_key = ?
AND event_id = ?
AND attempt_index = ?
AND delivery_state = 0
AND next_attempt_at_ms IS NOT NULL
AND next_attempt_at_ms <= ?
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
long networkLeaseUntilMs = nowMs + 60_000L;
Long claimedNextAttemptAtMs = nextAttemptAtMs == null
? networkLeaseUntilMs
: Math.max(nextAttemptAtMs, networkLeaseUntilMs);
setNullableLong(ps, 1, claimedNextAttemptAtMs);
ps.setLong(2, nowMs);
ps.setLong(3, nowMs);
ps.setString(4, row.getOutgoingMessageKey());
ps.setString(5, row.getEventId());
ps.setInt(6, row.getAttemptIndex());
ps.setLong(7, nowMs);
return ps.executeUpdate() == 1;
}
}
public DmDeliveryStateEntry mergeRemote(
String eventId,
int remoteState,
String remoteDeliveredLogin,
String remoteRoutesHash,
long nowMs
) throws SQLException {
try (Connection c = db.getConnection()) {
boolean previousAutoCommit = c.getAutoCommit();
c.setAutoCommit(false);
try {
DmDeliveryStateEntry current = getByEventIdForUpdate(c, eventId);
if (current == null) {
c.rollback();
return null;
}
MergeResult merged = merge(
current.getDeliveryState(),
current.getDeliveredServerLogin(),
current.getRecipientRoutesHash(),
remoteState,
remoteDeliveredLogin,
remoteRoutesHash
);
updateMutableState(c, current, merged.state(), merged.deliveredLogin(),
merged.routesHash(), current.getNextAttemptAtMs(), current.getLastError(), nowMs);
c.commit();
return getByEventId(eventId);
} catch (Exception e) {
try { c.rollback(); } catch (Exception ignored) {}
throw e;
} finally {
c.setAutoCommit(previousAutoCommit);
}
}
}
public DmDeliveryStateEntry updateAfterAttempt(
String eventId,
int attemptedState,
String deliveredServerLogin,
String routesHash,
Long nextAttemptAtMs,
String lastError,
long nowMs
) throws SQLException {
try (Connection c = db.getConnection()) {
boolean previousAutoCommit = c.getAutoCommit();
c.setAutoCommit(false);
try {
DmDeliveryStateEntry current = getByEventIdForUpdate(c, eventId);
if (current == null) {
c.rollback();
return null;
}
MergeResult merged = merge(
current.getDeliveryState(), current.getDeliveredServerLogin(), current.getRecipientRoutesHash(),
attemptedState, deliveredServerLogin, routesHash
);
Long effectiveNext = merged.state() == DmDeliveryStateEntry.ACCEPTED
? nextAttemptAtMs
: null;
updateMutableState(c, current, merged.state(), merged.deliveredLogin(),
merged.routesHash(), effectiveNext, lastError, nowMs);
c.commit();
return getByEventId(eventId);
} catch (Exception e) {
try { c.rollback(); } catch (Exception ignored) {}
throw e;
} finally {
c.setAutoCommit(previousAutoCommit);
}
}
}
public DmDeliveryStateEntry finishAtExpiry(String eventId, long nowMs) throws SQLException {
try (Connection c = db.getConnection()) {
String sql = """
UPDATE dm_delivery_state
SET delivery_state = CASE WHEN delivery_state = 0 THEN 3 ELSE delivery_state END,
delivered_server_login = CASE WHEN delivery_state = 0 THEN NULL ELSE delivered_server_login END,
next_attempt_at_ms = NULL,
last_error = CASE WHEN delivery_state = 0 THEN 'DELIVERY_EXPIRED' ELSE last_error END,
updated_at_ms = ?
WHERE event_id = ? AND delivery_state = 0
""";
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setLong(1, nowMs);
ps.setString(2, eventId);
ps.executeUpdate();
}
}
return getByEventId(eventId);
}
/** Read-only peer confirmed that at least one recipient server accepted the message. */
public DmDeliveryStateEntry markDeliveredFromPeer(String eventId, long nowMs) throws SQLException {
try (Connection c = db.getConnection();
PreparedStatement ps = c.prepareStatement("""
UPDATE dm_delivery_state
SET delivery_state = 1,
delivered_server_login = NULL,
next_attempt_at_ms = NULL,
last_error = NULL,
updated_at_ms = ?
WHERE event_id = ? AND delivery_state = 0
""")) {
ps.setLong(1, nowMs);
ps.setString(2, eventId);
ps.executeUpdate();
}
return getByEventId(eventId);
}
public int removeByBaseKey(String baseKey) throws SQLException {
try (Connection c = db.getConnection();
PreparedStatement ps = c.prepareStatement("DELETE FROM dm_delivery_state WHERE base_key = ?")) {
ps.setString(1, baseKey);
return ps.executeUpdate();
}
}
public int removeConversationBefore(String fromLogin, String toLogin, long boundaryTimeMs) throws SQLException {
String sql = """
DELETE FROM dm_delivery_state
WHERE created_at_ms < ?
AND ((LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?))
OR (LOWER(from_login) = LOWER(?) AND LOWER(to_login) = LOWER(?)))
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setLong(1, boundaryTimeMs);
ps.setString(2, fromLogin);
ps.setString(3, toLogin);
ps.setString(4, toLogin);
ps.setString(5, fromLogin);
return ps.executeUpdate();
}
}
public int removeMissingMessages() throws SQLException {
String sql = """
DELETE FROM dm_delivery_state d
WHERE NOT EXISTS (
SELECT 1 FROM signed_messages m WHERE m.message_key = d.outgoing_message_key
)
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
return ps.executeUpdate();
}
}
private DmDeliveryStateEntry getByEventIdForUpdate(Connection c, String eventId) throws SQLException {
try (PreparedStatement ps = c.prepareStatement(selectColumns() + " WHERE event_id = ? FOR UPDATE")) {
ps.setString(1, eventId);
try (ResultSet rs = ps.executeQuery()) {
return rs.next() ? mapRow(rs) : null;
}
}
}
private void updateMutableState(
Connection c,
DmDeliveryStateEntry current,
int state,
String deliveredLogin,
String routesHash,
Long nextAttemptAtMs,
String lastError,
long nowMs
) throws SQLException {
String sql = """
UPDATE dm_delivery_state
SET delivery_state = ?,
delivered_server_login = ?,
recipient_routes_hash = ?,
next_attempt_at_ms = ?,
last_error = ?,
updated_at_ms = ?
WHERE outgoing_message_key = ? AND event_id = ?
""";
try (PreparedStatement ps = c.prepareStatement(sql)) {
ps.setInt(1, state);
setNullableString(ps, 2, state == DmDeliveryStateEntry.DELIVERED_ONE ? normalize(deliveredLogin) : null);
setNullableString(ps, 3, normalizeBlank(routesHash));
setNullableLong(ps, 4, nextAttemptAtMs);
setNullableString(ps, 5, normalizeBlank(lastError));
ps.setLong(6, nowMs);
ps.setString(7, current.getOutgoingMessageKey());
ps.setString(8, current.getEventId());
ps.executeUpdate();
}
}
private MergeResult merge(
int localState,
String localLogin,
String localHash,
int remoteStateRaw,
String remoteLoginRaw,
String remoteHashRaw
) {
int remoteState = normalizeState(remoteStateRaw, remoteLoginRaw);
String remoteLogin = normalize(remoteLoginRaw);
String localNormalizedLogin = normalize(localLogin);
String localRoutesHash = normalizeBlank(localHash);
String remoteRoutesHash = normalizeBlank(remoteHashRaw);
String resultHash = remoteRoutesHash != null ? remoteRoutesHash : localRoutesHash;
if (localState == DmDeliveryStateEntry.DELIVERED_ONE || localState == DmDeliveryStateEntry.DELIVERED_ALL) {
return new MergeResult(DmDeliveryStateEntry.DELIVERED, localNormalizedLogin,
localRoutesHash != null ? localRoutesHash : remoteRoutesHash);
}
if (remoteState == DmDeliveryStateEntry.DELIVERED_ONE || remoteState == DmDeliveryStateEntry.DELIVERED_ALL) {
return new MergeResult(DmDeliveryStateEntry.DELIVERED, remoteLogin, resultHash);
}
if (localState == DmDeliveryStateEntry.FAILED_FINAL || remoteState == DmDeliveryStateEntry.FAILED_FINAL) {
return new MergeResult(DmDeliveryStateEntry.FAILED_FINAL, null, resultHash);
}
return new MergeResult(DmDeliveryStateEntry.PENDING_NONE, null, resultHash);
}
private int normalizeState(int state, String deliveredLogin) {
if (state < DmDeliveryStateEntry.PENDING_NONE || state > DmDeliveryStateEntry.FAILED_FINAL) {
return DmDeliveryStateEntry.PENDING_NONE;
}
if (state == DmDeliveryStateEntry.DELIVERED_ALL) return DmDeliveryStateEntry.DELIVERED;
return state;
}
private String selectColumns() {
return """
SELECT outgoing_message_key, event_id, base_key, from_login, to_login,
incoming_message_key, created_at_ms, delivery_expires_at_ms,
delivery_state, delivered_server_login, recipient_routes_hash,
attempt_index, next_attempt_at_ms, last_attempt_at_ms, last_error, updated_at_ms
FROM dm_delivery_state
""";
}
private DmDeliveryStateEntry mapRow(ResultSet rs) throws SQLException {
DmDeliveryStateEntry e = new DmDeliveryStateEntry();
e.setOutgoingMessageKey(rs.getString("outgoing_message_key"));
e.setEventId(rs.getString("event_id"));
e.setBaseKey(rs.getString("base_key"));
e.setFromLogin(rs.getString("from_login"));
e.setToLogin(rs.getString("to_login"));
e.setIncomingMessageKey(rs.getString("incoming_message_key"));
e.setCreatedAtMs(rs.getLong("created_at_ms"));
e.setDeliveryExpiresAtMs(rs.getLong("delivery_expires_at_ms"));
e.setDeliveryState(rs.getInt("delivery_state"));
e.setDeliveredServerLogin(rs.getString("delivered_server_login"));
e.setRecipientRoutesHash(rs.getString("recipient_routes_hash"));
e.setAttemptIndex(rs.getInt("attempt_index"));
long next = rs.getLong("next_attempt_at_ms");
e.setNextAttemptAtMs(rs.wasNull() ? null : next);
long last = rs.getLong("last_attempt_at_ms");
e.setLastAttemptAtMs(rs.wasNull() ? null : last);
e.setLastError(rs.getString("last_error"));
e.setUpdatedAtMs(rs.getLong("updated_at_ms"));
return e;
}
private static void setNullableString(PreparedStatement ps, int index, String value) throws SQLException {
if (isBlank(value)) ps.setNull(index, Types.VARCHAR);
else ps.setString(index, value.trim());
}
private static void setNullableLong(PreparedStatement ps, int index, Long value) throws SQLException {
if (value == null) ps.setNull(index, Types.BIGINT);
else ps.setLong(index, value);
}
private static String normalize(String value) {
if (value == null) return null;
String normalized = value.trim().toLowerCase();
return normalized.isEmpty() ? null : normalized;
}
private static String normalizeBlank(String value) {
if (value == null) return null;
String normalized = value.trim();
return normalized.isEmpty() ? null : normalized;
}
private static boolean isBlank(String value) {
return value == null || value.isBlank();
}
private record MergeResult(int state, String deliveredLogin, String routesHash) {}
}
@@ -0,0 +1,232 @@
package shine.db.dao;
import shine.db.DbController;
import shine.db.entities.DmSyncOutboxEntry;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Types;
import java.util.ArrayList;
import java.util.List;
/** Outbox событий DM для единственного второго access-сервера пользователя. */
public final class DmSyncOutboxDAO {
private static volatile DmSyncOutboxDAO instance;
private final DbController db = DbController.getInstance();
private DmSyncOutboxDAO() {}
public static DmSyncOutboxDAO getInstance() {
if (instance == null) {
synchronized (DmSyncOutboxDAO.class) {
if (instance == null) instance = new DmSyncOutboxDAO();
}
}
return instance;
}
public void upsert(
String ownerLogin,
String primaryMessageKey,
String eventId,
String secondaryMessageKey,
boolean synced,
long createdAtMs
) throws SQLException {
long now = System.currentTimeMillis();
String sql = """
INSERT INTO dm_sync_outbox (
owner_login, primary_message_key, event_id, secondary_message_key,
synced, created_at_ms, updated_at_ms
) VALUES (LOWER(?), ?, ?, ?, ?, ?, ?)
ON CONFLICT (owner_login, primary_message_key) DO UPDATE SET
event_id = EXCLUDED.event_id,
secondary_message_key = EXCLUDED.secondary_message_key,
synced = CASE
WHEN dm_sync_outbox.event_id = EXCLUDED.event_id
THEN dm_sync_outbox.synced OR EXCLUDED.synced
ELSE EXCLUDED.synced
END,
created_at_ms = CASE
WHEN dm_sync_outbox.event_id = EXCLUDED.event_id
THEN dm_sync_outbox.created_at_ms
ELSE EXCLUDED.created_at_ms
END,
updated_at_ms = EXCLUDED.updated_at_ms
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, ownerLogin);
ps.setString(2, primaryMessageKey);
ps.setString(3, eventId);
if (secondaryMessageKey == null || secondaryMessageKey.isBlank()) ps.setNull(4, Types.VARCHAR);
else ps.setString(4, secondaryMessageKey.trim());
ps.setBoolean(5, synced);
ps.setLong(6, createdAtMs);
ps.setLong(7, now);
ps.executeUpdate();
}
}
public List<DmSyncOutboxEntry> listUnsynced(String ownerLogin, int limit) throws SQLException {
return listUnsynced(ownerLogin, 0L, "", limit);
}
public List<DmSyncOutboxEntry> listUnsynced(
String ownerLogin, long afterCreatedAtMs, String afterPrimaryMessageKey, int limit
) throws SQLException {
String sql = """
SELECT owner_login, primary_message_key, event_id, secondary_message_key,
synced, created_at_ms, updated_at_ms
FROM dm_sync_outbox
WHERE LOWER(owner_login) = LOWER(?) AND synced = FALSE
AND (created_at_ms > ? OR (created_at_ms = ? AND primary_message_key > ?))
ORDER BY created_at_ms ASC, primary_message_key ASC
LIMIT ?
""";
List<DmSyncOutboxEntry> out = new ArrayList<>();
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, ownerLogin);
ps.setLong(2, Math.max(0L, afterCreatedAtMs));
ps.setLong(3, Math.max(0L, afterCreatedAtMs));
ps.setString(4, afterPrimaryMessageKey == null ? "" : afterPrimaryMessageKey);
ps.setInt(5, Math.max(1, limit));
try (ResultSet rs = ps.executeQuery()) {
while (rs.next()) out.add(mapRow(rs));
}
}
return out;
}
public boolean hasUnsyncedForServer(String serverLogin) throws SQLException {
String sql = """
SELECT 1
FROM dm_sync_outbox o
JOIN user_access_servers_current r
ON LOWER(r.user_login) = LOWER(o.owner_login)
WHERE o.synced = FALSE AND LOWER(r.server_login) = LOWER(?)
AND EXISTS (
SELECT 1
FROM user_access_servers_current peer
WHERE LOWER(peer.user_login) = LOWER(o.owner_login)
AND LOWER(peer.server_login) <> LOWER(?)
)
LIMIT 1
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, serverLogin);
ps.setString(2, serverLogin);
try (ResultSet rs = ps.executeQuery()) {
return rs.next();
}
}
}
public List<String> listUnsyncedOwnersForServer(String serverLogin) throws SQLException {
String sql = """
SELECT DISTINCT o.owner_login
FROM dm_sync_outbox o
JOIN user_access_servers_current r
ON LOWER(r.user_login) = LOWER(o.owner_login)
WHERE o.synced = FALSE AND LOWER(r.server_login) = LOWER(?)
AND EXISTS (
SELECT 1
FROM user_access_servers_current peer
WHERE LOWER(peer.user_login) = LOWER(o.owner_login)
AND LOWER(peer.server_login) <> LOWER(?)
)
ORDER BY o.owner_login
""";
List<String> result = new ArrayList<>();
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setString(1, serverLogin);
ps.setString(2, serverLogin);
try (ResultSet rs = ps.executeQuery()) {
while (rs.next()) result.add(rs.getString("owner_login"));
}
}
return result;
}
public int markSynced(String ownerLogin, String eventId) throws SQLException {
String sql = """
UPDATE dm_sync_outbox
SET synced = TRUE, updated_at_ms = ?
WHERE LOWER(owner_login) = LOWER(?) AND event_id = ?
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setLong(1, System.currentTimeMillis());
ps.setString(2, ownerLogin);
ps.setString(3, eventId);
return ps.executeUpdate();
}
}
public int markSynced(String ownerLogin, List<String> eventIds) throws SQLException {
if (eventIds == null || eventIds.isEmpty()) return 0;
int total = 0;
try (Connection c = db.getConnection();
PreparedStatement ps = c.prepareStatement("""
UPDATE dm_sync_outbox
SET synced = TRUE, updated_at_ms = ?
WHERE LOWER(owner_login) = LOWER(?) AND event_id = ?
""")) {
long now = System.currentTimeMillis();
for (String eventId : eventIds) {
if (eventId == null || eventId.isBlank()) continue;
ps.setLong(1, now);
ps.setString(2, ownerLogin);
ps.setString(3, eventId.trim());
ps.addBatch();
}
for (int changed : ps.executeBatch()) if (changed > 0) total += changed;
}
return total;
}
public int markAllUnsynced(String ownerLogin) throws SQLException {
String sql = """
UPDATE dm_sync_outbox
SET synced = FALSE, updated_at_ms = ?
WHERE LOWER(owner_login) = LOWER(?)
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
ps.setLong(1, System.currentTimeMillis());
ps.setString(2, ownerLogin);
return ps.executeUpdate();
}
}
public int markAllUnsynced() throws SQLException {
try (Connection c = db.getConnection();
PreparedStatement ps = c.prepareStatement("UPDATE dm_sync_outbox SET synced = FALSE, updated_at_ms = ?")) {
ps.setLong(1, System.currentTimeMillis());
return ps.executeUpdate();
}
}
public int removeMissingMessages() throws SQLException {
String sql = """
DELETE FROM dm_sync_outbox o
WHERE NOT EXISTS (
SELECT 1 FROM signed_messages m WHERE m.message_key = o.primary_message_key
)
""";
try (Connection c = db.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) {
return ps.executeUpdate();
}
}
private DmSyncOutboxEntry mapRow(ResultSet rs) throws SQLException {
DmSyncOutboxEntry e = new DmSyncOutboxEntry();
e.setOwnerLogin(rs.getString("owner_login"));
e.setPrimaryMessageKey(rs.getString("primary_message_key"));
e.setEventId(rs.getString("event_id"));
e.setSecondaryMessageKey(rs.getString("secondary_message_key"));
e.setSynced(rs.getBoolean("synced"));
e.setCreatedAtMs(rs.getLong("created_at_ms"));
e.setUpdatedAtMs(rs.getLong("updated_at_ms"));
return e;
}
}
@@ -0,0 +1,84 @@
package shine.db.entities;
public class DmDeliveryStateEntry {
public static final int ACCEPTED = 0;
public static final int DELIVERED = 1;
public static final int FAILED = 3;
// Совместимость со строками БД, созданными ранней экспериментальной миграцией.
public static final int PENDING_NONE = ACCEPTED;
public static final int DELIVERED_ONE = DELIVERED;
public static final int DELIVERED_ALL = 2;
public static final int FAILED_FINAL = FAILED;
private String outgoingMessageKey;
private String eventId;
private String baseKey;
private String fromLogin;
private String toLogin;
private String incomingMessageKey;
private long createdAtMs;
private long deliveryExpiresAtMs;
private int deliveryState;
private String deliveredServerLogin;
private String recipientRoutesHash;
private int attemptIndex;
private Long nextAttemptAtMs;
private Long lastAttemptAtMs;
private String lastError;
private long updatedAtMs;
public String getOutgoingMessageKey() { return outgoingMessageKey; }
public void setOutgoingMessageKey(String value) { this.outgoingMessageKey = value; }
public String getEventId() { return eventId; }
public void setEventId(String value) { this.eventId = value; }
public String getBaseKey() { return baseKey; }
public void setBaseKey(String value) { this.baseKey = value; }
public String getFromLogin() { return fromLogin; }
public void setFromLogin(String value) { this.fromLogin = value; }
public String getToLogin() { return toLogin; }
public void setToLogin(String value) { this.toLogin = value; }
public String getIncomingMessageKey() { return incomingMessageKey; }
public void setIncomingMessageKey(String value) { this.incomingMessageKey = value; }
public long getCreatedAtMs() { return createdAtMs; }
public void setCreatedAtMs(long value) { this.createdAtMs = value; }
public long getDeliveryExpiresAtMs() { return deliveryExpiresAtMs; }
public void setDeliveryExpiresAtMs(long value) { this.deliveryExpiresAtMs = value; }
public int getDeliveryState() { return deliveryState; }
public void setDeliveryState(int value) { this.deliveryState = value; }
public String getDeliveredServerLogin() { return deliveredServerLogin; }
public void setDeliveredServerLogin(String value) { this.deliveredServerLogin = value; }
public String getRecipientRoutesHash() { return recipientRoutesHash; }
public void setRecipientRoutesHash(String value) { this.recipientRoutesHash = value; }
public int getAttemptIndex() { return attemptIndex; }
public void setAttemptIndex(int value) { this.attemptIndex = value; }
public Long getNextAttemptAtMs() { return nextAttemptAtMs; }
public void setNextAttemptAtMs(Long value) { this.nextAttemptAtMs = value; }
public Long getLastAttemptAtMs() { return lastAttemptAtMs; }
public void setLastAttemptAtMs(Long value) { this.lastAttemptAtMs = value; }
public String getLastError() { return lastError; }
public void setLastError(String value) { this.lastError = value; }
public long getUpdatedAtMs() { return updatedAtMs; }
public void setUpdatedAtMs(long value) { this.updatedAtMs = value; }
public String deliveryStateCode() {
return switch (deliveryState) {
case DELIVERED_ONE, DELIVERED_ALL -> "delivered";
case FAILED_FINAL -> "failed";
default -> "accepted";
};
}
public static int parseStateCode(String value) {
if (value == null) return ACCEPTED;
return switch (value.trim().toLowerCase()) {
case "delivered", "delivered_one", "delivered_all" -> DELIVERED;
case "failed", "failed_final" -> FAILED;
default -> ACCEPTED;
};
}
public boolean isDelivered() {
return deliveryState == DELIVERED_ONE || deliveryState == DELIVERED_ALL;
}
}
@@ -0,0 +1,26 @@
package shine.db.entities;
public class DmSyncOutboxEntry {
private String ownerLogin;
private String primaryMessageKey;
private String eventId;
private String secondaryMessageKey;
private boolean synced;
private long createdAtMs;
private long updatedAtMs;
public String getOwnerLogin() { return ownerLogin; }
public void setOwnerLogin(String value) { this.ownerLogin = value; }
public String getPrimaryMessageKey() { return primaryMessageKey; }
public void setPrimaryMessageKey(String value) { this.primaryMessageKey = value; }
public String getEventId() { return eventId; }
public void setEventId(String value) { this.eventId = value; }
public String getSecondaryMessageKey() { return secondaryMessageKey; }
public void setSecondaryMessageKey(String value) { this.secondaryMessageKey = value; }
public boolean isSynced() { return synced; }
public void setSynced(boolean value) { this.synced = value; }
public long getCreatedAtMs() { return createdAtMs; }
public void setCreatedAtMs(long value) { this.createdAtMs = value; }
public long getUpdatedAtMs() { return updatedAtMs; }
public void setUpdatedAtMs(long value) { this.updatedAtMs = value; }
}
@@ -0,0 +1,139 @@
BEGIN;
CREATE TABLE IF NOT EXISTS dm_delivery_state (
outgoing_message_key TEXT PRIMARY KEY,
event_id TEXT NOT NULL UNIQUE,
base_key TEXT NOT NULL,
from_login TEXT NOT NULL,
to_login TEXT NOT NULL,
incoming_message_key TEXT NOT NULL,
created_at_ms BIGINT NOT NULL,
delivery_expires_at_ms BIGINT NOT NULL,
delivery_state INTEGER NOT NULL DEFAULT 0 CHECK (delivery_state IN (0, 1, 2, 3)),
delivered_server_login TEXT,
recipient_routes_hash TEXT,
attempt_index INTEGER NOT NULL DEFAULT 0,
next_attempt_at_ms BIGINT,
last_attempt_at_ms BIGINT,
last_error TEXT,
updated_at_ms BIGINT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_dm_delivery_state_due
ON dm_delivery_state(next_attempt_at_ms, delivery_state)
WHERE delivery_state = 0 AND next_attempt_at_ms IS NOT NULL;
CREATE INDEX IF NOT EXISTS idx_dm_delivery_state_base
ON dm_delivery_state(base_key, from_login);
CREATE TABLE IF NOT EXISTS dm_sync_outbox (
owner_login TEXT NOT NULL,
primary_message_key TEXT NOT NULL,
event_id TEXT NOT NULL,
secondary_message_key TEXT,
synced BOOLEAN NOT NULL DEFAULT FALSE,
created_at_ms BIGINT NOT NULL,
updated_at_ms BIGINT NOT NULL,
PRIMARY KEY (owner_login, primary_message_key)
);
CREATE INDEX IF NOT EXISTS idx_dm_sync_outbox_unsynced
ON dm_sync_outbox(owner_login, created_at_ms, primary_message_key)
WHERE synced = FALSE;
-- До v12 факт межсерверной доставки не сохранялся. Старые исходящие пары
-- считаем уже доставленными, чтобы обновление не вызвало повторную рассылку.
INSERT INTO dm_delivery_state (
outgoing_message_key, event_id, base_key, from_login, to_login,
incoming_message_key, created_at_ms, delivery_expires_at_ms,
delivery_state, delivered_server_login, recipient_routes_hash,
attempt_index, next_attempt_at_ms, last_attempt_at_ms, last_error, updated_at_ms
)
SELECT
outgoing.message_key,
outgoing.message_key || ':' || outgoing.revision_time_ms || ':' || outgoing.reencrypted_at_ms,
outgoing.base_key,
outgoing.from_login,
outgoing.to_login,
incoming.message_key,
outgoing.created_at_ms,
outgoing.created_at_ms + 3600000,
2,
NULL,
NULL,
9,
NULL,
NULL,
'MIGRATED_AS_DELIVERED',
CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
FROM signed_messages outgoing
JOIN signed_messages incoming
ON incoming.base_key = outgoing.base_key
AND incoming.message_type = outgoing.message_type - 1
WHERE outgoing.message_type IN (2, 4)
ON CONFLICT (outgoing_message_key) DO NOTHING;
-- Историю до v12 считаем подтверждённой, иначе само обновление вызовет
-- массовую повторную передачу. При замене сервера общий reset вернёт FALSE.
INSERT INTO dm_sync_outbox (
owner_login, primary_message_key, event_id, secondary_message_key,
synced, created_at_ms, updated_at_ms
)
SELECT
LOWER(outgoing.from_login),
outgoing.message_key,
outgoing.message_key || ':' || outgoing.revision_time_ms || ':' || outgoing.reencrypted_at_ms,
incoming.message_key,
TRUE,
outgoing.created_at_ms,
CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
FROM signed_messages outgoing
JOIN signed_messages incoming
ON incoming.base_key = outgoing.base_key
AND incoming.message_type = outgoing.message_type - 1
WHERE outgoing.message_type IN (2, 4)
ON CONFLICT (owner_login, primary_message_key) DO NOTHING;
-- Входящая копия синхронизируется отдельно между серверами получателя.
INSERT INTO dm_sync_outbox (
owner_login, primary_message_key, event_id, secondary_message_key,
synced, created_at_ms, updated_at_ms
)
SELECT
LOWER(incoming.to_login),
incoming.message_key,
incoming.message_key || ':' || incoming.revision_time_ms || ':' || incoming.reencrypted_at_ms,
NULL,
TRUE,
incoming.created_at_ms,
CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
FROM signed_messages incoming
WHERE incoming.message_type IN (1, 3)
ON CONFLICT (owner_login, primary_message_key) DO NOTHING;
-- Tombstone принадлежит обоим участникам и должен попасть на второй сервер
-- каждого из них. Для self-DM конфликт безопасно схлопывается.
INSERT INTO dm_sync_outbox (
owner_login, primary_message_key, event_id, secondary_message_key,
synced, created_at_ms, updated_at_ms
)
SELECT
LOWER(owner_login),
tombstone.message_key,
tombstone.message_key || ':' || tombstone.revision_time_ms || ':' || tombstone.reencrypted_at_ms,
NULL,
TRUE,
tombstone.created_at_ms,
CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
FROM signed_messages tombstone
CROSS JOIN LATERAL (VALUES (tombstone.from_login), (tombstone.to_login)) owners(owner_login)
WHERE tombstone.message_type IN (5, 6, 7, 8)
ON CONFLICT (owner_login, primary_message_key) DO NOTHING;
INSERT INTO db_schema_version (id, schema_version, updated_at_ms)
VALUES (1, 12, 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;
@@ -0,0 +1,46 @@
BEGIN;
-- Нормализация экспериментальной схемы доставки v12:
-- 0=accepted, 1=delivered, 3=failed. Старое delivered_all (2)
-- объединяется с delivered, поскольку ACK одного сервера теперь достаточен.
UPDATE dm_delivery_state
SET delivery_state = 1,
delivered_server_login = NULL,
next_attempt_at_ms = NULL,
updated_at_ms = CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
WHERE delivery_state = 2;
-- Для незавершённых записей перестраиваем только будущую очередь. Повторная
-- передача безопасна благодаря messageKey и идемпотентному приёму.
UPDATE dm_delivery_state
SET delivery_expires_at_ms = created_at_ms + 3600000,
attempt_index = CASE
WHEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT) >= created_at_ms + 3600000 THEN 4
WHEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT) >= created_at_ms + 1500000 THEN 3
WHEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT) >= created_at_ms + 300000 THEN 2
ELSE 1
END,
next_attempt_at_ms = CASE
WHEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT) >= created_at_ms + 3600000
THEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
WHEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT) >= created_at_ms + 1500000
THEN created_at_ms + 3600000
WHEN CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT) >= created_at_ms + 300000
THEN created_at_ms + 1500000
ELSE created_at_ms + 30000
END,
updated_at_ms = CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)
WHERE delivery_state = 0;
DROP INDEX IF EXISTS idx_dm_delivery_state_due;
CREATE INDEX idx_dm_delivery_state_due
ON dm_delivery_state(next_attempt_at_ms, delivery_state)
WHERE delivery_state = 0 AND next_attempt_at_ms IS NOT NULL;
INSERT INTO db_schema_version (id, schema_version, updated_at_ms)
VALUES (1, 13, 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;
@@ -22,7 +22,7 @@ CREATE TABLE IF NOT EXISTS db_schema_version (
); );
INSERT INTO db_schema_version (id, schema_version, updated_at_ms) INSERT INTO db_schema_version (id, schema_version, updated_at_ms)
VALUES (1, 11, CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT)) VALUES (1, 13, CAST(EXTRACT(EPOCH FROM clock_timestamp()) * 1000 AS BIGINT))
ON CONFLICT (id) DO UPDATE SET ON CONFLICT (id) DO UPDATE SET
schema_version = EXCLUDED.schema_version, schema_version = EXCLUDED.schema_version,
updated_at_ms = EXCLUDED.updated_at_ms; updated_at_ms = EXCLUDED.updated_at_ms;
@@ -808,6 +808,51 @@ CREATE TABLE IF NOT EXISTS dm_sync_peer_state (
CREATE INDEX IF NOT EXISTS idx_dm_sync_peer_state_owner CREATE INDEX IF NOT EXISTS idx_dm_sync_peer_state_owner
ON dm_sync_peer_state(owner_login); ON dm_sync_peer_state(owner_login);
-- Изменяемое состояние доставки исходящей пары. Подписанные блоки остаются
-- неизменяемыми; эта таблица описывает только сетевую доставку пары 1/2 или 3/4.
CREATE TABLE IF NOT EXISTS dm_delivery_state (
outgoing_message_key TEXT PRIMARY KEY,
event_id TEXT NOT NULL UNIQUE,
base_key TEXT NOT NULL,
from_login TEXT NOT NULL,
to_login TEXT NOT NULL,
incoming_message_key TEXT NOT NULL,
created_at_ms BIGINT NOT NULL,
delivery_expires_at_ms BIGINT NOT NULL,
delivery_state INTEGER NOT NULL DEFAULT 0 CHECK (delivery_state IN (0, 1, 2, 3)),
delivered_server_login TEXT,
recipient_routes_hash TEXT,
attempt_index INTEGER NOT NULL DEFAULT 0,
next_attempt_at_ms BIGINT,
last_attempt_at_ms BIGINT,
last_error TEXT,
updated_at_ms BIGINT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_dm_delivery_state_due
ON dm_delivery_state(next_attempt_at_ms, delivery_state)
WHERE delivery_state = 0 AND next_attempt_at_ms IS NOT NULL;
CREATE INDEX IF NOT EXISTS idx_dm_delivery_state_base
ON dm_delivery_state(base_key, from_login);
-- У одного пользователя теперь максимум один второй access-сервер, поэтому
-- достаточно одного флага ACK на событие. Курсор по времени больше не нужен.
CREATE TABLE IF NOT EXISTS dm_sync_outbox (
owner_login TEXT NOT NULL,
primary_message_key TEXT NOT NULL,
event_id TEXT NOT NULL,
secondary_message_key TEXT,
synced BOOLEAN NOT NULL DEFAULT FALSE,
created_at_ms BIGINT NOT NULL,
updated_at_ms BIGINT NOT NULL,
PRIMARY KEY (owner_login, primary_message_key)
);
CREATE INDEX IF NOT EXISTS idx_dm_sync_outbox_unsynced
ON dm_sync_outbox(owner_login, created_at_ms, primary_message_key)
WHERE synced = FALSE;
CREATE TABLE IF NOT EXISTS message_views_state ( CREATE TABLE IF NOT EXISTS message_views_state (
viewer_login TEXT NOT NULL, viewer_login TEXT NOT NULL,
to_bch_name TEXT NOT NULL, to_bch_name TEXT NOT NULL,
@@ -98,6 +98,7 @@ import server.logic.ws_protocol.JSON.messages.Net_CallInviteBroadcast_Handler;
import server.logic.ws_protocol.JSON.messages.Net_CallSignalToSession_Handler; import server.logic.ws_protocol.JSON.messages.Net_CallSignalToSession_Handler;
import server.logic.ws_protocol.JSON.messages.Net_DeleteConversation_Handler; import server.logic.ws_protocol.JSON.messages.Net_DeleteConversation_Handler;
import server.logic.ws_protocol.JSON.messages.Net_DeleteMessage_Handler; import server.logic.ws_protocol.JSON.messages.Net_DeleteMessage_Handler;
import server.logic.ws_protocol.JSON.messages.Net_GetDmDeliveryStatus_Handler;
import server.logic.ws_protocol.JSON.messages.Net_MarkAllUserSettingsUnsynced_Handler; import server.logic.ws_protocol.JSON.messages.Net_MarkAllUserSettingsUnsynced_Handler;
import server.logic.ws_protocol.JSON.messages.Net_DmSyncBatch_Handler; import server.logic.ws_protocol.JSON.messages.Net_DmSyncBatch_Handler;
import server.logic.ws_protocol.JSON.messages.Net_GetDirectMessages_Handler; import server.logic.ws_protocol.JSON.messages.Net_GetDirectMessages_Handler;
@@ -112,6 +113,7 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_CallInviteBroadcast_R
import server.logic.ws_protocol.JSON.messages.entyties.Net_CallSignalToSession_Request; import server.logic.ws_protocol.JSON.messages.entyties.Net_CallSignalToSession_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteConversation_Request; import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteConversation_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteMessage_Request; import server.logic.ws_protocol.JSON.messages.entyties.Net_DeleteMessage_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDmDeliveryStatus_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_MarkAllUserSettingsUnsynced_Request; import server.logic.ws_protocol.JSON.messages.entyties.Net_MarkAllUserSettingsUnsynced_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_DmSyncBatch_Request; import server.logic.ws_protocol.JSON.messages.entyties.Net_DmSyncBatch_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Request; import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Request;
@@ -220,6 +222,7 @@ public final class JsonHandlerRegistry {
Map.entry("DeleteMessage", new Net_DeleteMessage_Handler()), Map.entry("DeleteMessage", new Net_DeleteMessage_Handler()),
Map.entry("DeleteConversation", new Net_DeleteConversation_Handler()), Map.entry("DeleteConversation", new Net_DeleteConversation_Handler()),
Map.entry("DmSyncBatch", new Net_DmSyncBatch_Handler()), Map.entry("DmSyncBatch", new Net_DmSyncBatch_Handler()),
Map.entry("GetDmDeliveryStatus", new Net_GetDmDeliveryStatus_Handler()),
Map.entry("UserSettingsSyncBatch", new Net_UserSettingsSyncBatch_Handler()), Map.entry("UserSettingsSyncBatch", new Net_UserSettingsSyncBatch_Handler()),
Map.entry("MarkAllUserSettingsUnsynced", new Net_MarkAllUserSettingsUnsynced_Handler()), Map.entry("MarkAllUserSettingsUnsynced", new Net_MarkAllUserSettingsUnsynced_Handler()),
Map.entry("GetDirectMessages", new Net_GetDirectMessages_Handler()), Map.entry("GetDirectMessages", new Net_GetDirectMessages_Handler()),
@@ -310,6 +313,7 @@ public final class JsonHandlerRegistry {
Map.entry("DeleteMessage", Net_DeleteMessage_Request.class), Map.entry("DeleteMessage", Net_DeleteMessage_Request.class),
Map.entry("DeleteConversation", Net_DeleteConversation_Request.class), Map.entry("DeleteConversation", Net_DeleteConversation_Request.class),
Map.entry("DmSyncBatch", Net_DmSyncBatch_Request.class), Map.entry("DmSyncBatch", Net_DmSyncBatch_Request.class),
Map.entry("GetDmDeliveryStatus", Net_GetDmDeliveryStatus_Request.class),
Map.entry("UserSettingsSyncBatch", Net_UserSettingsSyncBatch_Request.class), Map.entry("UserSettingsSyncBatch", Net_UserSettingsSyncBatch_Request.class),
Map.entry("MarkAllUserSettingsUnsynced", Net_MarkAllUserSettingsUnsynced_Request.class), Map.entry("MarkAllUserSettingsUnsynced", Net_MarkAllUserSettingsUnsynced_Request.class),
Map.entry("GetDirectMessages", Net_GetDirectMessages_Request.class), Map.entry("GetDirectMessages", Net_GetDirectMessages_Request.class),
@@ -0,0 +1,12 @@
package server.logic.ws_protocol.JSON.messages;
import shine.db.entities.SignedMessageEntry;
public final class DmDeliveryIds {
private DmDeliveryIds() {}
public static String forEntry(SignedMessageEntry entry) {
if (entry == null) throw new IllegalArgumentException("EMPTY_MESSAGE");
return entry.getMessageKey() + ":" + entry.getRevisionTimeMs() + ":" + entry.getReencryptedAtMs();
}
}
@@ -0,0 +1,25 @@
package server.logic.ws_protocol.JSON.messages;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import server.logic.ws_protocol.JSON.ActiveConnectionsRegistry;
import server.logic.ws_protocol.JSON.ConnectionContext;
import server.logic.ws_protocol.JSON.push.WsEventSender;
import shine.db.entities.DmDeliveryStateEntry;
public final class DmDeliveryRealtime {
private static final ObjectMapper MAPPER = new ObjectMapper();
private DmDeliveryRealtime() {}
public static void notifySender(DmDeliveryStateEntry state) {
if (state == null || state.getFromLogin() == null || state.getFromLogin().isBlank()) return;
ObjectNode payload = MAPPER.createObjectNode();
payload.put("baseKey", state.getBaseKey());
payload.put("outgoingKey", state.getOutgoingMessageKey());
payload.put("deliveryState", state.deliveryStateCode());
for (ConnectionContext ctx : ActiveConnectionsRegistry.getInstance().getByLogin(state.getFromLogin())) {
WsEventSender.sendEvent(ctx, "DmDeliveryStateChanged", state.getOutgoingMessageKey(), payload);
}
}
}
@@ -1,7 +1,13 @@
package server.logic.ws_protocol.JSON.messages; package server.logic.ws_protocol.JSON.messages;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.entities.DmDeliveryStateEntry;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
import server.sync.DmDeliveryCoordinator;
import java.util.List;
public final class DmSyncApplySupport { public final class DmSyncApplySupport {
private DmSyncApplySupport() {} private DmSyncApplySupport() {}
@@ -47,4 +53,56 @@ public final class DmSyncApplySupport {
int messageType, int messageType,
SignedMessagesDAO.ApplyStatus status SignedMessagesDAO.ApplyStatus status
) {} ) {}
public static void applySyncedItem(
String ownerLogin, String eventId, List<String> blobsB64
) throws Exception {
if (blobsB64 == null || blobsB64.isEmpty() || blobsB64.size() > 2) {
throw new IllegalArgumentException("BAD_BLOB_COUNT");
}
if (blobsB64.size() == 1) {
ApplyResult result = applySyncedBlob(ownerLogin, blobsB64.get(0));
SignedMessageEntry stored = SignedMessagesDAO.getInstance().getByMessageKey(result.messageKey());
long createdAt = stored == null ? System.currentTimeMillis() : stored.getCreatedAtMs();
DmSyncOutboxDAO.getInstance().upsert(
ownerLogin, result.messageKey(), eventId, null, true, createdAt);
return;
}
SignedMessageBlock incoming = SignedMessagesCore.parseFromB64(blobsB64.get(0));
SignedMessageBlock outgoing = SignedMessagesCore.parseFromB64(blobsB64.get(1));
SignedMessagesCore.validatePair(incoming, outgoing);
SignedMessagesCore.verifyUsersAndSignature(incoming);
SignedMessagesCore.verifyUsersAndSignature(outgoing);
if (!outgoing.fromLogin.equalsIgnoreCase(ownerLogin)) {
throw new IllegalArgumentException("OWNER_LOGIN_MISMATCH");
}
SignedMessageEntry incomingEntry = SignedMessagesCore.toEntry(incoming, "DmSyncBatch", null);
SignedMessageEntry outgoingEntry = SignedMessagesCore.toEntry(outgoing, "DmSyncBatch", null);
if (incoming.isContentType()) {
SignedMessagesDAO.getInstance().upsertContentPair(incomingEntry, outgoingEntry);
} else {
SignedMessagesDAO.getInstance().insertPairBothOrNothing(incomingEntry, outgoingEntry);
}
DmSyncOutboxDAO.getInstance().upsert(
ownerLogin, outgoingEntry.getMessageKey(), eventId,
incomingEntry.getMessageKey(), true, outgoingEntry.getCreatedAtMs());
long now = System.currentTimeMillis();
long signedAt = Math.max(outgoing.timeMs,
Math.max(outgoing.revisionTimeMs, outgoing.reencryptedAtMs));
long acceptedAt = signedAt > 0L ? Math.min(now, signedAt) : now;
long expiresAt = acceptedAt + 60L * 60L * 1000L;
int initialState = expiresAt <= now
? DmDeliveryStateEntry.FAILED_FINAL
: DmDeliveryStateEntry.PENDING_NONE;
DmDeliveryStateEntry delivery = DmDeliveryStateDAO.getInstance().upsertPair(
outgoingEntry.getMessageKey(), eventId, outgoingEntry.getBaseKey(),
outgoingEntry.getFromLogin(), outgoingEntry.getToLogin(), incomingEntry.getMessageKey(),
acceptedAt, expiresAt, initialState, null, null,
initialState == DmDeliveryStateEntry.PENDING_NONE);
if (delivery != null && initialState == DmDeliveryStateEntry.PENDING_NONE) {
DmDeliveryCoordinator.assistReceivedPairAsync(eventId);
}
}
} }
@@ -10,6 +10,8 @@ import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import server.sync.DmFederationService; import server.sync.DmFederationService;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
public class Net_DeleteConversation_Handler implements JsonMessageHandler { public class Net_DeleteConversation_Handler implements JsonMessageHandler {
@@ -43,6 +45,8 @@ public class Net_DeleteConversation_Handler implements JsonMessageHandler {
SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters(); SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters();
if (status.applied()) { if (status.applied()) {
DmDeliveryStateDAO.getInstance().removeMissingMessages();
recordOutboxForLocalOwners(entry);
counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, block); counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, block);
} }
if (status.applied() && ctx != null && ctx.isAuthenticatedUser()) { if (status.applied() && ctx != null && ctx.isAuthenticatedUser()) {
@@ -63,4 +67,15 @@ public class Net_DeleteConversation_Handler implements JsonMessageHandler {
private boolean isBlank(String s) { private boolean isBlank(String s) {
return s == null || s.isBlank(); return s == null || s.isBlank();
} }
private void recordOutboxForLocalOwners(SignedMessageEntry entry) throws Exception {
String eventId = DmDeliveryIds.forEntry(entry);
if (server.sync.DmDeliveryCoordinator.isLocalAccessServer(entry.getFromLogin())) {
DmSyncOutboxDAO.getInstance().upsert(entry.getFromLogin(), entry.getMessageKey(), eventId, null, false, entry.getCreatedAtMs());
}
if (!entry.getToLogin().equalsIgnoreCase(entry.getFromLogin())
&& server.sync.DmDeliveryCoordinator.isLocalAccessServer(entry.getToLogin())) {
DmSyncOutboxDAO.getInstance().upsert(entry.getToLogin(), entry.getMessageKey(), eventId, null, false, entry.getCreatedAtMs());
}
}
} }
@@ -10,6 +10,8 @@ import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import server.sync.DmFederationService; import server.sync.DmFederationService;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
public class Net_DeleteMessage_Handler implements JsonMessageHandler { public class Net_DeleteMessage_Handler implements JsonMessageHandler {
@@ -43,6 +45,8 @@ public class Net_DeleteMessage_Handler implements JsonMessageHandler {
SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters(); SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters();
if (status.applied()) { if (status.applied()) {
DmDeliveryStateDAO.getInstance().removeByBaseKey(entry.getBaseKey());
recordOutboxForLocalOwners(entry);
counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, block); counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, block);
} }
if (status.applied() && ctx != null && ctx.isAuthenticatedUser()) { if (status.applied() && ctx != null && ctx.isAuthenticatedUser()) {
@@ -63,4 +67,15 @@ public class Net_DeleteMessage_Handler implements JsonMessageHandler {
private boolean isBlank(String s) { private boolean isBlank(String s) {
return s == null || s.isBlank(); return s == null || s.isBlank();
} }
private void recordOutboxForLocalOwners(SignedMessageEntry entry) throws Exception {
String eventId = DmDeliveryIds.forEntry(entry);
if (server.sync.DmDeliveryCoordinator.isLocalAccessServer(entry.getFromLogin())) {
DmSyncOutboxDAO.getInstance().upsert(entry.getFromLogin(), entry.getMessageKey(), eventId, null, false, entry.getCreatedAtMs());
}
if (!entry.getToLogin().equalsIgnoreCase(entry.getFromLogin())
&& server.sync.DmDeliveryCoordinator.isLocalAccessServer(entry.getToLogin())) {
DmSyncOutboxDAO.getInstance().upsert(entry.getToLogin(), entry.getMessageKey(), eventId, null, false, entry.getCreatedAtMs());
}
}
} }
@@ -8,8 +8,10 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_DmSyncBatch_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_DmSyncBatch_Response; import server.logic.ws_protocol.JSON.messages.entyties.Net_DmSyncBatch_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory; import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.dao.UserAccessServersCurrentDAO; import shine.db.dao.UserAccessServersCurrentDAO;
import shine.db.entities.DmSyncOutboxEntry;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
import shine.db.entities.UserAccessServerRouteEntry; import shine.db.entities.UserAccessServerRouteEntry;
import utils.config.AppConfig; import utils.config.AppConfig;
@@ -18,6 +20,7 @@ import java.util.ArrayList;
import java.util.Base64; import java.util.Base64;
import java.util.List; import java.util.List;
/** Pull синхронизация только событий synced=false с ACK предыдущей страницы. */
public class Net_DmSyncBatch_Handler implements JsonMessageHandler { public class Net_DmSyncBatch_Handler implements JsonMessageHandler {
private static final int DEFAULT_LIMIT = 500; private static final int DEFAULT_LIMIT = 500;
private static final int MAX_LIMIT = 500; private static final int MAX_LIMIT = 500;
@@ -46,13 +49,14 @@ public class Net_DmSyncBatch_Handler implements JsonMessageHandler {
long afterStoredAtMs = Math.max(0L, req.getAfterStoredAtMs() == null ? 0L : req.getAfterStoredAtMs()); long afterStoredAtMs = Math.max(0L, req.getAfterStoredAtMs() == null ? 0L : req.getAfterStoredAtMs());
String afterMessageKey = req.getAfterMessageKey() == null ? "" : req.getAfterMessageKey().trim(); String afterMessageKey = req.getAfterMessageKey() == null ? "" : req.getAfterMessageKey().trim();
SignedMessagesDAO.SyncBatch batch = SignedMessagesDAO.getInstance().listSyncBatch( if (req.getAckSyncIds() != null && req.getAckSyncIds().size() > MAX_LIMIT) {
ownerLogin, return NetExceptionResponseFactory.error(req, WireCodes.Status.BAD_REQUEST,
afterStoredAtMs, "TOO_MANY_ACKS", "ackSyncIds содержит слишком много элементов");
afterMessageKey, }
limit, DmSyncOutboxDAO outbox = DmSyncOutboxDAO.getInstance();
maxBytes outbox.markSynced(ownerLogin, req.getAckSyncIds());
); List<DmSyncOutboxEntry> batch = outbox.listUnsynced(
ownerLogin, afterStoredAtMs, afterMessageKey, limit);
Net_DmSyncBatch_Response resp = new Net_DmSyncBatch_Response(); Net_DmSyncBatch_Response resp = new Net_DmSyncBatch_Response();
resp.setOp(req.getOp()); resp.setOp(req.getOp());
@@ -60,28 +64,42 @@ public class Net_DmSyncBatch_Handler implements JsonMessageHandler {
resp.setStatus(WireCodes.Status.OK); resp.setStatus(WireCodes.Status.OK);
resp.setOwnerLogin(ownerLogin); resp.setOwnerLogin(ownerLogin);
resp.setLimit(limit); resp.setLimit(limit);
resp.setRawBytes(batch.rawBytes()); resp.setRawBytes(0);
resp.setHasMore(batch.hasMore()); resp.setHasMore(false);
resp.setNextStoredAtMs(afterStoredAtMs); resp.setNextStoredAtMs(afterStoredAtMs);
resp.setNextMessageKey(afterMessageKey); resp.setNextMessageKey(afterMessageKey);
List<Net_DmSyncBatch_Response.Item> items = new ArrayList<>(); List<Net_DmSyncBatch_Response.Item> items = new ArrayList<>();
Base64.Encoder encoder = Base64.getEncoder(); Base64.Encoder encoder = Base64.getEncoder();
for (SignedMessageEntry entry : batch.items()) { int rawBytes = 0;
Net_DmSyncBatch_Response.Item item = new Net_DmSyncBatch_Response.Item(); for (DmSyncOutboxEntry row : batch) {
item.setMessageKey(entry.getMessageKey()); SignedMessageEntry primary = SignedMessagesDAO.getInstance().getByMessageKey(row.getPrimaryMessageKey());
item.setBaseKey(entry.getBaseKey()); if (primary == null || primary.getRawBlock() == null) continue;
item.setTargetLogin(entry.getTargetLogin()); SignedMessageEntry secondary = null;
item.setFromLogin(entry.getFromLogin()); if (row.getSecondaryMessageKey() != null && !row.getSecondaryMessageKey().isBlank()) {
item.setToLogin(entry.getToLogin()); secondary = SignedMessagesDAO.getInstance().getByMessageKey(row.getSecondaryMessageKey());
item.setMessageType(entry.getMessageType()); if (secondary == null || secondary.getRawBlock() == null) continue;
item.setTimeMs(entry.getTimeMs());
item.setStoredAtMs(entry.getCreatedAtMs());
item.setBlobB64(encoder.encodeToString(entry.getRawBlock()));
items.add(item);
resp.setNextStoredAtMs(entry.getCreatedAtMs());
resp.setNextMessageKey(entry.getMessageKey());
} }
int itemBytes = primary.getRawBlock().length + (secondary == null ? 0 : secondary.getRawBlock().length);
if (!items.isEmpty() && rawBytes + itemBytes > maxBytes) {
resp.setHasMore(true);
break;
}
Net_DmSyncBatch_Response.Item item = new Net_DmSyncBatch_Response.Item();
item.setSyncId(row.getEventId());
item.setPrimaryMessageKey(row.getPrimaryMessageKey());
item.setStoredAtMs(row.getCreatedAtMs());
List<String> blobs = new ArrayList<>();
if (secondary != null) blobs.add(encoder.encodeToString(secondary.getRawBlock()));
blobs.add(encoder.encodeToString(primary.getRawBlock()));
item.setBlobsB64(blobs);
items.add(item);
rawBytes += itemBytes;
resp.setNextStoredAtMs(row.getCreatedAtMs());
resp.setNextMessageKey(row.getPrimaryMessageKey());
}
if (!resp.isHasMore() && batch.size() >= limit) resp.setHasMore(true);
resp.setRawBytes(rawBytes);
resp.setItems(items); resp.setItems(items);
return resp; return resp;
} }
@@ -89,9 +107,7 @@ public class Net_DmSyncBatch_Handler implements JsonMessageHandler {
private boolean isLocalAccessServer(String ownerLogin, String ownServerLogin) throws Exception { private boolean isLocalAccessServer(String ownerLogin, String ownServerLogin) throws Exception {
for (UserAccessServerRouteEntry route : UserAccessServersCurrentDAO.getInstance().listByUserLogin(ownerLogin)) { for (UserAccessServerRouteEntry route : UserAccessServersCurrentDAO.getInstance().listByUserLogin(ownerLogin)) {
if (route == null || route.getServerLogin() == null) continue; if (route == null || route.getServerLogin() == null) continue;
if (ownServerLogin.equals(normalize(route.getServerLogin()))) { if (ownServerLogin.equals(normalize(route.getServerLogin()))) return true;
return true;
}
} }
return false; return false;
} }
@@ -11,11 +11,15 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDirectMessages_Res
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory; import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.entities.DmDeliveryStateEntry;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Base64; import java.util.Base64;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
public class Net_GetDirectMessages_Handler implements JsonMessageHandler { public class Net_GetDirectMessages_Handler implements JsonMessageHandler {
private static final Logger log = LoggerFactory.getLogger(Net_GetDirectMessages_Handler.class); private static final Logger log = LoggerFactory.getLogger(Net_GetDirectMessages_Handler.class);
@@ -55,6 +59,11 @@ public class Net_GetDirectMessages_Handler implements JsonMessageHandler {
if (hasMore) { if (hasMore) {
page = new ArrayList<>(page.subList(0, limit)); page = new ArrayList<>(page.subList(0, limit));
} }
Map<String, DmDeliveryStateEntry> deliveryByKey = DmDeliveryStateDAO.getInstance()
.listByOutgoingMessageKeys(page.stream()
.filter(entry -> entry.getMessageType() == SignedMessageBlock.TYPE_OUTGOING_COPY)
.map(SignedMessageEntry::getMessageKey)
.collect(Collectors.toList()));
Net_GetDirectMessages_Response resp = new Net_GetDirectMessages_Response(); Net_GetDirectMessages_Response resp = new Net_GetDirectMessages_Response();
resp.setOp(req.getOp()); resp.setOp(req.getOp());
@@ -80,6 +89,10 @@ public class Net_GetDirectMessages_Handler implements JsonMessageHandler {
item.setCreatedAtMs(entry.getCreatedAtMs()); item.setCreatedAtMs(entry.getCreatedAtMs());
item.setReadAtMs(entry.getReadAtMs()); item.setReadAtMs(entry.getReadAtMs());
item.setBlobB64(Base64.getEncoder().encodeToString(entry.getRawBlock())); item.setBlobB64(Base64.getEncoder().encodeToString(entry.getRawBlock()));
DmDeliveryStateEntry delivery = deliveryByKey.get(entry.getMessageKey());
if (delivery != null) {
item.setDeliveryState(delivery.deliveryStateCode());
}
items.add(item); items.add(item);
} }
resp.setMessages(items); resp.setMessages(items);
@@ -0,0 +1,35 @@
package server.logic.ws_protocol.JSON.messages;
import server.logic.ws_protocol.JSON.ConnectionContext;
import server.logic.ws_protocol.JSON.entyties.Net_Request;
import server.logic.ws_protocol.JSON.entyties.Net_Response;
import server.logic.ws_protocol.JSON.handlers.JsonMessageHandler;
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDmDeliveryStatus_Request;
import server.logic.ws_protocol.JSON.messages.entyties.Net_GetDmDeliveryStatus_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.entities.DmDeliveryStateEntry;
/** Read-only проверка: доставил ли peer хотя бы одну копию получателю. */
public class Net_GetDmDeliveryStatus_Handler implements JsonMessageHandler {
@Override
public Net_Response handle(Net_Request baseRequest, ConnectionContext ctx) throws Exception {
Net_GetDmDeliveryStatus_Request req = (Net_GetDmDeliveryStatus_Request) baseRequest;
if (req.getMessageKey() == null || req.getMessageKey().isBlank()) {
return NetExceptionResponseFactory.error(
req, WireCodes.Status.BAD_REQUEST, "EMPTY_MESSAGE_KEY", "messageKey обязателен");
}
String messageKey = req.getMessageKey().trim();
DmDeliveryStateEntry state = DmDeliveryStateDAO.getInstance().getByOutgoingMessageKey(messageKey);
Net_GetDmDeliveryStatus_Response resp = new Net_GetDmDeliveryStatus_Response();
resp.setOp(req.getOp());
resp.setRequestId(req.getRequestId());
resp.setStatus(WireCodes.Status.OK);
resp.setMessageKey(messageKey);
resp.setKnown(state != null);
resp.setDelivered(state != null && state.isDelivered());
return resp;
}
}
@@ -12,6 +12,8 @@ import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import shine.db.DbController; import shine.db.DbController;
import shine.db.dao.UserSettingsDAO; import shine.db.dao.UserSettingsDAO;
import shine.db.dao.DmSyncOutboxDAO;
import server.sync.DmSyncWakeSignal;
import java.sql.Connection; import java.sql.Connection;
@@ -23,16 +25,21 @@ public class Net_MarkAllUserSettingsUnsynced_Handler implements JsonMessageHandl
Net_MarkAllUserSettingsUnsynced_Request req = (Net_MarkAllUserSettingsUnsynced_Request) baseRequest; Net_MarkAllUserSettingsUnsynced_Request req = (Net_MarkAllUserSettingsUnsynced_Request) baseRequest;
try (Connection c = DbController.getInstance().getConnection()) { try (Connection c = DbController.getInstance().getConnection()) {
int updated; int updated;
int dmUpdated;
if (req.getLogin() == null || req.getLogin().isBlank()) { if (req.getLogin() == null || req.getLogin().isBlank()) {
updated = UserSettingsDAO.getInstance().markAllUnsynced(c); updated = UserSettingsDAO.getInstance().markAllUnsynced(c);
dmUpdated = DmSyncOutboxDAO.getInstance().markAllUnsynced();
} else { } else {
updated = UserSettingsDAO.getInstance().markAllUnsynced(c, req.getLogin().trim()); updated = UserSettingsDAO.getInstance().markAllUnsynced(c, req.getLogin().trim());
dmUpdated = DmSyncOutboxDAO.getInstance().markAllUnsynced(req.getLogin().trim());
} }
DmSyncWakeSignal.request();
Net_MarkAllUserSettingsUnsynced_Response resp = new Net_MarkAllUserSettingsUnsynced_Response(); Net_MarkAllUserSettingsUnsynced_Response resp = new Net_MarkAllUserSettingsUnsynced_Response();
resp.setOp(req.getOp()); resp.setOp(req.getOp());
resp.setRequestId(req.getRequestId()); resp.setRequestId(req.getRequestId());
resp.setStatus(WireCodes.Status.OK); resp.setStatus(WireCodes.Status.OK);
resp.setUpdated(updated); resp.setUpdated(updated);
resp.setDmUpdated(dmUpdated);
return resp; return resp;
} catch (Exception e) { } catch (Exception e) {
log.error("MarkAllUserSettingsUnsynced failed", e); log.error("MarkAllUserSettingsUnsynced failed", e);
@@ -8,9 +8,12 @@ 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.messages.entyties.Net_ReceiveIncomingMessage_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory; import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import server.sync.DmDeliveryCoordinator;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
import java.util.Arrays;
import java.util.Base64; import java.util.Base64;
public class Net_ReceiveIncomingMessage_Handler implements JsonMessageHandler { public class Net_ReceiveIncomingMessage_Handler implements JsonMessageHandler {
@@ -40,6 +43,10 @@ public class Net_ReceiveIncomingMessage_Handler implements JsonMessageHandler {
return NetExceptionResponseFactory.error(req, status, code, "Сообщение не прошло проверку"); return NetExceptionResponseFactory.error(req, status, code, "Сообщение не прошло проверку");
} }
if (!DmDeliveryCoordinator.isLocalAccessServer(incoming.toLogin)) {
return NetExceptionResponseFactory.error(req, 403, "LOCAL_SERVER_NOT_ACCESS_SERVER", "Сервер не обслуживает получателя сообщения");
}
final SignedMessageEntry entry; final SignedMessageEntry entry;
try { try {
entry = SignedMessagesCore.toEntry(incoming, "ReceiveIncomingMessage", null); entry = SignedMessagesCore.toEntry(incoming, "ReceiveIncomingMessage", null);
@@ -53,16 +60,27 @@ public class Net_ReceiveIncomingMessage_Handler implements JsonMessageHandler {
SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters(); SignedMessagesRealtime.DeliveryCounters counters = new SignedMessagesRealtime.DeliveryCounters();
if (status.applied()) { if (status.applied()) {
counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, incoming); counters = SignedMessagesRealtime.deliverToRelevantSessions(entry, incoming);
server.sync.DmFederationService.fanOutIncomingToRecipientAccessServers(
incoming.toLogin,
req.getIncomingBlobB64().trim(),
req.getSourceServerLogin()
);
} }
if (status == SignedMessagesDAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) { if (status == SignedMessagesDAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) {
bounceConversationDeleteIfKnown(incoming.fromLogin, incoming.toLogin); bounceConversationDeleteIfKnown(incoming.fromLogin, incoming.toLogin);
return NetExceptionResponseFactory.error(req, 409, "BLOCKED_BY_CONVERSATION_TOMBSTONE", "Переписка уже удалена этой ревизией");
} }
SignedMessageEntry stored = SignedMessagesDAO.getInstance().getByMessageKey(entry.getMessageKey());
if (stored == null || !Arrays.equals(stored.getRawBlock(), entry.getRawBlock())) {
return NetExceptionResponseFactory.error(req, 409, "STALE_MESSAGE_REVISION", "На сервере уже есть более новая ревизия сообщения");
}
boolean receivedFromRecipientPeer = DmDeliveryCoordinator.sourceIsOtherAccessServer(
incoming.toLogin, req.getSourceServerLogin());
DmSyncOutboxDAO.getInstance().upsert(
incoming.toLogin,
entry.getMessageKey(),
DmDeliveryIds.forEntry(entry),
null,
receivedFromRecipientPeer,
entry.getCreatedAtMs()
);
Net_ReceiveIncomingMessage_Response resp = new Net_ReceiveIncomingMessage_Response(); Net_ReceiveIncomingMessage_Response resp = new Net_ReceiveIncomingMessage_Response();
resp.setOp(req.getOp()); resp.setOp(req.getOp());
resp.setRequestId(req.getRequestId()); resp.setRequestId(req.getRequestId());
@@ -8,10 +8,15 @@ import server.logic.ws_protocol.JSON.messages.entyties.Net_SendMessagePair_Reque
import server.logic.ws_protocol.JSON.messages.entyties.Net_SendMessagePair_Response; import server.logic.ws_protocol.JSON.messages.entyties.Net_SendMessagePair_Response;
import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory; import server.logic.ws_protocol.JSON.utils.NetExceptionResponseFactory;
import server.logic.ws_protocol.WireCodes; import server.logic.ws_protocol.WireCodes;
import server.sync.DmDeliveryCoordinator;
import server.sync.DmFederationService; import server.sync.DmFederationService;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.dao.SignedMessagesDAO; import shine.db.dao.SignedMessagesDAO;
import shine.db.entities.DmDeliveryStateEntry;
import shine.db.entities.SignedMessageEntry; import shine.db.entities.SignedMessageEntry;
import java.util.Arrays;
import java.util.Base64; import java.util.Base64;
public class Net_SendMessagePair_Handler implements JsonMessageHandler { public class Net_SendMessagePair_Handler implements JsonMessageHandler {
@@ -43,8 +48,10 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
SignedMessageEntry incomingEntry; SignedMessageEntry incomingEntry;
SignedMessageEntry outgoingEntry; SignedMessageEntry outgoingEntry;
boolean fromPeer = "ReceiveOutcomingMessage".equalsIgnoreCase(req.getOp())
|| !isBlank(req.getSourceServerLogin());
try { try {
String sourceApi = "SendMessagePair"; String sourceApi = fromPeer ? "ReceiveOutcomingMessage" : "SendMessagePair";
String originSessionId = (ctx != null && !isBlank(ctx.getSessionId())) ? ctx.getSessionId() : null; String originSessionId = (ctx != null && !isBlank(ctx.getSessionId())) ? ctx.getSessionId() : null;
incomingEntry = SignedMessagesCore.toEntry(incoming, sourceApi, originSessionId); incomingEntry = SignedMessagesCore.toEntry(incoming, sourceApi, originSessionId);
outgoingEntry = SignedMessagesCore.toEntry(outgoing, sourceApi, originSessionId); outgoingEntry = SignedMessagesCore.toEntry(outgoing, sourceApi, originSessionId);
@@ -79,15 +86,53 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
if (pairStatus == SignedMessagesDAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) { if (pairStatus == SignedMessagesDAO.ApplyStatus.BLOCKED_BY_CONVERSATION_TOMBSTONE) {
bounceConversationDeleteIfKnown(incoming.fromLogin, incoming.toLogin); bounceConversationDeleteIfKnown(incoming.fromLogin, incoming.toLogin);
return NetExceptionResponseFactory.error(req, 409, "BLOCKED_BY_CONVERSATION_TOMBSTONE", "Переписка уже удалена этой ревизией");
} }
if (pairStatus.applied() && ctx != null && ctx.isAuthenticatedUser()) { SignedMessageEntry storedOutgoing = SignedMessagesDAO.getInstance().getByMessageKey(outgoingEntry.getMessageKey());
DmFederationService.fanOutPair( if (storedOutgoing == null || !Arrays.equals(storedOutgoing.getRawBlock(), outgoingEntry.getRawBlock())) {
incoming.fromLogin, return NetExceptionResponseFactory.error(req, 409, "STALE_MESSAGE_REVISION", "На сервере уже есть более новая ревизия сообщения");
incoming.toLogin, }
req.getIncomingBlobB64().trim(),
req.getOutgoingBlobB64().trim() String eventId = DmDeliveryIds.forEntry(outgoingEntry);
long nowMs = System.currentTimeMillis();
long acceptedAtMs;
long signedAtMs = Math.max(outgoing.timeMs,
Math.max(outgoing.revisionTimeMs, outgoing.reencryptedAtMs));
acceptedAtMs = fromPeer && signedAtMs > 0L ? Math.min(nowMs, signedAtMs) : nowMs;
long expiresAtMs = acceptedAtMs + 60L * 60L * 1000L;
int initialState = DmDeliveryStateEntry.ACCEPTED;
DmDeliveryStateEntry delivery = DmDeliveryStateDAO.getInstance().upsertPair(
outgoingEntry.getMessageKey(), eventId, outgoingEntry.getBaseKey(),
outgoingEntry.getFromLogin(), outgoingEntry.getToLogin(), incomingEntry.getMessageKey(),
acceptedAtMs, expiresAtMs, initialState,
null, null, fromPeer && initialState == DmDeliveryStateEntry.PENDING_NONE
); );
DmSyncOutboxDAO.getInstance().upsert(
outgoingEntry.getFromLogin(), outgoingEntry.getMessageKey(), eventId,
incomingEntry.getMessageKey(), fromPeer, acceptedAtMs);
// Если этот же сервер также обслуживает получателя, его входящая копия
// имеет отдельный sync-флаг владельца-получателя.
if (DmDeliveryCoordinator.isLocalAccessServer(incomingEntry.getToLogin())) {
boolean incomingAlreadySynced = DmDeliveryCoordinator.sourceIsOtherAccessServer(
incomingEntry.getToLogin(), req.getSourceServerLogin());
DmSyncOutboxDAO.getInstance().upsert(
incomingEntry.getToLogin(), incomingEntry.getMessageKey(), DmDeliveryIds.forEntry(incomingEntry),
null, incomingAlreadySynced, acceptedAtMs);
}
if (delivery != null && pairStatus.applied()) {
if (fromPeer) {
DmDeliveryCoordinator.assistReceivedPairAsync(delivery.getEventId());
} else {
// Первая доставка выполняется до ответа клиенту. Два сервера
// получателя вызываются параллельно внутри координатора.
DmDeliveryCoordinator.processDueEntry(delivery);
delivery = DmDeliveryStateDAO.getInstance()
.getByOutgoingMessageKey(outgoingEntry.getMessageKey());
}
} }
Net_SendMessagePair_Response resp = new Net_SendMessagePair_Response(); Net_SendMessagePair_Response resp = new Net_SendMessagePair_Response();
@@ -99,6 +144,7 @@ public class Net_SendMessagePair_Handler implements JsonMessageHandler {
resp.setOutgoingKey(outgoingEntry.getMessageKey()); resp.setOutgoingKey(outgoingEntry.getMessageKey());
resp.setDeliveredWsSessions(inCounters.wsDelivered + outCounters.wsDelivered); resp.setDeliveredWsSessions(inCounters.wsDelivered + outCounters.wsDelivered);
resp.setDeliveredWebPushSessions(inCounters.pushDelivered + outCounters.pushDelivered); resp.setDeliveredWebPushSessions(inCounters.pushDelivered + outCounters.pushDelivered);
resp.setDeliveryState(delivery.deliveryStateCode());
return resp; return resp;
} }
@@ -2,12 +2,16 @@ package server.logic.ws_protocol.JSON.messages.entyties;
import server.logic.ws_protocol.JSON.entyties.Net_Request; import server.logic.ws_protocol.JSON.entyties.Net_Request;
import java.util.ArrayList;
import java.util.List;
public class Net_DmSyncBatch_Request extends Net_Request { public class Net_DmSyncBatch_Request extends Net_Request {
private String ownerLogin; private String ownerLogin;
private Long afterStoredAtMs; private Long afterStoredAtMs;
private String afterMessageKey; private String afterMessageKey;
private Integer limit; private Integer limit;
private Integer maxBytes; private Integer maxBytes;
private List<String> ackSyncIds = new ArrayList<>();
public String getOwnerLogin() { return ownerLogin; } public String getOwnerLogin() { return ownerLogin; }
public void setOwnerLogin(String ownerLogin) { this.ownerLogin = ownerLogin; } public void setOwnerLogin(String ownerLogin) { this.ownerLogin = ownerLogin; }
@@ -19,4 +23,8 @@ public class Net_DmSyncBatch_Request extends Net_Request {
public void setLimit(Integer limit) { this.limit = limit; } public void setLimit(Integer limit) { this.limit = limit; }
public Integer getMaxBytes() { return maxBytes; } public Integer getMaxBytes() { return maxBytes; }
public void setMaxBytes(Integer maxBytes) { this.maxBytes = maxBytes; } public void setMaxBytes(Integer maxBytes) { this.maxBytes = maxBytes; }
public List<String> getAckSyncIds() { return ackSyncIds; }
public void setAckSyncIds(List<String> ackSyncIds) {
this.ackSyncIds = ackSyncIds == null ? new ArrayList<>() : ackSyncIds;
}
} }
@@ -30,33 +30,20 @@ public class Net_DmSyncBatch_Response extends Net_Response {
public void setItems(List<Item> items) { this.items = items; } public void setItems(List<Item> items) { this.items = items; }
public static class Item { public static class Item {
private String messageKey; private String syncId;
private String baseKey; private String primaryMessageKey;
private String targetLogin;
private String fromLogin;
private String toLogin;
private int messageType;
private long timeMs;
private long storedAtMs; private long storedAtMs;
private String blobB64; private List<String> blobsB64 = new ArrayList<>();
public String getMessageKey() { return messageKey; } public String getSyncId() { return syncId; }
public void setMessageKey(String messageKey) { this.messageKey = messageKey; } public void setSyncId(String syncId) { this.syncId = syncId; }
public String getBaseKey() { return baseKey; } public String getPrimaryMessageKey() { return primaryMessageKey; }
public void setBaseKey(String baseKey) { this.baseKey = baseKey; } public void setPrimaryMessageKey(String primaryMessageKey) { this.primaryMessageKey = primaryMessageKey; }
public String getTargetLogin() { return targetLogin; }
public void setTargetLogin(String targetLogin) { this.targetLogin = targetLogin; }
public String getFromLogin() { return fromLogin; }
public void setFromLogin(String fromLogin) { this.fromLogin = fromLogin; }
public String getToLogin() { return toLogin; }
public void setToLogin(String toLogin) { this.toLogin = toLogin; }
public int getMessageType() { return messageType; }
public void setMessageType(int messageType) { this.messageType = messageType; }
public long getTimeMs() { return timeMs; }
public void setTimeMs(long timeMs) { this.timeMs = timeMs; }
public long getStoredAtMs() { return storedAtMs; } public long getStoredAtMs() { return storedAtMs; }
public void setStoredAtMs(long storedAtMs) { this.storedAtMs = storedAtMs; } public void setStoredAtMs(long storedAtMs) { this.storedAtMs = storedAtMs; }
public String getBlobB64() { return blobB64; } public List<String> getBlobsB64() { return blobsB64; }
public void setBlobB64(String blobB64) { this.blobB64 = blobB64; } public void setBlobsB64(List<String> blobsB64) {
this.blobsB64 = blobsB64 == null ? new ArrayList<>() : blobsB64;
}
} }
} }
@@ -42,6 +42,7 @@ public class Net_GetDirectMessages_Response extends Net_Response {
private long createdAtMs; private long createdAtMs;
private Long readAtMs; private Long readAtMs;
private String blobB64; private String blobB64;
private String deliveryState;
public String getMessageKey() { return messageKey; } public String getMessageKey() { return messageKey; }
public void setMessageKey(String messageKey) { this.messageKey = messageKey; } public void setMessageKey(String messageKey) { this.messageKey = messageKey; }
@@ -67,5 +68,7 @@ public class Net_GetDirectMessages_Response extends Net_Response {
public void setReadAtMs(Long readAtMs) { this.readAtMs = readAtMs; } public void setReadAtMs(Long readAtMs) { this.readAtMs = readAtMs; }
public String getBlobB64() { return blobB64; } public String getBlobB64() { return blobB64; }
public void setBlobB64(String blobB64) { this.blobB64 = blobB64; } public void setBlobB64(String blobB64) { this.blobB64 = blobB64; }
public String getDeliveryState() { return deliveryState; }
public void setDeliveryState(String deliveryState) { this.deliveryState = deliveryState; }
} }
} }
@@ -0,0 +1,10 @@
package server.logic.ws_protocol.JSON.messages.entyties;
import server.logic.ws_protocol.JSON.entyties.Net_Request;
public class Net_GetDmDeliveryStatus_Request extends Net_Request {
private String messageKey;
public String getMessageKey() { return messageKey; }
public void setMessageKey(String value) { this.messageKey = value; }
}
@@ -0,0 +1,16 @@
package server.logic.ws_protocol.JSON.messages.entyties;
import server.logic.ws_protocol.JSON.entyties.Net_Response;
public class Net_GetDmDeliveryStatus_Response extends Net_Response {
private String messageKey;
private boolean known;
private boolean delivered;
public String getMessageKey() { return messageKey; }
public void setMessageKey(String value) { this.messageKey = value; }
public boolean isKnown() { return known; }
public void setKnown(boolean value) { this.known = value; }
public boolean isDelivered() { return delivered; }
public void setDelivered(boolean value) { this.delivered = value; }
}
@@ -4,7 +4,10 @@ import server.logic.ws_protocol.JSON.entyties.Net_Response;
public class Net_MarkAllUserSettingsUnsynced_Response extends Net_Response { public class Net_MarkAllUserSettingsUnsynced_Response extends Net_Response {
private Integer updated; private Integer updated;
private Integer dmUpdated;
public Integer getUpdated() { return updated; } public Integer getUpdated() { return updated; }
public void setUpdated(Integer updated) { this.updated = updated; } public void setUpdated(Integer updated) { this.updated = updated; }
public Integer getDmUpdated() { return dmUpdated; }
public void setDmUpdated(Integer dmUpdated) { this.dmUpdated = dmUpdated; }
} }
@@ -8,6 +8,7 @@ public class Net_SendMessagePair_Response extends Net_Response {
private String outgoingKey; private String outgoingKey;
private int deliveredWsSessions; private int deliveredWsSessions;
private int deliveredWebPushSessions; private int deliveredWebPushSessions;
private String deliveryState;
public String getBaseKey() { return baseKey; } public String getBaseKey() { return baseKey; }
public void setBaseKey(String baseKey) { this.baseKey = baseKey; } public void setBaseKey(String baseKey) { this.baseKey = baseKey; }
@@ -19,4 +20,6 @@ public class Net_SendMessagePair_Response extends Net_Response {
public void setDeliveredWsSessions(int deliveredWsSessions) { this.deliveredWsSessions = deliveredWsSessions; } public void setDeliveredWsSessions(int deliveredWsSessions) { this.deliveredWsSessions = deliveredWsSessions; }
public int getDeliveredWebPushSessions() { return deliveredWebPushSessions; } public int getDeliveredWebPushSessions() { return deliveredWebPushSessions; }
public void setDeliveredWebPushSessions(int deliveredWebPushSessions) { this.deliveredWebPushSessions = deliveredWebPushSessions; } public void setDeliveredWebPushSessions(int deliveredWebPushSessions) { this.deliveredWebPushSessions = deliveredWebPushSessions; }
public String getDeliveryState() { return deliveryState; }
public void setDeliveryState(String deliveryState) { this.deliveryState = deliveryState; }
} }
@@ -0,0 +1,309 @@
package server.sync;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import server.logic.ws_protocol.JSON.messages.DmDeliveryRealtime;
import shine.db.dao.DmDeliveryStateDAO;
import shine.db.dao.DmSyncOutboxDAO;
import shine.db.dao.SignedMessagesDAO;
import shine.db.dao.UserAccessServersCurrentDAO;
import shine.db.entities.DmDeliveryStateEntry;
import shine.db.entities.SignedMessageEntry;
import shine.db.entities.UserAccessServerRouteEntry;
import utils.config.AppConfig;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.util.ArrayList;
import java.util.Base64;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.ThreadFactory;
/**
* Координатор доставки исходящей DM-пары. Первая попытка выполняется до ответа
* клиенту, последующие пятисекундным воркером.
*/
public final class DmDeliveryCoordinator {
private static final Logger log = LoggerFactory.getLogger(DmDeliveryCoordinator.class);
private static final long[] ATTEMPT_OFFSETS_MS = {
0L,
30_000L,
5L * 60_000L,
25L * 60_000L,
60L * 60_000L
};
private static final String SERVER_LOGIN_CONFIG = "server.SHiNE.login";
private static final RemoteDmSyncClient REMOTE = new RemoteDmSyncClient();
private static final DmDeliveryStateDAO DELIVERY_DAO = DmDeliveryStateDAO.getInstance();
private static final DmSyncOutboxDAO OUTBOX_DAO = DmSyncOutboxDAO.getInstance();
private static final SignedMessagesDAO MESSAGES_DAO = SignedMessagesDAO.getInstance();
private static final UserAccessServersCurrentDAO ACCESS_DAO = UserAccessServersCurrentDAO.getInstance();
private static final ExecutorService ASSIST_EXECUTOR = new ThreadPoolExecutor(
2, 2, 0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(500),
daemonThreadFactory("dm-peer-assist"),
new ThreadPoolExecutor.DiscardPolicy());
private static final ExecutorService RECIPIENT_EXECUTOR = new ThreadPoolExecutor(
4, 16, 60L, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(1000),
daemonThreadFactory("dm-recipient-delivery"),
new ThreadPoolExecutor.CallerRunsPolicy());
private DmDeliveryCoordinator() {}
public static List<DmDeliveryStateEntry> listDue(int limit) throws Exception {
return DELIVERY_DAO.listDue(System.currentTimeMillis(), Math.max(1, limit));
}
public static void processDueEntry(DmDeliveryStateEntry snapshot) {
processDueEntry(snapshot, true);
}
private static void processDueEntry(DmDeliveryStateEntry snapshot, boolean allowInitialHandoff) {
if (snapshot == null) return;
try {
long now = System.currentTimeMillis();
Long nextAttemptAtMs = nextAttemptAt(snapshot.getCreatedAtMs(), snapshot.getAttemptIndex() + 1);
if (!DELIVERY_DAO.claimAttempt(snapshot, now, nextAttemptAtMs)) return;
DmDeliveryStateEntry current = DELIVERY_DAO.getByEventId(snapshot.getEventId());
if (current == null || current.isDelivered()
|| current.getDeliveryState() == DmDeliveryStateEntry.FAILED) return;
boolean finalAttempt = snapshot.getAttemptIndex() >= ATTEMPT_OFFSETS_MS.length - 1
|| now >= current.getDeliveryExpiresAtMs();
// Перед попытками на 5-й, 25-й и 60-й минутах сначала спрашиваем
// второй сервер отправителя: возможно, он уже доставил сообщение.
if (current.getDeliveryState() == DmDeliveryStateEntry.ACCEPTED
&& (snapshot.getAttemptIndex() >= 2 || finalAttempt)) {
current = acceptPeerDeliveryStatus(current);
if (current != null && current.isDelivered()) {
DmDeliveryRealtime.notifySender(current);
return;
}
}
DmDeliveryStateEntry afterAttempt = attemptRecipientRoutes(
current, finalAttempt ? null : nextAttemptAtMs);
// После первой попытки сразу передаём полную пару второму серверу
// отправителя старой операцией ReceiveOutcomingMessage.
if (allowInitialHandoff && snapshot.getAttemptIndex() == 0 && afterAttempt != null) {
afterAttempt = handoffPairToSenderPeer(afterAttempt);
}
if (finalAttempt && afterAttempt != null && !afterAttempt.isDelivered()) {
afterAttempt = DELIVERY_DAO.finishAtExpiry(afterAttempt.getEventId(), System.currentTimeMillis());
}
DmDeliveryRealtime.notifySender(afterAttempt);
} catch (Exception e) {
log.warn("DM delivery attempt failed: eventId={}", snapshot.getEventId(), e);
}
}
public static void assistReceivedPairAsync(String eventId) {
if (eventId == null || eventId.isBlank()) return;
ASSIST_EXECUTOR.execute(() -> {
try {
DmDeliveryStateEntry row = DELIVERY_DAO.getByEventId(eventId);
if (row != null) processDueEntry(row, false);
} catch (Exception e) {
log.warn("DM peer assist failed: eventId={}", eventId, e);
}
});
}
private static DmDeliveryStateEntry attemptRecipientRoutes(
DmDeliveryStateEntry current,
Long nextAttemptAtMs
) throws Exception {
if (current == null) return null;
List<UserAccessServerRouteEntry> routes = cappedRoutes(current.getToLogin());
String routesHash = routesHash(routes);
String alreadyDelivered = normalize(current.getDeliveredServerLogin());
List<String> acceptedLogins = new ArrayList<>();
if (alreadyDelivered != null) acceptedLogins.add(alreadyDelivered);
SignedMessageEntry incoming = MESSAGES_DAO.getByMessageKey(current.getIncomingMessageKey());
if (incoming == null || incoming.getRawBlock() == null) {
return DELIVERY_DAO.updateAfterAttempt(
current.getEventId(), current.getDeliveryState(), current.getDeliveredServerLogin(),
routesHash, nextAttemptAtMs, "INCOMING_BLOB_NOT_FOUND", System.currentTimeMillis());
}
String incomingBlobB64 = Base64.getEncoder().encodeToString(incoming.getRawBlock());
String ownServerLogin = ownServerLogin();
java.util.concurrent.CompletionService<RouteAttempt> completion =
new java.util.concurrent.ExecutorCompletionService<>(RECIPIENT_EXECUTOR);
int submitted = 0;
for (UserAccessServerRouteEntry route : routes) {
String routeLogin = normalize(route.getServerLogin());
if (routeLogin == null || acceptedLogins.contains(routeLogin)) continue;
completion.submit(() -> {
try {
if (!routeLogin.equals(ownServerLogin)) {
REMOTE.receiveIncomingMessage(route.getServerUrl(), incomingBlobB64, ownServerLogin);
}
return new RouteAttempt(routeLogin, null);
} catch (Exception e) {
log.info("DM recipient server unavailable: messageKey={} server={}",
current.getOutgoingMessageKey(), routeLogin);
return new RouteAttempt(routeLogin, compactError(e));
}
});
submitted++;
}
String lastError = null;
for (int i = 0; i < submitted && acceptedLogins.isEmpty(); i++) {
RouteAttempt result = completion.take().get();
if (result.error() == null) {
acceptedLogins.add(result.serverLogin());
} else {
lastError = result.error();
}
}
int nextState = acceptedLogins.isEmpty()
? DmDeliveryStateEntry.ACCEPTED
: DmDeliveryStateEntry.DELIVERED;
String oneLogin = acceptedLogins.isEmpty() ? null : acceptedLogins.get(0);
return DELIVERY_DAO.updateAfterAttempt(
current.getEventId(), nextState, oneLogin, routesHash,
nextAttemptAtMs, lastError, System.currentTimeMillis());
}
private static DmDeliveryStateEntry handoffPairToSenderPeer(DmDeliveryStateEntry current) throws Exception {
UserAccessServerRouteEntry peer = senderPeer(current.getFromLogin());
if (peer == null) return current;
SignedMessageEntry incoming = MESSAGES_DAO.getByMessageKey(current.getIncomingMessageKey());
SignedMessageEntry outgoing = MESSAGES_DAO.getByMessageKey(current.getOutgoingMessageKey());
if (incoming == null || outgoing == null) return current;
try {
REMOTE.sendMessagePair(
peer.getServerUrl(),
Base64.getEncoder().encodeToString(incoming.getRawBlock()),
Base64.getEncoder().encodeToString(outgoing.getRawBlock()),
ownServerLogin());
OUTBOX_DAO.markSynced(current.getFromLogin(), current.getEventId());
} catch (Exception e) {
log.info("DM sender peer unavailable: eventId={} server={}", current.getEventId(), peer.getServerLogin());
}
return current;
}
private static DmDeliveryStateEntry acceptPeerDeliveryStatus(DmDeliveryStateEntry current) throws Exception {
UserAccessServerRouteEntry peer = senderPeer(current.getFromLogin());
if (peer == null) return current;
try {
RemoteDmSyncClient.RemoteDeliveryStatus remote = REMOTE.getDmDeliveryStatus(
peer.getServerUrl(), current.getOutgoingMessageKey());
if (remote.known() && remote.delivered()) {
return DELIVERY_DAO.markDeliveredFromPeer(current.getEventId(), System.currentTimeMillis());
}
return current;
} catch (Exception e) {
log.info("DM delivery status unavailable: eventId={} server={}", current.getEventId(), peer.getServerLogin());
return current;
}
}
private static UserAccessServerRouteEntry senderPeer(String senderLogin) throws Exception {
String own = ownServerLogin();
for (UserAccessServerRouteEntry route : cappedRoutes(senderLogin)) {
String routeLogin = normalize(route.getServerLogin());
if (routeLogin != null && !routeLogin.equals(own)) return route;
}
return null;
}
public static boolean isLocalAccessServer(String ownerLogin) throws Exception {
String own = ownServerLogin();
if (own == null) return false;
for (UserAccessServerRouteEntry route : cappedRoutes(ownerLogin)) {
if (own.equals(normalize(route.getServerLogin()))) return true;
}
return false;
}
public static boolean sourceIsOtherAccessServer(String ownerLogin, String sourceServerLogin) throws Exception {
String source = normalize(sourceServerLogin);
String own = ownServerLogin();
if (source == null || source.equals(own)) return false;
for (UserAccessServerRouteEntry route : cappedRoutes(ownerLogin)) {
if (source.equals(normalize(route.getServerLogin()))) return true;
}
return false;
}
public static String currentServerLogin() {
return ownServerLogin();
}
private static List<UserAccessServerRouteEntry> cappedRoutes(String ownerLogin) throws Exception {
Map<String, UserAccessServerRouteEntry> unique = new LinkedHashMap<>();
for (UserAccessServerRouteEntry route : ACCESS_DAO.listByUserLogin(ownerLogin)) {
if (route == null || route.getServerUrl() == null || route.getServerUrl().isBlank()) continue;
String login = normalize(route.getServerLogin());
if (login == null) continue;
unique.putIfAbsent(login, route);
if (unique.size() == 2) break;
}
return new ArrayList<>(unique.values());
}
private static String routesHash(List<UserAccessServerRouteEntry> routes) throws Exception {
List<String> logins = new ArrayList<>();
for (UserAccessServerRouteEntry route : routes) {
String login = normalize(route.getServerLogin());
if (login != null) logins.add(login);
}
logins.sort(String::compareTo);
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] hash = digest.digest(String.join("\n", logins).getBytes(StandardCharsets.UTF_8));
StringBuilder out = new StringBuilder(hash.length * 2);
for (byte b : hash) out.append(String.format("%02x", b & 0xff));
return out.toString();
}
private static Long nextAttemptAt(long createdAtMs, int nextAttemptIndex) {
if (nextAttemptIndex < 0 || nextAttemptIndex >= ATTEMPT_OFFSETS_MS.length) return null;
return createdAtMs + ATTEMPT_OFFSETS_MS[nextAttemptIndex];
}
private record RouteAttempt(String serverLogin, String error) {}
private static String ownServerLogin() {
return normalize(AppConfig.getInstance().getParam(SERVER_LOGIN_CONFIG));
}
private static String normalize(String value) {
if (value == null) return null;
String normalized = value.trim().toLowerCase(Locale.ROOT);
return normalized.isEmpty() ? null : normalized;
}
private static String compactError(Exception e) {
String text = String.valueOf(e == null ? "unknown" : e.getMessage());
return text.length() <= 500 ? text : text.substring(0, 500);
}
private static ThreadFactory daemonThreadFactory(String prefix) {
return new ThreadFactory() {
private int sequence;
@Override
public synchronized Thread newThread(Runnable r) {
Thread thread = new Thread(r, prefix + "-" + (++sequence));
thread.setDaemon(true);
return thread;
}
};
}
}
@@ -0,0 +1,17 @@
package server.sync;
import java.util.concurrent.atomic.AtomicBoolean;
public final class DmSyncWakeSignal {
private static final AtomicBoolean REQUESTED = new AtomicBoolean(false);
private DmSyncWakeSignal() {}
public static void request() {
REQUESTED.set(true);
}
public static boolean consume() {
return REQUESTED.getAndSet(false);
}
}
@@ -17,13 +17,19 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException; import java.util.concurrent.TimeoutException;
/** Клиент стабильных межсерверных DM-операций. */
public final class RemoteDmSyncClient { public final class RemoteDmSyncClient {
private static final ObjectMapper MAPPER = new ObjectMapper(); private static final ObjectMapper MAPPER = new ObjectMapper();
private static final HttpClient HTTP = HttpClient.newBuilder() private static final HttpClient HTTP = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(6)) .connectTimeout(Duration.ofSeconds(6))
.build(); .build();
public void sendMessagePair(String serverAddressRaw, String incomingBlobB64, String outgoingBlobB64, String sourceServerLogin) throws Exception { public void sendMessagePair(
String serverAddressRaw,
String incomingBlobB64,
String outgoingBlobB64,
String sourceServerLogin
) throws Exception {
String incomingJson = MAPPER.writeValueAsString(incomingBlobB64); String incomingJson = MAPPER.writeValueAsString(incomingBlobB64);
String outgoingJson = MAPPER.writeValueAsString(outgoingBlobB64); String outgoingJson = MAPPER.writeValueAsString(outgoingBlobB64);
String sourceServerLoginJson = toOptionalJsonField("sourceServerLogin", sourceServerLogin); String sourceServerLoginJson = toOptionalJsonField("sourceServerLogin", sourceServerLogin);
@@ -40,7 +46,11 @@ public final class RemoteDmSyncClient {
ensureOk("ReceiveOutcomingMessage", response); ensureOk("ReceiveOutcomingMessage", response);
} }
public void receiveIncomingMessage(String serverAddressRaw, String incomingBlobB64, String sourceServerLogin) throws Exception { public void receiveIncomingMessage(
String serverAddressRaw,
String incomingBlobB64,
String sourceServerLogin
) throws Exception {
String incomingJson = MAPPER.writeValueAsString(incomingBlobB64); String incomingJson = MAPPER.writeValueAsString(incomingBlobB64);
String sourceServerLoginJson = toOptionalJsonField("sourceServerLogin", sourceServerLogin); String sourceServerLoginJson = toOptionalJsonField("sourceServerLogin", sourceServerLogin);
JsonNode response = send(serverAddressRaw, """ JsonNode response = send(serverAddressRaw, """
@@ -55,17 +65,54 @@ public final class RemoteDmSyncClient {
ensureOk("ReceiveIncomingMessage", response); ensureOk("ReceiveIncomingMessage", response);
} }
public RemoteDeliveryStatus getDmDeliveryStatus(String serverAddressRaw, String messageKey) throws Exception {
String messageKeyJson = MAPPER.writeValueAsString(messageKey);
JsonNode response = send(serverAddressRaw, """
{
"op":"GetDmDeliveryStatus",
"requestId":%s,
"payload":{
"messageKey":%s
}
}
""".formatted("%s", messageKeyJson));
ensureOk("GetDmDeliveryStatus", response);
JsonNode payload = response.path("payload");
return new RemoteDeliveryStatus(
payload.path("messageKey").asText(messageKey),
payload.path("known").asBoolean(false),
payload.path("delivered").asBoolean(false)
);
}
public RemoteDmBatch dmSyncBatch( public RemoteDmBatch dmSyncBatch(
String serverAddressRaw, String serverAddressRaw,
String ownerLogin, String ownerLogin,
long afterStoredAtMs, long afterStoredAtMs,
String afterMessageKey, String afterMessageKey,
int limit, int limit,
int maxBytes int maxBytes,
List<String> ackSyncIds
) throws Exception {
try (RemoteSyncSession session = new RemoteSyncSession(serverAddressRaw)) {
return dmSyncBatch(session, ownerLogin, afterStoredAtMs, afterMessageKey,
limit, maxBytes, ackSyncIds);
}
}
public RemoteDmBatch dmSyncBatch(
RemoteSyncSession session,
String ownerLogin,
long afterStoredAtMs,
String afterMessageKey,
int limit,
int maxBytes,
List<String> ackSyncIds
) throws Exception { ) throws Exception {
String ownerLoginJson = MAPPER.writeValueAsString(ownerLogin); String ownerLoginJson = MAPPER.writeValueAsString(ownerLogin);
String afterMessageKeyJson = MAPPER.writeValueAsString(afterMessageKey == null ? "" : afterMessageKey); String afterMessageKeyJson = MAPPER.writeValueAsString(afterMessageKey == null ? "" : afterMessageKey);
JsonNode response = send(serverAddressRaw, """ String ackSyncIdsJson = MAPPER.writeValueAsString(ackSyncIds == null ? List.of() : ackSyncIds);
JsonNode response = session.send("""
{ {
"op":"DmSyncBatch", "op":"DmSyncBatch",
"requestId":%s, "requestId":%s,
@@ -74,10 +121,12 @@ public final class RemoteDmSyncClient {
"afterStoredAtMs":%d, "afterStoredAtMs":%d,
"afterMessageKey":%s, "afterMessageKey":%s,
"limit":%d, "limit":%d,
"maxBytes":%d "maxBytes":%d,
"ackSyncIds":%s
} }
} }
""".formatted("%s", ownerLoginJson, Math.max(0L, afterStoredAtMs), afterMessageKeyJson, limit, maxBytes)); """.formatted("%s", ownerLoginJson, Math.max(0L, afterStoredAtMs), afterMessageKeyJson,
limit, maxBytes, ackSyncIdsJson));
ensureOk("DmSyncBatch", response); ensureOk("DmSyncBatch", response);
JsonNode payload = response.path("payload"); JsonNode payload = response.path("payload");
@@ -85,11 +134,17 @@ public final class RemoteDmSyncClient {
JsonNode arr = payload.path("items"); JsonNode arr = payload.path("items");
if (arr.isArray()) { if (arr.isArray()) {
for (JsonNode item : arr) { for (JsonNode item : arr) {
List<String> blobs = new ArrayList<>();
JsonNode blobsNode = item.path("blobsB64");
if (blobsNode.isArray()) for (JsonNode blob : blobsNode) blobs.add(blob.asText(""));
if (blobs.isEmpty() && !item.path("blobB64").asText("").isBlank()) {
blobs.add(item.path("blobB64").asText(""));
}
items.add(new RemoteDmItem( items.add(new RemoteDmItem(
item.path("messageKey").asText(""), item.path("syncId").asText(""),
item.path("primaryMessageKey").asText(item.path("messageKey").asText("")),
item.path("storedAtMs").asLong(0L), item.path("storedAtMs").asLong(0L),
item.path("blobB64").asText("") blobs));
));
} }
} }
return new RemoteDmBatch( return new RemoteDmBatch(
@@ -106,9 +161,7 @@ public final class RemoteDmSyncClient {
{ {
"op":"DeleteMessage", "op":"DeleteMessage",
"requestId":%s, "requestId":%s,
"payload":{ "payload":{"blobB64":%s}
"blobB64":%s
}
} }
""".formatted("%s", blobJson)); """.formatted("%s", blobJson));
ensureOk("DeleteMessage", response); ensureOk("DeleteMessage", response);
@@ -120,9 +173,7 @@ public final class RemoteDmSyncClient {
{ {
"op":"DeleteConversation", "op":"DeleteConversation",
"requestId":%s, "requestId":%s,
"payload":{ "payload":{"blobB64":%s}
"blobB64":%s
}
} }
""".formatted("%s", blobJson)); """.formatted("%s", blobJson));
ensureOk("DeleteConversation", response); ensureOk("DeleteConversation", response);
@@ -132,24 +183,19 @@ public final class RemoteDmSyncClient {
String requestId = MAPPER.writeValueAsString("dm-sync-" + UUID.randomUUID()); String requestId = MAPPER.writeValueAsString("dm-sync-" + UUID.randomUUID());
String json = jsonTemplate.formatted(requestId); String json = jsonTemplate.formatted(requestId);
String wsUrl = RemoteBlockchainSyncClient.buildWsUrl(serverAddressRaw); String wsUrl = RemoteBlockchainSyncClient.buildWsUrl(serverAddressRaw);
if (wsUrl == null) { if (wsUrl == null) throw new IllegalArgumentException("Invalid server address: " + serverAddressRaw);
throw new IllegalArgumentException("Invalid server address: " + serverAddressRaw);
}
CompletableFuture<String> responseFuture = new CompletableFuture<>(); CompletableFuture<String> responseFuture = new CompletableFuture<>();
CountDownLatch openLatch = new CountDownLatch(1); CountDownLatch openLatch = new CountDownLatch(1);
SyncWsListener listener = new SyncWsListener(responseFuture, openLatch); SyncWsListener listener = new SyncWsListener(responseFuture, openLatch);
WebSocket webSocket = HTTP.newWebSocketBuilder() WebSocket webSocket = HTTP.newWebSocketBuilder()
.connectTimeout(Duration.ofSeconds(6)) .connectTimeout(Duration.ofSeconds(6))
.buildAsync(URI.create(wsUrl), listener) .buildAsync(URI.create(wsUrl), listener)
.get(8, TimeUnit.SECONDS); .get(8, TimeUnit.SECONDS);
if (!openLatch.await(8, TimeUnit.SECONDS)) { if (!openLatch.await(8, TimeUnit.SECONDS)) {
tryAbort(webSocket); tryAbort(webSocket);
throw new TimeoutException("WS open timeout"); throw new TimeoutException("WS open timeout");
} }
webSocket.sendText(json, true).get(8, TimeUnit.SECONDS); webSocket.sendText(json, true).get(8, TimeUnit.SECONDS);
String responseJson = responseFuture.get(12, TimeUnit.SECONDS); String responseJson = responseFuture.get(12, TimeUnit.SECONDS);
tryAbort(webSocket); tryAbort(webSocket);
@@ -157,9 +203,7 @@ public final class RemoteDmSyncClient {
} }
private String toOptionalJsonField(String fieldName, String value) throws Exception { private String toOptionalJsonField(String fieldName, String value) throws Exception {
if (value == null || value.isBlank()) { if (value == null || value.isBlank()) return "";
return "";
}
return ",\n \"" + fieldName + "\":" + MAPPER.writeValueAsString(value.trim()); return ",\n \"" + fieldName + "\":" + MAPPER.writeValueAsString(value.trim());
} }
@@ -172,27 +216,14 @@ public final class RemoteDmSyncClient {
} }
private static void tryAbort(WebSocket webSocket) { private static void tryAbort(WebSocket webSocket) {
try { try { webSocket.sendClose(WebSocket.NORMAL_CLOSURE, "ok"); } catch (Exception ignored) {}
webSocket.sendClose(WebSocket.NORMAL_CLOSURE, "ok"); try { webSocket.abort(); } catch (Exception ignored) {}
} catch (Exception ignored) {
}
try {
webSocket.abort();
} catch (Exception ignored) {
}
} }
public record RemoteDmBatch( public record RemoteDeliveryStatus(String messageKey, boolean known, boolean delivered) {}
long nextStoredAtMs, public record RemoteDmBatch(long nextStoredAtMs, String nextMessageKey, boolean hasMore, List<RemoteDmItem> items) {}
String nextMessageKey,
boolean hasMore,
List<RemoteDmItem> items
) {}
public record RemoteDmItem( public record RemoteDmItem(
String messageKey, String syncId, String primaryMessageKey, long storedAtMs, List<String> blobsB64
long storedAtMs,
String blobB64
) {} ) {}
private static final class SyncWsListener implements WebSocket.Listener { private static final class SyncWsListener implements WebSocket.Listener {
@@ -214,9 +245,7 @@ public final class RemoteDmSyncClient {
@Override @Override
public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) { public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) {
textBuffer.append(data); textBuffer.append(data);
if (last && !responseFuture.isDone()) { if (last && !responseFuture.isDone()) responseFuture.complete(textBuffer.toString());
responseFuture.complete(textBuffer.toString());
}
webSocket.request(1); webSocket.request(1);
return CompletableFuture.completedFuture(null); return CompletableFuture.completedFuture(null);
} }
@@ -230,16 +259,15 @@ public final class RemoteDmSyncClient {
@Override @Override
public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) { public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) {
if (!responseFuture.isDone()) { if (!responseFuture.isDone()) {
responseFuture.completeExceptionally(new IllegalStateException("WS closed before response: " + statusCode + " " + reason)); responseFuture.completeExceptionally(
new IllegalStateException("WS closed before response: " + statusCode + " " + reason));
} }
return CompletableFuture.completedFuture(null); return CompletableFuture.completedFuture(null);
} }
@Override @Override
public void onError(WebSocket webSocket, Throwable error) { public void onError(WebSocket webSocket, Throwable error) {
if (!responseFuture.isDone()) { if (!responseFuture.isDone()) responseFuture.completeExceptionally(error);
responseFuture.completeExceptionally(error);
}
openLatch.countDown(); openLatch.countDown();
} }
} }
@@ -0,0 +1,70 @@
package server.sync;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.WebSocket;
import java.nio.ByteBuffer;
import java.time.Duration;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/** Один последовательный WS-сеанс для синхронизации всех данных access-сервера. */
public final class RemoteSyncSession implements AutoCloseable {
private static final ObjectMapper MAPPER = new ObjectMapper();
private static final HttpClient HTTP = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(6)).build();
private final LinkedBlockingQueue<String> responses = new LinkedBlockingQueue<>();
private final WebSocket webSocket;
public RemoteSyncSession(String serverAddressRaw) throws Exception {
String wsUrl = RemoteBlockchainSyncClient.buildWsUrl(serverAddressRaw);
if (wsUrl == null) throw new IllegalArgumentException("Invalid server address: " + serverAddressRaw);
Listener listener = new Listener(responses);
webSocket = HTTP.newWebSocketBuilder().connectTimeout(Duration.ofSeconds(6))
.buildAsync(URI.create(wsUrl), listener).get(8, TimeUnit.SECONDS);
}
public synchronized JsonNode send(String jsonTemplate) throws Exception {
String requestId = MAPPER.writeValueAsString("access-sync-" + UUID.randomUUID());
webSocket.sendText(jsonTemplate.formatted(requestId), true).get(8, TimeUnit.SECONDS);
String json = responses.poll(12, TimeUnit.SECONDS);
if (json == null) throw new TimeoutException("WS response timeout");
return MAPPER.readTree(json);
}
@Override
public void close() {
try { webSocket.sendClose(WebSocket.NORMAL_CLOSURE, "ok"); } catch (Exception ignored) {}
}
private static final class Listener implements WebSocket.Listener {
private final LinkedBlockingQueue<String> responses;
private final StringBuilder text = new StringBuilder();
private Listener(LinkedBlockingQueue<String> responses) { this.responses = responses; }
@Override public void onOpen(WebSocket ws) { ws.request(1); }
@Override public CompletionStage<?> onText(WebSocket ws, CharSequence data, boolean last) {
text.append(data);
if (last) {
responses.offer(text.toString());
text.setLength(0);
}
ws.request(1);
return CompletableFuture.completedFuture(null);
}
@Override public CompletionStage<?> onBinary(WebSocket ws, ByteBuffer data, boolean last) {
ws.request(1);
return CompletableFuture.completedFuture(null);
}
@Override public void onError(WebSocket ws, Throwable error) {
responses.offer("{\"status\":500,\"code\":\"WS_ERROR\"}");
}
}
}
@@ -24,7 +24,13 @@ public final class RemoteUserSettingsSyncClient {
.build(); .build();
public void upsertUserSetting(String serverAddressRaw, UserSettingEntry entry, boolean syncDelivery) throws Exception { public void upsertUserSetting(String serverAddressRaw, UserSettingEntry entry, boolean syncDelivery) throws Exception {
JsonNode response = send(serverAddressRaw, """ try (RemoteSyncSession session = new RemoteSyncSession(serverAddressRaw)) {
upsertUserSetting(session, entry, syncDelivery);
}
}
public void upsertUserSetting(RemoteSyncSession session, UserSettingEntry entry, boolean syncDelivery) throws Exception {
JsonNode response = session.send("""
{ {
"op":"UpsertUserSetting", "op":"UpsertUserSetting",
"requestId":%s, "requestId":%s,
@@ -63,7 +69,20 @@ public final class RemoteUserSettingsSyncClient {
int limit, int limit,
int maxBytes int maxBytes
) throws Exception { ) throws Exception {
JsonNode response = send(serverAddressRaw, """ try (RemoteSyncSession session = new RemoteSyncSession(serverAddressRaw)) {
return userSettingsSyncBatch(session, ownerLogin, afterTimeMs, afterSettingKey, limit, maxBytes);
}
}
public RemoteUserSettingsBatch userSettingsSyncBatch(
RemoteSyncSession session,
String ownerLogin,
long afterTimeMs,
String afterSettingKey,
int limit,
int maxBytes
) throws Exception {
JsonNode response = session.send("""
{ {
"op":"UserSettingsSyncBatch", "op":"UserSettingsSyncBatch",
"requestId":%s, "requestId":%s,
@@ -2,165 +2,52 @@ package server.sync;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import server.logic.ws_protocol.JSON.messages.DmSyncApplySupport; import shine.db.entities.DmDeliveryStateEntry;
import shine.db.dao.DmSyncPeerStateDAO;
import shine.db.dao.UserAccessServersCurrentDAO;
import shine.db.entities.DmSyncPeerStateEntry;
import shine.db.entities.UserAccessServerRouteEntry;
import utils.config.AppConfig; import utils.config.AppConfig;
import java.util.LinkedHashSet; import java.util.concurrent.ArrayBlockingQueue;
import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicBoolean;
/** /** Пятисекундный воркер отвечает только за доставку DM получателю. */
* Периодическая догоняющая синхронизация личной переписки между access-серверами пользователя.
*/
public final class PeriodicDmSyncService { public final class PeriodicDmSyncService {
private static final Logger log = LoggerFactory.getLogger(PeriodicDmSyncService.class); private static final Logger log = LoggerFactory.getLogger(PeriodicDmSyncService.class);
private static final String SERVER_LOGIN_CONFIG = "server.SHiNE.login";
private static final AtomicBoolean STARTED = new AtomicBoolean(false); private static final AtomicBoolean STARTED = new AtomicBoolean(false);
private static final ScheduledExecutorService SCHEDULER =
private static final RemoteDmSyncClient REMOTE = new RemoteDmSyncClient(); Executors.newSingleThreadScheduledExecutor(daemonThreadFactory("dm-worker-dispatcher"));
private static final UserAccessServersCurrentDAO ACCESS_DAO = UserAccessServersCurrentDAO.getInstance(); private static final ThreadPoolExecutor DELIVERY_EXECUTOR = new ThreadPoolExecutor(
private static final DmSyncPeerStateDAO STATE_DAO = DmSyncPeerStateDAO.getInstance(); 4, 4, 0L, TimeUnit.MILLISECONDS, new ArrayBlockingQueue<>(500),
daemonThreadFactory("dm-delivery"), new ThreadPoolExecutor.DiscardPolicy());
private static final ScheduledExecutorService EXECUTOR = Executors.newSingleThreadScheduledExecutor(new ThreadFactory() {
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, "periodic-dm-sync");
t.setDaemon(true);
return t;
}
});
private PeriodicDmSyncService() {} private PeriodicDmSyncService() {}
public static void startOrLog() { public static void startOrLog() {
if (!isEnabled()) { if (!isEnabled()) {
log.info("Periodic DM sync disabled by dm.sync.enabled=false"); log.info("DM delivery worker disabled by dm.sync.enabled=false");
return; return;
} }
if (!STARTED.compareAndSet(false, true)) { if (!STARTED.compareAndSet(false, true)) return;
return; long pollSeconds = configLong("dm.worker.pollSeconds", 5L, 1L, 60L);
} SCHEDULER.scheduleWithFixedDelay(
long initialDelaySec = configLong("dm.sync.initialDelaySeconds", 60L, 0L, 3600L); PeriodicDmSyncService::tickSafe, 0L, pollSeconds, TimeUnit.SECONDS);
long periodHours = configLong("dm.sync.periodHours", 6L, 1L, 168L); log.info("DM delivery worker scheduled every {} seconds", pollSeconds);
EXECUTOR.scheduleWithFixedDelay(
PeriodicDmSyncService::runCycleSafe,
initialDelaySec,
TimeUnit.HOURS.toSeconds(periodHours),
TimeUnit.SECONDS
);
log.info("Periodic DM sync scheduled: first run in {} seconds, then every {} hours", initialDelaySec, periodHours);
} }
private static void runCycleSafe() { private static void tickSafe() {
try { try {
runCycle(); int limit = (int) configLong("dm.worker.dueLimit", 100L, 1L, 1000L);
for (DmDeliveryStateEntry row : DmDeliveryCoordinator.listDue(limit)) {
DELIVERY_EXECUTOR.execute(() -> DmDeliveryCoordinator.processDueEntry(row));
}
} catch (Exception e) { } catch (Exception e) {
log.error("Periodic DM sync failed unexpectedly", e); log.error("DM delivery dispatcher failed", e);
} }
} }
private static void runCycle() throws Exception {
String ownServerLogin = normalize(AppConfig.getInstance().getParam(SERVER_LOGIN_CONFIG));
if (ownServerLogin == null) {
log.warn("Periodic DM sync skipped: {} is empty", SERVER_LOGIN_CONFIG);
return;
}
List<String> ownersRaw = ACCESS_DAO.listUserLoginsByServerLogin(ownServerLogin);
Set<String> owners = new LinkedHashSet<>(ownersRaw);
if (owners.isEmpty()) {
log.info("Periodic DM sync skipped: no local access-server users for {}", ownServerLogin);
return;
}
int syncedPeers = 0;
int appliedEvents = 0;
for (String ownerLogin : owners) {
List<UserAccessServerRouteEntry> routes = ACCESS_DAO.listByUserLogin(ownerLogin);
for (UserAccessServerRouteEntry route : routes) {
if (route == null) continue;
String remoteLogin = normalize(route.getServerLogin());
String remoteUrl = route.getServerUrl();
if (remoteLogin == null || remoteUrl == null || remoteUrl.isBlank()) continue;
if (remoteLogin.equals(ownServerLogin)) continue;
try {
appliedEvents += syncOwnerFromRemote(ownerLogin, route);
syncedPeers++;
} catch (Exception e) {
STATE_DAO.updateError(ownerLogin, route.getServerLogin(), remoteUrl, String.valueOf(e));
log.warn("Periodic DM sync peer failed: owner={} remoteServer={} reason={}",
ownerLogin, route.getServerLogin(), String.valueOf(e));
}
}
}
log.info("Periodic DM sync cycle finished: owners={} syncedPeers={} appliedEvents={}",
owners.size(), syncedPeers, appliedEvents);
}
private static int syncOwnerFromRemote(String ownerLogin, UserAccessServerRouteEntry route) throws Exception {
int limit = (int) configLong("dm.sync.batchLimit", 500L, 1L, 500L);
int maxBytes = (int) configLong("dm.sync.batchMaxBytes", 3_000_000L, 64_000L, 5_000_000L);
int maxPages = (int) configLong("dm.sync.maxPagesPerPeer", 50L, 1L, 500L);
DmSyncPeerStateEntry state = STATE_DAO.getOrCreate(ownerLogin, route.getServerLogin(), route.getServerUrl());
long cursorStoredAtMs = state.getCursorStoredAtMs();
String cursorMessageKey = state.getCursorMessageKey() == null ? "" : state.getCursorMessageKey();
int applied = 0;
boolean bootstrapCompleted = false;
for (int page = 0; page < maxPages; page++) {
RemoteDmSyncClient.RemoteDmBatch batch = REMOTE.dmSyncBatch(
route.getServerUrl(),
ownerLogin,
cursorStoredAtMs,
cursorMessageKey,
limit,
maxBytes
);
for (RemoteDmSyncClient.RemoteDmItem item : batch.items()) {
if (item == null || item.blobB64() == null || item.blobB64().isBlank()) continue;
DmSyncApplySupport.ApplyResult result = DmSyncApplySupport.applySyncedBlob(ownerLogin, item.blobB64());
if (result.status().applied()) {
applied++;
}
}
cursorStoredAtMs = Math.max(cursorStoredAtMs, batch.nextStoredAtMs());
cursorMessageKey = batch.nextMessageKey() == null ? "" : batch.nextMessageKey();
bootstrapCompleted = !batch.hasMore();
STATE_DAO.updateSuccess(
ownerLogin,
route.getServerLogin(),
route.getServerUrl(),
cursorStoredAtMs,
cursorMessageKey,
bootstrapCompleted
);
if (!batch.hasMore() || batch.items().isEmpty()) {
break;
}
}
if (!bootstrapCompleted) {
log.info("Periodic DM sync peer paused by page limit: owner={} remoteServer={} maxPages={}",
ownerLogin, route.getServerLogin(), maxPages);
}
return applied;
}
private static boolean isEnabled() { private static boolean isEnabled() {
String raw = AppConfig.getInstance().getParam("dm.sync.enabled"); String raw = AppConfig.getInstance().getParam("dm.sync.enabled");
return raw == null || raw.isBlank() || Boolean.parseBoolean(raw.trim()); return raw == null || raw.isBlank() || Boolean.parseBoolean(raw.trim());
@@ -169,17 +56,18 @@ public final class PeriodicDmSyncService {
private static long configLong(String key, long defaultValue, long min, long max) { private static long configLong(String key, long defaultValue, long min, long max) {
String raw = AppConfig.getInstance().getParam(key); String raw = AppConfig.getInstance().getParam(key);
if (raw == null || raw.isBlank()) return defaultValue; if (raw == null || raw.isBlank()) return defaultValue;
try { try { return Math.max(min, Math.min(max, Long.parseLong(raw.trim()))); }
long parsed = Long.parseLong(raw.trim()); catch (Exception ignored) { return defaultValue; }
return Math.max(min, Math.min(max, parsed));
} catch (Exception ignored) {
return defaultValue;
}
} }
private static String normalize(String value) { private static ThreadFactory daemonThreadFactory(String prefix) {
if (value == null) return null; return new ThreadFactory() {
String s = value.trim().toLowerCase(Locale.ROOT); private int sequence;
return s.isEmpty() ? null : s; @Override public synchronized Thread newThread(Runnable r) {
Thread t = new Thread(r, prefix + "-" + (++sequence));
t.setDaemon(true);
return t;
}
};
} }
} }
@@ -2,6 +2,7 @@ package server.sync;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import server.logic.ws_protocol.JSON.messages.DmSyncApplySupport;
import shine.db.dao.UserAccessServersCurrentDAO; import shine.db.dao.UserAccessServersCurrentDAO;
import shine.db.dao.UserSettingsDAO; import shine.db.dao.UserSettingsDAO;
import shine.db.dao.UserSettingsSyncPeerStateDAO; import shine.db.dao.UserSettingsSyncPeerStateDAO;
@@ -14,6 +15,7 @@ import java.sql.Connection;
import java.util.LinkedHashSet; import java.util.LinkedHashSet;
import java.util.List; import java.util.List;
import java.util.Locale; import java.util.Locale;
import java.util.ArrayList;
import java.util.Set; import java.util.Set;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
@@ -26,6 +28,7 @@ public final class PeriodicUserSettingsSyncService {
private static final String SERVER_LOGIN_CONFIG = "server.SHiNE.login"; private static final String SERVER_LOGIN_CONFIG = "server.SHiNE.login";
private static final AtomicBoolean STARTED = new AtomicBoolean(false); private static final AtomicBoolean STARTED = new AtomicBoolean(false);
private static final RemoteUserSettingsSyncClient REMOTE = new RemoteUserSettingsSyncClient(); private static final RemoteUserSettingsSyncClient REMOTE = new RemoteUserSettingsSyncClient();
private static final RemoteDmSyncClient DM_REMOTE = new RemoteDmSyncClient();
private static final UserAccessServersCurrentDAO ACCESS_DAO = UserAccessServersCurrentDAO.getInstance(); private static final UserAccessServersCurrentDAO ACCESS_DAO = UserAccessServersCurrentDAO.getInstance();
private static final UserSettingsDAO SETTINGS_DAO = UserSettingsDAO.getInstance(); private static final UserSettingsDAO SETTINGS_DAO = UserSettingsDAO.getInstance();
private static final UserSettingsSyncPeerStateDAO STATE_DAO = UserSettingsSyncPeerStateDAO.getInstance(); private static final UserSettingsSyncPeerStateDAO STATE_DAO = UserSettingsSyncPeerStateDAO.getInstance();
@@ -55,6 +58,9 @@ public final class PeriodicUserSettingsSyncService {
TimeUnit.HOURS.toSeconds(periodHours), TimeUnit.HOURS.toSeconds(periodHours),
TimeUnit.SECONDS TimeUnit.SECONDS
); );
EXECUTOR.scheduleWithFixedDelay(
PeriodicUserSettingsSyncService::runRequestedCycleSafe,
5L, 5L, TimeUnit.SECONDS);
log.info("Periodic user settings sync scheduled: first run in {} seconds, then every {} hours", initialDelaySec, periodHours); log.info("Periodic user settings sync scheduled: first run in {} seconds, then every {} hours", initialDelaySec, periodHours);
} }
@@ -66,6 +72,10 @@ public final class PeriodicUserSettingsSyncService {
} }
} }
private static void runRequestedCycleSafe() {
if (DmSyncWakeSignal.consume()) runCycleSafe();
}
private static void runCycle() throws Exception { private static void runCycle() throws Exception {
String ownServerLogin = normalize(AppConfig.getInstance().getParam(SERVER_LOGIN_CONFIG)); String ownServerLogin = normalize(AppConfig.getInstance().getParam(SERVER_LOGIN_CONFIG));
if (ownServerLogin == null) { if (ownServerLogin == null) {
@@ -83,6 +93,7 @@ public final class PeriodicUserSettingsSyncService {
int syncedPeers = 0; int syncedPeers = 0;
int appliedItems = 0; int appliedItems = 0;
int pushedItems = 0; int pushedItems = 0;
int appliedDmItems = 0;
for (String ownerLogin : owners) { for (String ownerLogin : owners) {
List<UserAccessServerRouteEntry> routes = ACCESS_DAO.listByUserLogin(ownerLogin); List<UserAccessServerRouteEntry> routes = ACCESS_DAO.listByUserLogin(ownerLogin);
for (UserAccessServerRouteEntry route : routes) { for (UserAccessServerRouteEntry route : routes) {
@@ -96,6 +107,7 @@ public final class PeriodicUserSettingsSyncService {
SyncStats stats = syncOwnerWithRemote(ownerLogin, route); SyncStats stats = syncOwnerWithRemote(ownerLogin, route);
appliedItems += stats.applied(); appliedItems += stats.applied();
pushedItems += stats.pushed(); pushedItems += stats.pushed();
appliedDmItems += stats.appliedDm();
syncedPeers++; syncedPeers++;
} catch (Exception e) { } catch (Exception e) {
STATE_DAO.updateError(ownerLogin, route.getServerLogin(), remoteUrl, String.valueOf(e)); STATE_DAO.updateError(ownerLogin, route.getServerLogin(), remoteUrl, String.valueOf(e));
@@ -105,8 +117,8 @@ public final class PeriodicUserSettingsSyncService {
} }
} }
log.info("Periodic user settings sync cycle finished: owners={} syncedPeers={} appliedItems={} pushedItems={}", log.info("Periodic access-data sync finished: owners={} peers={} settingsApplied={} settingsPushed={} dmApplied={}",
owners.size(), syncedPeers, appliedItems, pushedItems); owners.size(), syncedPeers, appliedItems, pushedItems, appliedDmItems);
} }
private static SyncStats syncOwnerWithRemote(String ownerLogin, UserAccessServerRouteEntry route) throws Exception { private static SyncStats syncOwnerWithRemote(String ownerLogin, UserAccessServerRouteEntry route) throws Exception {
@@ -124,9 +136,11 @@ public final class PeriodicUserSettingsSyncService {
int pushed = 0; int pushed = 0;
boolean bootstrapCompleted = false; boolean bootstrapCompleted = false;
int appliedDm;
try (RemoteSyncSession session = new RemoteSyncSession(route.getServerUrl())) {
for (int page = 0; page < maxPages; page++) { for (int page = 0; page < maxPages; page++) {
RemoteUserSettingsSyncClient.RemoteUserSettingsBatch batch = REMOTE.userSettingsSyncBatch( RemoteUserSettingsSyncClient.RemoteUserSettingsBatch batch = REMOTE.userSettingsSyncBatch(
route.getServerUrl(), session,
ownerLogin, ownerLogin,
cursorTimeMs, cursorTimeMs,
cursorSettingKey, cursorSettingKey,
@@ -163,17 +177,53 @@ public final class PeriodicUserSettingsSyncService {
try (Connection c = getDbConnection()) { try (Connection c = getDbConnection()) {
List<UserSettingEntry> unsynced = SETTINGS_DAO.listUnsyncedByLogin(c, ownerLogin, limit); List<UserSettingEntry> unsynced = SETTINGS_DAO.listUnsyncedByLogin(c, ownerLogin, limit);
for (UserSettingEntry entry : unsynced) { for (UserSettingEntry entry : unsynced) {
REMOTE.upsertUserSetting(route.getServerUrl(), entry, true); REMOTE.upsertUserSetting(session, entry, true);
SETTINGS_DAO.markSynced(c, entry.getLogin(), entry.getSettingType(), entry.getSettingKey()); SETTINGS_DAO.markSynced(c, entry.getLogin(), entry.getSettingType(), entry.getSettingKey());
pushed++; pushed++;
} }
} }
int dmLimit = (int) configLong("dm.sync.batchLimit", 200L, 1L, 500L);
int dmMaxBytes = (int) configLong("dm.sync.batchMaxBytes", 3_000_000L, 64_000L, 5_000_000L);
int dmMaxPages = (int) configLong("dm.sync.maxPagesPerPeer", 20L, 1L, 500L);
appliedDm = syncDmInSameSession(session, ownerLogin, dmLimit, dmMaxBytes, dmMaxPages);
}
if (!bootstrapCompleted) { if (!bootstrapCompleted) {
log.info("Periodic user settings sync peer paused by page limit: owner={} remoteServer={} maxPages={}", log.info("Periodic user settings sync peer paused by page limit: owner={} remoteServer={} maxPages={}",
ownerLogin, route.getServerLogin(), maxPages); ownerLogin, route.getServerLogin(), maxPages);
} }
return new SyncStats(applied, pushed); return new SyncStats(applied, pushed, appliedDm);
}
private static int syncDmInSameSession(
RemoteSyncSession session, String ownerLogin, int limit, int maxBytes, int maxPages
) throws Exception {
long cursorMs = 0L;
String cursorKey = "";
List<String> acknowledgements = new ArrayList<>();
int applied = 0;
for (int page = 0; page < maxPages; page++) {
RemoteDmSyncClient.RemoteDmBatch batch = DM_REMOTE.dmSyncBatch(
session, ownerLogin, cursorMs, cursorKey, Math.min(limit, 500), maxBytes, acknowledgements);
acknowledgements = new ArrayList<>();
for (RemoteDmSyncClient.RemoteDmItem item : batch.items()) {
if (item == null || item.syncId() == null || item.syncId().isBlank()) continue;
DmSyncApplySupport.applySyncedItem(ownerLogin, item.syncId(), item.blobsB64());
acknowledgements.add(item.syncId());
applied++;
}
cursorMs = batch.nextStoredAtMs();
cursorKey = batch.nextMessageKey() == null ? "" : batch.nextMessageKey();
if (!batch.hasMore() || batch.items().isEmpty()) {
if (!acknowledgements.isEmpty()) {
DM_REMOTE.dmSyncBatch(session, ownerLogin, 0L, "",
Math.min(limit, 500), maxBytes, acknowledgements);
}
break;
}
}
return applied;
} }
private static java.sql.Connection getDbConnection() throws Exception { private static java.sql.Connection getDbConnection() throws Exception {
@@ -202,5 +252,5 @@ public final class PeriodicUserSettingsSyncService {
return s.isEmpty() ? null : s; return s.isEmpty() ? null : s;
} }
private record SyncStats(int applied, int pushed) {} private record SyncStats(int applied, int pushed, int appliedDm) {}
} }
@@ -34,13 +34,13 @@ sync.importUserProfileFromPartner.enabled=false
# ------------------------------------------------------------ # ------------------------------------------------------------
server.version=${projectVersion} server.version=${projectVersion}
# Межсерверная догоняющая синхронизация личных сообщений. # Доставка и межсерверная синхронизация личных сообщений.
dm.sync.enabled=true dm.sync.enabled=true
dm.sync.initialDelaySeconds=60 dm.worker.pollSeconds=5
dm.sync.periodHours=6 dm.worker.dueLimit=100
dm.sync.batchLimit=500 dm.sync.batchLimit=200
dm.sync.batchMaxBytes=3000000 dm.sync.batchMaxBytes=3000000
dm.sync.maxPagesPerPeer=50 dm.sync.maxPagesPerPeer=20
server.info.url= server.info.url=
server.info.physicalRegion= server.info.physicalRegion=
server.info.description= server.info.description=
+5 -3
View File
@@ -62,11 +62,12 @@
| `UpsertPushToken` | `12_Direct_Messages_Push_Calls_API.md` | регистрация WebPush-токена | | `UpsertPushToken` | `12_Direct_Messages_Push_Calls_API.md` | регистрация WebPush-токена |
| `SendTestWebPush` | `12_Direct_Messages_Push_Calls_API.md` | тестовая push-доставка | | `SendTestWebPush` | `12_Direct_Messages_Push_Calls_API.md` | тестовая push-доставка |
| `SendMessagePair` | `12_Direct_Messages_Push_Calls_API.md` | отправка пары входящий/исходящий DM | | `SendMessagePair` | `12_Direct_Messages_Push_Calls_API.md` | отправка пары входящий/исходящий DM |
| `ReceiveOutcomingMessage` | `12_Direct_Messages_Push_Calls_API.md` | алиас `SendMessagePair` | | `ReceiveOutcomingMessage` | `12_Direct_Messages_Push_Calls_API.md` | старая полная DM-пара для второго access-сервера отправителя |
| `ReceiveIncomingMessage` | `12_Direct_Messages_Push_Calls_API.md` | прием входящего DM-блока | | `ReceiveIncomingMessage` | `12_Direct_Messages_Push_Calls_API.md` | прием входящего DM-блока |
| `DeleteMessage` | `12_Direct_Messages_Push_Calls_API.md` | tombstone одного личного сообщения у обеих сторон | | `DeleteMessage` | `12_Direct_Messages_Push_Calls_API.md` | tombstone одного личного сообщения у обеих сторон |
| `DeleteConversation` | `12_Direct_Messages_Push_Calls_API.md` | tombstone удаления истории переписки | | `DeleteConversation` | `12_Direct_Messages_Push_Calls_API.md` | tombstone удаления истории переписки |
| `DmSyncBatch` | `12_Direct_Messages_Push_Calls_API.md` | межсерверная догоняющая синхронизация DM по курсору | | `DmSyncBatch` | `12_Direct_Messages_Push_Calls_API.md` | pull событий `synced=false` с ACK предыдущей страницы |
| `GetDmDeliveryStatus` | `12_Direct_Messages_Push_Calls_API.md` | read-only статус доставки по существующему `messageKey` |
| `UserSettingsSyncBatch` | `13_User_Settings_API.md` | межсерверная догоняющая синхронизация пользовательских настроек по курсору | | `UserSettingsSyncBatch` | `13_User_Settings_API.md` | межсерверная догоняющая синхронизация пользовательских настроек по курсору |
| `MarkAllUserSettingsUnsynced` | `13_User_Settings_API.md` | служебная пометка всех настроек как несинхронизированных | | `MarkAllUserSettingsUnsynced` | `13_User_Settings_API.md` | служебная пометка всех настроек как несинхронизированных |
| `GetDirectMessages` | `12_Direct_Messages_Push_Calls_API.md` | постраничная загрузка истории личного диалога | | `GetDirectMessages` | `12_Direct_Messages_Push_Calls_API.md` | постраничная загрузка истории личного диалога |
@@ -76,7 +77,8 @@
## Важные замечания ## Важные замечания
- `ReceiveOutcomingMessage` сейчас зарегистрирован как алиас того же handler/request-класса, что и `SendMessagePair`. - `ReceiveOutcomingMessage` зарегистрирован как алиас того же handler/request-класса, что и `SendMessagePair`, и сохраняет прежний межсерверный payload.
- Межсерверные DM-операции пока доверяют `sourceServerLogin`; отдельная межсерверная авторизация запланирована позднее.
- Отдельных HTTP endpoints для DM-файлов сейчас нет. - Отдельных HTTP endpoints для DM-файлов сейчас нет.
- Классы `Net_MarkChannelMessagesSeen_*` существуют в коде, но операция `MarkChannelMessagesSeen` не зарегистрирована в `JsonHandlerRegistry`, поэтому в публичный список API не входит. - Классы `Net_MarkChannelMessagesSeen_*` существуют в коде, но операция `MarkChannelMessagesSeen` не зарегистрирована в `JsonHandlerRegistry`, поэтому в публичный список API не входит.
- HTTP debug endpoints из `src/main/java/server/debug/` не входят в этот индекс WebSocket `op`; они описаны отдельно в `13_HTTP_Debug_API.md`. - HTTP debug endpoints из `src/main/java/server/debug/` не входят в этот индекс WebSocket `op`; они описаны отдельно в `13_HTTP_Debug_API.md`.
+102 -33
View File
@@ -107,12 +107,17 @@
"baseKey": "from|to|time|nonce", "baseKey": "from|to|time|nonce",
"incomingKey": "from|to|time|nonce|1", "incomingKey": "from|to|time|nonce|1",
"outgoingKey": "from|to|time|nonce|2", "outgoingKey": "from|to|time|nonce|2",
"deliveryState": "accepted",
"deliveredWsSessions": 1, "deliveredWsSessions": 1,
"deliveredWebPushSessions": 0 "deliveredWebPushSessions": 0
} }
} }
``` ```
Успешный `status=200` подтверждает локальное сохранение. До ответа клиенту сервер параллельно пробует оба актуальных сервера получателя. Поэтому `deliveryState` уже может быть `delivered`; если никто не ответил, возвращается `accepted`.
Возможные `deliveryState`: `accepted`, `delivered`, `failed`.
### Ошибки ### Ошибки
- `400 / BAD_FIELDS` — пустой `incomingBlobB64` или `outgoingBlobB64` - `400 / BAD_FIELDS` — пустой `incomingBlobB64` или `outgoingBlobB64`
@@ -147,7 +152,9 @@
} }
``` ```
`sourceServerLogin` необязателен. Если поле есть, сервер использует его как подсказку, чтобы не отправлять событие обратно серверу-источнику. `sourceServerLogin` необязателен для совместимости. Пока межсерверная авторизация отложена, это поле считается доверенным. Пользовательская подпись signed-блока проверяется всегда.
Успешный ответ содержит существующие `messageKey`, `baseKey` и счётчики realtime-доставки. Повтор уже сохранённой той же ревизии обрабатывается идемпотентно.
### Примечание ### Примечание
@@ -243,6 +250,7 @@
"reencryptedAtMs": 0, "reencryptedAtMs": 0,
"createdAtMs": 1774700001123, "createdAtMs": 1774700001123,
"readAtMs": 1774700001456, "readAtMs": 1774700001456,
"deliveryState": "delivered",
"blobB64": "BASE64_SIGNED_BLOCK" "blobB64": "BASE64_SIGNED_BLOCK"
} }
] ]
@@ -252,67 +260,109 @@
Для следующей страницы клиент должен передать `nextBeforeTimeMs` и `nextBeforeMessageKey` из предыдущего ответа. Для следующей страницы клиент должен передать `nextBeforeTimeMs` и `nextBeforeMessageKey` из предыдущего ответа.
## 8. `DmSyncBatch` Поле `deliveryState` заполняется для исходящих элементов `type=2`. У входящих `type=1` оно отсутствует/равно `null`.
Межсерверная операция для догоняющей синхронизации истории одного пользователя. В текущей реализации не требует авторизации сервера-источника, но удалённый сервер отдаёт данные только если сам является access-сервером `ownerLogin` по `user_access_servers_current`. ## 8. Межсерверные операции DM
### Запрос До отдельного этапа server-auth поля `sourceServerLogin` доверяются. Каждый `SHiNE_DM` всё равно заново проходит проверку формата и пользовательской подписи.
### 8.1. `ReceiveOutcomingMessage`
Прежняя операция передачи полной пары второму access-серверу отправителя. Delivery-state в межсерверный запрос не входит.
```json ```json
{ {
"op": "DmSyncBatch", "op": "ReceiveOutcomingMessage",
"requestId": "dm-sync-001", "requestId": "dm-peer-001",
"payload": { "payload": {
"ownerLogin": "alice", "incomingBlobB64": "BASE64_INCOMING",
"afterStoredAtMs": 1774700000000, "outgoingBlobB64": "BASE64_OUTGOING",
"afterMessageKey": "alice|bob|1774699999000|123456780|2", "sourceServerLogin": "server-a"
"limit": 500,
"maxBytes": 3000000
} }
} }
``` ```
`afterStoredAtMs` и `afterMessageKey` образуют курсор. Если курсора нет, сервер передаёт `0` и пустую строку. `limit` ограничен максимумом `500`. Ответ имеет прежний формат `SendMessagePair`. Принимающий сервер сохраняет пару идемпотентно и самостоятельно ставит локальную задачу доставки.
### Успешный ответ ### 8.2. `ReceiveIncomingMessage`
Прежняя операция передачи одной входящей копии серверу получателя или второму серверу самого получателя.
```json
{
"op": "ReceiveIncomingMessage",
"requestId": "dm-incoming-001",
"payload": {
"incomingBlobB64": "BASE64_INCOMING",
"sourceServerLogin": "server-a"
}
}
```
Успешный 2xx-ответ означает, что signed-блок проверен и находится в БД (новая или идемпотентно повторённая запись).
### 8.3. `DmSyncBatch`
Pull-синхронизация событий владельца с `synced=false`. `ackSyncIds` подтверждает версии событий, успешно сохранённые из предыдущего ответа.
```json ```json
{ {
"op": "DmSyncBatch", "op": "DmSyncBatch",
"requestId": "dm-sync-001", "requestId": "dm-sync-001",
"status": 200,
"ok": true,
"payload": { "payload": {
"ownerLogin": "alice", "ownerLogin": "alice",
"afterStoredAtMs": 0,
"afterMessageKey": "",
"limit": 500, "limit": 500,
"rawBytes": 84512, "maxBytes": 3000000,
"hasMore": true, "ackSyncIds": []
"nextStoredAtMs": 1774700100000, }
"nextMessageKey": "alice|bob|1774700000123|123456789|1", }
```
Ответ возвращает страницу несинхронизированных событий и курсор внутри текущего цикла:
```json
{
"nextStoredAtMs": 1774700001123,
"nextMessageKey": "alice|bob|1774700000123|123456789|2",
"hasMore": false,
"items": [ "items": [
{ {
"messageKey": "alice|bob|1774700000123|123456789|1", "syncId": "alice|bob|1774700000123|123456789|2:0:0",
"baseKey": "alice|bob|1774700000123|123456789", "primaryMessageKey": "alice|bob|1774700000123|123456789|2",
"targetLogin": "alice", "storedAtMs": 1774700001123,
"fromLogin": "bob", "blobsB64": ["BASE64_INCOMING", "BASE64_OUTGOING"]
"toLogin": "alice",
"messageType": 1,
"timeMs": 1774700000123,
"storedAtMs": 1774700100000,
"blobB64": "BASE64_SIGNED_BLOCK"
} }
] ]
}
```
Для пары порядок всегда incoming/outgoing; для входящей копии и tombstone массив содержит один blob. `syncId` — технический идентификатор конкретной ревизии только для ACK синхронизации, а не второй ID сообщения. После сохранения вызывающий сервер передаёт `syncId` в `ackSyncIds` следующего запроса. Только тогда источник ставит `synced=true`. Новый цикл начинается с `afterStoredAtMs=0`; полный сброс флагов поэтому повторно отдаёт всю историю.
### 8.4. `GetDmDeliveryStatus`
Единственная новая межсерверная операция доставки. Read-only проверка существующего `messageKey`; не изменяет состояние отвечающего сервера.
```json
{
"op": "GetDmDeliveryStatus",
"requestId": "dm-status-001",
"payload": {
"messageKey": "alice|bob|1774700000123|123456789|2"
} }
} }
``` ```
События в `items` идут по `storedAtMs ASC, messageKey ASC`. В пачке могут быть сообщения любых диалогов пользователя, read-receipt и delete/tombstone типов `5/6/7/8`. Ответ содержит `messageKey`, `known` и `delivered`. `delivered=true` означает доставку хотя бы одному серверу получателя.
Ошибки: Периодический процесс использует одно WS-соединение с peer: сначала синхронизирует настройки, затем вызывает `DmSyncBatch` до завершения страниц и ACK. Существующий `MarkAllUserSettingsUnsynced` также сбрасывает DM-флаги и возвращает `dmUpdated`; отдельной операции `MarkAllDmUnsynced` нет.
- `400 / EMPTY_OWNER_LOGIN` — не передан `ownerLogin` Основные ошибки межсерверных операций:
- `403 / LOCAL_SERVER_NOT_ACCESS_SERVER` — этот сервер не является access-сервером пользователя
- `500 / LOCAL_SERVER_NOT_CONFIGURED` — не настроен `server.SHiNE.login` - `400 / EMPTY_OWNER_LOGIN`, `EMPTY_MESSAGE_KEY`;
- `403 / LOCAL_SERVER_NOT_ACCESS_SERVER`;
- `500 / LOCAL_SERVER_NOT_CONFIGURED`.
## 9. `AckSessionDelivery` ## 9. `AckSessionDelivery`
@@ -355,6 +405,25 @@
Для типов `5/6/7/8` событие тоже приходит в таком же конверте, но логика применения определяется `messageType` и бинарным `blobB64`. Для типов `5/6/7/8` событие тоже приходит в таком же конверте, но логика применения определяется `messageType` и бинарным `blobB64`.
## 10.1. Событие `DmDeliveryStateChanged`
Сервер отправляет событие активным сессиям отправителя при изменении сетевого состояния исходящей пары.
```json
{
"op": "DmDeliveryStateChanged",
"event": true,
"status": 200,
"payload": {
"baseKey": "alice|bob|1774700000123|123456789",
"outgoingKey": "alice|bob|1774700000123|123456789|2",
"deliveryState": "delivered"
}
}
```
`failed` приходит после последней неудачной попытки через час. Событие прочтения остаётся отдельным существующим read-receipt.
## 11. `CallInviteBroadcast` ## 11. `CallInviteBroadcast`
Требует авторизации. Шлёт приглашение к звонку в активные сессии `toLogin`. Требует авторизации. Шлёт приглашение к звонку в активные сессии `toLogin`.
+4 -2
View File
@@ -77,15 +77,17 @@
Внутренний служебный запрос. Внутренний служебный запрос.
- помечает все настройки пользователя или все настройки сразу как `synced=false`; - помечает все настройки пользователя или все настройки сразу как `synced=false`;
- одновременно сбрасывает DM-outbox того же пользователя (или всех пользователей), чтобы операция замены/добавления access-сервера не оставила историю DM несинхронизированной;
- ответ дополнительно содержит `dmUpdated` — количество сброшенных DM-событий;
- нужен после добавления нового sync-сервера или при потере локальной БД. - нужен после добавления нового sync-сервера или при потере локальной БД.
## 5. Синхронизация ## 5. Синхронизация
Синхронизация настроек работает отдельно от DM. Настройки и DM имеют раздельные таблицы и правила ACK, но периодический процесс открывает один последовательный WS-сеанс с peer: сначала синхронизирует настройки, затем забирает `DmSyncBatch`.
- локальная запись создаётся с `synced=false`, если её ещё не подтвердил второй сервер; - локальная запись создаётся с `synced=false`, если её ещё не подтвердил второй сервер;
- если запись пришла с другого сервера, она сохраняется сразу как `synced=true`; - если запись пришла с другого сервера, она сохраняется сразу как `synced=true`;
- периодический sync раз в 6 часов проверяет несинхронизированные записи и догружает новые записи по курсору; - периодический sync раз в 6 часов проверяет несинхронизированные настройки, догружает их по курсору и затем забирает несинхронизированные DM;
- если появляется новый sync-сервер или локальная БД была потеряна, нужно пометить все настройки несинхронизированными и заново догрузить batch с нуля. - если появляется новый sync-сервер или локальная БД была потеряна, нужно пометить все настройки несинхронизированными и заново догрузить batch с нуля.
## 6. Текущий UI-кейс ## 6. Текущий UI-кейс
+27 -17
View File
@@ -32,9 +32,11 @@
### 3.1 Личные сообщения (DM) ### 3.1 Личные сообщения (DM)
- Все DM-блоки форматов типов `1/2` (текст) и `3/4` (read-receipt). - Все DM-блоки форматов типов `1/2` (текст) и `3/4` (read-receipt).
- Сервер-отправитель: при получении пары блоков от клиента перенаправляет их серверу получателя. - Сервер-отправитель: сохраняет пару и ставит асинхронную delivery-задачу.
- Сервер-получатель: сохраняет блоки в `signed_messages_v2`, доставляет в активные сессии. - Сервер-получатель: сохраняет входящий блок в `signed_messages`, затем доставляет его активным сессиям.
- Дедупликация по уникальному `message_key = from|to|timeMs|nonce|type`. - Дедупликация по уникальному `message_key = from|to|timeMs|nonce|type`.
- Репликация между двумя access-серверами пользователя работает через `dm_sync_outbox.synced`, без постоянного time-cursor.
- Полная актуальная схема: `docs/Personal_Messages/Доставка_и_синхронизация_DM.md`.
### 3.2 Блоки пользовательского блокчейна ### 3.2 Блоки пользовательского блокчейна
@@ -181,7 +183,9 @@ Full resync запускается только тогда, когда:
Настройка влияет именно на этап подготовки отсутствующей локальной цепочки во время periodic sync. Настройка влияет именно на этап подготовки отсутствующей локальной цепочки во время periodic sync.
## 5. Целевой протокол следующего этапа ## 5. Возможное развитие server-to-server транспорта
Этот раздел не описывает текущую DM-доставку. DM уже использует короткие one-shot WebSocket-вызовы, indexed outbox и ACK. Ниже остаётся возможное развитие постоянного транспорта и server-auth.
### 5.1 Межсерверное соединение ### 5.1 Межсерверное соединение
@@ -192,14 +196,15 @@ Full resync запускается только тогда, когда:
### 5.2 Доставка новых данных (push) ### 5.2 Доставка новых данных (push)
- При получении нового блока или DM сервер немедленно пушит его всем подключённым партнёрам. - При получении нового блока сервер может немедленно пушить его всем подключённым партнёрам.
- Партнёр подтверждает приём (ACK). Без ACK — повтор с backoff. - Партнёр подтверждает приём (ACK). Без ACK — повтор с backoff.
- DM использует отдельное расписание, описанное в `docs/Personal_Messages/Доставка_и_синхронизация_DM.md`.
### 5.3 Начальная синхронизация (backfill) ### 5.3 Начальная синхронизация (backfill)
- При первом подключении к партнёру серверы обмениваются «курсорами» состояния: - При первом подключении к партнёру серверы могут обмениваться курсорами состояния блокчейнов.
последний глобальный номер блока, последний известный DM-ключ.
- Сервер с более полной историей досылает недостающее партнёру. - Сервер с более полной историей досылает недостающее партнёру.
- DM time-cursor удалён: новый peer получает историю после сброса `dm_sync_outbox.synced=false`.
### 5.4 Разрешение конфликтов ### 5.4 Разрешение конфликтов
@@ -212,19 +217,21 @@ Full resync запускается только тогда, когда:
При отправке DM от пользователя A к пользователю B: При отправке DM от пользователя A к пользователю B:
1. Клиент A отправляет пару блоков на свой сервер X. 1. Клиент A отправляет пару блоков на свой сервер X.
2. Сервер X определяет, на каком сервере зарегистрирован пользователь B. 2. Сервер X валидирует и локально сохраняет пару.
- Сначала проверяет локально (если B зарегистрирован на X). 3. До ответа клиенту X параллельно отправляет входящую копию максимум двум актуальным `access_servers` B.
- Иначе читает PDA пользователя B из Solana и смотрит `access_servers`. 4. ACK хотя бы одного маршрута B означает терминальный `delivered`; второй сервер B догоняется собственной DM-синхронизацией.
- Выбирает первый доступный сервер из `access_servers` и перенаправляет туда DM. 5. После первой попытки X передаёт полную пару второму access-серверу A старым `ReceiveOutcomingMessage`, без delivery-state.
3. Сервер Y (из `access_servers` B) сохраняет и доставляет блоки. 6. При нулевой доставке выполняются повторы через 30 секунд, 5 минут, 25 минут и 1 час. Перед тремя последними X спрашивает peer A через `GetDmDeliveryStatus(messageKey)`.
7. После неудачной попытки через час ставится терминальный `failed`.
Кэш адресов серверов: обновляется раз в сессию (при ошибке соединения). Изменения routing B учитываются до терминального состояния, потому что список маршрутов перечитывается на каждой попытке.
## 7. Безопасность ## 7. Безопасность
- Все блоки подписаны ключами пользователя на клиенте — сервер не может подделать содержимое. - Все блоки подписаны ключами пользователя на клиенте — сервер не может подделать содержимое.
- Серверы не расшифровывают DM-контент (шифрование — задача следующего этапа). - Серверы не расшифровывают DM-контент; E2EE уже выполняется клиентами.
- При синхронизации каждый блок проходит валидацию подписи на принимающем сервере. - При синхронизации каждый блок проходит валидацию подписи на принимающем сервере.
- Межсерверная авторизация DM-операций пока отложена; `sourceServerLogin` временно считается доверенным.
## 8. Статус реализации ## 8. Статус реализации
@@ -240,15 +247,18 @@ Full resync запускается только тогда, когда:
| Обход Solana RPC через `sync.importUserProfileFromPartner.enabled` | ✅ Реализовано | | Обход Solana RPC через `sync.importUserProfileFromPartner.enabled` | ✅ Реализовано |
| Обычный `AddBlock` через `tmp_bch`/`write_check`/`write_pending` | ✅ Реализовано | | Обычный `AddBlock` через `tmp_bch`/`write_check`/`write_pending` | ✅ Реализовано |
| Межсерверный постоянный WebSocket-канал | Нужна реализация | | Межсерверный постоянный WebSocket-канал | Нужна реализация |
| Push новых DM партнёрам | Нужна реализация | | Асинхронная доставка DM на access-серверы получателя | ✅ Реализовано |
| Retry DM до 1 часа + UI-state | ✅ Реализовано |
| Репликация DM на второй access-сервер по `synced` | ✅ Реализовано |
| Read-only `GetDmDeliveryStatus` | ✅ Реализовано |
| Push блоков блокчейна партнёрам | ✅ Реализована базовая one-shot версия | | Push блоков блокчейна партнёрам | ✅ Реализована базовая one-shot версия |
| Periodic backfill отсутствующего хвоста | ✅ Реализовано | | Periodic backfill отсутствующего хвоста | ✅ Реализовано |
| Разрешение рассинхрона / divergence | ✅ Реализована базовая full-resync схема во время periodic sync | | Разрешение рассинхрона / divergence | ✅ Реализована базовая full-resync схема во время periodic sync |
| Startup recovery по `*.resync_pending` marker-file | ✅ Реализовано | | Startup recovery по `*.resync_pending` marker-file | ✅ Реализовано |
| Маршрутизация DM через access_servers | Нужна реализация (заглушка) | | Маршрутизация DM через один/два `access_servers` | ✅ Реализовано |
| Криптографическая server-to-server авторизация DM | Нужна реализация |
Текущая версия сервера уже умеет базовую синхронизацию блокчейнов между партнёрами. Текущая версия сервера умеет синхронизацию блокчейнов и DM. Постоянные server-to-server соединения не требуются для текущей one-shot WS-реализации; отдельной будущей задачей остаётся криптографическая авторизация DM-вызовов.
Не реализованы ещё DM-sync и постоянные server-to-server соединения.
Следующие отдельные шаги после текущего этапа: Следующие отдельные шаги после текущего этапа:
- отдельно проверить full-resync и startup-recovery на реальном тестовом прогоне после ручного удаления БД/файлов. - отдельно проверить full-resync и startup-recovery на реальном тестовом прогоне после ручного удаления БД/файлов.
+1
View File
@@ -6,6 +6,7 @@
- `docs/Personal_Messages/Протокол_DM_v1.md` — логика протокола, роли API, серверное поведение, routing по `access_servers` - `docs/Personal_Messages/Протокол_DM_v1.md` — логика протокола, роли API, серверное поведение, routing по `access_servers`
- `docs/Personal_Messages/Формат_DM_v1.md` — точный бинарный формат контейнера `SHiNE_DM` - `docs/Personal_Messages/Формат_DM_v1.md` — точный бинарный формат контейнера `SHiNE_DM`
- `docs/Personal_Messages/Доставка_и_синхронизация_DM.md` — состояния доставки, retry-воркер, репликация между двумя access-серверами и UI-статусы
- `docs/Personal_Messages/Технические_вставки_DM_v1.md` — формат специальных `<S:...>` вставок внутри plaintext DM после расшифровки - `docs/Personal_Messages/Технические_вставки_DM_v1.md` — формат специальных `<S:...>` вставок внутри plaintext DM после расшифровки
Исторический устаревший документ сохранён отдельно: Исторический устаревший документ сохранён отдельно:
@@ -0,0 +1,136 @@
# Доставка и синхронизация личных сообщений
## 1. Главный принцип
DM считается доставленным, когда signed-входящую копию сохранил хотя бы один актуальный access-сервер получателя. Доставка на оба сервера не требуется: сервер получателя самостоятельно синхронизирует входящее сообщение со своим вторым сервером.
Клиентский `status=200` от `SendMessagePair` означает, что собственный сервер сохранил пару. Это отдельный факт от доставки получателю.
Формат подписанного контейнера `SHiNE_DM` не меняется. Состояние доставки и флаг синхронизации — изменяемые серверные метаданные.
## 2. Идентификаторы
- `baseKey` связывает входящую и исходящую копии одного логического сообщения;
- `incomingKey` идентифицирует копию получателя;
- `outgoingKey`/`messageKey` идентифицирует копию отправителя и используется для проверки доставки;
- дополнительный публичный `eventId` для доставки не создаётся;
- `syncId` используется только внутри `DmSyncBatch` как ACK конкретной ревизии, потому что редактирование сохраняет прежний `messageKey`.
## 3. Состояния доставки
| Состояние | Смысл | Повторные попытки |
|---|---|---|
| `accepted` | пара сохранена сервером отправителя, но ни один сервер получателя ещё не подтвердил запись | да |
| `delivered` | хотя бы один сервер получателя подтвердил запись | нет |
| `failed` | последняя попытка через час завершилась без доставки | никогда |
Состояние `delivered` терминальное. Сервер отправителя не пытается отдельно добиться второго ACK получателя.
## 4. Обычная отправка
1. Клиент отправляет `SendMessagePair` на один свой access-сервер.
2. Сервер проверяет формат, пользователей, подписи и согласованность пары.
3. Сервер атомарно сохраняет входящую и исходящую копии.
4. В том же request-процессе сервер читает до двух актуальных маршрутов получателя.
5. Оба вызова `ReceiveIncomingMessage` запускаются параллельно.
6. Если хотя бы один вызов успешен, устанавливается `delivered`.
7. После этой попытки сервер передаёт полную пару своему второму access-серверу старой операцией `ReceiveOutcomingMessage`.
8. `SendMessagePair` возвращает существующие ключи и единственное новое поле `deliveryState`.
Успешный повтор уже сохранённого signed-блока считается ACK. Все операции должны быть идемпотентными.
## 5. Расписание повторов
Воркер запускается каждые 5 секунд и выбирает только due-строки по индексу. Он не сканирует всю таблицу сообщений.
Попытки привязаны к времени первоначального принятия:
| Номер | Время от старта | Сначала спросить второй сервер отправителя |
|---:|---:|---|
| 1 | сразу | нет |
| 2 | 30 секунд | нет |
| 3 | 5 минут | да |
| 4 | 25 минут | да |
| 5 | 1 час | да |
Из-за шага воркера повтор может начаться на несколько секунд позже указанного времени. Первая попытка выполняется немедленно и от воркера не зависит.
На 5-й, 25-й и 60-й минутах сервер сначала вызывает у второго сервера отправителя:
```json
{
"op": "GetDmDeliveryStatus",
"payload": {
"messageKey": "alice|bob|1774700000123|123456789|2"
}
}
```
Ответ:
```json
{
"messageKey": "alice|bob|1774700000123|123456789|2",
"known": true,
"delivered": true
}
```
Операция read-only. Если peer отвечает `delivered=true`, локальный сервер устанавливает `delivered` и не обращается к серверам получателя. Ошибка или отсутствие операции на старом peer не блокирует собственную попытку.
Перед каждой попыткой маршруты получателя заново читаются из `user_access_servers_current`. Изменение серверов учитывается только пока сообщение находится в `accepted`. После `delivered` или `failed` состояние больше не открывается.
Если последняя проверка и попытка через час не дали ACK, устанавливается `failed`, `next_attempt_at_ms` очищается и сообщение больше никогда автоматически не отправляется.
## 6. Старая межсерверная доставка
Форматы существующих операций не расширяются данными о результате доставки:
- `ReceiveOutcomingMessage` получает прежнюю пару `incomingBlobB64` + `outgoingBlobB64` и необязательный `sourceServerLogin`;
- `ReceiveIncomingMessage` получает один `incomingBlobB64` и необязательный `sourceServerLogin`.
Второй сервер отправителя после получения пары создаёт собственное локальное состояние доставки и самостоятельно пробует маршруты получателя. Серверы обмениваются результатом только через read-only `GetDmDeliveryStatus`.
## 7. Догоняющая синхронизация двух серверов пользователя
Для каждого владельца сервер ведёт outbox с флагом `synced`:
- исходящая пара отправителя — один элемент с двумя blob в порядке incoming/outgoing;
- входящая копия получателя — один элемент с одним blob;
- read-receipt и tombstone применяются теми же проверенными обработчиками.
`DmSyncBatch` возвращает только элементы с `synced=false`. Получатель проверяет signed-блоки, сохраняет их идемпотентно и в следующем запросе подтверждает `syncId` через `ackSyncIds`. Источник ставит `synced=true` только после ACK.
При обрыве соединения неподтверждённый элемент остаётся `synced=false` и безопасно приходит повторно. Элемент, полученный от peer, локально сразу считается синхронизированным, чтобы не образовалась петля.
Синхронизация настроек и DM проходит последовательно через один WS-сеанс: сначала настройки, затем все страницы DM. Если второй сервер был выключен, после включения он сам догружает пропущенные элементы.
Существующий `MarkAllUserSettingsUnsynced` также сбрасывает DM-флаги. После сброса история повторно передаётся как синхронизация, но старые сообщения получателю заново не отправляются: delivery-состояние создаётся с учётом их возраста и не открывает завершённую часовую очередь.
## 8. UI
| Вид | Значение |
|---|---|
| одна серая галочка | собственный сервер принял сообщение (`accepted`) |
| одна светлая галочка | хотя бы один сервер получателя принял сообщение (`delivered`) |
| две светлые галочки | пришёл существующий read-receipt |
| красный `!` и «Сообщение не доставлено» | окончательное состояние `failed` |
Read-receipt имеет приоритет над delivery-индикатором: если сообщение прочитано, оно заведомо было доставлено.
Сервер сообщает поздние изменения через `DmDeliveryStateChanged` с полями `outgoingKey`, `baseKey`, `deliveryState`. При повторном открытии чата то же состояние приходит в `GetDirectMessages`.
## 9. Нагрузка и отказоустойчивость
- воркер выбирает только due-записи по частичному индексу;
- терминальные строки не попадают в рабочую выборку;
- два маршрута первой попытки выполняются параллельно;
- ограниченный пул потоков и очередь защищают сервер при всплеске отправок;
- сетевой lease и сравнение версии строки предотвращают одновременную обработку одной задачи несколькими worker-потоками;
- повторная запись одного signed-блока безопасна;
- отсутствие второго сервера отправителя не мешает собственной доставке;
- отсутствие обоих серверов получателя завершает задачу через час.
## 10. Граница доверия
Межсерверная авторизация пока отложена. `sourceServerLogin` временно принимается на доверии, но каждый контейнер `SHiNE_DM` всё равно проходит проверку пользовательской подписи. `GetDmDeliveryStatus` сообщает только факт локального delivery-state и не изменяет данные.
@@ -26,6 +26,10 @@
- `docs/Personal_Messages/Технические_вставки_DM_v1.md` - `docs/Personal_Messages/Технические_вставки_DM_v1.md`
Изменяемое состояние доставки, retry-расписание и репликация между двумя access-серверами подробно описаны отдельно:
- `docs/Personal_Messages/Доставка_и_синхронизация_DM.md`
Устаревшая предыдущая версия сохранена отдельно: Устаревшая предыдущая версия сохранена отдельно:
- `docs/Personal_Messages/Спецификация_DM_v0.5_устаревшая.md` - `docs/Personal_Messages/Спецификация_DM_v0.5_устаревшая.md`
@@ -300,12 +304,17 @@
`sync_servers` не являются списком пользовательских серверов доставки DM. `sync_servers` не являются списком пользовательских серверов доставки DM.
### 7.3. Несколько серверов у отправителя и получателя ### 7.3. Два сервера у отправителя и получателя
Протокол должен поддерживать ситуацию, когда: Актуальный runtime исходит из ограничения:
- у отправителя несколько `access_servers`; - у одного пользователя не более двух `access_servers`;
- у получателя несколько `access_servers`; - следовательно, у локального сервера есть не более одного peer для репликации данных пользователя.
Протокол поддерживает ситуацию, когда:
- у отправителя один или два `access_servers`;
- у получателя один или два `access_servers`;
- часть серверов у сторон совпадает; - часть серверов у сторон совпадает;
- часть серверов уникальна. - часть серверов уникальна.
@@ -332,7 +341,14 @@
- `ReceiveIncomingMessage` — приём одной входящей копии, входящих редактирований и входящего read-receipt; - `ReceiveIncomingMessage` — приём одной входящей копии, входящих редактирований и входящего read-receipt;
- `ReceiveOutcomingMessage` — алиас `SendMessagePair`. - `ReceiveOutcomingMessage` — алиас `SendMessagePair`.
### 8.2. Новые методы, которые нужны ### 8.2. Межсерверная синхронизация
- `ReceiveOutcomingMessage` — прежняя полная пара второго access-сервера отправителя;
- `ReceiveIncomingMessage` — прежняя одиночная входящая копия;
- `DmSyncBatch` — pull событий с `synced=false` и ACK сохранённой предыдущей страницы;
- `GetDmDeliveryStatus` — единственная новая read-only проверка доставки по существующему `messageKey`.
Delivery-state между серверами не передаётся и не объединяется.
### 8.3. Серверный слой диалогов ### 8.3. Серверный слой диалогов
@@ -446,8 +462,12 @@ Request:
Правила: Правила:
- клиенту достаточно отправить пару на один любой доступный сервер; - клиенту достаточно отправить пару на один любой доступный сервер;
- успешный `status=200` означает локальное сохранение;
- первая сетевая попытка выполняется до ответа клиенту, параллельно для двух серверов получателя;
- сервер после принятия сам отвечает за дальнейшую межсерверную доставку. - сервер после принятия сам отвечает за дальнейшую межсерверную доставку.
Ответ сохраняет прежние `baseKey`, `incomingKey`, `outgoingKey` и счётчики доставки в клиентские сессии. Единственное новое поле — `deliveryState`: `accepted`, `delivered` или `failed`.
### 10.2. `ReceiveIncomingMessage` ### 10.2. `ReceiveIncomingMessage`
Назначение: Назначение:
@@ -469,7 +489,7 @@ Request:
} }
``` ```
`sourceServerLogin` необязателен и используется как best-effort подсказка, чтобы сервер при дальнейшей пересылке не отправлял то же событие обратно серверу-источнику. `sourceServerLogin` пока доверяется без отдельной межсерверной подписи. Сервер всё равно проверяет пользовательскую подпись самого `SHiNE_DM`. Успешный ответ сохраняет старые поля `messageKey`, `baseKey` и счётчики realtime-доставки.
### 10.3. `DeleteMessage` ### 10.3. `DeleteMessage`
@@ -558,6 +578,8 @@ UI-следствие для клиента:
- если точное время для старого сообщения неизвестно, но по более новым данным видно, что сообщение уже точно прочитано, UI может показывать его как прочитанное без точного времени; - если точное время для старого сообщения неизвестно, но по более новым данным видно, что сообщение уже точно прочитано, UI может показывать его как прочитанное без точного времени;
- read-receipt при этом остаётся отдельным DM-событием синхронизации, но в обычную историю страницы не подмешивается. - read-receipt при этом остаётся отдельным DM-событием синхронизации, но в обычную историю страницы не подмешивается.
Для исходящих элементов сервер также возвращает единственное поле `deliveryState`.
## 11. Межсерверная доставка ## 11. Межсерверная доставка
### 11.1. Клиентская сторона ### 11.1. Клиентская сторона
@@ -568,47 +590,23 @@ UI-следствие для клиента:
### 11.2. Серверная сторона ### 11.2. Серверная сторона
После принятия валидного события сервер должен отправлять его: После локального принятия пара получает изменяемое состояние `accepted`. Сервер сразу вызывает до двух текущих маршрутов получателя параллельно. Успех хотя бы одного маршрута переводит сообщение в терминальное `delivered`; ждать второй маршрут не требуется.
- на все серверы из `access_servers` отправителя; После первой попытки полная пара передаётся единственному второму access-серверу отправителя через прежний `ReceiveOutcomingMessage`. Peer самостоятельно пробует актуальные маршруты получателя; delivery-state в запрос не входит.
- на все серверы из `access_servers` получателя.
Если часть серверов совпадает, это допустимо. Идемпотентность обязательна. ACK означает запись в БД; повтор уже известной ревизии считается успешным ACK.
Если один и тот же сервер присутствует у обеих сторон, он не должен слать сообщение сам себе повторно, но обязан локально сохранить событие и доставить его в нужные пользовательские сессии.
Идемпотентность обязательна.
### 11.3. Догоняющая синхронизация истории ### 11.3. Догоняющая синхронизация истории
Для восстановления пропущенных DM-событий между access-серверами используется отдельная операция: Для восстановления пропущенных DM-событий между access-серверами используется pull-операция:
- `DmSyncBatch` - `DmSyncBatch`
Сервер-получатель синхронизации запрашивает у другого access-сервера историю одного пользователя по курсору: Второй сервер запрашивает outbox-события владельца с `synced=false`, сохраняет их и передаёт подтверждённые `syncId` в `ackSyncIds` следующего запроса. Источник ставит `synced=true` только после ACK. Каждый цикл начинается с начала списка; курсор нужен только для страниц текущего цикла.
- `ownerLogin`; Полная пара отправителя передаётся двумя blob. Входящая копия получателя и tombstone передаются одним blob.
- `afterStoredAtMs`;
- `afterMessageKey`;
- `limit`, максимум `500`;
- `maxBytes`, ограничение суммарного размера raw-блоков пачки.
Удалённый сервер отдаёт все DM-события, относящиеся к этому пользователю: Полная пара отправителя передаётся одним элементом с двумя blob в порядке incoming/outgoing. Входящая копия получателя и tombstone передаются одним blob.
- контентные копии и read-receipt по `target_login`;
- tombstone типов `5/6/7/8`, где пользователь участвует как `fromLogin` или `toLogin`.
Порядок пачки:
- `created_at_ms ASC`;
- `message_key ASC`.
Курсор хранится локально для пары:
- пользователь;
- удалённый access-сервер.
При первом добавлении сервера или отсутствии курсора синхронизация стартует с `0` и постепенно подтягивает всю доступную историю пачками.
При применении событий, полученных через `DmSyncBatch`, сервер: При применении событий, полученных через `DmSyncBatch`, сервер:
@@ -618,17 +616,17 @@ UI-следствие для клиента:
- не отправляет realtime/push-уведомления клиентам; - не отправляет realtime/push-уведомления клиентам;
- не запускает повторный fan-out, чтобы не создавать циклы. - не запускает повторный fan-out, чтобы не создавать циклы.
Плановый sync запускается фоном после старта WebSocket-сервера и повторяется раз в 6 часов. Синхронизация настроек и DM выполняется одним периодическим процессом и через один последовательный WS-сеанс с peer. Выборка DM использует частичный индекс по `synced=false`.
В текущей реализации межсерверная авторизация для `DmSyncBatch` ещё не включена. Сервер отдаёт пачку только если сам локально является access-сервером `ownerLogin` по актуальной таблице `user_access_servers_current`. В текущей реализации межсерверная авторизация DM ещё не включена. Принимающий сервер проверяет, что сам является access-сервером `ownerLogin`, и всегда проверяет пользовательские подписи signed-блоков.
### 11.4. Ошибки доставки ### 11.4. Ошибки доставки
Если часть серверов временно недоступна: Если ни один сервер получателя не подтвердил запись, попытки выполняются сразу, через 30 секунд, 5 минут, 25 минут и 1 час от первоначального принятия.
- это не должно отменять локальное принятие уже валидного сообщения; Перед попытками через 5 минут, 25 минут и час сервер read-only спрашивает peer отправителя через `GetDmDeliveryStatus(messageKey)`. Если peer уже доставил хотя бы на один сервер, сообщение считается доставленным.
- повторная доставка может делаться отдельным retry-механизмом;
- повторное получение того же события должно быть безопасным. После неудачной последней попытки устанавливается `failed`; дальнейших автоматических попыток и кнопки ручного повтора нет.
## 12. Хранение в БД ## 12. Хранение в БД
@@ -653,11 +651,13 @@ UI-следствие для клиента:
Сообщение об удалении переписки тоже хранится в БД, а старые сообщения до его времени из БД удаляются. Сообщение об удалении переписки тоже хранится в БД, а старые сообщения до его времени из БД удаляются.
Для догоняющей межсерверной синхронизации дополнительно используются: Изменяемая сетевая часть хранится отдельно:
- индекс по `target_login`, `created_at_ms`, `message_key`; - `dm_delivery_state` — состояние и расписание доставки исходящей пары;
- отдельные индексы по delete-событиям для `from_login` и `to_login`; - `dm_sync_outbox` — событие владельца и единственный флаг ACK `synced`;
- таблица `dm_sync_peer_state` с курсором чтения для пары `ownerLogin + remoteServerLogin`. - частичные индексы содержат только due/unsynced строки.
Legacy-таблица `dm_sync_peer_state` после миграции v12 физически остаётся для безопасной установки ZIP-накладки, но новым DM-кодом не используется.
## 13. Что обязательно должно измениться в коде относительно v0.5 ## 13. Что обязательно должно измениться в коде относительно v0.5
@@ -670,7 +670,7 @@ UI-следствие для клиента:
- межсерверная маршрутизация DM должна идти через `access_servers`; - межсерверная маршрутизация DM должна идти через `access_servers`;
- сервер должен добирать отсутствующих пользователей из Solana PDA до проверки подписи DM; - сервер должен добирать отсутствующих пользователей из Solana PDA до проверки подписи DM;
- при выборе актуальной версии должен учитываться `reencryptedAtMs`, если `revisionTimeMs` совпадает; - при выборе актуальной версии должен учитываться `reencryptedAtMs`, если `revisionTimeMs` совпадает;
- логика должна быть безопасна для нескольких серверов у каждой стороны. - логика должна быть безопасна для одного или двух серверов у каждой стороны.
## 14. Что в v1 пока не входит ## 14. Что в v1 пока не входит
@@ -678,4 +678,4 @@ UI-следствие для клиента:
- хранение отдельного `keyId` шифрования в DM; - хранение отдельного `keyId` шифрования в DM;
- ротация `clientKey`; - ротация `clientKey`;
- финальная конкретная UI-реализация массовой перешифровки; - финальная конкретная UI-реализация массовой перешифровки;
- межсерверная авторизация `DmSyncBatch`. - межсерверная авторизация DM-синхронизации и доставки.
+4 -1
View File
@@ -13,6 +13,9 @@
Логика протокола, API и поведение сервера описаны отдельно: Логика протокола, API и поведение сервера описаны отдельно:
- `docs/Personal_Messages/Протокол_DM_v1.md` - `docs/Personal_Messages/Протокол_DM_v1.md`
- `docs/Personal_Messages/Доставка_и_синхронизация_DM.md`
`deliveryState`, retry-времена и sync-флаг не входят в `SHiNE_DM` и не подписываются пользователем. Это изменяемые серверные метаданные, хранящиеся отдельно от raw-контейнера. Для корреляции доставки используется уже существующий `messageKey`; добавление delivery-воркера не меняет ни одного байта формата ниже.
## 1. Общие правила ## 1. Общие правила
@@ -337,4 +340,4 @@ ReadReceiptBody_v1_0
В версии DM v1 все типы `1..8` используют единый контейнер `SHiNE_DM`. В версии DM v1 все типы `1..8` используют единый контейнер `SHiNE_DM`.
Межсерверная операция `DmSyncBatch` не вводит новый байтовый формат DM. Она передаёт уже сохранённые raw-контейнеры `SHiNE_DM` в Base64 вместе с серверными метаданными курсора (`storedAtMs`, `messageKey`), а принимающий сервер заново проверяет подпись и применяет тот же контейнер по его `messageType`. Межсерверные операции `ReceiveOutcomingMessage`, `ReceiveIncomingMessage` и `DmSyncBatch` не вводят новый байтовый формат DM. Они передают уже сохранённые raw-контейнеры `SHiNE_DM` в Base64. `ackSyncIds` подтверждает только факт сохранения синхронизированной ревизии; delivery-state между серверами не передаётся. Принимающий сервер заново проверяет подпись и применяет контейнер по его `messageType`.
+12
View File
@@ -43,6 +43,7 @@ import {
deleteConversationMessagesBefore, deleteConversationMessagesBefore,
deleteSignedMessageByBaseKey, deleteSignedMessageByBaseKey,
markIncomingReadByBaseKey, markIncomingReadByBaseKey,
markOutgoingDeliveryState,
markOutgoingReadByBaseKey, markOutgoingReadByBaseKey,
normalizeDmChatId, normalizeDmChatId,
setContacts, setContacts,
@@ -1390,6 +1391,17 @@ async function init() {
}); });
}); });
authService.onEvent('DmDeliveryStateChanged', (evt) => {
const payload = evt?.payload || {};
const changed = markOutgoingDeliveryState({
outgoingKey: payload.outgoingKey,
baseKey: payload.baseKey,
deliveryState: payload.deliveryState,
});
if (!changed) return;
window.dispatchEvent(new CustomEvent('shine-dm-delivery-updated', { detail: payload }));
});
authService.onEvent('SignedMessageArrived', async (evt) => { authService.onEvent('SignedMessageArrived', async (evt) => {
const payload = evt?.payload || {}; const payload = evt?.payload || {};
const messageKey = String(payload.messageKey || '').trim(); const messageKey = String(payload.messageKey || '').trim();
+46 -3
View File
@@ -528,10 +528,28 @@ function resolveEffectiveReadState(messages, msg) {
function resolveDeliveryStatus(messages, msg) { function resolveDeliveryStatus(messages, msg) {
if (msg?.from !== 'out') return ''; if (msg?.from !== 'out') return '';
if (resolveEffectiveReadState(messages, msg).isRead) return '✓✓'; if (resolveEffectiveReadState(messages, msg).isRead) return '✓✓';
if (msg?.firstTick) return ''; const deliveryState = String(msg?.deliveryState || '').trim().toLowerCase();
if (deliveryState === 'delivered') return '✓';
if (deliveryState === 'failed') return '!';
if (deliveryState === 'accepted' || msg?.firstTick) return '✓';
return '…'; return '…';
} }
function resolveDeliveryTone(messages, msg) {
if (resolveEffectiveReadState(messages, msg).isRead) return 'read';
const deliveryState = String(msg?.deliveryState || '').trim().toLowerCase();
if (deliveryState === 'delivered') return 'delivered';
if (deliveryState === 'failed') return 'failed';
return 'accepted';
}
function resolveDeliveryNote(msg) {
if (msg?.from !== 'out') return '';
const deliveryState = String(msg?.deliveryState || '').trim().toLowerCase();
if (deliveryState === 'failed') return 'Сообщение не доставлено.';
return '';
}
function resolveMessageEditedTimeMs(msg) { function resolveMessageEditedTimeMs(msg) {
const revisionTimeMs = Number(msg?.revisionTimeMs || 0); const revisionTimeMs = Number(msg?.revisionTimeMs || 0);
if (!Number.isFinite(revisionTimeMs) || revisionTimeMs <= 0) return 0; if (!Number.isFinite(revisionTimeMs) || revisionTimeMs <= 0) return 0;
@@ -744,13 +762,21 @@ function renderLog(
const status = resolveDeliveryStatus(messages, msg); const status = resolveDeliveryStatus(messages, msg);
if (status) { if (status) {
const statusNode = document.createElement('span'); const statusNode = document.createElement('span');
statusNode.className = 'bubble-status'; statusNode.className = `bubble-status bubble-status--${resolveDeliveryTone(messages, msg)}`;
statusNode.textContent = status; statusNode.textContent = status;
metaNode.append(statusNode); metaNode.append(statusNode);
} }
bubble.append(metaNode); bubble.append(metaNode);
const deliveryNote = resolveDeliveryNote(msg);
if (deliveryNote) {
const deliveryNoteNode = document.createElement('div');
deliveryNoteNode.className = `bubble-delivery-note${String(msg?.deliveryState || '') === 'failed' ? ' bubble-delivery-note--failed' : ''}`;
deliveryNoteNode.textContent = deliveryNote;
bubble.append(deliveryNoteNode);
}
const editedAtMs = resolveMessageEditedTimeMs(msg); const editedAtMs = resolveMessageEditedTimeMs(msg);
if (editedAtMs > 0) { if (editedAtMs > 0) {
const editedNode = document.createElement('div'); const editedNode = document.createElement('div');
@@ -831,6 +857,7 @@ async function mergeDirectMessagesPage(chatId, payloadMessages) {
readAtMs: Number(item?.readAtMs || 0), readAtMs: Number(item?.readAtMs || 0),
rawBlobB64: blobB64, rawBlobB64: blobB64,
revisionTimeMs: Number(item?.revisionTimeMs || parsed?.revisionTimeMs || 0), revisionTimeMs: Number(item?.revisionTimeMs || parsed?.revisionTimeMs || 0),
deliveryState: String(item?.deliveryState || ''),
}); });
} catch (error) { } catch (error) {
addAppLogEntry({ addAppLogEntry({
@@ -1248,7 +1275,12 @@ export function render({ navigate, route, chrome }) {
focusInputToEnd(); focusInputToEnd();
}; };
const applyLocalRevision = async ({ localOutgoingBlobB64, fallbackMessageKey = '', fallbackBaseKey = '' }) => { const applyLocalRevision = async ({
localOutgoingBlobB64,
fallbackMessageKey = '',
fallbackBaseKey = '',
deliveryState = '',
}) => {
if (!localOutgoingBlobB64) return; if (!localOutgoingBlobB64) return;
try { try {
const parsed = authService.parseSignedMessageBlob(localOutgoingBlobB64); const parsed = authService.parseSignedMessageBlob(localOutgoingBlobB64);
@@ -1271,6 +1303,7 @@ export function render({ navigate, route, chrome }) {
unread: false, unread: false,
rawBlobB64: localOutgoingBlobB64, rawBlobB64: localOutgoingBlobB64,
revisionTimeMs: Number(parsed?.revisionTimeMs || 0), revisionTimeMs: Number(parsed?.revisionTimeMs || 0),
deliveryState,
deleted: Boolean(parsed?.deleted), deleted: Boolean(parsed?.deleted),
}); });
return true; return true;
@@ -1346,6 +1379,7 @@ export function render({ navigate, route, chrome }) {
markOutgoingSent(tempId, { markOutgoingSent(tempId, {
messageKey: result?.outgoingKey || '', messageKey: result?.outgoingKey || '',
baseKey: result?.baseKey || result?.localBaseKey || '', baseKey: result?.baseKey || result?.localBaseKey || '',
deliveryState: result?.deliveryState || 'accepted',
}); });
} }
@@ -1353,6 +1387,7 @@ export function render({ navigate, route, chrome }) {
localOutgoingBlobB64: result?.localOutgoingBlobB64 || '', localOutgoingBlobB64: result?.localOutgoingBlobB64 || '',
fallbackMessageKey: result?.outgoingKey || '', fallbackMessageKey: result?.outgoingKey || '',
fallbackBaseKey: result?.baseKey || result?.localBaseKey || '', fallbackBaseKey: result?.baseKey || result?.localBaseKey || '',
deliveryState: result?.deliveryState || 'accepted',
}); });
if (editing) { if (editing) {
@@ -1580,7 +1615,14 @@ export function render({ navigate, route, chrome }) {
void sendReadReceiptsForVisible(chatId); void sendReadReceiptsForVisible(chatId);
}; };
const handleDeliveryRefresh = () => {
preserveComposerSelection(input, () => {
renderChatLog({ scrollMode: 'preserve', markAsRead: false });
});
};
window.addEventListener('shine-chat-messages-updated', handleIncomingChatRefresh); window.addEventListener('shine-chat-messages-updated', handleIncomingChatRefresh);
window.addEventListener('shine-dm-delivery-updated', handleDeliveryRefresh);
wrap.append(historyLoader, log); wrap.append(historyLoader, log);
screen.append(wrap); screen.append(wrap);
@@ -1609,6 +1651,7 @@ export function render({ navigate, route, chrome }) {
hideUnreadSeparator({ rerender: false }); hideUnreadSeparator({ rerender: false });
stopAllTwemojiAnimations(); stopAllTwemojiAnimations();
window.removeEventListener('shine-chat-messages-updated', handleIncomingChatRefresh); window.removeEventListener('shine-chat-messages-updated', handleIncomingChatRefresh);
window.removeEventListener('shine-dm-delivery-updated', handleDeliveryRefresh);
boundScrollContainer?.removeEventListener('scroll', handleHistoryScroll); boundScrollContainer?.removeEventListener('scroll', handleHistoryScroll);
clearUnreadSeparatorHideTimer(); clearUnreadSeparatorHideTimer();
document.body.classList.remove('chat-topbar-overlay'); document.body.classList.remove('chat-topbar-overlay');
+2
View File
@@ -17,6 +17,7 @@ const DM_BLOB_PREVIEW_CACHE = new Map();
const DM_BLOB_PREVIEW_PENDING = new Map(); const DM_BLOB_PREVIEW_PENDING = new Map();
const dmAvatarSnapshotCache = new Map(); const dmAvatarSnapshotCache = new Map();
const dmAvatarPendingByLogin = new Map(); const dmAvatarPendingByLogin = new Map();
const SVG_CHEVRON = '<svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round" aria-hidden="true"><path d="M9 6l6 6-6 6"/></svg>';
const RELATION_ORDER = new Map([ const RELATION_ORDER = new Map([
['close_friend', 0], ['close_friend', 0],
@@ -166,6 +167,7 @@ function compareChatRows(a, b) {
export function render({ navigate, chrome }) { export function render({ navigate, chrome }) {
const screen = document.createElement('section'); const screen = document.createElement('section');
screen.className = 'stack dm-screen dm-list-screen'; screen.className = 'stack dm-screen dm-list-screen';
const login = String(state.session.login || '').trim();
const head = document.createElement('header'); const head = document.createElement('header');
head.className = 'dm-head'; head.className = 'dm-head';
head.innerHTML = ` head.innerHTML = `
+37 -1
View File
@@ -414,6 +414,7 @@ function persistMessageRecord(chatId, row) {
readAtMs: Number(row.readAtMs || 0), readAtMs: Number(row.readAtMs || 0),
readReceiptSent: Boolean(row.readReceiptSent), readReceiptSent: Boolean(row.readReceiptSent),
refBaseKey: String(row.refBaseKey || ''), refBaseKey: String(row.refBaseKey || ''),
deliveryState: String(row.deliveryState || ''),
ts: resolvedTs > 0 ? resolvedTs : Date.now(), ts: resolvedTs > 0 ? resolvedTs : Date.now(),
}).catch(() => {}); }).catch(() => {});
} }
@@ -450,6 +451,7 @@ export async function hydrateMessagesFromStore() {
readAtMs: Number(row.readAtMs || 0), readAtMs: Number(row.readAtMs || 0),
readReceiptSent: Boolean(row.readReceiptSent), readReceiptSent: Boolean(row.readReceiptSent),
refBaseKey: String(row.refBaseKey || ''), refBaseKey: String(row.refBaseKey || ''),
deliveryState: String(row.deliveryState || ''),
createdAtMs: Number(row.ts || 0), createdAtMs: Number(row.ts || 0),
}); });
}); });
@@ -531,7 +533,11 @@ export function addOutgoingPendingMessage(chatId, text) {
return tempId; return tempId;
} }
export function markOutgoingSent(tempId, { messageKey = '', baseKey = '' } = {}) { export function markOutgoingSent(tempId, {
messageKey = '',
baseKey = '',
deliveryState = 'accepted',
} = {}) {
if (!tempId) return; if (!tempId) return;
const keys = Object.keys(state.chats || {}); const keys = Object.keys(state.chats || {});
keys.forEach((chatId) => { keys.forEach((chatId) => {
@@ -541,6 +547,7 @@ export function markOutgoingSent(tempId, { messageKey = '', baseKey = '' } = {})
row.firstTick = true; row.firstTick = true;
row.messageKey = messageKey || row.messageKey || ''; row.messageKey = messageKey || row.messageKey || '';
row.baseKey = baseKey || row.baseKey || ''; row.baseKey = baseKey || row.baseKey || '';
row.deliveryState = String(deliveryState || row.deliveryState || 'accepted');
if (messageKey) { if (messageKey) {
state.knownMessageKeys[messageKey] = true; state.knownMessageKeys[messageKey] = true;
persistMessageRecord(chatId, row); persistMessageRecord(chatId, row);
@@ -549,6 +556,33 @@ export function markOutgoingSent(tempId, { messageKey = '', baseKey = '' } = {})
}); });
} }
export function markOutgoingDeliveryState({
outgoingKey = '',
baseKey = '',
deliveryState = '',
} = {}) {
const normalizedOutgoingKey = String(outgoingKey || '').trim();
const normalizedBaseKey = String(baseKey || '').trim();
let changed = false;
Object.keys(state.chats || {}).forEach((chatId) => {
getChatMessages(chatId).forEach((row) => {
if (row?.from !== 'out') return;
const matches = (
(normalizedOutgoingKey && String(row.messageKey || '') === normalizedOutgoingKey)
|| (normalizedBaseKey && String(row.baseKey || '') === normalizedBaseKey)
);
if (!matches) return;
row.deliveryState = String(deliveryState || row.deliveryState || 'accepted');
if (row.deliveryState === 'delivered') {
row.firstTick = true;
}
persistMessageRecord(chatId, row);
changed = true;
});
});
return changed;
}
export function markOutgoingReadByBaseKey(baseKey, readAtMs = 0) { export function markOutgoingReadByBaseKey(baseKey, readAtMs = 0) {
if (!baseKey) return; if (!baseKey) return;
const normalizedReadAtMs = Number(readAtMs || 0); const normalizedReadAtMs = Number(readAtMs || 0);
@@ -629,6 +663,7 @@ export function addSignedMessageToChat({
rawBlobB64 = '', rawBlobB64 = '',
refBaseKey = '', refBaseKey = '',
revisionTimeMs = 0, revisionTimeMs = 0,
deliveryState = '',
deleted = false, deleted = false,
} = {}) { } = {}) {
const normalizedChatId = normalizeDmChatId(chatId); const normalizedChatId = normalizeDmChatId(chatId);
@@ -663,6 +698,7 @@ export function addSignedMessageToChat({
row.messageType = Number(messageType || 0); row.messageType = Number(messageType || 0);
row.rawBlobB64 = String(rawBlobB64 || ''); row.rawBlobB64 = String(rawBlobB64 || '');
row.revisionTimeMs = nextRevision; row.revisionTimeMs = nextRevision;
row.deliveryState = String(deliveryState || existing?.deliveryState || (row.from === 'out' ? 'accepted' : ''));
row.unread = row.from === 'in' ? Boolean(unread) : false; row.unread = row.from === 'in' ? Boolean(unread) : false;
row.refBaseKey = String(refBaseKey || ''); row.refBaseKey = String(refBaseKey || '');
row.firstTick = row.from === 'out'; row.firstTick = row.from === 'out';
+23 -1
View File
@@ -1827,7 +1827,29 @@
.bubble-status { .bubble-status {
min-width: 14px; min-width: 14px;
text-align: right; text-align: right;
color: rgba(215, 230, 255, 0.92); color: rgba(148, 163, 184, 0.92);
}
.bubble-status--delivered,
.bubble-status--read {
color: rgba(215, 230, 255, 0.96);
text-shadow: 0 0 6px rgba(126, 180, 255, 0.36);
}
.bubble-status--failed {
color: rgba(255, 145, 145, 0.96);
}
.bubble-delivery-note {
max-width: 280px;
margin-top: 6px;
font-size: 10px;
line-height: 1.3;
color: rgba(255, 220, 142, 0.92);
}
.bubble-delivery-note--failed {
color: rgba(255, 145, 145, 0.96);
} }
.bubble.in { .bubble.in {