summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/tests/test_transfer_slots.py
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-09 14:28:40 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-09 14:28:40 +0200
commit7e2d078fe1870d256ae47781bee6ac4f454edf24 (patch)
tree05689049fd48f01ebbdb9995d5189cba15cee052 /packages/meshbay-node/tests/test_transfer_slots.py
parent813d18424ec57963bb56e6f40824a2db0ccce50d (diff)
parente6f895c473a0b19e7b186889c1836d3945bc880b (diff)
downloadmeshbay-7e2d078fe1870d256ae47781bee6ac4f454edf24.tar.gz
Merge branch 'fix/large-download-paths'
Concurrent-transfer limits, with the queue, the pause and the flag day. A node now caps how many transfers it runs at once (8 downloads, 8 uploads, node-wide) and how many one member may run in one group (2 by default, operator-signed). Beyond that the node answers "queued" and the client waits its turn, visibly, in the transfers panel — and a slot that frees starts whatever is next, skipping past a member who is at their own cap rather than letting them stall everyone behind them. Browsing is never subject to a slot: not the poster grid, not the covers, not opening a photo to look at it. That is structural — a transfer is what the transfers widget shows — and the exemption is bounded rather than open, at two files in flight per session, because an exemption with no bound is a leaseless branch under another name. Transfers can be cancelled, and now paused and resumed. A paused one holds nothing: its slot goes back at once and resuming rejoins the queue at the tail. Uploads survive the connection that started them and resume where the node stopped, asked for inside the seal rather than on a clear message. What they leave behind when they are abandoned is reaped, which closes a disk leak that predates this work. MNP 3.0 makes the lease compulsory and refuses 2.x at the handshake, with the desktop client checking `client.minimum` before connecting so an un-updated one says "update" instead of failing every connection in a protocol vocabulary. Fourteen defects were found on the way, eight of them by a person clicking Download and pasting a console — none of which 2075 tests could reach. Section 12 of ~/next/improve-downloads.md is that report, including the three this work introduced itself and the one that turned out to be caused by an instruction to hard-reload after each deployment. Node suite 1209 passed, hub suite 866 passed. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HCGdheDLxGReuKHga3BtST
Diffstat (limited to 'packages/meshbay-node/tests/test_transfer_slots.py')
-rw-r--r--packages/meshbay-node/tests/test_transfer_slots.py365
1 files changed, 365 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..7056e93
--- /dev/null
+++ b/packages/meshbay-node/tests/test_transfer_slots.py
@@ -0,0 +1,365 @@
+"""
+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_MISSED_GRANTS, MAX_QUEUED_PER_MEMBER, REASON_ABANDONED, 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)
+
+
+# ── the cycle the node's own log showed ─────────────────────────────────────
+
+def test_a_grant_is_not_requeued_for_ever(_=None):
+ """
+ A revoked grant went back in the queue, was granted again a millisecond
+ later because there was room, and was revoked again 30 s on. The node
+ logged the same two reclaims every 30 s for as long as it ran — minutes
+ after the transfers involved had finished.
+
+ Three chances, then it is closed and the peer told, which is what ends the
+ cycle. `test_the_counter_never_drifts` could not see this: nothing drifted,
+ the same lease simply never left.
+ """
+ s = _slots(node=4, per_member=4)
+ _open(s, "ghost", now=0.0)
+ now = 0.0
+ reasons = []
+ for _ in range(6):
+ now += GRANT_DEADLINE_SECS + 1
+ ended, _granted = s.sweep(now=now)
+ reasons += [r for _, r in ended]
+ assert reasons.count(REASON_NOT_TAKEN_UP) == MAX_MISSED_GRANTS - 1
+ assert reasons.count(REASON_ABANDONED) == 1
+ assert "ghost" not in s.leases, "the lease is still cycling"
+ assert s.queues[DOWNLOAD] == []
+
+
+def test_a_transfer_that_is_running_is_never_revoked(_=None):
+ """
+ The other half, and the one that mattered: nothing marked a lease used, so
+ `used` stayed False for a whole download and the sweeper revoked a grant
+ every 30 s while the file transferred at 20 MB/s.
+ """
+ s = _slots(node=2, per_member=2)
+ _open(s, "live", now=0.0)
+ now = 0.0
+ for _ in range(10):
+ now += GRANT_DEADLINE_SECS - 5
+ assert s.touch("live", now=now), "a granted lease refused a touch"
+ ended, _granted = s.sweep(now=now)
+ assert ended == [], f"a running transfer was revoked: {ended}"
+ assert s.leases["live"].state == "granted"
+
+
+def test_using_a_lease_forgives_its_earlier_misses(_=None):
+ """A slow start is not an abandoned one: a client that took two grants to
+ get going must not be closed on its third."""
+ s = _slots(node=2, per_member=2)
+ _open(s, "slow", now=0.0)
+ s.sweep(now=GRANT_DEADLINE_SECS + 1)
+ s.sweep(now=2 * GRANT_DEADLINE_SECS + 2)
+ assert s.leases["slow"].missed_grants == 2
+ s.touch("slow", now=2 * GRANT_DEADLINE_SECS + 3)
+ assert s.leases["slow"].missed_grants == 0