1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
|
"""
MeshBay Node daemon — main process.
Startup sequence:
1. Load config (~/.config/meshbay/node.toml)
2. Load or create keystore (Argon2id unlock)
3. Connect to hub: register → login → announce node
4. Fetch GEK bundle from hub (if group configured)
5. Start directory indexer (watchdog)
6. Start TCP+TLS chunk server on node.port
7. Start local web UI on node.ui_port (localhost only)
8. Run until SIGINT/SIGTERM
Usage:
meshbay-node # interactive password prompt
meshbay-node --config /path # custom config
meshbay-node init # write example config + create keystore
meshbay-node --calibrate-argon2 # benchmark Argon2id, suggest parameters
"""
import asyncio
import logging
import signal
import sys
from pathlib import Path
import uvicorn
from meshbay_node.config import Config, load_config, write_example_config
from meshbay_node.hub_client import HubClient, HubConfig
from meshbay_node.indexer import DirectoryIndexer
from meshbay_node.keystore import load_or_create_keystore
from meshbay_node.transport import ChunkServer
from meshbay_node.transport.quic_server import QuicChunkServer
from meshbay_node.ui import create_ui_app
log = logging.getLogger(__name__)
# ── Argon2id calibration ──────────────────────────────────────────────────────
def calibrate_argon2(target_ms: int = 500) -> None:
"""Benchmark Argon2id and suggest parameters targeting ~target_ms."""
import time
import os
from meshbay_common.crypto import derive_keystore_key
print(f"Calibrating Argon2id (target: {target_ms}ms) ...")
salt = os.urandom(16)
for mem in [65536, 131072, 262144, 524288]:
from cryptography.hazmat.primitives.kdf.argon2 import Argon2id
t0 = time.perf_counter()
Argon2id(salt=salt, length=32, iterations=3,
lanes=4, memory_cost=mem).derive(b"benchmark")
elapsed_ms = (time.perf_counter() - t0) * 1000
print(f" memory_cost={mem:>7} ({mem//1024:>4}MB): {elapsed_ms:.0f}ms", end="")
if abs(elapsed_ms - target_ms) < target_ms * 0.3:
print(" ← recommended")
else:
print()
print("Set memory_cost in meshbay_common/crypto.py: ARGON2_MEMORY_COST")
# ── Daemon ────────────────────────────────────────────────────────────────────
class NodeDaemon:
def __init__(self, config: Config):
self._config = config
self._state = {
"status": "starting",
"hub_url": config.hub.url,
"username": config.hub.username,
"groups": [g.name for g in config.groups],
"node_port": config.node.port,
"quic_port": config.node.quic_port,
"endpoint_hint": None,
"indexes": {},
}
self._tcp_server: ChunkServer | None = None
self._quic_server: QuicChunkServer | None = None
self._indexers: list[DirectoryIndexer] = []
self._tasks: list[asyncio.Task] = []
async def run(self) -> None:
log.info("MeshBay Node starting up")
# 1. Keystore
keys = load_or_create_keystore(
path=self._config.keystore.path,
unlock_file=self._config.keystore.unlock_file,
)
log.info("Keys loaded: %s", keys.pk_ed25519_b64[:16])
# 2. Hub connection
hub_cfg = HubConfig(
hub_url=self._config.hub.url,
username=self._config.hub.username,
password=self._config.hub.password,
)
async with HubClient(hub_cfg, keys) as hub:
session = await hub.startup(endpoint_hint=None)
self._state["endpoint_hint"] = session.node_id
# 3. Build per-group contexts
groups_ctx: dict[str, dict] = {}
for group_cfg in self._config.groups:
if not group_cfg.id or not group_cfg.shared_dir:
log.warning("Group %r missing id or shared_dir — skipping",
group_cfg.name)
continue
shared_root = Path(group_cfg.shared_dir).expanduser().resolve()
if not shared_root.exists():
log.warning("Shared dir not found: %s — skipping group %s",
shared_root, group_cfg.name)
continue
gek = None
if group_cfg.visibility == "private":
try:
gek = await hub.fetch_gek(group_cfg.id)
log.info("GEK loaded for group %s", group_cfg.id[:8])
except LookupError:
log.warning("No GEK for group %s — skipping", group_cfg.name)
continue
indexer = DirectoryIndexer(
root=shared_root,
group_id=group_cfg.id,
sk_node=keys.sk_ed25519,
gek=gek,
)
await indexer.start()
self._indexers.append(indexer)
self._state["indexes"][group_cfg.id] = indexer.index
log.info("Indexing group %s: %s (%d files)",
group_cfg.name, shared_root, indexer.index.count)
groups_ctx[group_cfg.id] = {
"gek": gek,
"shared_root": shared_root,
"index": indexer.index,
}
# 4. QUIC chunk server (primary transport, all groups on one port)
if groups_ctx:
first = next(iter(groups_ctx.values()))
self._quic_server = QuicChunkServer(
sk_node=keys.sk_ed25519,
hub_pk_pem=session.hub_pk_pem,
gek=first["gek"],
shared_root=first["shared_root"],
index=first["index"],
host="::",
port=self._config.node.quic_port,
groups=groups_ctx,
)
await self._quic_server.start()
log.info("QUIC server on port %d (%d groups)",
self._config.node.quic_port, len(groups_ctx))
# TCP+TLS server (fallback transport, same groups)
self._tcp_server = ChunkServer(
sk_node=keys.sk_ed25519,
hub_pk_pem=session.hub_pk_pem,
gek=first["gek"],
shared_root=first["shared_root"],
index=first["index"],
host="0.0.0.0",
port=self._config.node.port,
groups=groups_ctx,
)
await self._tcp_server.start()
log.info("TCP+TLS server on port %d", self._config.node.port)
# 5. Local web UI
ui_app = create_ui_app(self._state)
ui_cfg = uvicorn.Config(
ui_app,
host="127.0.0.1",
port=self._config.node.ui_port,
log_level="warning",
)
ui_server = uvicorn.Server(ui_cfg)
self._tasks.append(asyncio.create_task(ui_server.serve()))
log.info("Local UI at http://localhost:%d", self._config.node.ui_port)
self._state["status"] = "running"
log.info("Node ready — %d groups", len(groups_ctx))
# 6. Wait for shutdown
stop_event = asyncio.Event()
loop = asyncio.get_event_loop()
for sig in (signal.SIGINT, signal.SIGTERM):
loop.add_signal_handler(sig, stop_event.set)
await stop_event.wait()
await self._shutdown()
async def _shutdown(self) -> None:
log.info("Shutting down...")
self._state["status"] = "stopping"
for task in self._tasks:
task.cancel()
for indexer in self._indexers:
await indexer.stop()
if self._quic_server:
await self._quic_server.stop()
if self._tcp_server:
await self._tcp_server.stop()
log.info("Node stopped")
# ── Entry point ───────────────────────────────────────────────────────────────
def main() -> None:
import argparse
parser = argparse.ArgumentParser(description="MeshBay Node daemon")
parser.add_argument("command", nargs="?",
choices=["init", "calibrate-argon2"],
help="init: write example config | calibrate-argon2: benchmark")
parser.add_argument("--config", type=Path, default=None,
help="Config file path")
parser.add_argument("--log-level", default="INFO",
choices=["DEBUG", "INFO", "WARNING", "ERROR"])
args = parser.parse_args()
logging.basicConfig(
level=getattr(logging, args.log_level),
format="%(asctime)s %(levelname)-8s %(name)s: %(message)s",
)
if args.command == "init":
write_example_config()
print(f"Example config written. Edit it and run: meshbay-node")
return
if args.command == "calibrate-argon2":
calibrate_argon2()
return
cfg = load_config(args.config)
if not cfg.hub.username:
print("Error: hub.username not set in config. Run: meshbay-node init")
sys.exit(1)
daemon = NodeDaemon(cfg)
asyncio.run(daemon.run())
if __name__ == "__main__":
main()
|