Browse Source

Commit DM activity timestamps together with message state

master
evgeny 2 days ago
parent
commit
10de4f1c71
  1. 15
      src/dm/dm_core.c
  2. 2
      src/dm/dm_core.h

15
src/dm/dm_core.c

@ -207,13 +207,14 @@ struct dm_conv {
uint8_t content_key[DM_CONTENT_KEY_SIZE];
uint64_t last_out_seq;
uint64_t last_in_seq;
uint64_t last_ts;
};
static int dm_conv_load(struct dm_state* dm, const char* conv_id, struct dm_conv* c) {
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(dm->db,
"SELECT peer_node_id, peer_x25519, peer_ed25519, COALESCE(peer_name,''),"
" COALESCE(group_id,0), last_out_seq, last_in_seq FROM dm_conversations WHERE conv_id=?",
" COALESCE(group_id,0), last_out_seq, last_in_seq, last_ts FROM dm_conversations WHERE conv_id=?",
-1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation query failed: %s", DM_ID, sqlite3_errmsg(dm->db));
return -1;
@ -238,6 +239,7 @@ static int dm_conv_load(struct dm_state* dm, const char* conv_id, struct dm_conv
c->group_id = (uint64_t)sqlite3_column_int64(st, 4);
c->last_out_seq = (uint64_t)sqlite3_column_int64(st, 5);
c->last_in_seq = (uint64_t)sqlite3_column_int64(st, 6);
c->last_ts = (uint64_t)sqlite3_column_int64(st, 7);
/* Старые беседы могли быть авто-созданы с group_id=0 (маршрутизация только напрямую) —
* подтягиваем общую группу, чтобы DM шёл через неё (indirect/BGP). */
if (c->group_id == 0) {
@ -256,11 +258,11 @@ static int dm_conv_save(struct dm_state* dm, const struct dm_conv* c) {
if (sqlite3_prepare_v2(dm->db,
"INSERT INTO dm_conversations(conv_id,peer_node_id,peer_x25519,peer_ed25519,"
" peer_name,group_id,last_out_seq,last_in_seq,created_at,last_ts)"
" VALUES(?,?,?,?,?,?,?,?,?,0)"
" VALUES(?,?,?,?,?,?,?,?,?,?)"
" ON CONFLICT(conv_id) DO UPDATE SET"
" peer_x25519=excluded.peer_x25519, peer_ed25519=excluded.peer_ed25519,"
" peer_name=excluded.peer_name, group_id=excluded.group_id,"
" last_out_seq=excluded.last_out_seq, last_in_seq=excluded.last_in_seq",
" last_out_seq=excluded.last_out_seq, last_in_seq=excluded.last_in_seq, last_ts=excluded.last_ts",
-1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation save prepare failed: %s", DM_ID, sqlite3_errmsg(dm->db));
return -1;
@ -274,6 +276,7 @@ static int dm_conv_save(struct dm_state* dm, const struct dm_conv* c) {
sqlite3_bind_int64(st, 7, (sqlite3_int64)c->last_out_seq);
sqlite3_bind_int64(st, 8, (sqlite3_int64)c->last_in_seq);
sqlite3_bind_int64(st, 9, (sqlite3_int64)ntp_time_get_seconds(dm->inst));
sqlite3_bind_int64(st, 10, (sqlite3_int64)c->last_ts);
int rc = sqlite3_step(st) == SQLITE_DONE ? 0 : -1;
if (rc) DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conversation save failed conv=%s: %s", DM_ID, c->conv_id, sqlite3_errmsg(dm->db));
sqlite3_finalize(st);
@ -442,6 +445,9 @@ int dm_accept_message(struct UTUN_INSTANCE* inst, const uint8_t* msg, size_t len
if (dm_exec(dm, "BEGIN IMMEDIATE") != 0) return -1;
int inserted = dm_store_message(dm, conv_id, 0, msg, len);
if (seq > c.last_in_seq) c.last_in_seq = seq; /* Только статистика; не курсор dedup. */
uint64_t ts;
memcpy(&ts, msg + 16, 8);
if (ts > c.last_ts) c.last_ts = ts;
if (inserted < 0 || dm_conv_save(dm, &c) != 0 || dm_exec(dm, "COMMIT") != 0) {
dm_exec(dm, "ROLLBACK");
return -1;
@ -679,7 +685,7 @@ int dm_start(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
memcpy(c.peer_ed25519, peer_ed25519, 32);
if (peer_name) snprintf(c.peer_name, sizeof(c.peer_name), "%s", peer_name);
c.group_id = source_ch_id ? strtoull(source_ch_id, NULL, 10) : 0;
dm_derive_content_key(inst->my_keys.private_key, c.peer_x25519, c.content_key);
if (dm_derive_content_key(inst->my_keys.private_key, c.peer_x25519, c.content_key) != 0) return -1;
if (dm_conv_save(dm, &c) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: conv_save failed conv=%s", DM_ID, conv_id);
@ -722,6 +728,7 @@ int dm_send(struct UTUN_INSTANCE* inst, const char* conv_id, const char* content
sqlite3_stmt* st = NULL;
if (dm_exec(dm, "BEGIN IMMEDIATE") != 0) { u_free(body); return -1; }
c.last_out_seq = seq;
if (ts > c.last_ts) c.last_ts = ts;
if (dm_store_message(dm, conv_id, 1, body, len) != 1 || dm_conv_save(dm, &c) != 0) goto rollback;
int rc = sqlite3_prepare_v2(dm->db, "INSERT INTO dm_outbox(conv_id,seq,body) VALUES(?,?,?)", -1, &st, NULL);
if (rc == SQLITE_OK) {

2
src/dm/dm_core.h

@ -67,7 +67,7 @@ void dm_core_destroy(struct UTUN_INSTANCE* inst);
/* Начать беседу с пользователем (peer). source_ch_id — канал, где нашли target
* (общая группа; нужен для маршрутизации и подключения). Создаёт беседу если её
* ещё нет, выводит conv_id/ключи, пытается установить прямое соединение.
* ещё нет, выводит conv_id/ключи. Подключение для передачи медиа принадлежит media.
* Возвращает 0 при успехе, <0 при ошибке. */
int dm_start(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
const uint8_t peer_x25519[32], const uint8_t peer_ed25519[32],

Loading…
Cancel
Save