diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops.py | 22 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transfers.py | 29 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | 11 |
3 files changed, 53 insertions, 9 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/ops.py b/packages/meshbay-node/src/meshbay_node/ops.py index 81b92c7..088ee32 100644 --- a/packages/meshbay-node/src/meshbay_node/ops.py +++ b/packages/meshbay-node/src/meshbay_node/ops.py @@ -1433,16 +1433,28 @@ async def list_transfers(state: dict) -> dict: where it would be tempting to add one. """ webrtc = state.get("webrtc") - slots = getattr(webrtc, "_ctx", {}).get("_transfer_slots") if webrtc else None + ctx = getattr(webrtc, "_ctx", {}) if webrtc else {} + slots = ctx.get("_transfer_slots") if slots is None: from meshbay_node.transfers import ( DEFAULT_MAX_CONCURRENT, DEFAULT_MAX_PER_MEMBER, KINDS) # No pool built means nothing has transferred since the daemon started, # which is a real answer and not an error. - return {"pools": {k: {"in_use": 0, "cap": DEFAULT_MAX_CONCURRENT, - "per_member": DEFAULT_MAX_PER_MEMBER, - "queued": 0} for k in KINDS}, - "leases": []} + # + # The caps still have to be the operator's own. Reporting the module + # defaults here was worse than reporting nothing: `transfers set 2 2` + # answered "applied now", and `transfers show` immediately said 0/8 — + # a setting written, acknowledged and displayed wrong, which reads + # exactly like the hot-swap that did nothing for months. Found by + # running it, not by a test: the test asserted the defaults and so + # agreed with the bug. + return {"pools": { + k: {"in_use": 0, + "cap": int(ctx.get(f"max_concurrent_{k}s") + or DEFAULT_MAX_CONCURRENT), + "per_member": DEFAULT_MAX_PER_MEMBER, + "queued": 0} + for k in KINDS}, "leases": []} return slots.snapshot() diff --git a/packages/meshbay-node/src/meshbay_node/transfers.py b/packages/meshbay-node/src/meshbay_node/transfers.py index 533e985..dd5da5c 100644 --- a/packages/meshbay-node/src/meshbay_node/transfers.py +++ b/packages/meshbay-node/src/meshbay_node/transfers.py @@ -64,6 +64,10 @@ IDLE_TIMEOUT_SECS = 120.0 # Per account, per kind. Unbounded queues are how a node runs out of memory # politely; past this the client keeps the rest in its own list. MAX_QUEUED_PER_MEMBER = 32 +# How many times a lease may be granted and not taken up before it is closed +# rather than queued again. Without a bound the requeue is a permanent cycle, +# and a node logs the same reclaim every 30 s until it restarts. +MAX_MISSED_GRANTS = 3 # Why a lease ended, as it reaches the peer. REASON_DONE = "done" @@ -73,6 +77,7 @@ REASON_FAILED = "failed" REASON_SESSION_GONE = "session_gone" REASON_IDLE = "idle" REASON_NOT_TAKEN_UP = "not_taken_up" +REASON_ABANDONED = "abandoned" @dataclass @@ -92,6 +97,12 @@ class Lease: # first is a grant to revoke and pass on, the second a transfer to reclaim. used: bool = False last_seen: float = 0.0 + # How many grants this lease has been given and not taken up. Bounded + # because the requeue is otherwise a permanent cycle: revoked, put back, + # granted again a millisecond later because there is room, revoked 30 s + # later, for ever. Seen doing exactly that in a node's log, every 30 s, + # minutes after the transfers involved had finished. + missed_grants: int = 0 @property def member(self) -> tuple[str, str]: @@ -211,6 +222,7 @@ class TransferSlots: if lease is None or lease.state != "granted": return False lease.used = True + lease.missed_grants = 0 lease.last_seen = time.monotonic() if now is None else now return True @@ -267,11 +279,20 @@ class TransferSlots: continue if not lease.used and lease.granted_at is not None \ and now - lease.granted_at > GRANT_DEADLINE_SECS: - lease.state = "queued" + lease.missed_grants += 1 lease.granted_at = None - self.queues[lease.kind].append(lease.tr) - ended.append((lease, REASON_NOT_TAKEN_UP)) - requeued = True + if lease.missed_grants >= MAX_MISSED_GRANTS: + # It has had its chances. Closing it is what ends the cycle, + # and the peer is told so a client that is somehow still + # there can ask again from a clean state rather than hold a + # slot it has never once used. + self.leases.pop(lease.tr, None) + ended.append((lease, REASON_ABANDONED)) + else: + lease.state = "queued" + self.queues[lease.kind].append(lease.tr) + ended.append((lease, REASON_NOT_TAKEN_UP)) + requeued = True elif lease.used and now - lease.last_seen > IDLE_TIMEOUT_SECS: self.leases.pop(lease.tr, None) ended.append((lease, REASON_IDLE)) 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 fdab53d..6a75b52 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -3651,6 +3651,17 @@ class WebRTCPeerSession: async def _do_file_request(self, msg: dict) -> None: ctx = self._group_ctx() + # A chunk request is what "this transfer is alive" looks like. Nothing + # marked a lease used, so `used` stayed False for the whole download and + # the sweeper revoked the grant every 30 s as never-taken-up — while the + # file was transferring at 20 MB/s. Found in the node's own log, which + # repeated the same two reclaims every 30 s for as long as the daemon + # ran. + tr = msg.get("tr") + if tr: + slots = self._ctx.get("_transfer_slots") + if slots is not None: + slots.touch(str(tr)[:64]) file_id = msg["file_id"] chunk_index = msg["chunk_index"] entry = ctx["index"].get_entry(file_id) |