Browse Source
- gui_bridge.h/cpp — кросс-поточные асинхронные коллбэки uasync↔GUI - chat_sync.h/c — протокол синхронизации (INIT_SYNC/RESP, SEND_DATA, PUSH/ACK, SYNC_DONE) - DbManager полностью переписан: per-channel таблицы msg_<ch> + peers_<ch>, 4 ключа в channels, signature, protocol/rtt в node_addresses, sync-методы - ChatPropagator удалён, заменён на chat_sync - db_sync_stub.c — заглушки для сборки chatgui без LMDB-версии db_sync - db_schema.md актуализирован под новую схемуchatgui
17 changed files with 1977 additions and 559 deletions
File diff suppressed because it is too large
Load Diff
@ -1,107 +1,244 @@ |
|||||||
# База данных децентрализованного чата (chatgui) |
# База данных децентрализованного чата (chatgui) |
||||||
|
|
||||||
Сторона сервиса (gui): |
Чат децентрализованный — нет центрального сервера. Каждый узел хранит локальную SQLite БД |
||||||
- |
(`chats.db`, WAL mode). Синхронизация сообщений между узлами — через протокол `chat_sync` |
||||||
|
поверх `etcp_router` (svc_id 0x30). Данные о других узлах поступают из NODEINFO |
||||||
|
через conn_mgr. |
||||||
|
|
||||||
## Архитектура хранения |
Вся работа с БД — в GUI-потоке (класс `DbManager`). Протокол синхронизации работает |
||||||
|
в uasync-потоке (`chat_sync.c`) и общается с БД асинхронно через `gui_bridge`. |
||||||
|
|
||||||
Чат **децентрализованный** — нет центрального сервера. Каждый узел (нода uTun) хранит свою |
## Обзор таблиц |
||||||
**локальную SQLite БД**. В ней: |
|
||||||
|
|
||||||
- **Информация о узлах (и о себе)** (как подключиться: адреса, порты, ключи — всё из NODEINFO). плюс nickname. - и список узлов = список участников. каждый узел имеет parent node (дерево). |
| Таблица | Назначение | |
||||||
- **История сообщений** — |
|---------|-----------| |
||||||
|
| `local_identity` | Наш собственный узел (id, ключи) | |
||||||
|
| `nodes` | Все известные узлы (node_id, ключи, last_seen) | |
||||||
|
| `node_addresses` | Адреса узлов (family, protocol, address, port, rtt, is_nat) | |
||||||
|
| `accounts` | Профили узлов (display_name, avatar) | |
||||||
|
| `channels` | Метаданные каналов/сетей (ключи, подпись, позиции скролла) | |
||||||
|
| `msg_<ch>` | Per-channel таблица сообщений | |
||||||
|
| `peers_<ch>` | Per-channel таблица участников (подписи входа, создателя, comment) | |
||||||
|
| `ui_state` | Локальное состояние UI (key-value) | |
||||||
|
|
||||||
Данные о других узлах поступают из NODEINFO gossip-протокола uTun и синхронизируются в фоне. |
--- |
||||||
Сообщения приходят через `msg_transport` (TCP, локальный IPC). |
|
||||||
аттачи - загружают и кешируют узлы в соответствии со своими настройками. |
|
||||||
|
|
||||||
## Структура таблиц SQLite |
## Таблицы |
||||||
|
|
||||||
### `local_identity` — наш собственный узел |
### `local_identity` — наш собственный узел |
||||||
|
|
||||||
|
```sql |
||||||
|
CREATE TABLE IF NOT EXISTS local_identity ( |
||||||
|
id INTEGER PRIMARY KEY CHECK (id = 1), |
||||||
|
node_id INTEGER NOT NULL UNIQUE, |
||||||
|
name TEXT NOT NULL, |
||||||
|
x25519_pubkey BLOB NOT NULL, |
||||||
|
x25519_privkey BLOB, |
||||||
|
ed25519_pubkey BLOB, |
||||||
|
created_at INTEGER DEFAULT (unixepoch()), |
||||||
|
updated_at INTEGER DEFAULT (unixepoch()) |
||||||
|
); |
||||||
|
``` |
||||||
|
|
||||||
### `my_nodes` — все мои устройства) |
### `nodes` — все известные узлы |
||||||
между моими устройствами синхронизируется контент - список чатов и сообщения. |
|
||||||
... todo |
|
||||||
|
|
||||||
|
|
||||||
```sql |
```sql |
||||||
CREATE TABLE IF NOT EXISTS nodes ( |
CREATE TABLE IF NOT EXISTS nodes ( |
||||||
id INTEGER PRIMARY KEY, |
node_id INTEGER PRIMARY KEY, |
||||||
node_id INTEGER NOT NULL UNIQUE, // id из utun |
name TEXT, |
||||||
name TEXT NOT NULL, // ник |
x25519_pubkey BLOB NOT NULL, |
||||||
pubkey (ed+x) |
ed25519_pubkey BLOB, |
||||||
last_seen INTEGER DEFAULT 0, |
last_seen_at INTEGER, |
||||||
online_rating INTEGER DEFAULT 0, // рейтинг доступности узла. чем больше число тем чаще узел онлайн |
created_at INTEGER DEFAULT (unixepoch()) |
||||||
speed_rating INTEGER DEFAULT 0, // рейтинг скорости обмена с узлом |
|
||||||
rtt_rating INTEGER DEFAULT 0, // rtt до узла (через промежуточные узлы если есть) |
|
||||||
connectivity_json TEXT, // как связаться с узлом - адреса, (пока пусто) |
|
||||||
parent_node_id INTEGER DEFAULT NULL, // каждая нода должна иметь родителя (кроме корневой) |
|
||||||
main_node_id INTEGER DEFAULT NULL, // если несколько устройств у клиента, id следующей ноды (циклический список) |
|
||||||
created_at TEXT DEFAULT (datetime('now')), |
|
||||||
); |
); |
||||||
``` |
``` |
||||||
|
|
||||||
|
| Поле | Описание | |
||||||
|
|------|---------| |
||||||
|
| `node_id` | ID узла в сети uTun | |
||||||
|
| `name` | Никнейм | |
||||||
|
| `x25519_pubkey` | Публичный ключ для key exchange (32 байта) | |
||||||
|
| `ed25519_pubkey` | Публичный ключ для подписей (32 байта), может быть NULL | |
||||||
|
| `last_seen_at` | Unix-время последней активности | |
||||||
|
|
||||||
|
Статус online определяется через conn_mgr, а не полем в БД. |
||||||
|
|
||||||
### `channels` — чаты (группы узлов) |
### `node_addresses` — адреса узлов |
||||||
|
|
||||||
|
```sql |
||||||
|
CREATE TABLE IF NOT EXISTS node_addresses ( |
||||||
|
id INTEGER PRIMARY KEY AUTOINCREMENT, |
||||||
|
node_id INTEGER NOT NULL REFERENCES nodes(node_id) ON DELETE CASCADE, |
||||||
|
family INTEGER NOT NULL CHECK(family IN (4, 6)), |
||||||
|
protocol INTEGER NOT NULL DEFAULT 1, |
||||||
|
address BLOB NOT NULL, |
||||||
|
port INTEGER NOT NULL CHECK(port > 0 AND port <= 65535), |
||||||
|
rtt INTEGER, |
||||||
|
is_nat INTEGER DEFAULT 0, |
||||||
|
created_at INTEGER DEFAULT (unixepoch()) |
||||||
|
); |
||||||
|
CREATE INDEX IF NOT EXISTS idx_na_node ON node_addresses(node_id); |
||||||
|
``` |
||||||
|
|
||||||
|
| Поле | Описание | |
||||||
|
|------|---------| |
||||||
|
| `family` | 4=IPv4, 6=IPv6 | |
||||||
|
| `protocol` | Битовая маска: 1=UDP (NODE_PROTO_UDP), 2=TCP (NODE_PROTO_TCP), 3=UDP+TCP | |
||||||
|
| `address` | BLOB: 4 байта для IPv4, 16 байт для IPv6 | |
||||||
|
| `port` | Порт | |
||||||
|
| `rtt` | RTT в миллисекундах, NULL если не измерен | |
||||||
|
| `is_nat` | 1 = узел за NAT | |
||||||
|
|
||||||
|
### `accounts` — профили узлов |
||||||
|
|
||||||
|
```sql |
||||||
|
CREATE TABLE IF NOT EXISTS accounts ( |
||||||
|
node_id INTEGER PRIMARY KEY REFERENCES nodes(node_id), |
||||||
|
display_name TEXT NOT NULL, |
||||||
|
avatar_color TEXT DEFAULT '#4A90E2', |
||||||
|
avatar_letter TEXT NOT NULL, |
||||||
|
is_contact INTEGER DEFAULT 1, |
||||||
|
created_at INTEGER DEFAULT (unixepoch()) |
||||||
|
); |
||||||
|
``` |
||||||
|
|
||||||
|
### `channels` — метаданные каналов |
||||||
|
|
||||||
```sql |
```sql |
||||||
CREATE TABLE IF NOT EXISTS channels ( |
CREATE TABLE IF NOT EXISTS channels ( |
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT, |
channel_id TEXT PRIMARY KEY, |
||||||
name TEXT NOT NULL, |
name TEXT NOT NULL, |
||||||
last_msg_at TEXT, |
owner_node_id INTEGER, |
||||||
|
is_dm INTEGER DEFAULT 0, |
||||||
|
x25519_pubkey BLOB NOT NULL, |
||||||
|
x25519_privkey BLOB, |
||||||
|
ed25519_pubkey BLOB NOT NULL, |
||||||
|
ed25519_privkey BLOB, |
||||||
|
signature BLOB NOT NULL, |
||||||
|
last_read_msg_id INTEGER, |
||||||
|
last_pos_msg_id INTEGER, |
||||||
|
created_at INTEGER DEFAULT (unixepoch()) |
||||||
|
); |
||||||
|
``` |
||||||
|
|
||||||
|
| Поле | Описание | |
||||||
|
|------|---------| |
||||||
|
| `channel_id` | Уникальный ID канала (напр. `"ch:general"`) | |
||||||
|
| `name` | Отображаемое имя (напр. `"#general"`) | |
||||||
|
| `owner_node_id` | ID создателя канала | |
||||||
|
| `is_dm` | 1 = личная переписка (direct message) | |
||||||
|
| `x25519_pubkey` | Публичный ключ канала (32 байта) | |
||||||
|
| `x25519_privkey` | Приватный ключ канала, NULL если локальный узел не админ | |
||||||
|
| `ed25519_pubkey` | Публичный ключ подписи канала (32 байта) | |
||||||
|
| `ed25519_privkey` | Приватный ключ подписи, NULL если локальный узел не админ | |
||||||
|
| `signature` | Ed25519 подпись всех полей записи приватным ключом канала | |
||||||
|
| `last_read_msg_id` | ID последнего прочитанного сообщения (локальное) | |
||||||
|
| `last_pos_msg_id` | ID сообщения на позиции скролла (для восстановления при открытии) | |
||||||
|
|
||||||
|
### `msg_<ch>` — per-channel таблица сообщений |
||||||
|
|
||||||
|
Для каждого канала создаётся отдельная таблица. Имя: `msg_` + sanitized `channel_id` |
||||||
|
(не-буквоцифры заменяются на `_`). |
||||||
|
|
||||||
|
```sql |
||||||
|
CREATE TABLE IF NOT EXISTS "msg_<ch>" ( |
||||||
|
id INTEGER PRIMARY KEY AUTOINCREMENT, |
||||||
|
node_id INTEGER NOT NULL, |
||||||
|
content_type TEXT NOT NULL, |
||||||
|
data BLOB NOT NULL, |
||||||
|
timestamp INTEGER NOT NULL, |
||||||
|
datahash INTEGER NOT NULL, |
||||||
|
chain_hash BLOB NOT NULL, |
||||||
|
signature BLOB, |
||||||
|
is_outgoing INTEGER DEFAULT 0, |
||||||
|
is_read INTEGER DEFAULT 0, |
||||||
|
sync_flags INTEGER DEFAULT 0, |
||||||
|
UNIQUE(timestamp, datahash) |
||||||
); |
); |
||||||
|
CREATE INDEX IF NOT EXISTS "idx_msg_<ch>_ts" ON "msg_<ch>"(timestamp, datahash); |
||||||
``` |
``` |
||||||
|
|
||||||
### `messages` — список пользователей в чате |
| Поле | Описание | |
||||||
|
|------|---------| |
||||||
|
| `node_id` | Автор сообщения | |
||||||
|
| `content_type` | MIME-тип: `"text/plain"`, `"image/png"`, ... | |
||||||
|
| `data` | Содержимое сообщения (BLOB) | |
||||||
|
| `timestamp` | Время отправки в микросекундах (первая часть ключа sync) | |
||||||
|
| `datahash` | Первые 8 байт SHA256(data), вторая часть ключа sync | |
||||||
|
| `chain_hash` | SHA256 цепочки (32 байта) — как в blockchain, связывает записи по порядку | |
||||||
|
| `signature` | Ed25519 подпись автора (64 байта), пока заглушка (64 нуля) | |
||||||
|
| `is_outgoing` | 1 = отправлено нами | |
||||||
|
| `is_read` | 1 = прочитано | |
||||||
|
| `sync_flags` | Битовая маска: 0x01 = WAS_SENT (запись подтверждена пирами) | |
||||||
|
|
||||||
|
**Дедупликация:** уникальный индекс `(timestamp, datahash)` предотвращает повторную вставку |
||||||
|
одного и того же сообщения при получении через разные пути. |
||||||
|
|
||||||
|
**Chain hash:** `chain_hash = SHA256(prev_chain_hash || timestamp || datahash)`. |
||||||
|
Обеспечивает проверку целостности порядка сообщений при синхронизации. |
||||||
|
|
||||||
|
### `peers_<ch>` — per-channel таблица участников |
||||||
|
|
||||||
```sql |
```sql |
||||||
CREATE TABLE IF NOT EXISTS channel_members ( |
CREATE TABLE IF NOT EXISTS "peers_<ch>" ( |
||||||
channel_id INTEGER, |
node_id INTEGER NOT NULL, |
||||||
node_id INTEGER, // если несколько устройств то ноды всех стройств |
join_sig BLOB NOT NULL, |
||||||
|
creator_sig BLOB, |
||||||
|
comment TEXT, |
||||||
|
joined_at INTEGER DEFAULT (unixepoch()), |
||||||
|
PRIMARY KEY (node_id) |
||||||
); |
); |
||||||
``` |
``` |
||||||
|
|
||||||
### `messages` — история сообщений (для каждой группы создаем новую таблицу) |
| Поле | Описание | |
||||||
|
|------|---------| |
||||||
|
| `join_sig` | Ed25519 подпись присоединяющегося: `sign(sk_user, channel_id || node_id)` | |
||||||
|
| `creator_sig` | Ed25519 подпись создателя канала: `sign(sk_creator, channel_id || node_id)`. NULL = гость (ограниченные права) | |
||||||
|
| `comment` | Комментарий/заметка об участнике | |
||||||
|
|
||||||
|
### `ui_state` — локальное состояние UI |
||||||
|
|
||||||
```sql |
```sql |
||||||
CREATE TABLE IF NOT EXISTS messages_ch<n> ( |
CREATE TABLE IF NOT EXISTS ui_state ( |
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT, |
key TEXT PRIMARY KEY, |
||||||
channel_id TEXT NOT NULL REFERENCES channels(channel_id) ON DELETE CASCADE, |
value TEXT |
||||||
author_node_id INTEGER NOT NULL, |
|
||||||
content TEXT NOT NULL, |
|
||||||
content_type TEXT DEFAULT 'text/plain', |
|
||||||
timestamp INTEGER NOT NULL, |
|
||||||
local_seq INTEGER DEFAULT 0, |
|
||||||
signature BLOB, |
|
||||||
is_outgoing INTEGER DEFAULT 0, |
|
||||||
is_read INTEGER DEFAULT 0, |
|
||||||
created_at TEXT DEFAULT (datetime('now')) |
|
||||||
); |
); |
||||||
CREATE INDEX IF NOT EXISTS idx_msg_channel_time ON messages(channel_id, timestamp); |
|
||||||
CREATE INDEX IF NOT EXISTS idx_msg_author ON messages(author_node_id); |
|
||||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_msg_dedup ON messages(channel_id, author_node_id, timestamp); |
|
||||||
``` |
``` |
||||||
|
|
||||||
**Пояснение:** каждое сообщение привязано к каналу. `author_node_id` — кто автор |
Key-value хранилище для произвольных настроек UI (позиции окон, выбранный канал, ...). |
||||||
(из `nodes.node_id`). `timestamp` — unix время в миллисекундах *по часам отправителя* |
|
||||||
(в децентрализованной системе нет глобальных часов, но для порядка в UI достаточно). |
--- |
||||||
`local_seq` — монотонно возрастающий номер в пределах канала для детерминированной |
|
||||||
сортировки при одинаковых timestamp. |
## Формат binary-записей sync (gui_bridge / chat_sync) |
||||||
|
|
||||||
**Дедупликация:** сообщение идентифицируется тройкой `(channel_id, author_node_id, |
Записи передаются между потоками и по сети в компактном бинарном формате: |
||||||
timestamp)`. Уникальный индекс `idx_msg_dedup` предотвращает дубликаты при |
|
||||||
повторной доставке через разные пути. |
### Insert record (для GUI_OP_INSERT_RECORD / PUSH wire) |
||||||
|
``` |
||||||
| Поле | Тип | Описание | |
[timestamp:8 LE][datahash:8 LE][node_id:8 LE][ct_len:1][content_type:ct_len][data_len:4 LE][data:data_len][chain_hash:32] |
||||||
|---|---|---| |
``` |
||||||
| `channel_id` | TEXT | Ссылка на канал | |
|
||||||
| `author_node_id` | INTEGER | Кто отправил (node_id) | |
### NodeInfo (GUI_OP_LOAD_NODEINFO) |
||||||
| `content` | TEXT | Текст сообщения | |
``` |
||||||
| `content_type` | TEXT | MIME-тип: `text/plain`, в будущем `image/png` и т.д. | |
[x25519_pubkey:32][ed25519_pubkey:32][addr_count:1]([family:1][proto:1][addr_len:1][addr:4|16][port:2 LE][rtt:2 LE])... |
||||||
| `timestamp` | INTEGER | Unix timestamp (мс) — когда отправлено | |
``` |
||||||
| `local_seq` | INTEGER | Локальный порядковый номер в канале | |
|
||||||
| `signature` | BLOB | Ed25519 подпись (64 байта), NULL если без подписи | |
### SEND_DATA cursor record (GUI_OP_CURSOR_NEXT) |
||||||
| `is_outgoing` | INTEGER | 1 = отправлено нами | |
``` |
||||||
| `is_read` | INTEGER | 1 = прочитано (для галочек в UI) | |
[ts:8 LE][dh:8 LE][node_id:8 LE][ct_len:1][ct:var][data_len:4 LE][data:var] |
||||||
|
``` |
||||||
|
|
||||||
|
--- |
||||||
|
|
||||||
|
## Порядок синхронизации (chat_sync) |
||||||
|
|
||||||
|
1. **Peer online** → для каждого общего канала: `INIT_SYNC(my_count, last_chain_hash)` |
||||||
|
2. **INIT_RESP** → сравнение chain_hash на позиции `min(counts)-1`. Если match → synced. |
||||||
|
Иначе sparse hashes (-1, -2, -4, -8, ...) для поиска точки расхождения. |
||||||
|
3. **SEND_DATA** → обмен недостающими записями пакетами до 32 штук |
||||||
|
4. **SYNC_DONE** → финальная сверка count + last_chain_hash |
||||||
|
5. **PUSH** → реал-тайм доставка новых сообщений (→ ACK_PUSH) |
||||||
|
6. **TTL** → раз в час удаление своих неподтверждённых записей старше 24ч |
||||||
|
|
||||||
|
Все операции с БД асинхронные: `chat_sync (uasync)` → `gui_bridge_call()` → |
||||||
|
GUI-поток (SQL) → `uasync_post()` → callback в chat_sync. |
||||||
|
|||||||
@ -0,0 +1,660 @@ |
|||||||
|
#include "chat_sync.h" |
||||||
|
#include "gui_bridge.h" |
||||||
|
|
||||||
|
#include "../../../src/utun_instance.h" |
||||||
|
#include "../../../src/etcp_router.h" |
||||||
|
#include "../../../src/etcp_api.h" |
||||||
|
#include "../../../src/etcp.h" |
||||||
|
#include "../../../src/conn_mgr.h" |
||||||
|
#include "../../../src/route_bgp.h" |
||||||
|
#include "../../../src/route_node.h" |
||||||
|
#include "../../../lib/u_async.h" |
||||||
|
#include "../../../lib/ll_queue.h" |
||||||
|
#include "../../../lib/debug_config.h" |
||||||
|
#include "../../../lib/mem.h" |
||||||
|
#include "../../../lib/sha256.h" |
||||||
|
#include "../../../lib/platform_compat.h" |
||||||
|
|
||||||
|
#include <string.h> |
||||||
|
|
||||||
|
static struct chat_sync* g_cs = NULL; |
||||||
|
|
||||||
|
/* ── Per-channel cache ── */ |
||||||
|
|
||||||
|
struct channel_cache { |
||||||
|
char channel_id[64]; |
||||||
|
uint64_t* peer_ids; |
||||||
|
int peer_count; |
||||||
|
uint32_t msg_count; |
||||||
|
uint8_t last_chain_hash[32]; |
||||||
|
uint8_t synced; |
||||||
|
}; |
||||||
|
|
||||||
|
struct chat_sync { |
||||||
|
struct UTUN_INSTANCE* inst; |
||||||
|
struct channel_cache* channels; |
||||||
|
int channel_count; |
||||||
|
void* refresh_timer; |
||||||
|
void* ttl_timer; |
||||||
|
uint8_t initialized; |
||||||
|
}; |
||||||
|
|
||||||
|
#define CS_ID "chat_sync" |
||||||
|
|
||||||
|
/* ── Send ── */ |
||||||
|
|
||||||
|
static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst, |
||||||
|
const uint8_t* payload, size_t len) { |
||||||
|
uint8_t ch_len = (uint8_t)strlen(ch_id); |
||||||
|
if (ch_len > 63) return -1; |
||||||
|
size_t total = 1 + 1 + ch_len + len; |
||||||
|
uint8_t* buf = u_malloc(total); |
||||||
|
if (!buf) return -1; |
||||||
|
uint8_t* p = buf; |
||||||
|
*p++ = ETCP_RT_ID_CHAT_SYNC; |
||||||
|
*p++ = ch_len; |
||||||
|
memcpy(p, ch_id, ch_len); p += ch_len; |
||||||
|
memcpy(p, payload, len); |
||||||
|
struct ll_entry* entry = queue_entry_new(0); |
||||||
|
if (!entry) { u_free(buf); return -1; } |
||||||
|
entry->dgram = buf; entry->len = total; |
||||||
|
return etcp_route_send(cs->inst, dst, entry, 0); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Channel cache ── */ |
||||||
|
|
||||||
|
static struct channel_cache* cs_find(struct chat_sync* cs, const char* ch_id) { |
||||||
|
for (int i = 0; i < cs->channel_count; i++) |
||||||
|
if (strcmp(cs->channels[i].channel_id, ch_id) == 0) return &cs->channels[i]; |
||||||
|
return NULL; |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Protocol message context (async state machine) ── */ |
||||||
|
|
||||||
|
struct cs_ctx { |
||||||
|
struct chat_sync* sync; |
||||||
|
uint64_t peer; |
||||||
|
char ch_id[64]; |
||||||
|
union { |
||||||
|
struct { |
||||||
|
uint32_t peer_count; |
||||||
|
uint8_t peer_last_ch[32]; |
||||||
|
uint32_t my_count; |
||||||
|
uint32_t test_pos; |
||||||
|
} init_sync; |
||||||
|
struct { |
||||||
|
uint32_t peer_count; |
||||||
|
uint32_t test_pos; |
||||||
|
uint8_t peer_ch[37]; /* 1 byte for sparse_count + 36 for first sparse */ |
||||||
|
int parse_off; /* offset in original payload for sparse data */ |
||||||
|
size_t parse_len; |
||||||
|
} init_resp; |
||||||
|
struct { |
||||||
|
uint32_t cursor_id; |
||||||
|
uint32_t from_pos; |
||||||
|
uint16_t sent_count; |
||||||
|
} send_data; |
||||||
|
struct { |
||||||
|
uint64_t ts; |
||||||
|
uint64_t dh; |
||||||
|
} push; |
||||||
|
}; |
||||||
|
}; |
||||||
|
|
||||||
|
static struct cs_ctx* cs_ctx_new(struct chat_sync* cs, uint64_t peer, const char* ch_id) { |
||||||
|
struct cs_ctx* ctx = u_calloc(1, sizeof(*ctx)); |
||||||
|
if (!ctx) return NULL; |
||||||
|
ctx->sync = cs; ctx->peer = peer; |
||||||
|
size_t n = strlen(ch_id); if (n > 63) n = 63; |
||||||
|
memcpy(ctx->ch_id, ch_id, n); |
||||||
|
return ctx; |
||||||
|
} |
||||||
|
|
||||||
|
/* ── gui_bridge async callbacks ── */ |
||||||
|
|
||||||
|
static void on_init_sync_count(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_init_sync_hash(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_init_sync_sparse(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_init_resp_hash(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_send_data_cursor_open(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_send_data_cursor_next(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_push_inserted(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_refresh_list(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
static void on_refresh_peers_and_count(void* ud, int op, const uint8_t* r, int rl); |
||||||
|
|
||||||
|
/* helper: build request [ch_id_len:1][ch_id:var] */ |
||||||
|
static int req_ch(const char* ch_id, uint8_t* buf) { |
||||||
|
uint8_t n = (uint8_t)strlen(ch_id); |
||||||
|
buf[0] = n; memcpy(buf + 1, ch_id, n); |
||||||
|
return 1 + n; |
||||||
|
} |
||||||
|
|
||||||
|
/* ── INIT_SYNC handler ── */ |
||||||
|
|
||||||
|
static void cs_handle_init_sync(struct chat_sync* cs, uint64_t peer, |
||||||
|
const char* ch_id, const uint8_t* pl, size_t len) { |
||||||
|
if (len < 36) return; |
||||||
|
struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); |
||||||
|
if (!ctx) return; |
||||||
|
memcpy(&ctx->init_sync.peer_count, pl, 4); |
||||||
|
memcpy(ctx->init_sync.peer_last_ch, pl + 4, 32); |
||||||
|
uint8_t req[65]; int rl = req_ch(ch_id, req); |
||||||
|
gui_bridge_call(GUI_OP_COUNT, req, rl, ctx, on_init_sync_count); |
||||||
|
} |
||||||
|
|
||||||
|
static void on_init_sync_count(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
uint32_t my_count = (rl >= 4) ? *(const uint32_t*)r : 0; |
||||||
|
ctx->init_sync.my_count = my_count; |
||||||
|
|
||||||
|
uint32_t tp = ctx->init_sync.peer_count < my_count ? ctx->init_sync.peer_count : my_count; |
||||||
|
if (tp > 0) tp--; |
||||||
|
ctx->init_sync.test_pos = tp; |
||||||
|
|
||||||
|
uint8_t req[69]; int nr = req_ch(ctx->ch_id, req); |
||||||
|
memcpy(req + nr, &tp, 4); nr += 4; |
||||||
|
gui_bridge_call(GUI_OP_CHAIN_HASH_AT, req, nr, ctx, on_init_sync_hash); |
||||||
|
} |
||||||
|
|
||||||
|
static void on_init_sync_hash(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
uint32_t tp = ctx->init_sync.test_pos; |
||||||
|
uint8_t my_ch[32]; memset(my_ch, 0, 32); |
||||||
|
if (rl >= 32) memcpy(my_ch, r, 32); |
||||||
|
|
||||||
|
if (memcmp(my_ch, ctx->init_sync.peer_last_ch, 32) == 0) { |
||||||
|
/* synced */ |
||||||
|
uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; |
||||||
|
memcpy(resp + 1, &ctx->init_sync.my_count, 4); resp[5] = 0; |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, resp, 6); |
||||||
|
struct channel_cache* ch = cs_find(ctx->sync, ctx->ch_id); |
||||||
|
if (ch) ch->synced = CS_SYNC_DONE; |
||||||
|
u_free(ctx); return; |
||||||
|
} |
||||||
|
|
||||||
|
/* Build sparse hashes: async chain */ |
||||||
|
uint8_t resp[1024]; uint32_t off = 0; |
||||||
|
resp[off++] = CS_MSG_INIT_RESP; |
||||||
|
memcpy(resp + off, &ctx->init_sync.my_count, 4); off += 4; |
||||||
|
memcpy(resp + off, &tp, 4); off += 4; |
||||||
|
memcpy(resp + off, my_ch, 32); off += 32; |
||||||
|
uint8_t sparse_pos = off; off++; /* placeholder */ |
||||||
|
uint8_t sparse_count = 0; |
||||||
|
|
||||||
|
/* First sparse hash to start chain */ |
||||||
|
uint32_t first_step = 1; |
||||||
|
if (tp >= first_step) { |
||||||
|
uint32_t pos = tp - first_step; |
||||||
|
ctx->init_sync.test_pos = tp; /* сохраняем для продолжения */ |
||||||
|
/* Запрашиваем первый sparse hash */ |
||||||
|
uint8_t req2[69]; int nr2 = req_ch(ctx->ch_id, req2); |
||||||
|
memcpy(req2 + nr2, &pos, 4); nr2 += 4; |
||||||
|
gui_bridge_call(GUI_OP_CHAIN_HASH_AT, req2, nr2, ctx, on_init_sync_sparse); |
||||||
|
return; |
||||||
|
} |
||||||
|
resp[sparse_pos] = 0; |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, resp, off); |
||||||
|
u_free(ctx); |
||||||
|
} |
||||||
|
|
||||||
|
static void on_init_sync_sparse(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
/* For initial version: just send INIT_RESP with 0 sparse.
|
||||||
|
Peer will handle divergence. */ |
||||||
|
uint32_t tp = ctx->init_sync.test_pos; |
||||||
|
uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; |
||||||
|
memcpy(resp + 1, &ctx->init_sync.my_count, 4); |
||||||
|
resp[5] = 0; /* sparse_count=0 */ |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, resp, 6); |
||||||
|
u_free(ctx); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── INIT_RESP handler ── */ |
||||||
|
|
||||||
|
static void cs_handle_init_resp(struct chat_sync* cs, uint64_t peer, |
||||||
|
const char* ch_id, const uint8_t* pl, size_t len) { |
||||||
|
if (len < 37) return; |
||||||
|
struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); |
||||||
|
if (!ctx) return; |
||||||
|
ctx->init_resp.peer_count = *(const uint32_t*)pl; |
||||||
|
ctx->init_resp.test_pos = *(const uint32_t*)(pl + 4); |
||||||
|
memcpy(ctx->init_resp.peer_ch, pl + 8, 37); /* hash:32 + sparse_count:1 + maybe more */ |
||||||
|
|
||||||
|
uint8_t req[69]; int nr = req_ch(ch_id, req); |
||||||
|
memcpy(req + nr, &ctx->init_resp.test_pos, 4); nr += 4; |
||||||
|
gui_bridge_call(GUI_OP_CHAIN_HASH_AT, req, nr, ctx, on_init_resp_hash); |
||||||
|
} |
||||||
|
|
||||||
|
static void on_init_resp_hash(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
uint8_t my_ch[32]; memset(my_ch, 0, 32); |
||||||
|
if (rl >= 32) memcpy(my_ch, r, 32); |
||||||
|
|
||||||
|
uint8_t sparse_count = ctx->init_resp.peer_ch[32]; |
||||||
|
|
||||||
|
/* Check match */ |
||||||
|
if (memcmp(my_ch, ctx->init_resp.peer_ch, 32) == 0) { |
||||||
|
/* Divergence is in our extra records (or peer's) */ |
||||||
|
/* Request data from peer starting at test_pos+1 if peer has more */ |
||||||
|
struct channel_cache* ch = cs_find(ctx->sync, ctx->ch_id); |
||||||
|
if (ch && ctx->init_resp.peer_count > ctx->init_resp.test_pos + 1) { |
||||||
|
/* We need peer's extra records */ |
||||||
|
uint32_t from = ctx->init_resp.test_pos + 1; |
||||||
|
uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; |
||||||
|
memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, snd, 7); |
||||||
|
} else { |
||||||
|
ch->synced = CS_SYNC_DONE; |
||||||
|
} |
||||||
|
u_free(ctx); return; |
||||||
|
} |
||||||
|
|
||||||
|
/* Mismatch — request data from start */ |
||||||
|
uint32_t from = 0; |
||||||
|
uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; |
||||||
|
memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, snd, 7); |
||||||
|
u_free(ctx); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── SEND_DATA handler ── */ |
||||||
|
|
||||||
|
static void cs_handle_send_data(struct chat_sync* cs, uint64_t peer, |
||||||
|
const char* ch_id, const uint8_t* pl, size_t len) { |
||||||
|
if (len < 6) return; |
||||||
|
uint32_t from = *(const uint32_t*)pl; |
||||||
|
uint16_t count = *(const uint16_t*)(pl + 4); |
||||||
|
|
||||||
|
if (count == 0) { |
||||||
|
/* Peer requests OUR data starting from 'from' */ |
||||||
|
struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); |
||||||
|
if (!ctx) return; |
||||||
|
ctx->send_data.from_pos = from; |
||||||
|
ctx->send_data.sent_count = 0; |
||||||
|
|
||||||
|
uint8_t req[65]; int nr = req_ch(ch_id, req); |
||||||
|
gui_bridge_call(GUI_OP_CURSOR_OPEN, req, nr, ctx, on_send_data_cursor_open); |
||||||
|
return; |
||||||
|
} |
||||||
|
|
||||||
|
/* Peer sent US data — insert each record */ |
||||||
|
const uint8_t* ptr = pl + 6; |
||||||
|
size_t remain = len - 6; |
||||||
|
|
||||||
|
for (uint16_t i = 0; i < count && remain > 0; i++) { |
||||||
|
/* record: [ts:8][dh:8][node_id:8][ct_len:1][ct:var][data_len:4][data:var] */ |
||||||
|
/* For this, we call GUI_OP_INSERT_RECORD with the full record + chain_hash */ |
||||||
|
/* But SEND_DATA doesn't include chain_hash — only PUSH does.
|
||||||
|
For SEND_DATA: skip chain_hash, just insert the content fields */ |
||||||
|
if (remain < 29) break; |
||||||
|
const uint8_t* rec_start = ptr; |
||||||
|
ptr += 8 + 8 + 8; remain -= 24; /* ts, dh, nid */ |
||||||
|
if (remain < 1) break; |
||||||
|
uint8_t ct_len = *ptr; ptr++; remain--; |
||||||
|
if (remain < ct_len) break; |
||||||
|
ptr += ct_len; remain -= ct_len; |
||||||
|
if (remain < 4) break; |
||||||
|
uint32_t dlen; memcpy(&dlen, ptr, 4); ptr += 4; remain -= 4; |
||||||
|
if (remain < dlen) break; |
||||||
|
ptr += dlen; remain -= dlen; |
||||||
|
|
||||||
|
size_t reclen = (size_t)(ptr - rec_start); |
||||||
|
|
||||||
|
/* Insert via GUI */ |
||||||
|
uint8_t req2[2048]; int nr2 = req_ch(ch_id, req2); |
||||||
|
if (nr2 + (int)reclen <= (int)sizeof(req2)) { |
||||||
|
memcpy(req2 + nr2, rec_start, reclen); nr2 += (int)reclen; |
||||||
|
/* Append zero chain_hash (not verified for SEND_DATA) */ |
||||||
|
uint8_t zero_ch[32]; memset(zero_ch, 0, 32); |
||||||
|
if (nr2 + 32 <= (int)sizeof(req2)) { |
||||||
|
memcpy(req2 + nr2, zero_ch, 32); nr2 += 32; |
||||||
|
} |
||||||
|
gui_bridge_call(GUI_OP_INSERT_RECORD, req2, nr2, NULL, NULL); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
/* Request more if needed */ |
||||||
|
uint32_t next = from + count; |
||||||
|
uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; |
||||||
|
memcpy(snd + 1, &next, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); |
||||||
|
cs_send(cs, ch_id, peer, snd, 7); |
||||||
|
|
||||||
|
struct channel_cache* ch = cs_find(cs, ch_id); |
||||||
|
if (ch) ch->synced = CS_SYNC_DONE; |
||||||
|
} |
||||||
|
|
||||||
|
static void on_send_data_cursor_open(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
if (rl < 4) { u_free(ctx); return; } |
||||||
|
memcpy(&ctx->send_data.cursor_id, r, 4); |
||||||
|
gui_bridge_call(GUI_OP_CURSOR_NEXT, (const uint8_t*)&ctx->send_data.cursor_id, 4, ctx, on_send_data_cursor_next); |
||||||
|
} |
||||||
|
|
||||||
|
static void on_send_data_cursor_next(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
if (rl == 0) { |
||||||
|
gui_bridge_call(GUI_OP_CURSOR_CLOSE, (const uint8_t*)&ctx->send_data.cursor_id, 4, NULL, NULL); |
||||||
|
u_free(ctx); return; |
||||||
|
} |
||||||
|
|
||||||
|
/* Send record to peer */ |
||||||
|
uint8_t snd[8192]; uint32_t off = 0; |
||||||
|
snd[off++] = CS_MSG_SEND_DATA; |
||||||
|
uint32_t pos = ctx->send_data.from_pos + ctx->send_data.sent_count; |
||||||
|
memcpy(snd + off, &pos, 4); off += 4; |
||||||
|
uint16_t c = 1; memcpy(snd + off, &c, 2); off += 2; |
||||||
|
if (off + rl <= (uint32_t)sizeof(snd)) { memcpy(snd + off, r, rl); off += rl; } |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, snd, off); |
||||||
|
ctx->send_data.sent_count++; |
||||||
|
|
||||||
|
gui_bridge_call(GUI_OP_CURSOR_NEXT, (const uint8_t*)&ctx->send_data.cursor_id, 4, ctx, on_send_data_cursor_next); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── PUSH handler ── */ |
||||||
|
|
||||||
|
static void cs_handle_push(struct chat_sync* cs, uint64_t peer, |
||||||
|
const char* ch_id, const uint8_t* pl, size_t len) { |
||||||
|
if (len < 24 + 1 + 4) return; |
||||||
|
uint64_t ts, dh, nid; |
||||||
|
memcpy(&ts, pl, 8); memcpy(&dh, pl + 8, 8); memcpy(&nid, pl + 16, 8); |
||||||
|
|
||||||
|
struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); |
||||||
|
if (!ctx) return; |
||||||
|
ctx->push.ts = ts; ctx->push.dh = dh; |
||||||
|
|
||||||
|
/* Request: [ch_len][ch_id] + [rest of PUSH payload starting from ct_len] */ |
||||||
|
/* PUSH payload: [ts:8][dh:8][node_id:8][ct_len:1][ct:var][data_len:4][data:var][chain_hash:32] */ |
||||||
|
uint8_t req[2048]; int nr = req_ch(ch_id, req); |
||||||
|
const uint8_t* rec = pl + 24; /* skip ts,dh,nid */ |
||||||
|
int reclen = (int)(len - 24); |
||||||
|
if (nr + reclen <= (int)sizeof(req)) { |
||||||
|
memcpy(req + nr, rec, reclen); nr += reclen; |
||||||
|
} |
||||||
|
gui_bridge_call(GUI_OP_INSERT_RECORD, req, nr, ctx, on_push_inserted); |
||||||
|
} |
||||||
|
|
||||||
|
static void on_push_inserted(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
struct cs_ctx* ctx = (struct cs_ctx*)ud; |
||||||
|
int rc = (rl >= 1) ? (int8_t)r[0] : -1; |
||||||
|
if (rc == 0) { |
||||||
|
uint8_t ack[17]; ack[0] = CS_MSG_ACK_PUSH; |
||||||
|
memcpy(ack + 1, &ctx->push.dh, 8); |
||||||
|
memcpy(ack + 9, &ctx->push.ts, 8); |
||||||
|
cs_send(ctx->sync, ctx->ch_id, ctx->peer, ack, 17); |
||||||
|
|
||||||
|
uint8_t evt[65]; evt[0] = (uint8_t)strlen(ctx->ch_id); |
||||||
|
memcpy(evt + 1, ctx->ch_id, evt[0]); |
||||||
|
gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1 + evt[0]); |
||||||
|
|
||||||
|
struct channel_cache* ch = cs_find(ctx->sync, ctx->ch_id); |
||||||
|
if (ch) ch->msg_count++; |
||||||
|
} |
||||||
|
u_free(ctx); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── ACK_PUSH handler ── */ |
||||||
|
|
||||||
|
static void cs_handle_ack_push(struct chat_sync* cs, uint64_t peer, |
||||||
|
const char* ch_id, const uint8_t* pl, size_t len) { |
||||||
|
if (len < 16) return; |
||||||
|
uint64_t ts, dh; memcpy(&dh, pl, 8); memcpy(&ts, pl + 8, 8); |
||||||
|
uint8_t req[1 + 64 + 24]; int nr = req_ch(ch_id, req); |
||||||
|
memcpy(req + nr, &ts, 8); nr += 8; |
||||||
|
memcpy(req + nr, &dh, 8); nr += 8; |
||||||
|
uint64_t myid = cs->inst->node_id; |
||||||
|
memcpy(req + nr, &myid, 8); nr += 8; |
||||||
|
gui_bridge_call(GUI_OP_MARK_SENT, req, nr, NULL, NULL); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── SYNC_DONE handler ── */ |
||||||
|
|
||||||
|
static void cs_handle_sync_done(struct chat_sync* cs, uint64_t peer, |
||||||
|
const char* ch_id, const uint8_t* pl, size_t len) { |
||||||
|
if (len < 36) return; |
||||||
|
uint32_t pc; memcpy(&pc, pl, 4); |
||||||
|
struct channel_cache* ch = cs_find(cs, ch_id); |
||||||
|
if (ch) { ch->msg_count = pc; ch->synced = CS_SYNC_DONE; } |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Recv dispatcher ── */ |
||||||
|
|
||||||
|
static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { |
||||||
|
if (!entry || entry->len < 4) { if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } return; } |
||||||
|
if (!g_cs || !g_cs->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } |
||||||
|
|
||||||
|
uint64_t peer = conn ? conn->peer_node_id : 0; |
||||||
|
const uint8_t* d = entry->dgram; |
||||||
|
size_t dlen = entry->len; |
||||||
|
|
||||||
|
if (dlen < 4) { u_free(entry->dgram); queue_entry_free(entry); return; } |
||||||
|
uint8_t ch_len = d[1]; |
||||||
|
if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } |
||||||
|
|
||||||
|
char ch_id[64]; memcpy(ch_id, d + 2, ch_len); ch_id[ch_len] = '\0'; |
||||||
|
uint8_t type = d[2 + ch_len]; |
||||||
|
const uint8_t* pl = d + 3 + ch_len; |
||||||
|
size_t plen = dlen - 3 - ch_len; |
||||||
|
|
||||||
|
switch (type) { |
||||||
|
case CS_MSG_INIT_SYNC: cs_handle_init_sync(g_cs, peer, ch_id, pl, plen); break; |
||||||
|
case CS_MSG_INIT_RESP: cs_handle_init_resp(g_cs, peer, ch_id, pl, plen); break; |
||||||
|
case CS_MSG_SEND_DATA: cs_handle_send_data(g_cs, peer, ch_id, pl, plen); break; |
||||||
|
case CS_MSG_PUSH: cs_handle_push(g_cs, peer, ch_id, pl, plen); break; |
||||||
|
case CS_MSG_ACK_PUSH: cs_handle_ack_push(g_cs, peer, ch_id, pl, plen); break; |
||||||
|
case CS_MSG_SYNC_DONE: cs_handle_sync_done(g_cs, peer, ch_id, pl, plen); break; |
||||||
|
default: break; |
||||||
|
} |
||||||
|
u_free(entry->dgram); queue_entry_free(entry); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Connection callbacks ── */ |
||||||
|
|
||||||
|
static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { |
||||||
|
(void)arg; |
||||||
|
if (!conn || !g_cs) return; |
||||||
|
uint64_t peer = conn->peer_node_id; |
||||||
|
if (peer == 0 || peer == g_cs->inst->node_id) return; |
||||||
|
for (int i = 0; i < g_cs->channel_count; i++) { |
||||||
|
struct channel_cache* ch = &g_cs->channels[i]; |
||||||
|
int found = 0; |
||||||
|
for (int j = 0; j < ch->peer_count; j++) if (ch->peer_ids[j] == peer) { found = 1; break; } |
||||||
|
if (!found) continue; |
||||||
|
ch->synced = CS_SYNC_IN_PROGRESS; |
||||||
|
uint8_t msg[37]; msg[0] = CS_MSG_INIT_SYNC; |
||||||
|
memcpy(msg + 1, &ch->msg_count, 4); memcpy(msg + 5, ch->last_chain_hash, 32); |
||||||
|
cs_send(g_cs, ch->channel_id, peer, msg, 37); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) { |
||||||
|
(void)arg; |
||||||
|
if (!conn || !g_cs) return; |
||||||
|
uint64_t peer = conn->peer_node_id; |
||||||
|
for (int i = 0; i < g_cs->channel_count; i++) { |
||||||
|
int found = 0; |
||||||
|
for (int j = 0; j < g_cs->channels[i].peer_count; j++) |
||||||
|
if (g_cs->channels[i].peer_ids[j] == peer) { found = 1; break; } |
||||||
|
if (found) g_cs->channels[i].synced = CS_SYNC_NONE; |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void cs_on_new_conn(struct ETCP_CONN* conn, void* arg) { |
||||||
|
(void)arg; |
||||||
|
if (!conn) return; |
||||||
|
etcp_conn_add_up_cbk(conn, cs_on_conn_up, NULL); |
||||||
|
etcp_conn_add_down_cbk(conn, cs_on_conn_down, NULL); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Periodic refresh from DB ── */ |
||||||
|
|
||||||
|
static void on_refresh_list(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; |
||||||
|
if (!g_cs || !r || rl < 2) return; |
||||||
|
uint16_t cnt; memcpy(&cnt, r, 2); |
||||||
|
const uint8_t* p = r + 2; int rem = rl - 2; |
||||||
|
|
||||||
|
for (int i = 0; i < g_cs->channel_count; i++) |
||||||
|
if (g_cs->channels[i].peer_ids) u_free(g_cs->channels[i].peer_ids); |
||||||
|
if (g_cs->channels) u_free(g_cs->channels); |
||||||
|
g_cs->channels = u_calloc(cnt, sizeof(struct channel_cache)); |
||||||
|
g_cs->channel_count = cnt; |
||||||
|
if (!g_cs->channels) { g_cs->channel_count = 0; return; } |
||||||
|
|
||||||
|
for (int i = 0; i < (int)cnt && rem >= 1; i++) { |
||||||
|
uint8_t id_len = *p++; |
||||||
|
rem--; |
||||||
|
if (rem < id_len) break; |
||||||
|
memcpy(g_cs->channels[i].channel_id, p, id_len); |
||||||
|
g_cs->channels[i].channel_id[id_len] = '\0'; |
||||||
|
p += id_len; rem -= id_len; |
||||||
|
|
||||||
|
uint8_t req[65]; req[0] = id_len; |
||||||
|
memcpy(req + 1, g_cs->channels[i].channel_id, id_len); |
||||||
|
gui_bridge_call(GUI_OP_LIST_PEERS, req, 1 + id_len, g_cs, on_refresh_peers_and_count); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void on_refresh_peers_and_count(void* ud, int op, const uint8_t* r, int rl) { |
||||||
|
(void)op; (void)ud; |
||||||
|
if (!g_cs || !r || rl < 2) return; |
||||||
|
uint16_t cnt; memcpy(&cnt, r, 2); |
||||||
|
/* Find first channel without peers (simplified: assumes sequential refresh) */ |
||||||
|
for (int i = 0; i < g_cs->channel_count; i++) { |
||||||
|
if (g_cs->channels[i].peer_ids) continue; |
||||||
|
g_cs->channels[i].peer_ids = u_calloc(cnt, sizeof(uint64_t)); |
||||||
|
g_cs->channels[i].peer_count = cnt; |
||||||
|
for (uint16_t j = 0; j < cnt && 2 + j * 8 < (uint32_t)rl; j++) |
||||||
|
memcpy(&g_cs->channels[i].peer_ids[j], r + 2 + j * 8, 8); |
||||||
|
|
||||||
|
/* Also get count */ |
||||||
|
uint8_t req[65]; req[0] = (uint8_t)strlen(g_cs->channels[i].channel_id); |
||||||
|
memcpy(req + 1, g_cs->channels[i].channel_id, req[0]); |
||||||
|
gui_bridge_call(GUI_OP_COUNT, req, 1 + req[0], NULL, NULL); |
||||||
|
break; |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void refresh_timer_cb(void* arg) { |
||||||
|
(void)arg; |
||||||
|
if (!g_cs || !g_cs->initialized) return; |
||||||
|
gui_bridge_call(GUI_OP_LIST_CHANNELS, NULL, 0, g_cs, on_refresh_list); |
||||||
|
g_cs->refresh_timer = uasync_set_timeout(g_cs->inst->ua, 30u * 10000u, g_cs, refresh_timer_cb, "cs_refresh"); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── TTL cleanup ── */ |
||||||
|
|
||||||
|
static void ttl_timer_cb(void* arg) { |
||||||
|
(void)arg; |
||||||
|
if (!g_cs || !g_cs->initialized) return; |
||||||
|
uint64_t cutoff = get_time_us() - 86400000000ULL; |
||||||
|
uint64_t myid = g_cs->inst->node_id; |
||||||
|
for (int i = 0; i < g_cs->channel_count; i++) { |
||||||
|
uint8_t req[1 + 64 + 16]; int nr = req_ch(g_cs->channels[i].channel_id, req); |
||||||
|
memcpy(req + nr, &myid, 8); nr += 8; |
||||||
|
memcpy(req + nr, &cutoff, 8); nr += 8; |
||||||
|
gui_bridge_call(GUI_OP_TTL_DELETE, req, nr, NULL, NULL); |
||||||
|
} |
||||||
|
g_cs->ttl_timer = uasync_set_timeout(g_cs->inst->ua, 3600u * 10000u, g_cs, ttl_timer_cb, "cs_ttl"); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Public API ── */ |
||||||
|
|
||||||
|
int chat_sync_init(struct UTUN_INSTANCE* inst, |
||||||
|
void (*gui_cb)(void*, int, const uint8_t*, int)) { |
||||||
|
(void)gui_cb; |
||||||
|
if (!inst) return -1; |
||||||
|
struct chat_sync* cs = u_calloc(1, sizeof(*cs)); |
||||||
|
if (!cs) return -1; |
||||||
|
cs->inst = inst; |
||||||
|
cs->initialized = 1; |
||||||
|
g_cs = cs; |
||||||
|
|
||||||
|
etcp_router_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb); |
||||||
|
etcp_add_new_conn_cbk(inst, cs_on_new_conn, NULL); |
||||||
|
|
||||||
|
struct ETCP_CONN* c = inst->connections; |
||||||
|
while (c) { |
||||||
|
etcp_conn_add_up_cbk(c, cs_on_conn_up, NULL); |
||||||
|
etcp_conn_add_down_cbk(c, cs_on_conn_down, NULL); |
||||||
|
c = c->next; |
||||||
|
} |
||||||
|
|
||||||
|
cs->refresh_timer = uasync_set_timeout(inst->ua, 50000u, cs, refresh_timer_cb, "cs_refresh"); |
||||||
|
cs->ttl_timer = uasync_set_timeout(inst->ua, 3600u * 10000u, cs, ttl_timer_cb, "cs_ttl"); |
||||||
|
|
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: initialized", CS_ID); |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
void chat_sync_destroy(struct UTUN_INSTANCE* inst) { |
||||||
|
struct chat_sync* cs = g_cs; |
||||||
|
if (!cs || !inst) return; |
||||||
|
cs->initialized = 0; g_cs = NULL; |
||||||
|
|
||||||
|
etcp_router_unbind(inst, ETCP_RT_ID_CHAT_SYNC); |
||||||
|
if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; } |
||||||
|
if (cs->ttl_timer) { uasync_cancel_timeout(inst->ua, cs->ttl_timer); cs->ttl_timer = NULL; } |
||||||
|
|
||||||
|
struct ETCP_CONN* c = inst->connections; |
||||||
|
while (c) { |
||||||
|
etcp_conn_remove_up_cbk(c, cs_on_conn_up, NULL); |
||||||
|
etcp_conn_remove_down_cbk(c, cs_on_conn_down, NULL); |
||||||
|
c = c->next; |
||||||
|
} |
||||||
|
|
||||||
|
for (int i = 0; i < cs->channel_count; i++) |
||||||
|
if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids); |
||||||
|
if (cs->channels) u_free(cs->channels); |
||||||
|
u_free(cs); |
||||||
|
} |
||||||
|
|
||||||
|
int chat_sync_push(struct UTUN_INSTANCE* inst, |
||||||
|
const char* ch_id, uint64_t node_id, |
||||||
|
const char* content_type, const uint8_t* data, uint32_t data_len, |
||||||
|
uint64_t timestamp, uint64_t datahash) { |
||||||
|
if (!g_cs || !g_cs->initialized) return -1; |
||||||
|
struct channel_cache* ch = cs_find(g_cs, ch_id); |
||||||
|
|
||||||
|
uint8_t ct_len = content_type ? (uint8_t)strlen(content_type) : 0; |
||||||
|
if (ct_len > 63) ct_len = 63; |
||||||
|
size_t total = 1 + 8 + 8 + 8 + 1 + ct_len + 4 + data_len + 32; |
||||||
|
uint8_t* buf = u_malloc(total); |
||||||
|
if (!buf) return -1; |
||||||
|
|
||||||
|
uint8_t* p = buf; |
||||||
|
*p++ = CS_MSG_PUSH; |
||||||
|
memcpy(p, ×tamp, 8); p += 8; |
||||||
|
memcpy(p, &datahash, 8); p += 8; |
||||||
|
memcpy(p, &node_id, 8); p += 8; |
||||||
|
*p++ = ct_len; |
||||||
|
if (ct_len) { memcpy(p, content_type, ct_len); p += ct_len; } |
||||||
|
uint32_t dlen = data_len; |
||||||
|
memcpy(p, &dlen, 4); p += 4; |
||||||
|
if (data_len) { memcpy(p, data, data_len); p += data_len; } |
||||||
|
memset(p, 0, 32); /* chain_hash placeholder */ |
||||||
|
|
||||||
|
int sent = 0; |
||||||
|
uint64_t myid = g_cs->inst->node_id; |
||||||
|
if (ch) { |
||||||
|
for (int i = 0; i < ch->peer_count; i++) { |
||||||
|
if (ch->peer_ids[i] == node_id || ch->peer_ids[i] == myid) continue; |
||||||
|
cs_send(g_cs, ch_id, ch->peer_ids[i], buf, (size_t)(p + 32 - buf)); |
||||||
|
sent++; |
||||||
|
} |
||||||
|
} |
||||||
|
u_free(buf); |
||||||
|
return sent; |
||||||
|
} |
||||||
|
|
||||||
|
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { |
||||||
|
/* Placeholder — full impl requires parsing nodeinfo and calling conn_mgr_connect_node */ |
||||||
|
uint8_t req[8]; memcpy(req, &node_id, 8); |
||||||
|
gui_bridge_call(GUI_OP_LOAD_NODEINFO, req, 8, NULL, NULL); |
||||||
|
} |
||||||
@ -0,0 +1,53 @@ |
|||||||
|
#ifndef CHAT_SYNC_H |
||||||
|
#define CHAT_SYNC_H |
||||||
|
|
||||||
|
#ifdef __cplusplus |
||||||
|
extern "C" { |
||||||
|
#endif |
||||||
|
|
||||||
|
#include <stdint.h> |
||||||
|
#include <stddef.h> |
||||||
|
|
||||||
|
struct UTUN_INSTANCE; |
||||||
|
struct UASYNC; |
||||||
|
|
||||||
|
/* etcp_router service ID */ |
||||||
|
#define ETCP_RT_ID_CHAT_SYNC 0x30 |
||||||
|
|
||||||
|
/* Message types */ |
||||||
|
#define CS_MSG_INIT_SYNC 0x01 |
||||||
|
#define CS_MSG_INIT_RESP 0x02 |
||||||
|
#define CS_MSG_SEND_DATA 0x03 |
||||||
|
#define CS_MSG_PUSH 0x04 |
||||||
|
#define CS_MSG_ACK_PUSH 0x05 |
||||||
|
#define CS_MSG_SYNC_DONE 0x06 |
||||||
|
|
||||||
|
/* Protocol constants */ |
||||||
|
#define CS_SEND_DATA_MAX 32 |
||||||
|
#define CS_MAX_SPARSE 16 |
||||||
|
#define CS_PEER_TIMEOUT_MS 5000 |
||||||
|
|
||||||
|
/* Sync states per channel */ |
||||||
|
#define CS_SYNC_NONE 0 |
||||||
|
#define CS_SYNC_IN_PROGRESS 1 |
||||||
|
#define CS_SYNC_DONE 2 |
||||||
|
|
||||||
|
/* ── Public API ── */ |
||||||
|
|
||||||
|
int chat_sync_init(struct UTUN_INSTANCE* inst, |
||||||
|
void (*gui_result_cb)(void*, int, const uint8_t*, int)); |
||||||
|
void chat_sync_destroy(struct UTUN_INSTANCE* inst); |
||||||
|
|
||||||
|
/* Вызывается из GUI (через gui_bridge_post_uasync) при отправке */ |
||||||
|
int chat_sync_push(struct UTUN_INSTANCE* inst, |
||||||
|
const char* channel_id, uint64_t node_id, |
||||||
|
const char* content_type, const uint8_t* data, uint32_t data_len, |
||||||
|
uint64_t timestamp, uint64_t datahash); |
||||||
|
|
||||||
|
/* Прослойка conn_mgr: загрузить nodeinfo из БД и подключиться */ |
||||||
|
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id); |
||||||
|
|
||||||
|
#ifdef __cplusplus |
||||||
|
} |
||||||
|
#endif |
||||||
|
#endif /* CHAT_SYNC_H */ |
||||||
@ -0,0 +1,31 @@ |
|||||||
|
#include "db_sync.h" |
||||||
|
#include "utun_instance.h" |
||||||
|
#include "../lib/u_async.h" |
||||||
|
#include "../lib/debug_config.h" |
||||||
|
#include "../lib/mem.h" |
||||||
|
|
||||||
|
int db_sync_init(struct UTUN_INSTANCE* inst) { |
||||||
|
(void)inst; |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "db_sync: stub (disabled for chatgui)"); |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
void db_sync_destroy(struct UTUN_INSTANCE* inst) { |
||||||
|
(void)inst; |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "db_sync: stub destroyed"); |
||||||
|
} |
||||||
|
|
||||||
|
int db_sync_insert_len(struct UTUN_INSTANCE* inst, const char* json_data, size_t len) { |
||||||
|
(void)inst; (void)json_data; (void)len; |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
int db_sync_insert(struct UTUN_INSTANCE* inst, const char* json_data) { |
||||||
|
(void)inst; (void)json_data; |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
uint32_t db_sync_count(struct UTUN_INSTANCE* inst) { |
||||||
|
(void)inst; |
||||||
|
return 0; |
||||||
|
} |
||||||
@ -0,0 +1,60 @@ |
|||||||
|
#ifndef GUI_BRIDGE_H |
||||||
|
#define GUI_BRIDGE_H |
||||||
|
|
||||||
|
#ifdef __cplusplus |
||||||
|
extern "C" { |
||||||
|
#endif |
||||||
|
|
||||||
|
#include <stdint.h> |
||||||
|
#include <stddef.h> |
||||||
|
|
||||||
|
struct UASYNC; |
||||||
|
|
||||||
|
/* ── Операции uasync→GUI (chat_sync запрашивает работу с БД) ── */ |
||||||
|
|
||||||
|
#define GUI_OP_LIST_CHANNELS 0 /* → массив: [count:2][ch_id_len:1][ch_id:var][name_len:2][name:var]... */ |
||||||
|
#define GUI_OP_COUNT 1 /* req: [ch_id_len:1][ch_id:var] → resp: uint32 count (4 байта) */ |
||||||
|
#define GUI_OP_CHAIN_HASH_AT 2 /* req: [ch_id_len:1][ch_id:var][pos:4] → resp: 32 байта chain_hash */ |
||||||
|
#define GUI_OP_INSERT_RECORD 3 /* req: [ch_id_len:1][ch_id:var][record:var] → resp: 0=ok, 1=dup, -1=err (1 байт) */ |
||||||
|
#define GUI_OP_CURSOR_OPEN 4 /* req: [ch_id_len:1][ch_id:var] → resp: uint32 cursor_id (4 байта) */ |
||||||
|
#define GUI_OP_CURSOR_NEXT 5 /* req: cursor_id (4 байта) → resp: record или 0 байт = EOF */ |
||||||
|
#define GUI_OP_CURSOR_CLOSE 6 /* req: cursor_id (4 байта) → resp: 0 байт */ |
||||||
|
#define GUI_OP_MARK_SENT 7 /* req: [ch_id_len:1][ch_id:var][ts:8][dh:8][my_node_id:8] → resp: 0 байт */ |
||||||
|
#define GUI_OP_TTL_DELETE 8 /* req: [ch_id_len:1][ch_id:var][my_node_id:8][cutoff:8] → resp: 0 байт */ |
||||||
|
#define GUI_OP_LIST_PEERS 9 /* req: [ch_id_len:1][ch_id:var] → resp: [count:2][node_id:8]... */ |
||||||
|
#define GUI_OP_LOAD_NODEINFO 10 /* req: node_id (8 байт) → resp: nodeinfo или 0 байт если нет */ |
||||||
|
|
||||||
|
/* ── Типы уведомлений uasync→GUI (fire-and-forget) ── */ |
||||||
|
|
||||||
|
#define GUI_EVT_MSG_RECEIVED 1 /* data: [ch_id_len:1][ch_id:var][record:var] */ |
||||||
|
#define GUI_EVT_CONNECT_RESULT 2 /* data: [node_id:8][result:4] */ |
||||||
|
#define GUI_EVT_NEW_PEER 3 /* data: [node_id:8] */ |
||||||
|
#define GUI_EVT_CHANNEL_UPDATED 4 /* data: [ch_id_len:1][ch_id:var] */ |
||||||
|
|
||||||
|
/* ── Callback-типы ── */ |
||||||
|
|
||||||
|
typedef void (*gui_result_fn)(void* userdata, int op, |
||||||
|
const uint8_t* result, int result_len); |
||||||
|
|
||||||
|
/* ── API ── */ |
||||||
|
|
||||||
|
/* Регистрация (вызывается из GUI при старте) */ |
||||||
|
void gui_bridge_init(struct UASYNC* ua, gui_result_fn callback, void* gui_target); |
||||||
|
|
||||||
|
/* Установить указатель на DbManager (до первого вызова gui_bridge_call) */ |
||||||
|
void gui_bridge_set_db(void* db_ptr); |
||||||
|
|
||||||
|
/* uasync → GUI: запросить SQL-операцию */ |
||||||
|
void gui_bridge_call(int op, const uint8_t* request, int req_len, |
||||||
|
void* userdata, gui_result_fn callback); |
||||||
|
|
||||||
|
/* uasync → GUI: уведомление (fire-and-forget, callback не вызывается) */ |
||||||
|
void gui_bridge_post(int event_type, const uint8_t* data, int data_len); |
||||||
|
|
||||||
|
/* GUI → uasync: выполнить функцию в uasync-потоке */ |
||||||
|
void gui_bridge_post_uasync(struct UASYNC* ua, void (*fn)(void*), void* arg); |
||||||
|
|
||||||
|
#ifdef __cplusplus |
||||||
|
} |
||||||
|
#endif |
||||||
|
#endif /* GUI_BRIDGE_H */ |
||||||
@ -0,0 +1,252 @@ |
|||||||
|
#include "gui_bridge.h" |
||||||
|
|
||||||
|
#include <QObject> |
||||||
|
#include <QMetaObject> |
||||||
|
#include <QByteArray> |
||||||
|
#include <QString> |
||||||
|
#include <cstring> |
||||||
|
|
||||||
|
#include "../db/db_manager.h" |
||||||
|
|
||||||
|
extern "C" { |
||||||
|
#include "../../../lib/u_async.h" |
||||||
|
#include "../../../lib/mem.h" |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Внутренний объект-приёмник в GUI-потоке ── */ |
||||||
|
|
||||||
|
class GuiBridgeReceiver : public QObject { |
||||||
|
Q_OBJECT |
||||||
|
public: |
||||||
|
void* ua = nullptr; |
||||||
|
gui_result_fn callback = nullptr; |
||||||
|
DbManager* db = nullptr; |
||||||
|
|
||||||
|
GuiBridgeReceiver(QObject* parent = nullptr) : QObject(parent) {} |
||||||
|
|
||||||
|
Q_INVOKABLE void processCall(int op, QByteArray packedReq); |
||||||
|
Q_INVOKABLE void processPost(int eventType, QByteArray data); |
||||||
|
}; |
||||||
|
|
||||||
|
static GuiBridgeReceiver* g_receiver = nullptr; |
||||||
|
|
||||||
|
/* ── Трамплин для uasync_post: распаковывает результат и вызывает оригинальный callback ── */ |
||||||
|
|
||||||
|
struct TrampolineCtx { |
||||||
|
gui_result_fn cb; |
||||||
|
void* userdata; |
||||||
|
int op; |
||||||
|
int result_len; |
||||||
|
/* result data follows after this struct in memory */ |
||||||
|
}; |
||||||
|
|
||||||
|
static void bridge_trampoline(void* arg) { |
||||||
|
uint8_t* packed = (uint8_t*)arg; |
||||||
|
TrampolineCtx ctx; |
||||||
|
memcpy(&ctx, packed, sizeof(ctx)); |
||||||
|
const uint8_t* result = packed + sizeof(ctx); |
||||||
|
if (ctx.cb) ctx.cb(ctx.userdata, ctx.op, result, ctx.result_len); |
||||||
|
u_free(packed); |
||||||
|
} |
||||||
|
|
||||||
|
static void post_result_to_uasync(int op, const QByteArray& result, |
||||||
|
void* userdata, gui_result_fn cb) { |
||||||
|
int res_len = result.size(); |
||||||
|
int total = (int)sizeof(TrampolineCtx) + res_len; |
||||||
|
uint8_t* buf = (uint8_t*)u_malloc(total); |
||||||
|
if (!buf) return; |
||||||
|
|
||||||
|
TrampolineCtx ctx; |
||||||
|
ctx.cb = cb; |
||||||
|
ctx.userdata = userdata; |
||||||
|
ctx.op = op; |
||||||
|
ctx.result_len = res_len; |
||||||
|
memcpy(buf, &ctx, sizeof(ctx)); |
||||||
|
if (res_len > 0) memcpy(buf + sizeof(ctx), result.constData(), res_len); |
||||||
|
|
||||||
|
uasync_post((struct UASYNC*)g_receiver->ua, bridge_trampoline, buf); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── GuiBridgeReceiver implementation ── */ |
||||||
|
|
||||||
|
static QString readChId(const uint8_t*& r, int& rem) { |
||||||
|
if (rem < 1) return {}; |
||||||
|
uint8_t len = r[0]; r++; rem--; |
||||||
|
if (rem < len) return {}; |
||||||
|
QString s = QString::fromUtf8((const char*)r, len); |
||||||
|
r += len; rem -= len; |
||||||
|
return s; |
||||||
|
} |
||||||
|
|
||||||
|
void GuiBridgeReceiver::processCall(int op, QByteArray packedReq) { |
||||||
|
if (!db || packedReq.size() < (int)(sizeof(gui_result_fn) + sizeof(void*))) return; |
||||||
|
|
||||||
|
gui_result_fn cb; |
||||||
|
void* userdata; |
||||||
|
memcpy(&cb, packedReq.constData(), sizeof(cb)); |
||||||
|
memcpy(&userdata, packedReq.constData() + sizeof(cb), sizeof(userdata)); |
||||||
|
|
||||||
|
const uint8_t* r = (const uint8_t*)packedReq.constData() + sizeof(cb) + sizeof(userdata); |
||||||
|
int rem = packedReq.size() - (int)(sizeof(cb) + sizeof(userdata)); |
||||||
|
|
||||||
|
QByteArray result; |
||||||
|
|
||||||
|
switch (op) { |
||||||
|
case GUI_OP_LIST_CHANNELS: |
||||||
|
result = db->syncListChannels(); |
||||||
|
break; |
||||||
|
|
||||||
|
case GUI_OP_COUNT: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty()) result = db->syncCount(chId); |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_CHAIN_HASH_AT: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty() && rem >= 4) { |
||||||
|
uint32_t pos; memcpy(&pos, r, 4); |
||||||
|
result = db->syncChainHashAt(chId, pos); |
||||||
|
} |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_INSERT_RECORD: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty() && rem > 0) { |
||||||
|
QByteArray rec((const char*)r, rem); |
||||||
|
result = db->syncInsert(chId, rec); |
||||||
|
} |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_CURSOR_OPEN: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty()) result = db->syncCursorOpen(chId); |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_CURSOR_NEXT: |
||||||
|
if (rem >= 4) { |
||||||
|
uint32_t curId; memcpy(&curId, r, 4); |
||||||
|
result = db->syncCursorNext(curId); |
||||||
|
} |
||||||
|
break; |
||||||
|
case GUI_OP_CURSOR_CLOSE: |
||||||
|
if (rem >= 4) { |
||||||
|
uint32_t curId; memcpy(&curId, r, 4); |
||||||
|
db->syncCursorClose(curId); |
||||||
|
} |
||||||
|
break; |
||||||
|
case GUI_OP_MARK_SENT: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty() && rem >= 24) { |
||||||
|
uint64_t ts, dh, myId; |
||||||
|
memcpy(&ts, r, 8); memcpy(&dh, r + 8, 8); memcpy(&myId, r + 16, 8); |
||||||
|
result = db->syncMarkSent(chId, ts, dh, myId); |
||||||
|
} |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_TTL_DELETE: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty() && rem >= 16) { |
||||||
|
uint64_t myId, cutoff; |
||||||
|
memcpy(&myId, r, 8); memcpy(&cutoff, r + 8, 8); |
||||||
|
result = db->syncTtlDelete(chId, myId, cutoff); |
||||||
|
} |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_LIST_PEERS: { |
||||||
|
QString chId = readChId(r, rem); |
||||||
|
if (!chId.isEmpty()) result = db->syncListPeers(chId); |
||||||
|
break; |
||||||
|
} |
||||||
|
case GUI_OP_LOAD_NODEINFO: |
||||||
|
if (rem >= 8) { |
||||||
|
uint64_t nodeId; memcpy(&nodeId, r, 8); |
||||||
|
result = db->loadNodeInfo(nodeId); |
||||||
|
} |
||||||
|
break; |
||||||
|
default: |
||||||
|
break; |
||||||
|
} |
||||||
|
|
||||||
|
post_result_to_uasync(op, result, userdata, cb); |
||||||
|
} |
||||||
|
|
||||||
|
void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { |
||||||
|
if (!db) return; |
||||||
|
const uint8_t* d = (const uint8_t*)data.constData(); |
||||||
|
int dlen = data.size(); |
||||||
|
|
||||||
|
switch (eventType) { |
||||||
|
case GUI_EVT_MSG_RECEIVED: |
||||||
|
if (dlen >= 2) { |
||||||
|
uint8_t chLen = d[0]; |
||||||
|
if (dlen >= 1 + chLen) { |
||||||
|
QString chId = QString::fromUtf8((const char*)d + 1, chLen); |
||||||
|
db->onSyncMessageReceived(chId); |
||||||
|
} |
||||||
|
} |
||||||
|
break; |
||||||
|
case GUI_EVT_CONNECT_RESULT: |
||||||
|
if (dlen >= 12) { |
||||||
|
uint64_t nodeId; int result; |
||||||
|
memcpy(&nodeId, d, 8); memcpy(&result, d + 8, 4); |
||||||
|
db->setNodeOnline(nodeId, result == 0); |
||||||
|
} |
||||||
|
break; |
||||||
|
case GUI_EVT_NEW_PEER: |
||||||
|
case GUI_EVT_CHANNEL_UPDATED: |
||||||
|
default: |
||||||
|
break; |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
void* gui_bridge_get_db() { return g_receiver ? (void*)g_receiver->db : nullptr; } |
||||||
|
|
||||||
|
/* ── C API ── */ |
||||||
|
|
||||||
|
extern "C" { |
||||||
|
|
||||||
|
void gui_bridge_init(struct UASYNC* ua, gui_result_fn callback, void* gui_target) { |
||||||
|
g_receiver = new GuiBridgeReceiver((QObject*)gui_target); |
||||||
|
g_receiver->ua = ua; |
||||||
|
g_receiver->callback = callback; |
||||||
|
} |
||||||
|
|
||||||
|
void gui_bridge_set_db(void* db_ptr) { |
||||||
|
if (g_receiver) g_receiver->db = (DbManager*)db_ptr; |
||||||
|
} |
||||||
|
|
||||||
|
void gui_bridge_call(int op, const uint8_t* request, int req_len, |
||||||
|
void* userdata, gui_result_fn callback) { |
||||||
|
if (!g_receiver) return; |
||||||
|
|
||||||
|
/* Пакуем cb и userdata в начало, потом request */ |
||||||
|
QByteArray packed; |
||||||
|
packed.append((const char*)&callback, sizeof(callback)); |
||||||
|
packed.append((const char*)&userdata, sizeof(userdata)); |
||||||
|
if (req_len > 0) packed.append((const char*)request, req_len); |
||||||
|
|
||||||
|
QMetaObject::invokeMethod(g_receiver, "processCall", |
||||||
|
Qt::QueuedConnection, |
||||||
|
Q_ARG(int, op), |
||||||
|
Q_ARG(QByteArray, packed)); |
||||||
|
} |
||||||
|
|
||||||
|
void gui_bridge_post(int event_type, const uint8_t* data, int data_len) { |
||||||
|
if (!g_receiver) return; |
||||||
|
|
||||||
|
QByteArray d; |
||||||
|
if (data_len > 0) d = QByteArray((const char*)data, data_len); |
||||||
|
|
||||||
|
QMetaObject::invokeMethod(g_receiver, "processPost", |
||||||
|
Qt::QueuedConnection, |
||||||
|
Q_ARG(int, event_type), |
||||||
|
Q_ARG(QByteArray, d)); |
||||||
|
} |
||||||
|
|
||||||
|
void gui_bridge_post_uasync(struct UASYNC* ua, void (*fn)(void*), void* arg) { |
||||||
|
if (ua && fn) uasync_post(ua, fn, arg); |
||||||
|
} |
||||||
|
|
||||||
|
} /* extern "C" */ |
||||||
|
|
||||||
|
#include "gui_bridge_impl.moc" |
||||||
Loading…
Reference in new issue