aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/daemon.py
blob: 93ba3c4ac1b0cabc2dd63f0f0842ad6cc8da7f9f (plain) (blame)
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()