summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/bundle_store.py
blob: c65b94d7ae245b088d34d57966cdebdcabec71c6 (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
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
"""
Bundle store — SQLite-backed storage for GEK bundles and keypair bundles.

GEK bundles: ECIES-wrapped GEK targeted at a specific user's X25519 key.
Keypair bundles: AES-GCM encrypted (Ed25519 + X25519) private keys, encrypted
with the user's password-derived bundle_key. Opaque to the node. An optional
second copy (bundle_enc_recovery) is wrapped under the account's recovery key
instead, so a forgotten passphrase does not strand the identity — see
docs/auth-confirm.md §4.3.

Both are stored and served over the P2P DataChannel during MNP handshake.
"""

import logging
from pathlib import Path

import aiosqlite

log = logging.getLogger(__name__)

_SCHEMA_GEK = """\
CREATE TABLE IF NOT EXISTS gek_bundles (
    group_id    TEXT NOT NULL,
    user_id     TEXT NOT NULL,
    pk_eph_b64  TEXT NOT NULL,
    nonce_b64   TEXT NOT NULL,
    wrapped_b64 TEXT NOT NULL,
    stored_at   TEXT NOT NULL DEFAULT (datetime('now')),
    PRIMARY KEY (group_id, user_id)
);
"""

# Chat epoch keys. Wrapped to the node's own X25519 key, exactly as the node's
# copy of the group key is — never stored raw.
#
# That is the whole basis of the claim chat encryption makes: "unreadable to
# someone who obtains the node's storage without the keystore password". A
# plaintext table beside chat.db would collapse it to nothing, silently, and it
# is the obvious thing to write. `test_chat_key_storage.py` reads the file back
# and refuses to find the live key in it.
#
# Rows are kept, never replaced: opening a new epoch must not make the history
# of the old one unreadable to the members who could already read it, which is
# the difference between an epoch and a rotation.
_SCHEMA_CHAT_EPOCHS = """\
CREATE TABLE IF NOT EXISTS chat_epochs (
    group_id    TEXT NOT NULL,
    epoch       INTEGER NOT NULL,
    pk_eph_b64  TEXT NOT NULL,
    nonce_b64   TEXT NOT NULL,
    wrapped_b64 TEXT NOT NULL,
    created_at  TEXT NOT NULL DEFAULT (datetime('now')),
    PRIMARY KEY (group_id, epoch)
);
"""

_SCHEMA_KEYPAIR = """\
CREATE TABLE IF NOT EXISTS keypair_bundles (
    user_id             TEXT PRIMARY KEY,
    bundle_enc          TEXT NOT NULL,
    bundle_enc_recovery TEXT,
    stored_at           TEXT NOT NULL DEFAULT (datetime('now'))
);
"""

# Per-account blobs the node holds and cannot read — playlists today
# (docs/playlists.md §3.3). The same shape as a keypair bundle, with a
# different payload, so this adds no trust boundary the node was not already
# on the wrong side of for this same account.
#
# Two deliberate differences from keypair_bundles, each a mistake avoided
# rather than a preference:
#
#   `blob_enc` is a BLOB, not base64 TEXT. A keypair bundle is a few hundred
#   bytes and nobody will ever notice the 33% tax; a playlist is hundreds of
#   kilobytes, where a third is not a rounding error.
#
#   `kind` is a namespace, not an enum: "playlists" is the manifest and
#   "playlist:<uuid>" is one playlist's tracks. That is what lets one playlist
#   be rewritten without re-uploading the whole collection, and it costs no
#   second table. webrtc_server.py validates the shape.
#
# blob_enc_recovery is declared now and left NULL — the same column, on the
# same kind of table, that keypair_bundles needed a PRAGMA migration for one
# release late. CREATE TABLE IF NOT EXISTS never adds a column.
_SCHEMA_USER_BLOBS = """\
CREATE TABLE IF NOT EXISTS user_blobs (
    user_id           TEXT NOT NULL,
    kind              TEXT NOT NULL,
    rev               INTEGER NOT NULL,
    blob_enc          BLOB NOT NULL,
    blob_enc_recovery BLOB,
    stored_at         TEXT NOT NULL DEFAULT (datetime('now')),
    PRIMARY KEY (user_id, kind)
);
"""


class BundleStore:
    def __init__(self, db_path: Path):
        self._db_path = db_path
        self._db: aiosqlite.Connection | None = None

    async def open(self) -> None:
        self._db_path.parent.mkdir(parents=True, exist_ok=True)
        self._db = await aiosqlite.connect(str(self._db_path))
        await self._db.execute(_SCHEMA_GEK)
        await self._db.execute(_SCHEMA_KEYPAIR)
        await self._db.execute(_SCHEMA_CHAT_EPOCHS)
        await self._db.execute(_SCHEMA_USER_BLOBS)
        await self._migrate_keypair_recovery()
        await self._db.commit()

    async def _migrate_keypair_recovery(self) -> None:
        """
        Add bundle_enc_recovery to a keypair_bundles table created before it
        existed. SQLite has no ADD COLUMN IF NOT EXISTS, so check the columns
        first — this table is node-only and has no Alembic history.
        """
        assert self._db
        async with self._db.execute("PRAGMA table_info(keypair_bundles)") as cur:
            cols = {row[1] for row in await cur.fetchall()}
        if "bundle_enc_recovery" not in cols:
            await self._db.execute(
                "ALTER TABLE keypair_bundles ADD COLUMN bundle_enc_recovery TEXT")

    async def store(
        self,
        group_id: str,
        user_id: str,
        pk_eph_b64: str,
        nonce_b64: str,
        wrapped_b64: str,
    ) -> None:
        assert self._db
        await self._db.execute(
            "INSERT OR REPLACE INTO gek_bundles "
            "(group_id, user_id, pk_eph_b64, nonce_b64, wrapped_b64, stored_at) "
            "VALUES (?, ?, ?, ?, ?, datetime('now'))",
            (group_id, user_id, pk_eph_b64, nonce_b64, wrapped_b64),
        )
        await self._db.commit()

    async def fetch(self, group_id: str, user_id: str) -> dict | None:
        assert self._db
        async with self._db.execute(
            "SELECT pk_eph_b64, nonce_b64, wrapped_b64 FROM gek_bundles "
            "WHERE group_id = ? AND user_id = ?",
            (group_id, user_id),
        ) as cursor:
            row = await cursor.fetchone()
            if not row:
                return None
            return {
                "pk_eph_b64": row[0],
                "nonce_b64": row[1],
                "wrapped_b64": row[2],
            }

    async def store_keypair(
        self,
        user_id: str,
        bundle_enc: str,
        bundle_enc_recovery: str | None = None,
    ) -> None:
        """
        Store the passphrase-wrapped keypair bundle, and optionally a second
        copy wrapped under the account's recovery key.

        A call that omits bundle_enc_recovery — a plain re-backup, or a
        passphrase-change re-wrap (docs/auth-confirm.md §3.2) — must not erase a
        recovery copy already stored, so the upsert keeps the existing value
        when the new one is None.
        """
        assert self._db
        await self._db.execute(
            "INSERT INTO keypair_bundles "
            "(user_id, bundle_enc, bundle_enc_recovery, stored_at) "
            "VALUES (?, ?, ?, datetime('now')) "
            "ON CONFLICT(user_id) DO UPDATE SET "
            "  bundle_enc = excluded.bundle_enc, "
            "  bundle_enc_recovery = COALESCE(excluded.bundle_enc_recovery, "
            "                                 keypair_bundles.bundle_enc_recovery), "
            "  stored_at = excluded.stored_at",
            (user_id, bundle_enc, bundle_enc_recovery),
        )
        await self._db.commit()

    async def fetch_keypair(self, user_id: str) -> dict | None:
        assert self._db
        async with self._db.execute(
            "SELECT bundle_enc, bundle_enc_recovery FROM keypair_bundles "
            "WHERE user_id = ?",
            (user_id,),
        ) as cursor:
            row = await cursor.fetchone()
            if not row:
                return None
            return {"bundle_enc": row[0], "bundle_enc_recovery": row[1]}

    async def delete_keypair(self, user_id: str) -> bool:
        """
        Drop someone's keypair bundle at their own request.

        Backing keys up here is what lets a second browser recover them with the
        password — and it is also what puts a PBKDF2-protected blob on every node
        whose group they join (finding C4). Someone who does not need the first
        should be able to withdraw the second, and not merely stop adding to it.
        """
        assert self._db
        cur = await self._db.execute(
            "DELETE FROM keypair_bundles WHERE user_id = ?", (user_id,))
        await self._db.commit()
        return cur.rowcount > 0

    # ── Per-account blobs (playlists) ────────────────────────────────────
    #
    # Opaque throughout: the node stores bytes it cannot read, returns them, and
    # never looks inside. Every method takes `user_id` from the caller, which
    # takes it from the authenticated session and never from a message body —
    # a user_id off the wire would let any member read or overwrite any other
    # member's blob.

    async def store_user_blob(
        self,
        user_id: str,
        kind: str,
        rev: int,
        blob_enc: bytes,
    ) -> None:
        """
        Write one blob, replacing whatever was there.

        Older revisions are not kept. The client is the authority on what the
        merged state is (docs/playlists.md §6.4) and holds its own copy, so a
        node keeping history would buy nothing and would mean the node deciding
        which revision is current — which is exactly what it must not do.
        """
        assert self._db
        await self._db.execute(
            "INSERT INTO user_blobs (user_id, kind, rev, blob_enc, stored_at) "
            "VALUES (?, ?, ?, ?, datetime('now')) "
            "ON CONFLICT(user_id, kind) DO UPDATE SET "
            "  rev = excluded.rev, "
            "  blob_enc = excluded.blob_enc, "
            "  stored_at = excluded.stored_at",
            (user_id, kind, rev, blob_enc),
        )
        await self._db.commit()

    async def fetch_user_blob(self, user_id: str, kind: str) -> dict | None:
        assert self._db
        async with self._db.execute(
            "SELECT rev, blob_enc FROM user_blobs WHERE user_id = ? AND kind = ?",
            (user_id, kind),
        ) as cursor:
            row = await cursor.fetchone()
            if not row:
                return None
            return {"rev": row[0], "blob_enc": row[1]}

    async def list_user_blobs(self, user_id: str) -> list[dict]:
        """
        Which blobs exist and at what revision — never a payload.

        A client that has lost its local state (a cache clear, a new device)
        cannot otherwise discover which playlists exist: their kinds carry
        client-generated UUIDs, and guessing is not a plan.
        """
        assert self._db
        async with self._db.execute(
            "SELECT kind, rev FROM user_blobs WHERE user_id = ? ORDER BY kind",
            (user_id,),
        ) as cursor:
            return [{"kind": r[0], "rev": r[1]} for r in await cursor.fetchall()]

    async def delete_user_blob(self, user_id: str, kind: str) -> bool:
        assert self._db
        cur = await self._db.execute(
            "DELETE FROM user_blobs WHERE user_id = ? AND kind = ?",
            (user_id, kind))
        await self._db.commit()
        return cur.rowcount > 0

    async def user_blob_total_bytes(self, user_id: str) -> int:
        """
        What this account is already using, so a store can refuse to take it
        over the per-account cap. An unbounded write primitive pointed at
        somebody else's disk needs a number that is actually checked.
        """
        assert self._db
        async with self._db.execute(
            "SELECT COALESCE(SUM(LENGTH(blob_enc)), 0) FROM user_blobs "
            "WHERE user_id = ?", (user_id,)) as cursor:
            row = await cursor.fetchone()
            return int(row[0]) if row else 0

    # ── Chat epoch keys ──────────────────────────────────────────────────

    async def store_chat_epoch(
        self,
        group_id: str,
        epoch: int,
        pk_eph_b64: str,
        nonce_b64: str,
        wrapped_b64: str,
    ) -> None:
        """
        Record one epoch key, wrapped to the node's own key.

        `INSERT OR IGNORE`, not `REPLACE`: an epoch's key is written once and is
        then the only way to read the messages sent under it. Overwriting one —
        which a retry, or two callers racing to open the same epoch, would do —
        would destroy that history with no error anywhere.
        """
        assert self._db
        await self._db.execute(
            "INSERT OR IGNORE INTO chat_epochs "
            "(group_id, epoch, pk_eph_b64, nonce_b64, wrapped_b64) "
            "VALUES (?, ?, ?, ?, ?)",
            (group_id, epoch, pk_eph_b64, nonce_b64, wrapped_b64),
        )
        await self._db.commit()

    async def fetch_chat_epochs(self, group_id: str) -> list[dict]:
        """Every epoch this group has had, oldest first."""
        assert self._db
        async with self._db.execute(
            "SELECT epoch, pk_eph_b64, nonce_b64, wrapped_b64 FROM chat_epochs "
            "WHERE group_id = ? ORDER BY epoch", (group_id,)
        ) as cur:
            rows = await cur.fetchall()
        return [{"epoch": r[0], "pk_eph_b64": r[1], "nonce_b64": r[2],
                 "wrapped_b64": r[3]} for r in rows]

    async def latest_chat_epoch(self, group_id: str) -> int:
        """The highest epoch number, or 0 when the group has none yet."""
        assert self._db
        async with self._db.execute(
            "SELECT MAX(epoch) FROM chat_epochs WHERE group_id = ?",
            (group_id,)
        ) as cur:
            row = await cur.fetchone()
        return int(row[0] or 0)

    async def close(self) -> None:
        if self._db:
            await self._db.close()
            self._db = None