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
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
|
"""Who gets in: invitations and invitation links, the join request and its
pairing code, and device linking."""
import base64
import logging
import re
import time
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
from meshbay_common import MNP_VERSION
from meshbay_common.adminop import OP_INVITE_CANCEL, OP_INVITE_CREATE, OP_INVITE_LINK_CREATE
from meshbay_common.crypto import wrap_gek_aes
from meshbay_common.device import (
DEVICE_TTL,
device_add_transcript,
device_hello_transcript,
device_request_transcript,
)
from meshbay_common.join import JOIN_TTL, ROLE_MEMBER, ROLE_OPERATOR, join_transcript
from meshbay_common.protocol import MNP
from meshbay_node import ops
from meshbay_node.roster import KIND_ACCOUNT, KIND_LINK
from meshbay_node.transport.webrtc.channel import _get_remote_ip
log = logging.getLogger("meshbay_node.transport.webrtc_server")
# An invitation link's handle, as `roster.create_link_invite` mints it.
_INVITE_ID_RE = re.compile(r"[0-9a-f]{32}")
# Pairing codes carry 40 bits and are single-use, but a connection must not be
# allowed to sit there guessing. Failures are audited, so a grind is visible.
MAX_JOIN_ATTEMPTS = 5
# Per-connection limits alone would not bind an attacker who can open connections
# at will — and the adversary who can mint tokens for any account is the hub. So
# failed pairings are also counted node-wide over a window.
MAX_JOIN_FAILURES_WINDOW = 20
JOIN_FAILURE_WINDOW = 600 # seconds
class AdmissionMixin:
def _do_invite_create(self, msg: dict) -> None:
"""
Issue a one-time pairing code for someone the operator wants to admit.
Replaces the old invite path, where the inviter fetched the invitee's
public key from the hub and wrapped the group key for whatever came back
(H3). The node now needs nothing but a name: it will wrap the key itself,
later, for a key the invitee proves they hold.
"""
roster = self._ctx.get("roster")
if roster is None:
self._send({"type": "error", "detail": "Roster not available"})
return
invitee_id = msg.get("user_id", "")
group_id = msg.get("group_id") or self._group_id
if not invitee_id or not group_id:
self._send({"type": "error", "detail": "Missing user_id or group_id"})
return
if group_id != self._group_id:
self._send({"type": "error", "detail": "Wrong group for this session"})
return
if not self._has_admin_authority():
self._send({
"type": "error",
"detail": "No operator paired — run `meshbay-node operator pair`",
})
return
self._issue_admin_challenge(OP_INVITE_CREATE, invitee_id, {
"group_id": group_id,
"user_id": invitee_id,
"username": str(msg.get("username", ""))[:64],
})
def _do_invite_link_create(self, msg: dict) -> None:
"""
Issue a code bound to no account, for an invitation link — into the
group this connection authenticated to, and no other: a link names its
group, so the operator signs for exactly that one (docs/MESHBAY_DESIGN.md §3.4).
"""
group_id = self._group_id or ""
if not group_id:
self._send({"type": "error", "detail": "No group on this connection"})
return
if msg.get("group_id") and msg["group_id"] != group_id:
self._send({"type": "error", "detail": "Wrong group for this session"})
return
if not self._has_admin_authority():
self._send({
"type": "error",
"detail": "No operator paired — run `meshbay-node operator pair`",
})
return
self._issue_admin_challenge(
OP_INVITE_LINK_CREATE, f"link:{group_id}", {"group_id": group_id})
def _do_invite_cancel(self, msg: dict) -> None:
"""Take back an unredeemed link of this group, by its handle."""
group_id = self._group_id or ""
invite_id = str(msg.get("invite_id", ""))
if not group_id:
self._send({"type": "error", "detail": "No group on this connection"})
return
if not _INVITE_ID_RE.fullmatch(invite_id):
self._send({"type": "error", "detail": "Not an invitation id"})
return
if not self._has_admin_authority():
self._send({"type": "error", "detail": "No authorized key for this"})
return
self._issue_admin_challenge(
OP_INVITE_CANCEL, invite_id, {"group_id": group_id, "invite_id": invite_id})
def _join_refuse(self, reason: str, audit_detail: str = "") -> None:
self._join_attempts += 1
# Node-wide window, shared across connections: reconnecting must not reset
# the budget.
now = time.time()
failures = [t for t in self._ctx.get("join_failures", [])
if now - t < JOIN_FAILURE_WINDOW]
failures.append(now)
self._ctx["join_failures"] = failures
self._audit_join("join_refused", audit_detail or reason)
self._send({
"type": MNP.JOIN_RESULT,
"v": MNP_VERSION,
"ok": False,
"reason": reason,
})
def _audit_join(self, event: str, detail: str) -> None:
audit = self._ctx.get("audit_store")
if not audit:
return
self._remote_ip = self._remote_ip or _get_remote_ip(self._pc)
self._spawn(audit.log_event(
user_id=self._user_id or getattr(self, "_pending_sub", "unknown"),
event=event,
ip=self._remote_ip,
username=self._username or getattr(self, "_pending_username", ""),
group_id=self._group_id or getattr(self, "_pending_group", "") or "",
detail=detail,
))
async def _do_join_request(self, msg: dict) -> None:
"""
Pin an identity, or recognise one already pinned.
The client signs its own Ed25519 and X25519 keys together with the node's
nonce, so the identity key vouches for the encryption key — that is what
will make it safe for the node to wrap the GEK for a key that arrived over
the wire instead of one fetched from the hub's directory (H3).
A first pairing needs a one-time code, which the hub never sees. Afterwards
the pin is the credential and a changed key is refused outright, the same
rule the client applies to `pk_node` (11.5.8).
"""
roster = self._ctx.get("roster")
if roster is None:
self._send({"type": "error", "detail": "Roster not available"})
return
if self._join_attempts >= MAX_JOIN_ATTEMPTS:
self._send({"type": "error", "detail": "Too many attempts"})
return
now = time.time()
recent = [t for t in self._ctx.get("join_failures", [])
if now - t < JOIN_FAILURE_WINDOW]
if len(recent) >= MAX_JOIN_FAILURES_WINDOW:
self._audit_join("join_throttled", f"{len(recent)} failures in window")
self._send({"type": "error", "detail": "Pairing temporarily locked"})
return
user_id = self._user_id or getattr(self, "_pending_sub", "")
username = self._username or getattr(self, "_pending_username", "")
if not user_id:
self._send({"type": "error", "detail": "Handshake required"})
return
pk_ed_b64 = msg.get("pk_ed25519", "")
pk_x_b64 = msg.get("pk_x25519", "")
code = msg.get("code", "")
ts = msg.get("ts", 0)
try:
pk_ed_raw = base64.b64decode(pk_ed_b64)
pk_x_raw = base64.b64decode(pk_x_b64)
if len(pk_ed_raw) != 32 or len(pk_x_raw) != 32:
raise ValueError
pk_ed = Ed25519PublicKey.from_public_bytes(pk_ed_raw)
except Exception:
self._join_refuse("invalid_keys")
return
if not isinstance(ts, int) or abs(time.time() - ts) > JOIN_TTL:
self._join_refuse("stale_request")
return
# An empty group_id means operator pairing, which is node-wide. Anything
# else must be the group this connection authenticated to — a signature
# obtained for one group must not name another.
group_id = msg.get("group_id", "") or ""
session_group = self._group_id or getattr(self, "_pending_group", "") or ""
if group_id and group_id != session_group:
self._join_refuse("group_mismatch")
return
transcript = join_transcript(
node_pk_b64=self._node_pk_b64(),
group_id=group_id,
user_id=user_id,
pk_ed25519_b64=pk_ed_b64,
pk_x25519_b64=pk_x_b64,
nonce_node=self._nonce_node,
ts=ts,
)
try:
sig = base64.b64decode(msg.get("sig", ""))
except Exception:
self._join_refuse("invalid_signature_encoding")
return
if not self._verify_sig(pk_ed, transcript, sig):
self._join_refuse("signature_invalid")
return
# One person may hold several devices here — a browser and a desktop
# client are two keys on one account. So the question is not "is this
# THE key" but "is this ONE OF this account's live devices".
device = await roster.find_device(user_id, pk_ed_b64)
if device and device["pk_x25519"] != pk_x_b64:
# The Ed25519 key is pinned but arrives with a different encryption
# key. The join transcript signs both together, so this is either a
# client that regenerated half its identity or something splicing
# two messages; either way the pair is not the one admitted.
self._join_refuse(
"key_changed",
f"pinned x25519={device['pk_x25519'][:16]} presented={pk_x_b64[:16]}")
return
known = device
if not known and await roster.list_devices(user_id):
# The account is known here but this key is not one of its devices.
# Not an error to shout about: it is a second browser or a new
# client, and the way in is a device-add approved by a device that
# is already trusted — no operator, no new invitation code.
self._join_refuse(
"unknown_device",
f"presented={pk_ed_b64[:16]} — approve it from a device already "
f"paired with this node")
return
if known:
# This group's own row first; then the join message's group_id (empty
# on the node-wide first connect); then the operator's node-wide row,
# which is where an operator opening any group finds their authority.
member = (await roster.get_member(session_group, user_id)
or await roster.get_member(group_id, user_id)
or await roster.get_member("", user_id))
if not member and self._group_join_policy(session_group) == "open":
await roster.set_member(
group_id=session_group, user_id=user_id, role=ROLE_MEMBER,
status="active", approved_by="open-join",
)
member = await roster.get_member(session_group, user_id)
# A pending invite means the operator explicitly re-invited this
# person — require the code even if they already have a member
# row (e.g. they left and were re-invited, or were revoked then
# re-invited). Without this gate a stale roster row lets them
# back in without proving they received the new code.
pending_invite = any(
i["kind"] == KIND_ACCOUNT
and i["user_id"] == user_id
and i["group_id"] in (session_group, "")
for i in await roster.list_invites())
# Or they bring a link for this group: somebody already pinned here
# through another group, which is the ordinary case for a link, or
# somebody removed from it and invited back. Only when they are not
# an active member — a member opening the group leaves the link for
# whoever it was meant for.
active = bool(member) and member.get("status") == "active"
if pending_invite or (code and not active):
if not code:
self._join_refuse("code_required")
return
invite = await roster.consume_invite(code, user_id, session_group)
if not invite:
self._join_refuse("code_invalid")
return
await roster.set_member(
group_id=invite["group_id"], user_id=user_id,
role=invite["role"], status="active",
approved_by=invite["created_by"],
)
self._audit_join(
"join_pinned",
f"group={invite['group_id'][:8]} role={invite['role']} "
f"via={'link' if invite['kind'] == KIND_LINK else 'code'} "
"(device already known)")
member = (await roster.get_member(session_group, user_id)
or await roster.get_member(invite["group_id"], user_id))
if not member:
self._join_refuse("not_authorized_for_group")
return
await self._join_ok(
user_id, pk_x_raw, session_group,
role=member["role"] if member else "",
recognised=True,
)
return
if not code:
if self._group_join_policy(session_group) == "open":
# An open-join group admits anyone the hub calls a member, so a
# code would protect nothing — the hub can walk in through the
# front door. Pin what turns up and say so in the audit log.
await self._pin_and_admit(
roster, user_id, username, pk_ed_b64, pk_x_b64,
group_id=session_group, role=ROLE_MEMBER,
approved_by="open-join", via="tofu")
await self._join_ok(user_id, pk_x_raw, session_group,
role=ROLE_MEMBER, recognised=False)
return
self._join_refuse("code_required")
return
invite = await roster.consume_invite(code, user_id, session_group)
if not invite:
self._join_refuse("code_invalid")
return
await self._pin_and_admit(
# The name comes from the invitation, not from the token: the hub does
# not put a username claim in a JWT, so pinning from the session alone
# left the roster nameless and `member revoke <name>` unable to match.
roster, user_id, invite["username"] or username, pk_ed_b64, pk_x_b64,
group_id=invite["group_id"], role=invite["role"],
approved_by=invite["created_by"],
via="link" if invite["kind"] == KIND_LINK else "code")
# The roster row comes from the invitation; the key comes from the
# connection. An operator pairing is node-wide (empty group), but they
# redeemed the code while opening a group and expect to read it — and
# is_authorized() already grants an operator every group on this node.
await self._join_ok(user_id, pk_x_raw, session_group or invite["group_id"],
role=invite["role"], recognised=False)
async def _do_device_request(self, msg: dict) -> None:
"""
A new device files itself as pending, bound to a code it displays.
Served in the pre-proof window: by construction the caller holds no key
this node knows, so there is nothing yet to prove. Filing is inert —
nothing is admitted until an existing device countersigns.
"""
roster = self._ctx.get("roster")
if roster is None or not self._user_id or not self._nonce_node:
self._send({"type": "error", "detail": "Not ready for a device request"})
return
if not self._spend_device_attempt():
return
pk_ed_b64 = str(msg.get("pk_ed25519", ""))
pk_x_b64 = str(msg.get("pk_x25519", ""))
code_hash = str(msg.get("code_hash", ""))
if not (pk_ed_b64 and pk_x_b64 and code_hash):
self._send({"type": "error", "detail": "Missing device keys or code"})
return
# The account must already be known here. Anti-spam rather than a
# security boundary: the filing key is unpinned by construction, so this
# bounds the table, not the trust.
existing = await roster.list_devices(self._user_id)
if not existing:
self._send({"type": "error",
"detail": "This account has no device on this node yet — "
"an invitation code is what admits the first"})
return
if len(existing) >= roster.MAX_DEVICES_PER_USER:
self._send({"type": "error",
"detail": f"Already {len(existing)} devices, which is the "
f"limit. Revoke one first."})
return
ts = int(msg.get("ts", 0))
if abs(time.time() - ts) > DEVICE_TTL:
self._send({"type": "error", "detail": "Device request expired"})
return
transcript = device_request_transcript(
node_pk_b64=self._node_pk_b64(), user_id=self._user_id,
pk_ed25519_b64=pk_ed_b64, pk_x25519_b64=pk_x_b64,
code_hash=code_hash, nonce_node=self._nonce_node, ts=ts)
try:
pk_ed = Ed25519PublicKey.from_public_bytes(base64.b64decode(pk_ed_b64))
sig = base64.b64decode(msg.get("sig", ""))
except Exception:
self._send({"type": "error", "detail": "Invalid device key encoding"})
return
if not self._verify_sig(pk_ed, transcript, sig):
# Proof of possession, and nothing more: this says the caller holds
# the keys, never that they belong to this account.
self._send({"type": "error", "detail": "Device signature invalid"})
return
ttl = int(self._ctx.get("device_request_ttl") or 3600)
expires = await roster.file_device_request(
user_id=self._user_id, username=self._username or "",
pk_ed25519=pk_ed_b64, pk_x25519=pk_x_b64,
code_hash=code_hash, ttl=ttl)
self._audit("device_request", f"{pk_ed_b64[:16]}")
log.info("Device request filed for %s (%s)", self._user_id[:8],
pk_ed_b64[:16])
self._send({"type": MNP.DEVICE_REQUEST_ACK, "v": MNP_VERSION,
"expires_at": expires})
async def _do_device_lookup(self, msg: dict) -> None:
"""
List this account's pending device requests, each with its code hash.
**The node never learns the code**, which is what makes it unable to
substitute a key. It answers with candidates; the approver recomputes
`sha256(code ‖ keys)` for each and keeps the one that matches. A node
offering fabricated keys would have to produce a hash matching
`sha256(code ‖ fabricated)` — and it does not know the code.
An earlier version of this took the hash from the client and looked the
request up by it. That is circular: the client cannot compute the hash
without already knowing the keys it is asking about.
"""
roster = self._ctx.get("roster")
if roster is None or not self._user_id:
self._send({"type": "error", "detail": "Roster not available"})
return
pending = await roster.list_device_requests(self._user_id)
self._send({
"type": MNP.DEVICE_LOOKUP_RESULT, "v": MNP_VERSION,
"requests": [
{"pk_ed25519": r["pk_ed25519"], "pk_x25519": r["pk_x25519"],
"code_hash": r["code_hash"], "created_at": r["created_at"]}
for r in pending
],
})
async def _do_device_add(self, msg: dict) -> None:
"""
Admit a device, countersigned by one this node already pinned.
The whole control is in `_verify_device_signer`: the signature must
verify against a **live device of this same account**. The hub holds no
user keys and so cannot produce one.
"""
roster = self._ctx.get("roster")
if roster is None or not self._user_id or not self._nonce_node:
self._send({"type": "error", "detail": "Not ready to add a device"})
return
if not self._spend_device_attempt():
return
pk_ed_b64 = str(msg.get("pk_ed25519", ""))
pk_x_b64 = str(msg.get("pk_x25519", ""))
ts = int(msg.get("ts", 0))
if not (pk_ed_b64 and pk_x_b64):
self._send({"type": "error", "detail": "Missing device keys"})
return
if abs(time.time() - ts) > DEVICE_TTL:
self._send({"type": "error", "detail": "Approval expired"})
return
transcript = device_add_transcript(
node_pk_b64=self._node_pk_b64(), user_id=self._user_id,
pk_ed25519_b64=pk_ed_b64, pk_x25519_b64=pk_x_b64,
nonce_node=self._nonce_node, ts=ts)
signer = await self._verify_device_signer(roster, transcript,
msg.get("sig", ""))
if signer is None:
self._audit("device_add_refused", pk_ed_b64[:16])
self._send({"type": "error",
"detail": "Not signed by a device already paired here"})
return
devices = await roster.list_devices(self._user_id)
if len(devices) >= roster.MAX_DEVICES_PER_USER:
self._send({"type": "error", "detail": "Device limit reached"})
return
# Spend the request. Single use: an approval cannot be replayed, and a
# code that was used is gone whatever else happens next.
code_hash = str(msg.get("code_hash", ""))
if code_hash and not await roster.take_device_request(
code_hash, self._user_id):
self._send({"type": "error",
"detail": "That request is no longer pending"})
return
# The countersignature is **kept**, with the two fields needed to rebuild
# what it signed. Until 2026-09-07 it was verified here and thrown away,
# leaving only `added_by_pk` — which says *which* key approved and
# proves nothing to anyone else. `device_add_transcript` binds
# `nonce_node`, this connection's handshake nonce, so a stored signature
# without it is still unverifiable; that is why all three go in.
#
# This is what lets another member check for themselves that this device
# belongs to an account whose earlier device they have already pinned,
# instead of taking the node's word (Tier 2, docs/MESHBAY_DESIGN.md §3.3).
await roster.pin_identity(
user_id=self._user_id, username=self._username or "",
pk_ed25519=pk_ed_b64, pk_x25519=pk_x_b64, via="device",
label=str(msg.get("label", ""))[:64], added_by_pk=signer,
add_sig=str(msg.get("sig", "")),
add_nonce=base64.b64encode(self._nonce_node).decode(),
add_ts=ts)
self._audit("device_added", f"{pk_ed_b64[:16]} by {signer[:16]}")
log.info("Device added for %s: %s (approved by %s)",
self._user_id[:8], pk_ed_b64[:16], signer[:16])
self._send({"type": MNP.DEVICE_ADD_ACK, "v": MNP_VERSION,
"pk_ed25519": pk_ed_b64})
async def _do_device_hello(self, msg: dict) -> None:
"""
Learn which of this account's devices is on this connection.
The handshake authenticates a *group membership* (the GEK-HMAC) and an
*account* (the hub's token). It has never authenticated a device, and
while one account meant one key that was the same statement. It stopped
being so on 2026-08-18, and `_load_pinned_pk` — which resolves the
account's oldest live device — has been standing in for the real answer
ever since, including as the recorded uploader of every file.
What is checked, in order: the key is a live device *of this account* in
the node's own roster (never a token claim — that is
`docs/MESHBAY_DESIGN.md` §3.2's rule), the timestamp is fresh, and the
signature verifies over a transcript naming this node, this group and
this connection's nonce. A key that is merely well-formed proves nothing.
Idempotent for the same key, refused for a different one: a connection
does not get to change device half way through, which would let one
session's uploads be attributed to two.
"""
roster = self._ctx.get("roster")
if roster is None or not self._user_id:
self._send({"type": "error", "detail": "Roster not available"})
return
if not self._spend_device_attempt():
return
pk_ed_b64 = str(msg.get("pk_ed25519", ""))
ts = int(msg.get("ts", 0) or 0)
if not pk_ed_b64:
self._send({"type": "error", "detail": "Missing device key"})
return
if self._device_confirmed and pk_ed_b64 != self._pinned_pk:
self._send({"type": "error",
"detail": "This connection is already another device"})
return
if abs(time.time() - ts) > DEVICE_TTL:
self._send({"type": "error", "detail": "Stale device_hello"})
return
device = await roster.find_device(self._user_id, pk_ed_b64)
if device is None:
self._audit("device_hello_refused", pk_ed_b64[:16])
self._send({"type": "error",
"detail": "Not a device paired here"})
return
transcript = device_hello_transcript(
node_pk_b64=self._node_pk_b64(), group_id=self._group_id or "",
user_id=self._user_id, pk_ed25519_b64=pk_ed_b64,
nonce_node=self._nonce_node, ts=ts)
try:
pk = Ed25519PublicKey.from_public_bytes(base64.b64decode(pk_ed_b64))
except Exception:
self._send({"type": "error", "detail": "Unreadable device key"})
return
try:
sig = base64.b64decode(msg.get("sig", ""))
except Exception:
sig = b""
if not self._verify_sig(pk, transcript, sig):
self._audit("device_hello_refused", pk_ed_b64[:16])
self._send({"type": "error", "detail": "Signature verification failed"})
return
self._pinned_pk = pk_ed_b64
self._device_confirmed = True
log.info("Device identified on connection: user=%s device=%s",
self._user_id[:8], pk_ed_b64[:16])
self._send({"type": MNP.DEVICE_HELLO_ACK, "v": MNP_VERSION,
"pk_ed25519": pk_ed_b64})
async def _do_device_list(self, msg: dict) -> None:
"""This account's devices. Anyone may read their own, nobody else's."""
roster = self._ctx.get("roster")
if roster is None or not self._user_id:
self._send({"type": "error", "detail": "Roster not available"})
return
devices = await roster.list_devices(self._user_id)
pending = await roster.pending_device_requests(self._user_id)
self._send({
"type": MNP.DEVICE_LIST_RESULT, "v": MNP_VERSION,
"pending": pending,
"devices": [
{"pk_ed25519": d["pk_ed25519"], "label": d.get("label", ""),
"pinned_at": d["pinned_at"], "pinned_via": d["pinned_via"],
"added_by_pk": d.get("added_by_pk", ""),
"is_this_one": d["pk_ed25519"] == self._pinned_pk}
for d in devices
],
})
async def _do_device_revoke(self, msg: dict) -> None:
"""
Retire one of this account's devices — a lost laptop.
Countersigned like an addition, by a live device of the same account.
The last one cannot go: an account with no device on this node can only
return through an operator's invitation code, and doing that to yourself
by accident is not a mistake worth allowing.
"""
roster = self._ctx.get("roster")
if roster is None or not self._user_id or not self._nonce_node:
self._send({"type": "error", "detail": "Not ready"})
return
if not self._spend_device_attempt():
return
target = str(msg.get("pk_ed25519", ""))
ts = int(msg.get("ts", 0))
if not target:
self._send({"type": "error", "detail": "Missing device key"})
return
if abs(time.time() - ts) > DEVICE_TTL:
self._send({"type": "error", "detail": "Request expired"})
return
victim = await roster.find_device(self._user_id, target)
if victim is None:
self._send({"type": "error", "detail": "No such device"})
return
transcript = device_add_transcript(
node_pk_b64=self._node_pk_b64(), user_id=self._user_id,
pk_ed25519_b64=target, pk_x25519_b64=victim["pk_x25519"],
nonce_node=self._nonce_node, ts=ts)
signer = await self._verify_device_signer(roster, transcript,
msg.get("sig", ""))
if signer is None:
self._send({"type": "error",
"detail": "Not signed by a device already paired here"})
return
if len(await roster.list_devices(self._user_id)) <= 1:
self._send({"type": "error",
"detail": "This is your only device here — removing it "
"would need an operator code to come back"})
return
await roster.revoke_device(self._user_id, target)
# A revoked device holds every chat key it ever received — a lost laptop
# reads the group's chat until the epoch moves.
await self._new_chat_epoch(self._group_id or "", "device_revoke")
self._audit("device_revoked", f"{target[:16]} by {signer[:16]}")
log.info("Device revoked for %s: %s", self._user_id[:8], target[:16])
self._send({"type": MNP.DEVICE_ADD_ACK, "v": MNP_VERSION,
"revoked": target})
async def _verify_device_signer(self, roster, transcript: bytes,
sig_b64: str) -> str | None:
"""
The pinned key that signed this, or None.
Every live device of the account is tried, because any of them may
approve. A revoked one is not in the list — that is the point of marking
rather than deleting: a lost laptop must stop being able to admit its
replacement.
"""
try:
sig = base64.b64decode(sig_b64)
except Exception:
return None
for device in await roster.list_devices(self._user_id):
try:
pk = Ed25519PublicKey.from_public_bytes(
base64.b64decode(device["pk_ed25519"]))
except Exception:
continue
if self._verify_sig(pk, transcript, sig):
return device["pk_ed25519"]
return None
def _spend_device_attempt(self) -> bool:
"""
Bound guessing on this connection, as the join path does.
A code is 40 bits, single use and bound to the keys it names, so this is
depth rather than the control — but an unbounded loop over the lookup is
still a free oracle, and a burst of failures belongs in the audit log.
"""
self._device_attempts = getattr(self, "_device_attempts", 0) + 1
if self._device_attempts > 5:
self._audit("device_attempts_exceeded", str(self._device_attempts))
self._send({"type": "error",
"detail": "Too many device attempts on this connection"})
return False
return True
def _group_join_policy(self, group_id: str) -> str:
"""
Admission policy for a group, read from the node's own configuration.
Never from the hub: a hub that could declare a group open would be handed
the key to it (docs/MESHBAY_DESIGN.md §3.4).
"""
gctx = (self._ctx.get("groups") or {}).get(group_id) or {}
return gctx.get("join_policy", "invite")
async def _pin_and_admit(
self, roster, user_id: str, username: str, pk_ed_b64: str, pk_x_b64: str,
*, group_id: str, role: str, approved_by: str, via: str,
) -> None:
await roster.pin_identity(
user_id=user_id, username=username,
pk_ed25519=pk_ed_b64, pk_x25519=pk_x_b64, via=via,
)
await roster.set_member(
group_id=group_id, user_id=user_id, role=role,
status="active", approved_by=approved_by,
)
if role == ROLE_OPERATOR:
self._ctx["has_admin_authority"] = True
log.info("Identity pinned (%s): user=%s role=%s", via, user_id[:8], role)
self._audit_join("join_pinned", f"role={role} via={via}")
async def _join_ok(
self, user_id: str, pk_x_raw: bytes, group_id: str,
*, role: str, recognised: bool,
) -> None:
"""
Answer a join, wrapping the group key for the key the caller just proved.
This is the H3 fix. The inviter used to fetch the invitee's public key from
the hub and wrap the GEK for whatever came back, so a hub that answered
with its own key was handed the group key by an honest member following the
protocol exactly. The node now wraps for a key that arrived from its owner
over an authenticated channel, bound to a pinned identity.
"""
reply = {
"type": MNP.JOIN_RESULT,
"v": MNP_VERSION,
"ok": True,
"recognised": recognised,
"role": role,
}
roster = self._ctx["roster"]
if group_id and not await roster.is_authorized(group_id, user_id):
# Pinned on this node, but not admitted to this group. Hub membership
# alone must not produce a key.
reply["gek"] = False
reply["reason"] = "not_authorized_for_group"
self._send(reply)
self._audit_join("join_no_gek", f"group={group_id[:8]} not authorized")
return
gctx = (self._ctx.get("groups") or {}).get(group_id) or {}
gek = gctx.get("gek")
if not gek:
reply["gek"] = False
reply["reason"] = "no_gek"
self._send(reply)
return
bundle = wrap_gek_aes(gek, pk_x_raw)
reply["gek"] = True
reply["pk_eph_b64"] = bundle["pk_eph_b64"]
reply["nonce_b64"] = bundle["nonce_b64"]
reply["wrapped_b64"] = bundle["wrapped_b64"]
self._send(reply)
self._audit_join("gek_wrapped", f"group={group_id[:8]}")
async def _admin_exec_invite_create(
self, pending: dict, transcript: bytes, sig: bytes,
) -> None:
# Node operator only. A group admin who does not run the node has no
# authority over who this node admits (deny by default). Delegation is
# designed but deferred — see docs/MESHBAY_DESIGN.md §3.4.
if not await self._verify_admin_sig(transcript, sig):
self._send({"type": "error", "detail": "Signature verification failed"})
self._audit("admin_auth_failed", f"invite_create:{pending['subject'][:16]}")
return
payload = pending["payload"]
try:
result = await self._run_op(
ops.create_invite,
payload["group_id"],
payload.get("username", ""),
user_id=payload["user_id"],
created_by=self._user_id or "",
)
except ops.OpError as e:
self._send({"type": "error", "detail": e.message})
return
self._audit("invite_create", f"target={payload['user_id'][:8]}")
self._send({
"type": MNP.INVITE_RESULT,
"v": MNP_VERSION,
"code": result["code"],
"expires_at": result["expires_at"],
"user_id": result["user_id"],
"username": result.get("username", ""),
})
async def _admin_exec_invite_link_create(
self, pending: dict, transcript: bytes, sig: bytes,
) -> None:
if not await self._verify_admin_sig(transcript, sig):
self._send({"type": "error", "detail": "Signature verification failed"})
self._audit("admin_auth_failed", "invite_link_create")
return
try:
result = await self._run_op(
ops.create_link_invite, pending["payload"]["group_id"],
created_by=self._user_id or "")
except ops.OpError as e:
self._send({"type": "error", "detail": e.message})
return
self._audit("invite_link_create", f"invite={result['invite_id'][:8]}")
self._send({
"type": MNP.INVITE_LINK_RESULT,
"v": MNP_VERSION,
"code": result["code"],
"invite_id": result["invite_id"],
"expires_at": result["expires_at"],
"group_id": result["group_id"],
})
async def _admin_exec_invite_cancel(
self, pending: dict, transcript: bytes, sig: bytes,
) -> None:
if not await self._verify_admin_sig(transcript, sig):
self._send({"type": "error", "detail": "Signature verification failed"})
self._audit("admin_auth_failed", f"invite_cancel:{pending['subject'][:8]}")
return
payload = pending["payload"]
try:
await self._run_op(ops.cancel_invite, payload["group_id"], payload["invite_id"])
except ops.OpError as e:
self._send({"type": "error", "detail": e.message})
return
self._audit("invite_cancel", f"invite={payload['invite_id'][:8]}")
self._send({"type": "ack", "v": MNP_VERSION, "detail": "invite_cancelled",
"invite_id": payload["invite_id"]})
|