aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py120
-rw-r--r--packages/meshbay-node/tests/test_roster_pairing.py82
2 files changed, 202 insertions, 0 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
index 81db0e9..eab1cac 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -59,6 +59,7 @@ from meshbay_common.adminop import (
OP_DIR_DELETE,
OP_FILE_DELETE,
OP_INVITE_CREATE,
+ OP_MEMBER_REVOKE,
admin_transcript,
)
from meshbay_common.crypto import pk_to_b64, wrap_gek_aes
@@ -171,6 +172,10 @@ def _extract_dtls_fingerprint(sdp: str) -> bytes:
STREAM_SEGMENT_SIZE = 256 * 1024
+# What a client may ask for in one go, and how long the node waits for it to ask
+# again before deciding nobody is watching any more.
+STREAM_MAX_CREDIT = 256
+STREAM_CREDIT_TIMEOUT = 120
_H264_PROFILES = {"Baseline": "42", "Main": "4d", "High": "64", "High 10": "6e"}
@@ -289,6 +294,9 @@ class WebRTCPeerSession:
# Set from the roster: the key this node pinned for this account. Never
# from the JWT — the hub picks what goes in there.
self._pinned_pk: str = ""
+ # Flow control for video: how many segments the client says it can take.
+ self._stream_credit = 0
+ self._stream_credit_evt = asyncio.Event()
self._gek_challenge: bytes | None = None
# Same value as the GEK challenge, but kept for the life of the connection:
# a join_request is signed over it, and it must stay verifiable after the
@@ -366,12 +374,16 @@ class WebRTCPeerSession:
self._do_admin_response(msg)
elif mtype == MNP.INVITE_CREATE:
self._do_invite_create(msg)
+ elif mtype == MNP.MEMBER_REVOKE:
+ self._do_member_revoke(msg)
elif mtype == MNP.KEYPAIR_BUNDLE_STORE:
asyncio.ensure_future(self._do_keypair_bundle_store(msg))
elif mtype == MNP.KEYPAIR_BUNDLE_DELETE:
asyncio.ensure_future(self._do_keypair_bundle_delete())
elif mtype == MNP.STREAM_REQUEST:
asyncio.ensure_future(self._stream_video(msg))
+ elif mtype == MNP.STREAM_MORE:
+ self._grant_stream_credit(msg)
else:
log.warning("Unknown MNP message type on DataChannel: %s", mtype)
except Exception as e:
@@ -999,6 +1011,69 @@ class WebRTCPeerSession:
self._audit("dir_delete", rel)
self._send({"type": MNP.DIR_DELETE_ACK, "v": MNP_VERSION, "dir": rel})
+ def _do_member_revoke(self, msg: dict) -> None:
+ """
+ Stop serving the group key to someone, at the operator's request.
+
+ The same authority as an invite, and the same reason: the roster decides
+ who this node serves, so only a key the node pinned as an operator may
+ change it. Membership on the hub is not consulted — the hub can remove
+ someone from a group, and that stops them reaching the node at all, but
+ it cannot make the node forget them.
+ """
+ user_id = str(msg.get("user_id", "")).strip()
+ if not user_id:
+ self._send({"type": "error", "detail": "Missing user_id"})
+ return
+ if user_id == self._user_id:
+ # Removing yourself from your own node is not a member operation;
+ # it would leave the group with nobody able to invite.
+ self._send({"type": "error", "detail": "Cannot revoke yourself"})
+ return
+ if not self._has_admin_authority():
+ self._send({"type": "error", "detail": "No authorized key for this"})
+ return
+ self._issue_admin_challenge(OP_MEMBER_REVOKE, user_id)
+
+ async def _admin_exec_member_revoke(
+ self, pending: dict, transcript: bytes, sig: bytes,
+ ) -> None:
+ user_id = pending["subject"]
+ if not await self._verify_admin_sig(transcript, sig):
+ self._send({"type": "error", "detail": "Signature verification failed"})
+ self._audit("admin_auth_failed", f"member_revoke:{user_id[:8]}")
+ return
+
+ roster = self._ctx.get("roster")
+ if roster is None:
+ self._send({"type": "error", "detail": "Roster not available"})
+ return
+
+ group_id = self._group_id or ""
+ if not await roster.set_status(group_id, user_id, "revoked"):
+ self._send({"type": "error", "detail": "Not a member of this group"})
+ return
+
+ # Anyone connected right now keeps the key they already unwrapped; what
+ # they lose is the next one. Rotating it is the operator's call, and the
+ # ack says so rather than implying this undid anything already read.
+ peer = self._peer_registry().get(user_id)
+ if peer is not None:
+ try:
+ await peer.close()
+ except Exception:
+ pass
+
+ log.info("Member revoked by %s: user=%s group=%s",
+ self._user_id[:8], user_id[:8], group_id[:8] or "-")
+ self._audit("member_revoke", user_id)
+ self._send({
+ "type": MNP.MEMBER_REVOKE_ACK, "v": MNP_VERSION,
+ "user_id": user_id,
+ "reminder": "they still hold the current group key — rotate it with "
+ "meshbay-node gek-init",
+ })
+
async def _do_keypair_bundle_delete(self) -> None:
"""
Withdraw our own key backup from this node.
@@ -1542,6 +1617,9 @@ class WebRTCPeerSession:
elif pending["op"] == OP_DIR_DELETE:
asyncio.ensure_future(
self._admin_exec_dir_delete(pending, transcript, sig_bytes))
+ elif pending["op"] == OP_MEMBER_REVOKE:
+ asyncio.ensure_future(
+ self._admin_exec_member_revoke(pending, transcript, sig_bytes))
elif pending["op"] == OP_INVITE_CREATE:
asyncio.ensure_future(
self._admin_exec_invite_create(pending, transcript, sig_bytes))
@@ -1634,6 +1712,37 @@ class WebRTCPeerSession:
"file_id": file_id,
})
+ def _grant_stream_credit(self, msg: dict) -> None:
+ """The client has room for more segments."""
+ try:
+ n = int(msg.get("n", 1))
+ except (TypeError, ValueError):
+ n = 1
+ self._stream_credit += max(0, min(n, STREAM_MAX_CREDIT))
+ self._stream_credit_evt.set()
+
+ async def _await_stream_credit(self) -> bool:
+ """
+ Block until the client has room. False if it stopped asking.
+
+ Without this the node hands ffmpeg's entire output to the channel as
+ fast as it is produced, and the browser holds a four gigabyte film in a
+ JavaScript array while MediaSource consumes it a segment at a time.
+ """
+ while self._stream_credit <= 0:
+ self._stream_credit_evt.clear()
+ try:
+ await asyncio.wait_for(self._stream_credit_evt.wait(),
+ timeout=STREAM_CREDIT_TIMEOUT)
+ except asyncio.TimeoutError:
+ log.info("Stream stalled: no credit from peer=%s",
+ (self._user_id or "?")[:8])
+ return False
+ if self._channel is None or self._channel.readyState != "open":
+ return False
+ self._stream_credit -= 1
+ return True
+
async def _stream_video(self, msg: dict) -> None:
"""Stream a video file as fMP4 segments via MSE-compatible output."""
# One ffmpeg per request with no cap lets any member exhaust the node's
@@ -1692,9 +1801,20 @@ class WebRTCPeerSession:
"duration": duration,
})
+ # A client that says nothing gets the old behaviour, which is why this
+ # defaults to unlimited rather than to zero: a stream that waits for
+ # credit from a peer that will never send any is a stream that hangs.
+ try:
+ self._stream_credit = int(msg.get("credits", 0) or 0)
+ except (TypeError, ValueError):
+ self._stream_credit = 0
+ paced = self._stream_credit > 0
+
index = 0
try:
while True:
+ if paced and not await self._await_stream_credit():
+ break
data = await proc.stdout.read(STREAM_SEGMENT_SIZE)
if not data:
break
diff --git a/packages/meshbay-node/tests/test_roster_pairing.py b/packages/meshbay-node/tests/test_roster_pairing.py
index 88e426d..435cc76 100644
--- a/packages/meshbay-node/tests/test_roster_pairing.py
+++ b/packages/meshbay-node/tests/test_roster_pairing.py
@@ -914,3 +914,85 @@ async def test_someone_elses_signature_does_not_remove_it(tmp_path, roster):
assert _last(session).get("detail") == "Signature verification failed"
assert (tmp_path / "shared" / "theirs").is_dir(), (
"a member's signature removed a directory — only the operator may")
+
+
+# ── Removing a member ────────────────────────────────────────────────────────
+
+async def test_revoking_needs_an_operator_signature(tmp_path, roster):
+ from meshbay_common.adminop import admin_transcript
+
+ sk_op, pk_op, pk_x_op = _keypair()
+ await roster.pin_identity("grenet", "grenet", pk_op, pk_x_op, "code")
+ await roster.set_member("", "grenet", ROLE_OPERATOR, "active", "local-cli")
+
+ sk_m, pk_m, pk_x_m = _keypair()
+ await roster.pin_identity("victim", "victim", pk_m, pk_x_m, "code")
+ await roster.set_member("g1", "victim", ROLE_MEMBER, "active", "grenet")
+
+ session = _session(tmp_path, roster, group_id="g1", gek=generate_gek())
+ session._admin_ops = {}
+ session._ctx["has_admin_authority"] = True
+ session._ctx["peers"] = {}
+
+ transcript = admin_transcript(
+ op="member_revoke", node_pk_b64=session._node_pk_b64(), group_id="g1",
+ subject="victim", nonce=b"\x33" * 32, ts=int(time.time()))
+
+ # A member's own signature is not enough.
+ await session._admin_exec_member_revoke(
+ {"op": "member_revoke", "subject": "victim"}, transcript,
+ sk_m.sign(transcript))
+ assert _last(session).get("detail") == "Signature verification failed"
+ assert (await roster.get_member("g1", "victim"))["status"] == "active"
+
+ # The operator's is.
+ await session._admin_exec_member_revoke(
+ {"op": "member_revoke", "subject": "victim"}, transcript,
+ sk_op.sign(transcript))
+ assert _last(session)["type"] == "member_revoke_ack"
+ assert (await roster.get_member("g1", "victim"))["status"] == "revoked"
+
+
+async def test_revoking_is_confined_to_the_group_it_was_asked_for(tmp_path, roster):
+ """
+ A node hosting two groups must not lose someone from both. Their pinned
+ identity survives as well — forgetting a key is `member unpin`, and saying
+ "remove them" should not silently do it.
+ """
+ from meshbay_common.adminop import admin_transcript
+
+ sk_op, pk_op, pk_x_op = _keypair()
+ await roster.pin_identity("grenet", "grenet", pk_op, pk_x_op, "code")
+ await roster.set_member("", "grenet", ROLE_OPERATOR, "active", "local-cli")
+
+ _, pk_m, pk_x_m = _keypair()
+ await roster.pin_identity("both", "both", pk_m, pk_x_m, "code")
+ await roster.set_member("g1", "both", ROLE_MEMBER, "active", "grenet")
+ await roster.set_member("g2", "both", ROLE_MEMBER, "active", "grenet")
+
+ session = _session(tmp_path, roster, group_id="g1", gek=generate_gek())
+ session._admin_ops = {}
+ session._ctx["has_admin_authority"] = True
+ session._ctx["peers"] = {}
+
+ transcript = admin_transcript(
+ op="member_revoke", node_pk_b64=session._node_pk_b64(), group_id="g1",
+ subject="both", nonce=b"\x44" * 32, ts=int(time.time()))
+ await session._admin_exec_member_revoke(
+ {"op": "member_revoke", "subject": "both"}, transcript,
+ sk_op.sign(transcript))
+
+ assert (await roster.get_member("g1", "both"))["status"] == "revoked"
+ assert (await roster.get_member("g2", "both"))["status"] == "active"
+ assert await roster.get_identity("both") is not None, (
+ "the pinned identity was dropped; that is `member unpin`, not this")
+
+
+async def test_an_operator_cannot_revoke_themselves(tmp_path, roster):
+ """It would leave the group with nobody able to invite or remove."""
+ session = _session(tmp_path, roster, group_id="g1", gek=generate_gek())
+ session._admin_ops = {}
+ session._ctx["has_admin_authority"] = True
+
+ session._do_member_revoke({"user_id": session._user_id})
+ assert _last(session).get("detail") == "Cannot revoke yourself"