summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/tests
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-08 14:20:57 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-08 14:20:57 +0200
commitbdeffa448cdde9680fbf7bdda036746b101ddf75 (patch)
tree04d90f4d5bf7de00e508a23519add4a6c951d656 /packages/meshbay-node/tests
parent038066c43caa8ee76dd1e04761234271e4e67ecd (diff)
downloadmeshbay-bdeffa448cdde9680fbf7bdda036746b101ddf75.tar.gz
feat(node): transfer leases, pools and a queue for downloads and uploads
Step 2 of ~/next/improve-downloads.md. A download is invisible to the node: it is a series of independent `file_req` messages, with nothing saying one started or ended, so there is nothing to count and nothing to cap. The lease is that missing object. `meshbay_node/transfers.py` holds the decisions and has no asyncio and no transport in it, on purpose. The failure modes this has to survive — a slot the node never gets back, a client waiting on a grant the node has forgotten — are races through a DataChannel and unprovable there; here the clock is a parameter and every method returns what changed, so the caller does the I/O and the tests drive the worst case directly. What it decides: - two pools, downloads and uploads, separate from the stream pool: different resources with different costs, and merging them makes both caps meaningless; - per-member cap checked *before* the node-wide one, so a member at their own limit queues behind their own transfers rather than holding a slot a second member has none of. Per account across their devices, or the cap becomes a function of how many tabs somebody opens; - a queue that skips a member at their cap instead of waiting for them — granting strictly in arrival order lets one member's limit stall everyone; - `tr` drawn by the client and idempotent, which is what makes a reconnect safe; - bounded per member, because unbounded queues are how a node runs out of memory politely. Every way a slot comes back, with the session teardown as the one that matters (a closed tab, a quit browser and a dead network all arrive at `shutdown_tasks`, and none of them needs a timer): explicit close, session gone, a grant nobody took up in 30 s passed to the next in line, and a granted transfer silent for 120 s reclaimed with its peer told, so a widget can offer a resume rather than sit on a lie. `GET /api/transfers` is the operator's window: when somebody reports a transfer stuck at waiting, it is the only thing that says whether the node ever had them in a queue — a log cannot, when the symptom is that nothing is happening. It carries no filename and no path, which a test pins, because this is exactly where one would be tempting. Three things found while writing it, two of them mine: - the randomised property test rejected `in_use <= cap` at once, and it was right to: lowering a cap never interrupts a running transfer, so the count legitimately sits above the new value. The invariant is that a *new* grant never happens past the cap; - the sweeper was started with `self._spawn`, which ties a task to one session's set. It died with whichever peer opened the first transfer, and every other peer's abandoned lease then stopped being reclaimed — a node that fills up over days with nothing in the log. It belongs to the node now, with its strong reference on the transport context; - the pools are node-wide while `_peer_registry` is per group (finding H1), so a slot freed in one group can grant one in another and the peer to notify is not in the notifier's registry. Silently wrong in the first version. Nothing enforces a lease yet: `file_req` is untouched, no client asks, and the node grants everything. That is step 4's flag day, and this lands alone. 1148 node, 793 hub, 0 failed. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HCGdheDLxGReuKHga3BtST
Diffstat (limited to 'packages/meshbay-node/tests')
-rw-r--r--packages/meshbay-node/tests/test_transfer_slots.py308
-rw-r--r--packages/meshbay-node/tests/test_transfer_slots_wire.py298
2 files changed, 606 insertions, 0 deletions
diff --git a/packages/meshbay-node/tests/test_transfer_slots.py b/packages/meshbay-node/tests/test_transfer_slots.py
new file mode 100644
index 0000000..a0ef381
--- /dev/null
+++ b/packages/meshbay-node/tests/test_transfer_slots.py
@@ -0,0 +1,308 @@
+"""
+Transfer slots: the caps, the queue, and every way a slot can be lost.
+
+The requirement this is written against is not "a cap exists". It is that
+**nobody stays stuck** — neither a slot the node never gets back, which fills
+the node and queues everyone for ever, nor a transfer a client shows as waiting
+that the node has already forgotten.
+
+`TransferSlots` has no asyncio and no transport in it precisely so that those
+failures can be driven here instead of through a DataChannel, where they are
+rare, timing-dependent and unprovable. The clock is passed in, so the two
+timeouts are exercised without a test that sleeps for two minutes.
+
+`test_the_counter_never_drifts` is the one that matters most: a stuck slot is a
+race by nature, "it works now" is not evidence against a race, and the earlier
+flow-control bugs in this repo (window_leak.mjs, the discarded segment that
+leaked a slot per discard) were all found by forcing the worst case rather than
+by reasoning about it.
+"""
+
+import random
+
+import pytest
+
+from meshbay_node.transfers import (
+ DOWNLOAD, GRANT_DEADLINE_SECS, IDLE_TIMEOUT_SECS, KINDS,
+ MAX_QUEUED_PER_MEMBER, REASON_IDLE, REASON_NOT_TAKEN_UP, TransferSlots,
+ UPLOAD,
+)
+
+
+def _slots(node=8, per_member=2) -> TransferSlots:
+ s = TransferSlots()
+ s.caps = {k: node for k in KINDS}
+ s.per_member = {k: per_member for k in KINDS}
+ return s
+
+
+def _open(s, tr, *, session="s1", user="u1", group="g1", kind=DOWNLOAD, now=0.0):
+ lease, err = s.open(tr=tr, kind=kind, session_key=session, user_id=user,
+ group_id=group, now=now)
+ assert not err, err
+ return lease
+
+
+# ── the caps ────────────────────────────────────────────────────────────────
+
+def test_a_member_is_held_to_their_own_cap_first(_=None):
+ s = _slots(node=8, per_member=2)
+ assert _open(s, "a").state == "granted"
+ assert _open(s, "b").state == "granted"
+ assert _open(s, "c").state == "queued", (
+ "a third transfer for one member must queue even though the node has "
+ "six free slots — otherwise one member takes the node")
+
+
+def test_the_member_cap_spans_their_devices(_=None):
+ """Per account, not per connection: two browsers and a desktop client
+ signed in as the same person share the two slots, or the cap becomes a
+ function of how many tabs somebody opens."""
+ s = _slots(per_member=2)
+ _open(s, "a", session="laptop")
+ _open(s, "b", session="phone")
+ assert _open(s, "c", session="desktop").state == "queued"
+
+
+def test_the_node_cap_holds_across_members(_=None):
+ s = _slots(node=3, per_member=2)
+ _open(s, "a", user="u1")
+ _open(s, "b", user="u1")
+ _open(s, "c", user="u2")
+ assert _open(s, "d", user="u2").state == "queued"
+ assert s.in_use(DOWNLOAD) == 3
+
+
+def test_downloads_and_uploads_have_separate_pools(_=None):
+ s = _slots(node=2, per_member=2)
+ _open(s, "a", kind=DOWNLOAD)
+ _open(s, "b", kind=DOWNLOAD)
+ assert _open(s, "c", kind=UPLOAD).state == "granted", (
+ "a full download pool must not stop an upload")
+
+
+# ── the queue ───────────────────────────────────────────────────────────────
+
+def test_a_freed_slot_goes_to_whoever_was_waiting(_=None):
+ s = _slots(node=1, per_member=2)
+ _open(s, "a", user="u1")
+ queued = _open(s, "b", user="u2")
+ assert queued.state == "queued"
+ _, granted = s.close("a")
+ assert [x.tr for x in granted] == ["b"]
+ assert s.leases["b"].state == "granted"
+
+
+def test_a_member_at_their_cap_is_skipped_not_waited_for(_=None):
+ """Granting strictly in arrival order lets one member's own limit stall
+ every other member behind them."""
+ s = _slots(node=3, per_member=2)
+ _open(s, "a", user="u1")
+ _open(s, "b", user="u1")
+ hog = _open(s, "c", user="u1") # u1 is at their cap
+ other = _open(s, "d", user="u2") # arrives later
+ assert hog.state == "queued"
+ assert other.state == "granted", "u2 was made to wait behind u1's own limit"
+
+
+def test_position_is_reported_from_the_queue_itself(_=None):
+ s = _slots(node=1, per_member=8)
+ _open(s, "a")
+ b, c = _open(s, "b"), _open(s, "c")
+ assert (s.ahead_of(b), s.ahead_of(c)) == (0, 1)
+
+
+def test_a_member_cannot_queue_without_end(_=None):
+ s = _slots(node=1, per_member=1)
+ _open(s, "granted")
+ for i in range(MAX_QUEUED_PER_MEMBER):
+ _open(s, f"q{i}")
+ lease, err = s.open(tr="one-too-many", kind=DOWNLOAD, session_key="s1",
+ user_id="u1", group_id="g1")
+ assert lease is None and err == "too_many_queued"
+
+
+# ── every way a slot comes back (§5.1) ──────────────────────────────────────
+
+def test_closing_returns_the_slot(_=None):
+ s = _slots(node=1)
+ _open(s, "a")
+ s.close("a")
+ assert s.in_use(DOWNLOAD) == 0
+
+
+def test_losing_the_session_returns_everything_it_held(_=None):
+ """The primary reclaim, and the reason a lease is scoped to a connection:
+ a closed tab, a quit browser and a dropped network all arrive here, and
+ none of them needs a timer."""
+ s = _slots(node=8, per_member=8)
+ _open(s, "a", session="doomed")
+ _open(s, "b", session="doomed")
+ _open(s, "c", session="other")
+ gone, _ = s.release_session("doomed")
+ assert sorted(x.tr for x in gone) == ["a", "b"]
+ assert s.in_use(DOWNLOAD) == 1
+
+
+def test_a_queued_lease_dies_with_its_session_too(_=None):
+ s = _slots(node=1, per_member=8)
+ _open(s, "a", session="s1")
+ _open(s, "waiting", session="doomed")
+ s.release_session("doomed")
+ assert "waiting" not in s.leases
+ assert s.queues[DOWNLOAD] == []
+
+
+def test_a_grant_nobody_takes_up_is_passed_on(_=None):
+ s = _slots(node=1, per_member=8)
+ _open(s, "a", now=0.0)
+ _open(s, "b", now=0.0)
+ ended, granted = s.sweep(now=GRANT_DEADLINE_SECS + 1)
+ assert [(x.tr, r) for x, r in ended] == [("a", REASON_NOT_TAKEN_UP)]
+ assert [x.tr for x in granted] == ["b"], "the slot was not passed on"
+ assert s.leases["a"].state == "queued", "the abandoned one goes to the tail"
+
+
+def test_a_transfer_that_started_is_not_mistaken_for_an_abandoned_grant(_=None):
+ s = _slots(node=1)
+ _open(s, "a", now=0.0)
+ s.touch("a", now=1.0)
+ ended, _ = s.sweep(now=GRANT_DEADLINE_SECS + 2)
+ assert ended == [], "a transfer that is running was revoked"
+
+
+def test_a_transfer_that_goes_quiet_is_reclaimed(_=None):
+ s = _slots(node=1)
+ _open(s, "a", now=0.0)
+ s.touch("a", now=1.0)
+ ended, _ = s.sweep(now=1.0 + IDLE_TIMEOUT_SECS + 1)
+ assert [(x.tr, r) for x, r in ended] == [("a", REASON_IDLE)]
+ assert "a" not in s.leases
+
+
+def test_activity_keeps_a_slow_transfer_alive(_=None):
+ """A slow reader is not an absent one. The idle clock follows the lease's
+ own activity, not the wall since it started."""
+ s = _slots(node=1)
+ _open(s, "a", now=0.0)
+ t = 0.0
+ for _ in range(10):
+ t += IDLE_TIMEOUT_SECS - 1
+ s.touch("a", now=t)
+ assert s.sweep(now=t)[0] == []
+ assert "a" in s.leases
+
+
+# ── idempotence, which is what makes a reconnect safe ───────────────────────
+
+def test_reopening_the_same_transfer_does_not_charge_twice(_=None):
+ s = _slots(node=8, per_member=2)
+ first = _open(s, "a")
+ again = _open(s, "a")
+ assert again is first
+ assert s.in_use(DOWNLOAD) == 1
+
+
+def test_another_session_cannot_adopt_a_lease(_=None):
+ s = _slots()
+ _open(s, "a", session="mine")
+ lease, err = s.open(tr="a", kind=DOWNLOAD, session_key="theirs",
+ user_id="u1", group_id="g1")
+ assert lease is None and err == "not_your_transfer"
+
+
+# ── caps changed live ───────────────────────────────────────────────────────
+
+def test_raising_a_cap_starts_what_was_waiting(_=None):
+ s = _slots(node=1, per_member=8)
+ _open(s, "a")
+ _open(s, "b")
+ granted = s.set_caps(node={DOWNLOAD: 4})
+ assert [x.tr for x in granted] == ["b"]
+
+
+def test_lowering_a_cap_does_not_interrupt_anything(_=None):
+ s = _slots(node=4, per_member=4)
+ for tr in "abcd":
+ _open(s, tr)
+ s.set_caps(node={DOWNLOAD: 1})
+ assert s.in_use(DOWNLOAD) == 4, "a running transfer was taken away"
+ assert _open(s, "e").state == "queued"
+
+
+# ── the property that matters (§5.3) ────────────────────────────────────────
+
+@pytest.mark.parametrize("seed", range(25))
+def test_the_counter_never_drifts(seed):
+ """
+ Random open/close/drop/sweep/resize, checked after every single step.
+
+ A leaked slot is a race, and a test that reasons about the happy path
+ agrees with a broken implementation by construction. Two invariants, both
+ of which a real leak breaks: what the pool says is in use is exactly the
+ set of granted leases, and no queue entry names a lease that no longer
+ exists — the second being how "waiting for ever behind a ghost" starts.
+ """
+ rng = random.Random(seed)
+ s = _slots(node=rng.randint(1, 4), per_member=rng.randint(1, 3))
+ sessions = [f"s{i}" for i in range(4)]
+ users = ["u1", "u2", "u3"]
+ live: list[str] = []
+ now = 0.0
+ counter = 0
+
+ for _ in range(400):
+ before = {k: s.in_use(k) for k in KINDS}
+ member_before = {(k, m): s.member_in_use(k, m)
+ for k in KINDS
+ for m in {x.member for x in s.leases.values()}}
+ now += rng.uniform(0.0, 40.0)
+ action = rng.choice(
+ ["open", "open", "open", "close", "touch", "drop", "sweep", "caps"])
+ if action == "open":
+ counter += 1
+ tr = f"t{counter}"
+ lease, err = s.open(
+ tr=tr, kind=rng.choice(KINDS), session_key=rng.choice(sessions),
+ user_id=rng.choice(users), group_id="g1", now=now)
+ if lease is not None:
+ live.append(tr)
+ elif action == "close" and live:
+ s.close(live.pop(rng.randrange(len(live))), now=now)
+ elif action == "touch" and live:
+ s.touch(rng.choice(live), now=now)
+ elif action == "drop":
+ s.release_session(rng.choice(sessions), now=now)
+ elif action == "sweep":
+ s.sweep(now=now)
+ elif action == "caps":
+ s.set_caps(node={rng.choice(KINDS): rng.randint(1, 5)}, now=now)
+ live = [tr for tr in live if tr in s.leases]
+
+ for kind in KINDS:
+ granted = [x for x in s.leases.values()
+ if x.kind == kind and x.state == "granted"]
+ assert s.in_use(kind) == len(granted)
+ # Not `in_use <= cap`: lowering a cap never interrupts a transfer
+ # that is running, so the count legitimately sits above the new
+ # value until those finish. What must never happen is a *new* grant
+ # while the pool is at or over its cap -- so the count may fall or
+ # hold, and may only rise while there was room.
+ assert s.in_use(kind) <= max(s.caps[kind], before[kind]), (
+ f"{kind}: {before[kind]} -> {s.in_use(kind)} granted with a cap "
+ f"of {s.caps[kind]} — a slot was handed out past the cap")
+ for tr in s.queues[kind]:
+ assert tr in s.leases, "a queue entry outlived its lease"
+ assert s.leases[tr].state == "queued"
+ for member in {x.member for x in granted}:
+ assert s.member_in_use(kind, member) <= max(
+ s.per_member[kind], member_before.get((kind, member), 0))
+
+ # And at the end: drop every session and nothing may be left holding
+ # anything. A slot that survives the last connection is a slot nothing can
+ # ever release.
+ for session in sessions:
+ s.release_session(session, now=now)
+ assert s.leases == {}
+ assert all(q == [] for q in s.queues.values())
+ assert all(s.in_use(k) == 0 for k in KINDS)
diff --git a/packages/meshbay-node/tests/test_transfer_slots_wire.py b/packages/meshbay-node/tests/test_transfer_slots_wire.py
new file mode 100644
index 0000000..afa4582
--- /dev/null
+++ b/packages/meshbay-node/tests/test_transfer_slots_wire.py
@@ -0,0 +1,298 @@
+"""
+Transfer leases over the session, rather than over `TransferSlots` alone.
+
+test_transfer_slots.py proves the decisions; this proves the seam. Both exist
+because the seam is where this repo's defects have actually lived — a reply
+routed by arrival order, a session popped from a dict without its work being
+stopped, a slot released by a `finally` nobody reached.
+
+Three things can only be checked here:
+
+ - the handlers answer under the right shape, and refuse another connection's
+ transfer id;
+ - **losing the connection gives everything back.** That is the primary
+ reclaim, and it is a hook (`shutdown_tasks`) rather than a timer, so a test
+ of the pool alone would never touch it;
+ - a slot freed by one peer is *announced* to the peer waiting on it. A grant
+ nobody hears about is precisely the "stuck at waiting" report the design
+ exists to prevent, and it would look correct in the pool.
+"""
+
+import pytest
+
+# Every test drives a message handler, and in the node a message handler always
+# runs inside the event loop: `_do_transfer_open` starts the sweeper task there.
+# Calling these synchronously tested a situation that cannot happen and failed
+# on "no current event loop" the moment the sweeper stopped being faked.
+pytestmark = pytest.mark.asyncio
+
+from meshbay_common.protocol import MNP
+from meshbay_node.transfers import DOWNLOAD, UPLOAD
+from meshbay_node.transport.webrtc_server import WebRTCPeerSession
+
+
+class _Session(WebRTCPeerSession):
+ """A session with the DataChannel replaced by a list, and nothing else."""
+
+ def __init__(self, ctx, *, key, user, group="g1"):
+ self._ctx = ctx
+ self._registry_key = key
+ self._user_id = user
+ self._group_id = group
+ self.sent: list[dict] = []
+
+ def _send(self, msg):
+ self.sent.append(msg)
+
+ def _spawn(self, coro): # pragma: no cover - not used by these tests
+ coro.close()
+ return None
+
+ def last(self, mtype=MNP.TRANSFER_STATE):
+ return next(m for m in reversed(self.sent) if m.get("type") == mtype)
+
+
+@pytest.fixture
+def ctx():
+ """A transport context, with the sweeper stopped on the way out.
+
+ A task left running past the end of its test is a warning in the next one
+ and a hang in the worst case; the sweeper is started on demand by design, so
+ tearing it down is the test's job.
+ """
+ c: dict = {"_peers": {}}
+ yield c
+ task = c.get("_transfer_sweeper")
+ if task is not None:
+ task.cancel()
+
+
+def _join(ctx, key, user, group="g1") -> _Session:
+ s = _Session(ctx, key=key, user=user, group=group)
+ ctx["_peers"][key] = s
+ return s
+
+
+async def test_a_granted_transfer_is_answered_as_granted(ctx):
+ peer = _join(ctx, "s1", "alice")
+ peer._do_transfer_open({"tr": "t1", "kind": DOWNLOAD, "bytes": 10})
+ reply = peer.last()
+ assert reply["state"] == "granted"
+ assert reply["tr"] == "t1"
+ assert reply["kind"] == DOWNLOAD
+ assert reply["used"] == 1 and reply["cap"] >= 1
+
+
+async def test_a_queued_transfer_is_told_how_many_are_ahead(ctx):
+ peer = _join(ctx, "s1", "alice")
+ peer._slots().per_member[DOWNLOAD] = 1
+ peer._do_transfer_open({"tr": "t1"})
+ peer._do_transfer_open({"tr": "t2"})
+ peer._do_transfer_open({"tr": "t3"})
+ assert [m["state"] for m in peer.sent] == ["granted", "queued", "queued"]
+ assert peer.sent[-1]["ahead"] == 1
+
+
+async def test_the_reply_carries_no_name_and_no_path(ctx):
+ """A lease holds neither, and `transfer_state` stays in clear — so this is
+ the message where a filename would quietly become metadata on the wire."""
+ peer = _join(ctx, "s1", "alice")
+ peer._do_transfer_open({"tr": "t1", "name": "Some Saga.mkv",
+ "path": "/srv/films"})
+ assert set(peer.last()) <= {
+ "type", "v", "tr", "state", "kind", "used", "cap", "node_used",
+ "node_cap", "ahead", "reason"}
+
+
+async def test_closing_frees_the_slot(ctx):
+ peer = _join(ctx, "s1", "alice")
+ peer._do_transfer_open({"tr": "t1"})
+ peer._do_transfer_close({"tr": "t1", "reason": "done"})
+ assert peer.last()["state"] == "closed"
+ assert peer._slots().in_use(DOWNLOAD) == 0
+
+
+async def test_one_peer_cannot_close_anothers_transfer(ctx):
+ """A denial of service one random id away, otherwise."""
+ alice = _join(ctx, "s1", "alice")
+ bob = _join(ctx, "s2", "bob")
+ alice._do_transfer_open({"tr": "t1"})
+ bob._do_transfer_close({"tr": "t1"})
+ assert bob.last("error")["code"] == "not_your_transfer"
+ assert "t1" in alice._slots().leases
+
+
+async def test_one_peer_cannot_open_on_anothers_id(ctx):
+ alice = _join(ctx, "s1", "alice")
+ bob = _join(ctx, "s2", "bob")
+ alice._do_transfer_open({"tr": "t1"})
+ bob._do_transfer_open({"tr": "t1"})
+ assert bob.last("error")["code"] == "not_your_transfer"
+
+
+async def test_a_transfer_with_no_id_is_refused(ctx):
+ peer = _join(ctx, "s1", "alice")
+ peer._do_transfer_open({"kind": DOWNLOAD})
+ assert peer.last("error")["code"] == "bad_transfer_id"
+
+
+# ── the reclaim that matters ────────────────────────────────────────────────
+
+async def test_losing_the_connection_gives_everything_back(ctx):
+ peer = _join(ctx, "s1", "alice")
+ peer._do_transfer_open({"tr": "t1"})
+ peer._do_transfer_open({"tr": "t2", "kind": UPLOAD})
+ peer._release_transfers()
+ slots = peer._slots()
+ assert slots.leases == {}
+ assert slots.in_use(DOWNLOAD) == 0 and slots.in_use(UPLOAD) == 0
+
+
+async def test_the_freed_slot_reaches_the_peer_that_was_waiting(ctx):
+ """
+ The seam this file exists for. In the pool, granting is correct; if the
+ grant is not pushed, the waiting client sits on "waiting" for ever with a
+ node that believes it is streaming — and every unit test still passes.
+ """
+ alice = _join(ctx, "s1", "alice")
+ bob = _join(ctx, "s2", "bob")
+ alice._slots().caps[DOWNLOAD] = 1
+
+ alice._do_transfer_open({"tr": "a1"})
+ bob._do_transfer_open({"tr": "b1"})
+ assert bob.last()["state"] == "queued"
+
+ alice._release_transfers()
+ assert bob.last()["state"] == "granted", (
+ "bob was granted the slot and never told")
+ assert bob.last()["tr"] == "b1"
+
+
+async def test_a_grant_crossing_groups_still_reaches_its_peer(ctx):
+ """The pools are node-wide and `_peer_registry` is per group (finding H1),
+ so the peer to notify is not necessarily in the notifier's own registry."""
+ groups = {"g1": {"_peers": {}}, "g2": {"_peers": {}}}
+ ctx = {"groups": groups}
+ alice = _Session(ctx, key="s1", user="alice", group="g1")
+ groups["g1"]["_peers"]["s1"] = alice
+ bob = _Session(ctx, key="s2", user="bob", group="g2")
+ groups["g2"]["_peers"]["s2"] = bob
+
+ alice._slots().caps[DOWNLOAD] = 1
+ alice._do_transfer_open({"tr": "a1"})
+ bob._do_transfer_open({"tr": "b1"})
+ assert bob.last()["state"] == "queued"
+
+ alice._release_transfers()
+ assert bob.last()["state"] == "granted", (
+ "a slot freed in one group never reached the peer waiting in another")
+
+
+async def test_raising_the_cap_notifies_who_it_starts(ctx):
+ """`set_capacity` arrives from the loopback API, with no session behind it —
+ the grants it produces still have to be pushed."""
+ from meshbay_node.transport.webrtc_server import WebRTCTransport
+
+ transport = WebRTCTransport.__new__(WebRTCTransport)
+ transport._ctx = ctx
+ peer = _join(ctx, "s1", "alice")
+ peer._slots().caps[DOWNLOAD] = 1
+ peer._do_transfer_open({"tr": "t1"})
+ peer._do_transfer_open({"tr": "t2"})
+ assert peer.last()["state"] == "queued"
+
+ transport.set_capacity(max_concurrent_downloads=4)
+ assert peer.last()["state"] == "granted" and peer.last()["tr"] == "t2"
+
+
+async def test_the_operator_can_see_the_queue(ctx):
+ """`GET /api/transfers` is the answer to "was this peer ever queued", which
+ a log line cannot give when the symptom is that nothing is happening."""
+ from meshbay_node import ops
+
+ peer = _join(ctx, "s1", "alice")
+ peer._slots().caps[DOWNLOAD] = 1
+ peer._do_transfer_open({"tr": "t1", "bytes": 5})
+ peer._do_transfer_open({"tr": "t2", "bytes": 7})
+
+ class _T:
+ _ctx = ctx
+
+ snapshot = await ops.list_transfers({"webrtc": _T()})
+ assert snapshot["pools"][DOWNLOAD]["in_use"] == 1
+ assert snapshot["pools"][DOWNLOAD]["queued"] == 1
+ assert {x["state"] for x in snapshot["leases"]} == {"granted", "queued"}
+ assert all("name" not in x and "path" not in x for x in snapshot["leases"])
+
+
+async def test_asking_before_anything_has_transferred_is_not_an_error():
+ from meshbay_node import ops
+
+ class _T:
+ _ctx: dict = {}
+
+ snapshot = await ops.list_transfers({"webrtc": _T()})
+ assert snapshot["leases"] == []
+ assert snapshot["pools"][DOWNLOAD]["in_use"] == 0
+
+
+# ── the sweeper's lifetime ──────────────────────────────────────────────────
+
+async def test_the_sweeper_outlives_the_session_that_started_it(ctx):
+ """
+ It was started with `self._spawn`, which ties a task to one session's set —
+ so it was cancelled the moment that peer left, and every other peer's
+ abandoned lease stopped being reclaimed. Nothing else would have noticed:
+ the node simply fills up over days.
+ """
+ import asyncio
+
+ from meshbay_node.transport import webrtc_server as ws
+
+ alice = _join(ctx, "s1", "alice")
+ bob = _join(ctx, "s2", "bob")
+ # The real _spawn, so the session genuinely owns what it starts.
+ alice._tasks = set()
+ alice._spawn = ws.WebRTCPeerSession._spawn.__get__(alice)
+ bob._tasks = set()
+ bob._spawn = ws.WebRTCPeerSession._spawn.__get__(bob)
+
+ alice._do_transfer_open({"tr": "a1"})
+ bob._do_transfer_open({"tr": "b1"})
+ sweeper = ctx["_transfer_sweeper"]
+ assert sweeper is not None and not sweeper.done()
+
+ # Alice leaves, exactly as shutdown_tasks does it.
+ alice._release_transfers()
+ for task in list(alice._tasks):
+ task.cancel()
+ await asyncio.gather(*alice._tasks, return_exceptions=True)
+ await asyncio.sleep(0)
+
+ assert not sweeper.done(), (
+ "the sweeper died with the session that happened to start it; bob's "
+ "lease would never be reclaimed")
+ sweeper.cancel()
+
+
+async def test_the_sweeper_stops_when_the_last_lease_goes(ctx):
+ """An idle node must run no timer — the reason this is started on demand
+ rather than at boot."""
+ import asyncio
+
+ from meshbay_node.transport import webrtc_server as ws
+
+ peer = _join(ctx, "s1", "alice")
+ peer._tasks = set()
+ peer._spawn = ws.WebRTCPeerSession._spawn.__get__(peer)
+
+ original = ws.TRANSFER_SWEEP_SECS
+ ws.TRANSFER_SWEEP_SECS = 0.01
+ try:
+ peer._do_transfer_open({"tr": "t1"})
+ peer._do_transfer_close({"tr": "t1"})
+ sweeper = ctx["_transfer_sweeper"]
+ await asyncio.wait_for(sweeper, timeout=2)
+ assert ctx.get("_transfer_sweeper") is None
+ finally:
+ ws.TRANSFER_SWEEP_SECS = original