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
|
"""Transfer slots and leases: what a member opens before a download or an
upload, what the node tells everyone when one frees, and the sweep that
reclaims leases a vanished peer left behind."""
import asyncio
import logging
from typing import TYPE_CHECKING
from meshbay_common import MNP_VERSION
from meshbay_common.protocol import MNP
from meshbay_node import transfers as transfers_mod
from meshbay_node.transfers import TransferSlots
from meshbay_node.transport.webrtc.limits import LEASE_GRANTED, LEASE_NONE, LEASE_QUEUED
log = logging.getLogger("meshbay_node.transport.webrtc_server")
if TYPE_CHECKING: # annotations only: the facade imports this module
from meshbay_node.transport.webrtc_server import WebRTCPeerSession
# How often transfer leases are swept. Nothing depends on it being
# prompt -- the session teardown is the reclaim that matters and is
# immediate; this catches peers that vanished without the connection
# noticing, so it trades latency for a timer that hardly ever runs.
TRANSFER_SWEEP_SECS = 15
class TransferMixin:
def _slots(self) -> "TransferSlots":
"""The node's transfer pools, shared across every peer and every group.
On the transport context, not the session: it counts the node's
transfers, not one browser's. Built once, for the same reason the
transcode semaphore is — rebuilding it per call would hand every caller
its own budget and cap nothing at all.
"""
slots = self._ctx.get("_transfer_slots")
if slots is None:
slots = TransferSlots()
n = self._ctx.get("max_concurrent_downloads")
u = self._ctx.get("max_concurrent_uploads")
if n:
slots.caps[transfers_mod.DOWNLOAD] = int(n)
if u:
slots.caps[transfers_mod.UPLOAD] = int(u)
self._ctx["_transfer_slots"] = slots
log.info("transfer: %s", slots.summary())
# Refreshed from the group context rather than only at construction: a
# node serves several groups, each with its own signed cap, and the
# pools are built by whichever group happens to transfer first.
limits = self._group_ctx().get("transfer_limits")
if limits and self._group_id:
slots.group_limits[self._group_id] = dict(limits)
return slots
def _lease_of(self, tr) -> str:
"""What the `tr` on a request actually is, from the node's own record.
`tr` is drawn by the client (§5.5) and arrives on every chunk request
and every upload chunk. It was read as a boolean: *present* meant "this
is a leased transfer", and nothing asked whether the node had ever
granted such a lease — so any non-empty string skipped the leaseless
ceiling and every cap the operator set. `touch()` has always answered
exactly this question (`False` if it is not granted) and its answer was
discarded.
Three outcomes, because they deserve different treatment:
- **`granted`** — a live lease of *this session*, and the transfer is
under the caps it was granted against. The session is checked as well
as the id: a lease belongs to a connection, and touching somebody
else's would refresh their idle timer.
- **`queued`** — the node has this lease and has not granted it. The
client is jumping its own queue; refused, and no shipped client does
it (the transfer store awaits the grant before it reads a byte).
- **`none`** — the node has no such lease. Deliberately *not* a refusal:
it is what a reconnect looks like from here, where the session's
leases died with the old connection and the client is re-opening
them, and it is what a client that never asked looks like. Both are
then bounded by the leaseless ceiling instead — which is the residual
§5.5 already states: a client that lies gets that bound's worth of
files at a time, not the whole library.
"""
slots = self._ctx.get("_transfer_slots")
if slots is None:
return LEASE_NONE
lease = slots.leases.get(tr)
if lease is None or lease.session_key != self._registry_key:
return LEASE_NONE
if lease.state != "granted":
return LEASE_QUEUED
slots.touch(tr)
return LEASE_GRANTED
def _note_unleased(self, tr: str) -> None:
"""Say once that this connection transferred outside its lease.
Once per session, not per chunk: the interesting fact is that it
happened, and a per-chunk line would bury it under itself. The bound is
the leaseless ceiling either way; this is what makes the residual
visible to the operator rather than merely stated in a document.
"""
if self._unleased_noted:
return
self._unleased_noted = True
log.info("transfer: %s sent chunk requests under an unknown lease %s "
"— bounded by the leaseless ceiling",
(self._user_id or "?")[:8], str(tr)[:8])
self._audit("transfer_unleased", str(tr)[:16])
def _transfer_state_msg(self, lease, state: str, reason: str = "") -> dict:
slots = self._slots()
out = {
"type": MNP.TRANSFER_STATE,
"v": MNP_VERSION,
"tr": lease.tr,
"state": state,
"kind": lease.kind,
"used": slots.member_in_use(lease.kind, lease.member),
# `member_cap`, not the node-wide default: this group's own limit is
# what `_has_room` enforces and what the handshake ack announces, so
# reading the default here would have the widget contradicting both
# — "1 of 2" in a group where the operator signed 5, or two slots
# offered in a group limited to one.
"cap": slots.member_cap(lease.kind, lease.member),
"node_used": slots.in_use(lease.kind),
"node_cap": slots.caps.get(lease.kind,
transfers_mod.DEFAULT_MAX_CONCURRENT),
}
if state == "queued":
out["ahead"] = slots.ahead_of(lease)
if reason:
out["reason"] = reason
return out
def _notify_transfer(self, lease, state: str, reason: str = "") -> None:
"""Push a lease's state to the connection that owns it.
By session key, never by account: a lease belongs to one connection, and
telling a member's other device that *its* transfer was granted is how a
queue starts lying.
"""
session = self._peer_registry().get(lease.session_key)
for candidate in ([session] if session else
self._sessions_everywhere(lease.session_key)):
try:
candidate._send(self._transfer_state_msg(lease, state, reason))
except Exception:
pass
def _sessions_everywhere(self, session_key: str) -> list["WebRTCPeerSession"]:
"""The session with this key, whichever group it is in.
`_peer_registry` is per group (finding H1) and the pools are node-wide,
so a slot freed in one group can grant one in another: the peer to tell
is not necessarily in this session's own registry.
"""
groups = self._ctx.get("groups")
registries = ([g.get("_peers", {}) for g in groups.values()]
if groups else [self._ctx.get("_peers", {})])
return [reg[session_key] for reg in registries if session_key in reg]
def _announce(self, granted: list, ended: list | None = None) -> None:
for lease, reason in (ended or []):
self._notify_transfer(
lease, "queued" if lease.state == "queued" else "closed", reason)
for lease in granted:
self._notify_transfer(lease, "granted")
def _do_transfer_open(self, msg: dict) -> None:
tr = str(msg.get("tr") or "")[:64]
kind = str(msg.get("kind") or transfers_mod.DOWNLOAD)
if not tr:
self._send({"type": "error", "detail": "Missing transfer id",
"code": "bad_transfer_id"})
return
slots = self._slots()
try:
nbytes = int(msg.get("bytes") or 0)
chunks = int(msg.get("chunks") or 0)
except (TypeError, ValueError):
self._send({"type": "error", "detail": "Invalid transfer size",
"code": "bad_transfer_size", "tr": tr})
return
lease, err = slots.open(
tr=tr, kind=kind, session_key=self._registry_key,
user_id=self._user_id or "", group_id=self._group_id or "",
bytes=nbytes, chunks=chunks)
if err:
self._send({"type": "error", "detail": err, "code": err, "tr": tr})
return
self._send(self._transfer_state_msg(lease, lease.state))
# INFO, not DEBUG. This is the line that answers "did the client ever
# ask for a slot, and what was it told" when somebody reports a stuck
# transfer — and a whole afternoon was spent concluding "the node saw
# nothing" from a journal that could not have shown it. One line per
# transfer is not a volume problem; turning the root logger up to DEBUG
# to see it is, because aiortc logs every SCTP chunk.
log.info("transfer: open %s %s -> %s (%s)",
kind, tr[:8], lease.state, slots.summary())
self._ensure_transfer_sweeper()
def _do_transfer_close(self, msg: dict) -> None:
tr = str(msg.get("tr") or "")[:64]
reason = str(msg.get("reason") or transfers_mod.REASON_DONE)[:32]
slots = self._slots()
held = slots.leases.get(tr)
if held is not None and held.session_key != self._registry_key:
# Closing somebody else's transfer would be a denial of service one
# random id away.
self._send({"type": "error", "detail": "not_your_transfer",
"code": "not_your_transfer", "tr": tr})
return
lease, granted = slots.close(tr, reason)
if lease is not None:
self._send(self._transfer_state_msg(lease, "closed", reason))
self._announce(granted)
def _release_transfers(self) -> None:
"""Give back everything this connection held. Called from teardown."""
slots = self._ctx.get("_transfer_slots")
if slots is None:
return
gone, granted = slots.release_session(self._registry_key)
if gone:
log.info("transfer: session gone, released %d (%s)",
len(gone), slots.summary())
self._announce(granted)
def _ensure_transfer_sweeper(self) -> None:
"""Start the maintenance task, once, and only while it has work.
It reclaims what a session teardown cannot see — a grant nobody took up,
a transfer that went quiet — and logs the one line that answers "was
this peer ever in a queue" when somebody reports a stuck transfer. It
stops when the last lease goes, so an idle node runs no timer.
"""
running = self._ctx.get("_transfer_sweeper")
if running is not None and not running.done():
return
ctx = self._ctx
async def _sweep_loop() -> None:
while True:
await asyncio.sleep(TRANSFER_SWEEP_SECS)
slots = ctx.get("_transfer_slots")
if slots is None or not slots.leases:
return
ended, granted = slots.sweep()
for lease, reason in ended:
log.info("transfer: reclaimed %s (%s)", lease.tr[:8], reason)
self._announce(granted, ended)
log.debug("transfer: %s", slots.summary())
# Deliberately NOT `self._spawn`, which is otherwise the only way to
# start a task here. `_spawn` ties a task to *this session's* set, and
# `shutdown_tasks` cancels those when the peer leaves — so the sweeper
# would die with whichever connection happened to open the first
# transfer, and every other peer's abandoned lease would then never be
# reclaimed. It belongs to the node, so the strong reference that keeps
# it off the garbage collector lives on the transport context; the rule
# `_spawn` exists for (asyncio holds only a weak reference) is satisfied
# by that reference, not by which set it is in.
task = asyncio.ensure_future(_sweep_loop())
ctx["_transfer_sweeper"] = task
def _finished(done: asyncio.Task) -> None:
if ctx.get("_transfer_sweeper") is done:
ctx["_transfer_sweeper"] = None
if not done.cancelled() and done.exception() is not None:
# Nothing awaits this task, so an exception here would otherwise
# be swallowed and idle leases would silently stop being
# reclaimed — the failure mode is a node that fills up over days.
log.error("transfer: sweeper died: %r", done.exception())
task.add_done_callback(_finished)
|