Skip to content

Commit 12ca463

Browse files
committed
migration 'chat_prekey' in lmdb binary ( node relay ), and fixing some perfomance issue in chat e2ee.
1 parent 12afd78 commit 12ca463

11 files changed

Lines changed: 303 additions & 48 deletions

File tree

‎src/tsarchain/network/node_logic/chat_state.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,8 @@
1212
def init_chat_state(self) -> None:
1313
self.chat_lock = threading.RLock()
1414
self.chat_presence_pub = {}
15+
self.chat_presence_ts = {}
1516
self.chat_spend_pub = {}
16-
self.chat_prekeys = {}
1717
self.chat_mailbox = {}
1818
self.chat_global_count = 0
1919
self.chat_presence_seen = set()

‎src/tsarchain/network/rpc/docs/USER_RPC.MD‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -288,16 +288,17 @@ Provides decentralized, end-to-end encrypted messaging via TsarChain's Signal-co
288288
The node operates purely as a **Blind Relay & Ephemeral Mailbox Router** with zero ability to decrypt ciphertexts and zero permanent disk persistence for transient chat messages:
289289

290290
#### 1. Prekey Bundle Management (`chat_prekeys`)
291-
* **Data Structure**: In-memory dictionary `self.chat_prekeys[addr]`.
291+
* **Data Structure**: Dedicated LMDB binary database `LMDB_CHAT_PREKEYS` (`data/node/chat_prekeys`).
292292
* **Bundle Components**:
293293
- `ik` *(Identity Key)*: User's long-term public X25519 identity key.
294294
- `spk` *(Signed Prekey)*: Semi-static ephemeral public key authenticated via ECDSA Secp256k1.
295295
- `sig`: Identity digital signature over the canonical payload `TSAR-SPK|<spk_bytes>|<spend_pub_bytes>`.
296296
- `opk_list` *(One-Time Prekeys Pool)*: Pool of single-use ephemeral keys for maximum Forward Secrecy.
297+
* **Compact Binary Encoding**: Packed into a zero-overhead raw binary structure (`<QB` header, 32B IK, 32B SPK, signature bytes, and contiguous 32B OPKs) for >60% storage reduction.
297298
* **Atomic Key Consumption (FIFO Pop)**:
298-
- When an initiator calls `CHAT_GET_PREKEY`, the node consumes the top OPK via `lst.pop(0)`.
299+
- When an initiator calls `CHAT_GET_PREKEY`, the node consumes the top OPK via `lst.pop(0)` and atomically updates the record in LMDB.
299300
- A consumed OPK is **never issued to any other peer**, ensuring each new cryptographic session begins with unique ephemeral entropy.
300-
* **Capacity Constraints**: Capped at `CFG.CHAT_OPK_MAX_STORED` (default: 100 OPKs per address) to prevent memory exhaustion.
301+
* **Capacity Management**: Backed by `CFG.LMDB_PREKEYS_SIZE_MAX = 250MB` with rate limiting and PoW verification on publish endpoints.
301302

302303
#### 2. Offline Message Mailbox Queue Management (`chat_mailbox`)
303304
* **Data Structure**: In-memory FIFO queue per recipient: `self.chat_mailbox[addr] = collections.deque()`.

‎src/tsarchain/network/rpc/user_rpc/category/chat.py‎

Lines changed: 37 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -106,27 +106,30 @@ def chat_register(self, message, pow_obj, base_identity, addr, *,
106106
return spk_err
107107

108108
now = time.time()
109+
now_int = int(now)
109110
pid = secrets.token_hex(16)
110111
with self.chat_lock:
111112
self.chat_spend_pub[addr_s] = spend_pk
112113
self.chat_presence_pub[addr_s] = chat_pub
114+
if hasattr(self, "chat_presence_ts"):
115+
self.chat_presence_ts[addr_s] = now_int
113116
if hasattr(self, "record_presence_seen"):
114117
self.record_presence_seen(pid)
115118
else:
116119
self.chat_presence_seen.add(pid)
117-
b = self.chat_prekeys.get(addr_s) or {}
120+
b = self.get_prekey_bundle(addr_s)
118121
if "ik" not in b:
119122
b["ik"] = chat_pub
120-
b["ts"] = int(now)
121-
b["ts"] = int(now)
123+
b["ts"] = now_int
124+
b["ts"] = now_int
122125
if spk_valid:
123126
b["spk"] = spk_reg
124127
b["sig"] = sig_reg
125128
if opk_reg and len(opk_reg) == 64:
126129
b.setdefault("opk_list", []).append(opk_reg)
127-
self.chat_prekeys[addr_s] = b
130+
self.put_prekey_bundle(addr_s, b)
128131

129-
pres = {"pid": pid, "address": addr_s, "pubkey": chat_pub, "spend_pub": spend_pk, "presence_sig": presence_sig, "ts": int(now), "hops": 0}
132+
pres = {"pid": pid, "address": addr_s, "pubkey": chat_pub, "spend_pub": spend_pk, "presence_sig": presence_sig, "ts": now_int, "hops": 0}
130133
self.relay_presence_async(pres, exclude=addr)
131134
return {"type": "CHAT_REGISTERED", "address": addr_s, "pubkey": chat_pub}
132135
except Exception as exc:
@@ -175,10 +178,15 @@ def chat_lookup_pub(self, message, pow_obj, base_identity, *,
175178
return pow_resp
176179
pubhex = self.chat_presence_pub.get(addr_s)
177180
last_seen = None
178-
b = self.chat_prekeys.get(addr_s) or {}
179-
ts_field = b.get("ts")
180-
if isinstance(ts_field, (int, float)):
181-
last_seen = int(ts_field)
181+
if hasattr(self, "chat_presence_ts"):
182+
last_seen = self.chat_presence_ts.get(addr_s)
183+
if last_seen is None:
184+
b = self.get_prekey_bundle(addr_s)
185+
ts_field = b.get("ts")
186+
if isinstance(ts_field, (int, float)):
187+
last_seen = int(ts_field)
188+
if hasattr(self, "chat_presence_ts"):
189+
self.chat_presence_ts[addr_s] = last_seen
182190

183191
return {"type": "CHAT_PUBKEY", "address": addr_s, "pubkey": pubhex, "found": bool(pubhex), "last_seen": last_seen}
184192

@@ -237,18 +245,21 @@ def chat_presence(self, message, pow_obj, base_identity, addr, *,
237245
return pow_resp
238246

239247
pid = message.get("pid") or secrets.token_hex(16)
248+
now_int = int(time.time())
240249
with self.chat_lock:
241250
self.chat_presence_pub[addr_s] = pubhex
242251
self.chat_spend_pub[addr_s] = spend_pk
252+
if hasattr(self, "chat_presence_ts"):
253+
self.chat_presence_ts[addr_s] = now_int
243254
if hasattr(self, "record_presence_seen"):
244255
self.record_presence_seen(pid)
245256
else:
246257
self.chat_presence_seen.add(pid)
247-
b = self.chat_prekeys.get(addr_s) or {}
258+
b = self.get_prekey_bundle(addr_s)
248259
if "ik" not in b:
249260
b["ik"] = pubhex
250-
b["ts"] = int(time.time())
251-
self.chat_prekeys[addr_s] = b
261+
b["ts"] = now_int
262+
self.put_prekey_bundle(addr_s, b)
252263

253264
message["hops"] = hops + 1
254265
self.relay_presence_async(message, exclude=addr)
@@ -294,16 +305,15 @@ def chat_publish_prekeys(self, message, pow_obj, base_identity, *,
294305
if not sig_ok.get("spk"):
295306
log.warning("[chat_publish_prekeys] Bad SPK signature for %s", addr_s)
296307
return {"error":"bad_spk_sig"}
308+
now_int = int(time.time())
297309
with self.chat_lock:
298-
rec = self.chat_prekeys.get(addr_s) or {}
299-
rec.update({"ik": ik, "spk": spk, "sig": sig, "ts": int(time.time())})
300-
if isinstance(opk, str) and len(opk)==64:
301-
lst = rec.setdefault("opk_list", [])
302-
lst.append(opk)
303-
if len(lst) > CFG.CHAT_OPK_MAX_STORED:
304-
# keep it from getting bloated, only keep the latest OPK
305-
rec["opk_list"] = lst[-CFG.CHAT_OPK_MAX_STORED:]
306-
self.chat_prekeys[addr_s] = rec
310+
if hasattr(self, "chat_presence_ts"):
311+
self.chat_presence_ts[addr_s] = now_int
312+
rec = self.get_prekey_bundle(addr_s)
313+
rec.update({"ik": ik, "spk": spk, "sig": sig, "ts": now_int})
314+
if isinstance(opk, str) and len(opk) == 64:
315+
rec.setdefault("opk_list", []).append(opk)
316+
self.put_prekey_bundle(addr_s, rec)
307317

308318
return {"type":"CHAT_PUBLISH_PREKEYS"}
309319

@@ -313,12 +323,16 @@ def chat_get_prekey(self, message, *,
313323
client_ip, is_miner_sender, **kwargs):
314324
addr_s = (message.get("address") or "").strip().lower()
315325
with self.chat_lock:
316-
b = self.chat_prekeys.get(addr_s) or {}
326+
b = self.get_prekey_bundle(addr_s)
317327
if not b or ("ik" not in b or "spk" not in b or "sig" not in b):
318328
return {"error":"no_bundle"}
319329
lst = b.get("opk_list") or []
320330
opk = lst.pop(0) if lst else None
321-
self.chat_prekeys[addr_s] = b
331+
if lst:
332+
b["opk_list"] = lst
333+
else:
334+
b.pop("opk_list", None)
335+
self.put_prekey_bundle(addr_s, b)
322336
sp = self.chat_spend_pub.get(addr_s)
323337

324338
return {"type":"CHAT_PREKEY_BUNDLE","bundle":{"ik": b["ik"], "spk": b["spk"], "sig": b["sig"], "opk": opk, "spend_pub": sp}}

‎src/tsarchain/network/rpc_helper/chat.py‎

Lines changed: 132 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,22 +6,150 @@
66
import json
77
import time
88
import socket
9+
import struct
910
import threading
1011
import collections
12+
from concurrent.futures import ThreadPoolExecutor
1113

1214
# ---------------- Local Project ----------------
1315
from ...utils import config as CFG
1416
from .base import NetworkHandlerProxy
1517
from ..protocol import send_message, recv_message,build_envelope, SecureChannel
18+
from ...storage.kv import get as kv_get, put as kv_put, delete as kv_delete
1619

1720
# ---------------- Logger ----------------
1821
from ...utils.tsar_logging import get_ctx_logger
1922
log = get_ctx_logger("tsarchain.network.rpc_helper.chat")
2023

24+
_presence_executor: ThreadPoolExecutor | None = None
25+
_presence_executor_lock = threading.Lock()
26+
27+
28+
def get_presence_executor() -> ThreadPoolExecutor:
29+
global _presence_executor
30+
if _presence_executor is None:
31+
with _presence_executor_lock:
32+
if _presence_executor is None:
33+
_presence_executor = ThreadPoolExecutor(max_workers=4, thread_name_prefix="presence_relay")
34+
return _presence_executor
35+
36+
37+
def encode_prekey_bundle(bundle: dict) -> bytes:
38+
if not isinstance(bundle, dict):
39+
return b""
40+
ts = int(bundle.get("ts", 0) or 0)
41+
ik = bundle.get("ik")
42+
spk = bundle.get("spk")
43+
sig = bundle.get("sig")
44+
opk_list = bundle.get("opk_list") or []
45+
46+
flags = 0
47+
body = bytearray()
48+
if isinstance(ik, str) and len(ik) == 64:
49+
try:
50+
body.extend(bytes.fromhex(ik))
51+
flags |= 0x01
52+
except ValueError:
53+
pass
54+
if isinstance(spk, str) and len(spk) == 64:
55+
try:
56+
body.extend(bytes.fromhex(spk))
57+
flags |= 0x02
58+
except ValueError:
59+
pass
60+
if isinstance(sig, str) and sig:
61+
try:
62+
sig_bytes = bytes.fromhex(sig)
63+
if len(sig_bytes) <= 65535:
64+
body.extend(struct.pack("<H", len(sig_bytes)))
65+
body.extend(sig_bytes)
66+
flags |= 0x04
67+
except ValueError:
68+
pass
69+
70+
opk_bytes = bytearray()
71+
if isinstance(opk_list, list):
72+
for o in opk_list:
73+
if isinstance(o, str) and len(o) == 64:
74+
try:
75+
opk_bytes.extend(bytes.fromhex(o))
76+
except ValueError:
77+
pass
78+
opk_count = len(opk_bytes) // 32
79+
opk_header = struct.pack("<I", opk_count)
80+
header = struct.pack("<QB", ts, flags)
81+
return bytes(header + body + opk_header + opk_bytes)
82+
83+
84+
def decode_prekey_bundle(raw: bytes) -> dict:
85+
if not raw or len(raw) < 9:
86+
return {}
87+
try:
88+
ts, flags = struct.unpack_from("<QB", raw, 0)
89+
offset = 9
90+
ik = None
91+
spk = None
92+
sig = None
93+
if flags & 0x01:
94+
if offset + 32 <= len(raw):
95+
ik = raw[offset:offset+32].hex()
96+
offset += 32
97+
if flags & 0x02:
98+
if offset + 32 <= len(raw):
99+
spk = raw[offset:offset+32].hex()
100+
offset += 32
101+
if flags & 0x04:
102+
if offset + 2 <= len(raw):
103+
sig_len = struct.unpack_from("<H", raw, offset)[0]
104+
offset += 2
105+
if offset + sig_len <= len(raw):
106+
sig = raw[offset:offset+sig_len].hex()
107+
offset += sig_len
108+
opk_list = []
109+
if offset + 4 <= len(raw):
110+
opk_count = struct.unpack_from("<I", raw, offset)[0]
111+
offset += 4
112+
for _ in range(opk_count):
113+
if offset + 32 <= len(raw):
114+
opk_list.append(raw[offset:offset+32].hex())
115+
offset += 32
116+
res: dict = {"ts": ts}
117+
if ik:
118+
res["ik"] = ik
119+
if spk:
120+
res["spk"] = spk
121+
if sig:
122+
res["sig"] = sig
123+
if opk_list:
124+
res["opk_list"] = opk_list
125+
return res
126+
except Exception:
127+
return {}
128+
21129

22130
# ------------------------------ P2P Chat ------------------------------
23131

24132
class ChatHandler(NetworkHandlerProxy):
133+
def get_prekey_bundle(self, addr: str) -> dict:
134+
if not addr:
135+
return {}
136+
with self.chat_lock:
137+
raw = kv_get("chat_prekeys", addr.encode("utf-8"))
138+
return decode_prekey_bundle(raw) if raw else {}
139+
140+
def put_prekey_bundle(self, addr: str, bundle: dict) -> None:
141+
if not addr:
142+
return
143+
with self.chat_lock:
144+
raw = encode_prekey_bundle(bundle)
145+
kv_put("chat_prekeys", addr.encode("utf-8"), raw)
146+
147+
def delete_prekey_bundle(self, addr: str) -> None:
148+
if not addr:
149+
return
150+
with self.chat_lock:
151+
kv_delete("chat_prekeys", addr.encode("utf-8"))
152+
25153
def send_to_peer(self, peer: tuple[str,int], payload: dict) -> None:
26154
if not isinstance(peer, tuple) or len(peer) != 2:
27155
raise ValueError("bad peer")
@@ -42,7 +170,10 @@ def send_to_peer(self, peer: tuple[str,int], payload: dict) -> None:
42170

43171

44172
def relay_presence_async(self, pres: dict, exclude=None) -> None:
45-
threading.Thread(target=self._relay_presence, args=(pres, exclude), daemon=True).start()
173+
try:
174+
get_presence_executor().submit(self._relay_presence, pres, exclude)
175+
except Exception:
176+
pass
46177

47178

48179
def mailbox_put(self, addr, item, ttl_s, per_addr_max, global_max):

‎src/tsarchain/storage/kv.py‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ def get_db_path(name: str) -> str:
2626
"state": CFG.LMDB_STATE_DIR,
2727
"graffiti": CFG.LMDB_GRAFFITI_DIR,
2828
"mempool": CFG.LMDB_MEMPOOL_DIR,
29+
"chat_prekeys": CFG.LMDB_CHAT_PREKEYS,
2930

3031
# Keys and secrets
3132
"node_secrets": CFG.LMDB_KEYS_DIR,
@@ -47,11 +48,13 @@ def _init_native_store(name: str = "chain"):
4748
return _native_store
4849

4950
drive_override = os.getenv("TSAR_STORAGE_DRIVE_TYPE")
51+
map_size_init = int(CFG.LMDB_PREKEYS_SIZE_INIT) if name == "chat_prekeys" else int(CFG.LMDB_MAP_SIZE_INIT)
52+
map_size_max = int(CFG.LMDB_PREKEYS_SIZE_MAX) if name == "chat_prekeys" else int(CFG.LMDB_MAP_SIZE_MAX)
5053
store = _native_open_storage(
5154
"lmdb",
5255
path,
53-
map_size_init=int(CFG.LMDB_MAP_SIZE_INIT),
54-
map_size_max=int(CFG.LMDB_MAP_SIZE_MAX),
56+
map_size_init=map_size_init,
57+
map_size_max=map_size_max,
5558
pretty_json=False,
5659
drive_type=drive_override,
5760
)

‎src/tsarchain/utils/config.py‎

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -109,10 +109,13 @@
109109
LMDB_STATE_DIR = "data/node/state" # LMDB environment path for state
110110
LMDB_GRAFFITI_DIR = "data/node/graffiti" # LMDB environment path for graffiti
111111
LMDB_MEMPOOL_DIR = "data/node/mempool" # LMDB environment path for mempool
112+
LMDB_CHAT_PREKEYS = "data/node/chat_prekeys" # LMDB environment path for chat prekeys
112113

113-
LMDB_MAP_SIZE_INIT = 4 * 1024 * 1024 # initial LMDB map size (4 MB)
114-
LMDB_MAP_SIZE_MAX = 64 * 1024 * 1024 * 1024 # upper LMDB map cap (64 GB)
115-
KV_ITER_CHUNK = 512 # number of entries per chunk when iterating prefix scans (LMDB)
114+
LMDB_MAP_SIZE_INIT = 4 * 1024 * 1024 # initial LMDB map size (4 MB)
115+
LMDB_MAP_SIZE_MAX = 64 * 1024 * 1024 * 1024 # upper LMDB map cap (64 GB)
116+
LMDB_PREKEYS_SIZE_INIT = 4 * 1024 * 1024 # initial chat prekeys LMDB size (4 MB)
117+
LMDB_PREKEYS_SIZE_MAX = 250 * 1024 * 1024 # max chat prekeys LMDB size (250 MB)
118+
KV_ITER_CHUNK = 512 # number of entries per chunk when iterating prefix scans (LMDB)
116119

117120

118121
# ---- KEYS & SECRETS DATABASE PATHS (LMDB) ----
@@ -527,8 +530,7 @@
527530
CHAT_OPK_MIN_THRESHOLD = 5 # minimum one-time pre-keys kept ready
528531
CHAT_OPK_REFILL_COUNT = 20 # number of pre-keys generated when refilling
529532
CHAT_SPK_ROTATE_INTERVAL_S = 24 * 3600 # seconds between signed pre-key rotations
530-
CHAT_OPK_MAX_STORED = 200 # hard cap untuk jumlah OPK yang disimpan node per alamat
531-
CHAT_HISTORY_MAX_PER_PEER = 200 # maksimum entri riwayat chat per pasangan alamat
533+
CHAT_HISTORY_MAX_PER_PEER = 200 # Maximum paired chat history entries (stored on client)
532534

533535

534536
# =============================================================================
@@ -594,11 +596,11 @@
594596

595597
# ---- CHAT LOOKUP THROTTLING ----
596598
CHAT_LOOKUP_RL_IP_BURST = 20 # lookup pubkey chat per IP
597-
CHAT_LOOKUP_RL_IP_WINDOW_S = 10 # jendela waktu limiter lookup pubkey
598-
CHAT_LOOKUP_RL_BACKOFF_S = 5 # backoff setelah limiter lookup pubkey kena
599-
CHAT_LOOKUP_RL_ADDR_BURST = 10 # limiter lookup pubkey per alamat
600-
CHAT_LOOKUP_RL_ADDR_WINDOW_S = 10 # jendela waktu limiter per alamat
601-
CHAT_LOOKUP_RL_ADDR_BACKOFF_S = 8 # backoff setelah limiter per alamat kena
599+
CHAT_LOOKUP_RL_IP_WINDOW_S = 10 # seconds window for pubkey lookup limiter
600+
CHAT_LOOKUP_RL_BACKOFF_S = 5 # backoff after pubkey lookup limiter is triggered
601+
CHAT_LOOKUP_RL_ADDR_BURST = 10 # pubkey lookup limiter per address
602+
CHAT_LOOKUP_RL_ADDR_WINDOW_S = 10 # seconds window for per-address pubkey lookup limiter
603+
CHAT_LOOKUP_RL_ADDR_BACKOFF_S = 8 # backoff after per-address pubkey lookup limiter is triggered
602604

603605

604606
# ---- USER RPC THROTTLING ----

‎tests/unit/conftest.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ def isolate_unit_test_storage(request, tmp_path, monkeypatch):
3434
monkeypatch.setattr(CFG, "LMDB_STATE_DIR", str(mock_data / "node/state"))
3535
monkeypatch.setattr(CFG, "LMDB_GRAFFITI_DIR", str(mock_data / "node/graffiti"))
3636
monkeypatch.setattr(CFG, "LMDB_MEMPOOL_DIR", str(mock_data / "node/mempool"))
37+
monkeypatch.setattr(CFG, "LMDB_CHAT_PREKEYS", str(mock_data / "node/chat_prekeys"))
3738

3839
# 2. Patch web database and cache paths
3940
monkeypatch.setattr(CFG, "WEB_DATABASE_PATH", str(mock_data / "web"))

0 commit comments

Comments
 (0)