From ebe9a750bb7b44a115871f003a27ba08d5e876f8 Mon Sep 17 00:00:00 2001 From: Ben Papillon Date: Mon, 28 Sep 2026 14:01:18 -0700 Subject: [PATCH 1/5] skip the refund for a hold with no lease id --- src/schematic/leases/check.py | 9 +++++---- .../leases/redis_reservation_store.py | 8 ++++++-- src/schematic/leases/reservation_store.py | 7 ++++++- tests/leases/test_redis_reservation_store.py | 18 ++++++++++++++++++ tests/leases/test_reservation_store.py | 16 ++++++++++++++++ 5 files changed, 51 insertions(+), 7 deletions(-) diff --git a/src/schematic/leases/check.py b/src/schematic/leases/check.py index 8b2bdeb..8261375 100644 --- a/src/schematic/leases/check.py +++ b/src/schematic/leases/check.py @@ -259,11 +259,12 @@ async def failure(reason: str) -> "CheckResult": # claims whatever slice of the add landed and refunds it; a None says # nothing landed, so refund the debit directly. Both are pinned to the # lease the debit landed on (the record carries that id, so consume - # pins to it too), never to the acquired one. If the undo itself fails, - # accept the bounded leak: the slice comes back at lease expiry, which - # beats risking a double refund. + # pins to it too), never to the acquired one; a debit that names no + # lease is not refunded at all, as consume would not either. If the + # undo itself fails, accept the bounded leak: the slice comes back at + # lease expiry, which beats risking a double refund. try: - if await deps.reservations.consume(record.id, 0) is None: + if await deps.reservations.consume(record.id, 0) is None and reserve.lease_id: await deps.lease_store.refund(resolved_company.id, credit_id, credit_cost, reserve.lease_id) except Exception as undo_err: log.warning( diff --git a/src/schematic/leases/redis_reservation_store.py b/src/schematic/leases/redis_reservation_store.py index 0ed9473..32b4376 100644 --- a/src/schematic/leases/redis_reservation_store.py +++ b/src/schematic/leases/redis_reservation_store.py @@ -181,12 +181,16 @@ async def consume(self, reservation_id: str, credits_consumed: float) -> Optiona consumed = clamp_consumption(credits_consumed, reserved) refund = reserved - consumed - if refund > 0: + lease_id = raw.get("leaseId") + # A hold with no lease id cannot be pinned, and an empty pin disables + # the lease check, so the refund would land on whatever lease holds the + # slot now. Skip it: the slice comes back when its lease expires. + if refund > 0 and lease_id: # The lease store owns the lease hash, which keeps this cross-key # write out of a single Lua script. Pinned to the reservation's # lease so a hold carved out of an expired lease cannot inflate a # successor's balance. - await self._lease_store.refund(company_id, credit_type_id, refund, raw.get("leaseId")) + await self._lease_store.refund(company_id, credit_type_id, refund, lease_id) return consumed async def reserved_credits(self, company_id: str, credit_type_id: str) -> float: diff --git a/src/schematic/leases/reservation_store.py b/src/schematic/leases/reservation_store.py index 121631d..d4521dd 100644 --- a/src/schematic/leases/reservation_store.py +++ b/src/schematic/leases/reservation_store.py @@ -37,6 +37,11 @@ async def consume(self, reservation_id: str, credits_consumed: float) -> Optiona remainder is refunded to the lease (pinned to the reservation's lease), and the clamped figure is returned. A crash between the claim and the refund loses the refund, never double-refunds. + + A reservation with no lease id is claimed but not refunded: with + nothing to pin to, the refund would land on whichever lease holds the + slot now and could inflate a successor. The slice comes back when its + lease expires. Every SDK on a shared Redis must agree on this. """ @abc.abstractmethod @@ -84,7 +89,7 @@ async def consume(self, reservation_id: str, credits_consumed: float) -> Optiona return None consumed = clamp_consumption(credits_consumed, reservation.credits_reserved) refund = reservation.credits_reserved - consumed - if refund > 0: + if refund > 0 and reservation.lease_id: await self._lease_store.refund( reservation.company_id, reservation.credit_type_id, diff --git a/tests/leases/test_redis_reservation_store.py b/tests/leases/test_redis_reservation_store.py index 4034877..4a2b8ee 100644 --- a/tests/leases/test_redis_reservation_store.py +++ b/tests/leases/test_redis_reservation_store.py @@ -156,6 +156,24 @@ async def test_a_stale_lease_hold_never_inflates_its_successor( assert await reservations.get("res_1") is None +async def test_a_hold_with_no_lease_id_is_not_refunded( + leases: RedisLeaseStore, reservations: RedisReservationStore, frozen_clock: VirtualClock +) -> None: + # An empty pin disables the lease check in the refund script, so refunding + # would credit whichever lease holds the slot now. Skip it, as every other + # SDK sharing this Redis does. + await leases.try_reserve("co_1", "ct_1", 100) + await reservations.add(make_reservation(lease_id="", expires_at=frozen_clock() + 60)) + assert await reservations.consume("res_1", 30) == 30 + assert await _balance(leases) == 900 + assert await reservations.reserved_credits("co_1", "ct_1") == 0 + + await leases.try_reserve("co_1", "ct_1", 100) + await reservations.add(make_reservation(id="res_2", lease_id="", expires_at=frozen_clock() - 0.001)) + assert await reservations.sweep_expired() == 1 + assert await _balance(leases) == 800 + + async def test_reserved_credits_sums_open_holds( leases: RedisLeaseStore, reservations: RedisReservationStore, frozen_clock: VirtualClock ) -> None: diff --git a/tests/leases/test_reservation_store.py b/tests/leases/test_reservation_store.py index 10c5d2f..9fed4a1 100644 --- a/tests/leases/test_reservation_store.py +++ b/tests/leases/test_reservation_store.py @@ -131,3 +131,19 @@ async def test_reserved_credits_drops_a_consumed_hold( assert await reservations.reserved_credits("co_1", "ct_1") == 100 await reservations.consume("res_1", 30) assert await reservations.reserved_credits("co_1", "ct_1") == 0 + + +async def test_a_hold_with_no_lease_id_is_not_refunded( + leases: InMemoryLeaseStore, reservations: InMemoryReservationStore, clock: VirtualClock +) -> None: + # With nothing to pin to, a refund would land on whichever lease holds the + # slot now. The slice waits for its lease to expire instead. + await leases.try_reserve("co_1", "ct_1", 100) + await reservations.add(make_reservation(lease_id="", expires_at=clock() + 60)) + assert await reservations.consume("res_1", 30) == 30 + assert await _balance(leases) == 900 + + await leases.try_reserve("co_1", "ct_1", 100) + await reservations.add(make_reservation(id="res_2", lease_id="", expires_at=clock() - 0.001)) + assert await reservations.sweep_expired() == 1 + assert await _balance(leases) == 800 From 853ce6f327fed843cd5d194035e62086d6ef8817 Mon Sep 17 00:00:00 2001 From: Ben Papillon Date: Mon, 28 Sep 2026 14:02:46 -0700 Subject: [PATCH 2/5] round the server reservation quantity up --- src/schematic/client.py | 14 +++++++++----- tests/custom/test_client.py | 21 +++++++++++---------- 2 files changed, 20 insertions(+), 15 deletions(-) diff --git a/src/schematic/client.py b/src/schematic/client.py index 72b4be7..5fe331e 100644 --- a/src/schematic/client.py +++ b/src/schematic/client.py @@ -299,7 +299,7 @@ def _build_preflight(options: Optional[CheckFlagOptions]) -> Optional[PreflightR def _preflight_quantity(usage: float) -> int: """Cast a usage onto the integer the preflight body carries. - A hold can be sized from a fractional usage, but the API's preflight usage + A caller can pass a fractional usage, but the API's preflight usage is an integer. A preflight asks an upper-bound question ("would this action be allowed?"), so a fraction rounds up: the check must not pass on less usage than the operation is about to record. @@ -935,7 +935,7 @@ def failure(reason: str) -> CheckResult: flag_key, options, reason, self._resolve_default(flag_key, _check_options_to_flag_options(options)), ) - if not _is_valid_quantity(options.usage): + if options.usage is None or not _is_valid_quantity(options.usage): self.logger.error( f"Server reservation: invalid usage {options.usage!r} for flag {flag_key}; " "must be a finite, non-negative number" @@ -953,7 +953,9 @@ def failure(reason: str) -> CheckResult: flag_key, company=company, user=user, - quantity=options.usage, + # Whole units, like the local lease path: the settle bills + # ceil(actual), so a fractional hold would come up short. + quantity=_preflight_quantity(options.usage), expires_at=dt.datetime.now(dt.timezone.utc) + dt.timedelta(seconds=self._reservation_ttl), **_reservation_request_kwargs(options), ) @@ -1778,7 +1780,7 @@ def failure(reason: str) -> CheckResult: flag_key, options, reason, self._resolve_default(flag_key, _check_options_to_flag_options(options)), ) - if not _is_valid_quantity(options.usage): + if options.usage is None or not _is_valid_quantity(options.usage): self.logger.error( f"Server reservation: invalid usage {options.usage!r} for flag {flag_key}; " "must be a finite, non-negative number" @@ -1796,7 +1798,9 @@ def failure(reason: str) -> CheckResult: flag_key, company=company, user=user, - quantity=options.usage, + # Whole units, like the local lease path: the settle bills + # ceil(actual), so a fractional hold would come up short. + quantity=_preflight_quantity(options.usage), expires_at=dt.datetime.now(dt.timezone.utc) + dt.timedelta(seconds=self._reservation_ttl), **_reservation_request_kwargs(options), ) diff --git a/tests/custom/test_client.py b/tests/custom/test_client.py index 981c221..c74720d 100644 --- a/tests/custom/test_client.py +++ b/tests/custom/test_client.py @@ -1800,13 +1800,14 @@ def test_mints_a_fresh_idempotency_key_per_check(self): second = self.schematic.features.check_and_reserve_flag.call_args.kwargs["idempotency_key"] self.assertNotEqual(first, second) - def test_a_fractional_usage_sizes_the_hold_and_rounds_the_preflight_up(self): - self.schematic.check("inference", company={"id": "co_1"}, options=CheckOptions(usage=0.5)) + def test_a_fractional_usage_rounds_the_hold_and_the_preflight_up(self): + self.schematic.check("inference", company={"id": "co_1"}, options=CheckOptions(usage=2.5)) kwargs = self.schematic.features.check_and_reserve_flag.call_args.kwargs - self.assertEqual(kwargs["quantity"], 0.5) - # The hold takes the fraction; the preflight's usage is an integer, and - # rounding it down would ask about less usage than is about to land. - self.assertEqual(kwargs["preflight"], PreflightRequestBody(usage=1)) + # The settle bills whole units, so a 2.5 hold would take 2.5 * rate and + # the track would bill 3 * rate. Rounding the preflight down would ask + # about less usage than is about to land. + self.assertEqual(kwargs["quantity"], 3) + self.assertEqual(kwargs["preflight"], PreflightRequestBody(usage=3)) def test_an_integral_float_usage_reaches_the_preflight_unchanged(self): self.schematic.check( @@ -2268,11 +2269,11 @@ async def test_mints_a_fresh_idempotency_key_per_check(self): second = self.client.features.check_and_reserve_flag.call_args.kwargs["idempotency_key"] assert first != second - async def test_a_fractional_usage_sizes_the_hold_and_rounds_the_preflight_up(self): - await self.client.check("inference", company={"id": "co_1"}, options=CheckOptions(usage=0.5)) + async def test_a_fractional_usage_rounds_the_hold_and_the_preflight_up(self): + await self.client.check("inference", company={"id": "co_1"}, options=CheckOptions(usage=2.5)) kwargs = self.client.features.check_and_reserve_flag.call_args.kwargs - assert kwargs["quantity"] == 0.5 - assert kwargs["preflight"] == PreflightRequestBody(usage=1) + assert kwargs["quantity"] == 3 + assert kwargs["preflight"] == PreflightRequestBody(usage=3) async def test_a_reservation_ttl_above_the_cap_is_clamped(self): client = _async_server_client(credit_leases=CreditLeaseConfig(default_reservation_ttl=7200.0)) From f30fe7b3d12a87607dbdcd5e2dca988e31f4dda3 Mon Sep 17 00:00:00 2001 From: Ben Papillon Date: Mon, 28 Sep 2026 14:03:19 -0700 Subject: [PATCH 3/5] hold a joined extend to the joiner's own need --- src/schematic/leases/lease_manager.py | 15 ++++++---- tests/leases/test_lease_manager.py | 41 +++++++++++++++++++++++++++ 2 files changed, 51 insertions(+), 5 deletions(-) diff --git a/src/schematic/leases/lease_manager.py b/src/schematic/leases/lease_manager.py index 56f79ef..856d710 100644 --- a/src/schematic/leases/lease_manager.py +++ b/src/schematic/leases/lease_manager.py @@ -341,12 +341,17 @@ async def _maybe_extend( # The flight asked for at least what we need: every # watermark-driven joiner, and any check the tranche covers. # One wire call serves all of them, which is the point of - # single-flight. + # single-flight. But the flight re-checks against its starter's + # need, not ours: if a sibling's extend landed first it may + # have skipped the wire call and left less than we need, so + # hold its result to our own requirement before taking it. if additional_amount <= (inflight.requested_additional or 0.0): - return joined - # It asked for less. Go round again to re-read the slot it just - # moved, so what we ask for next is sized against the balance - # it left rather than the one we started from. + if joined is None or not self._needs_extend(joined, resolved, required_credits): + return joined + # It asked for less, or left us short. Go round again to + # re-read the slot it just moved, so what we ask for next is + # sized against the balance it left rather than the one we + # started from. joins_left -= 1 continue return await self._single_flight( diff --git a/tests/leases/test_lease_manager.py b/tests/leases/test_lease_manager.py index d570280..a9f32fa 100644 --- a/tests/leases/test_lease_manager.py +++ b/tests/leases/test_lease_manager.py @@ -321,6 +321,47 @@ async def test_a_follow_up_does_not_inherit_another_callers_smaller_follow_up( assert entry_a is not None and entry_a.local_remaining_credits >= 28_000 +async def test_a_joiner_left_short_by_a_skipped_flight_extends_for_itself( + clock: VirtualClock, monkeypatch: Any +) -> None: + # A water-mark flight re-checks the slot against its starter's need. A + # sibling's extend lands first and lifts the slot just past the water mark, + # so the flight skips the wire call. A joiner needing more than that must + # not take the result: its retry reserve would fail with the credits still + # on the server. + manager, store, wire = _make_manager(clock) + await _drawn_down_lease(store, clock) + wire.extend_responses.append({"lease": {"granted_total": 2100, "expires_at": clock() + 600}}) + + live = store.get + reads = 0 + gate = asyncio.Event() + + async def staged_get(company_id: str, credit_type_id: str) -> Optional[LeaseState]: + nonlocal reads + reads += 1 + if reads == 2: + # The flight's re-check: a sibling pod's extend lands first. + await gate.wait() + await store.extend("co_1", "ct_1", 1100, clock() + 600, "lse_1") + return await live(company_id, credit_type_id) + + monkeypatch.setattr(store, "get", staged_get) + + watermark = asyncio.ensure_future(manager.maybe_extend("co_1", "ct_1")) + await _settle() + joiner = asyncio.ensure_future(manager.maybe_extend("co_1", "ct_1", 400)) + await _settle() + + gate.set() + first = await watermark + joined = await joiner + + assert first is not None and first.local_remaining_credits == 300 + assert len(wire.extend_calls) == 1 + assert joined is not None and joined.local_remaining_credits >= 400 + + async def test_the_follow_up_never_chains(clock: VirtualClock) -> None: # A company whose balance cannot reach the request would otherwise spin: # the follow-up resolves short and the caller's retry reports insufficient From b726326ffd11b76bc4a5144f9d086a655e6b4b22 Mon Sep 17 00:00:00 2001 From: Ben Papillon Date: Mon, 28 Sep 2026 14:05:58 -0700 Subject: [PATCH 4/5] cap a check's lease joins at its start deadline --- src/schematic/leases/check.py | 8 +++-- src/schematic/leases/lease_manager.py | 36 ++++++++++++++++---- tests/leases/test_lease_manager.py | 48 +++++++++++++++++++++++++++ 3 files changed, 83 insertions(+), 9 deletions(-) diff --git a/src/schematic/leases/check.py b/src/schematic/leases/check.py index 8261375..7ec65a4 100644 --- a/src/schematic/leases/check.py +++ b/src/schematic/leases/check.py @@ -95,6 +95,10 @@ async def check_with_lease( log = deps.logger mode = options.on_acquire_failure or "fail-closed" usage = options.usage + # One budget for the whole check. Each wait on another caller's acquire or + # extend is capped against this, not against a fresh timeout, so a check + # cannot take its timeout once per step. + deadline = None if options.timeout is None else time.monotonic() + options.timeout # A malformed usage must never reach the stores, and the caller asked for a # contract for exactly this case, so resolve it through that rather than @@ -206,7 +210,7 @@ async def failure(reason: str) -> "CheckResult": # The caller's per-check timeout governs the lease wire calls, the same way # it governs the plain check's. - lease = await deps.manager.acquire_if_needed(resolved_company.id, credit_id, options.timeout) + lease = await deps.manager.acquire_if_needed(resolved_company.id, credit_id, options.timeout, deadline) if lease is None: return await failure("lease_acquire_failed") @@ -221,7 +225,7 @@ async def failure(reason: str) -> "CheckResult": if reserve is None: # Pass the cost as required_credits so a single large request # extends even while the ratio sits above the water mark. - await deps.manager.maybe_extend(resolved_company.id, credit_id, credit_cost, options.timeout) + await deps.manager.maybe_extend(resolved_company.id, credit_id, credit_cost, options.timeout, deadline) reserve = await deps.lease_store.try_reserve(resolved_company.id, credit_id, credit_cost) except Exception as err: log.error(f"Lease check: reserve against {resolved_company.id}/{credit_id} failed: {err}") diff --git a/src/schematic/leases/lease_manager.py b/src/schematic/leases/lease_manager.py index 856d710..4919b79 100644 --- a/src/schematic/leases/lease_manager.py +++ b/src/schematic/leases/lease_manager.py @@ -189,13 +189,18 @@ def sweep_interval(self) -> float: return self._config.sweep_interval or DEFAULT_SWEEP_INTERVAL async def acquire_if_needed( - self, company_id: str, credit_type_id: str, timeout: Optional[float] = None + self, + company_id: str, + credit_type_id: str, + timeout: Optional[float] = None, + deadline: Optional[float] = None, ) -> Optional[LeaseState]: """The slot's live lease, acquiring one over the wire if none is live. ``timeout`` governs the wire call this caller starts. A caller that - joins an in-flight acquire rides the first caller's timeout, since - there is one shared call to time out. + joins an in-flight acquire waits on the first caller's call, but only + until ``deadline`` (a ``time.monotonic()`` instant), and then resolves + to no lease while the call runs on for everybody else. """ try: existing = await self._lease_store.get(company_id, credit_type_id) @@ -214,7 +219,15 @@ async def acquire_if_needed( key = lease_key(company_id, credit_type_id) inflight = self._inflight_acquire.get(key) if inflight is not None: - return await asyncio.shield(inflight.task) + joined = await self._join_within(inflight.task, deadline) + if joined is _JOIN_TIMED_OUT: + logger.debug( + "Acquire in flight for %s/%s outlasted the caller's deadline; not waiting on it", + company_id, + credit_type_id, + ) + return None + return joined return await self._single_flight( self._inflight_acquire, key, self._acquire(company_id, credit_type_id, timeout) ) @@ -269,6 +282,7 @@ async def maybe_extend( credit_type_id: str, required_credits: Optional[float] = None, timeout: Optional[float] = None, + deadline: Optional[float] = None, ) -> Optional[LeaseState]: """Extend the slot's lease when the local view warrants it. @@ -283,8 +297,11 @@ async def maybe_extend( and fail its post-extend retry with credits still sitting on the server. A flight it finds on the way back is only joined if that one covers the shortfall too; a smaller one is waited out, never inherited. + + Waits on other callers' flights end at ``deadline`` (a + ``time.monotonic()`` instant), or ``timeout`` from now without one. """ - return await self._maybe_extend(company_id, credit_type_id, required_credits, timeout) + return await self._maybe_extend(company_id, credit_type_id, required_credits, timeout, deadline) async def _maybe_extend( self, @@ -292,12 +309,17 @@ async def _maybe_extend( credit_type_id: str, required_credits: Optional[float], timeout: Optional[float], + deadline: Optional[float] = None, ) -> Optional[LeaseState]: # A joiner waits on someone else's wire call, which runs on whatever # timeout ITS caller set (a background refresh uses the client # default). So the wait is capped at this caller's own timeout: a check - # with 200ms to spend must not sit behind a 30s extend. - join_deadline = None if timeout is None else time.monotonic() + timeout + # with 200ms to spend must not sit behind a 30s extend. A check passes + # the deadline it set when it started, so the time it already spent + # acquiring and reserving comes out of the same 200ms. + join_deadline = deadline + if join_deadline is None and timeout is not None: + join_deadline = time.monotonic() + timeout # Joins are budgeted, extends of our own are not: a caller may wait out # flights that ask for too little, but once the budget runs out it # issues its own single extend rather than joining again. Without the diff --git a/tests/leases/test_lease_manager.py b/tests/leases/test_lease_manager.py index a9f32fa..ddad932 100644 --- a/tests/leases/test_lease_manager.py +++ b/tests/leases/test_lease_manager.py @@ -417,6 +417,54 @@ async def test_a_joiner_gives_up_on_a_flight_that_outlasts_its_own_timeout( assert entry is not None and entry.granted_amount == 2000 +async def test_a_joiner_gives_up_on_an_acquire_at_its_deadline(clock: VirtualClock) -> None: + manager, _store, wire = _make_manager(clock) + gate = asyncio.Event() + original = wire.acquire + + async def slow_acquire(*args: Any, **kwargs: Any) -> LeaseGrant: + await gate.wait() + return await original(*args, **kwargs) + + wire.acquire = slow_acquire # type: ignore[method-assign] + wire.acquire_responses.append(_lease(clock)) + + flight = asyncio.ensure_future(manager.acquire_if_needed("co_1", "ct_1")) + await _settle() + + started = time.monotonic() + impatient = await manager.acquire_if_needed("co_1", "ct_1", deadline=started + 0.05) + assert impatient is None + assert time.monotonic() - started < 1 + + gate.set() + entry = await flight + assert entry is not None and entry.lease_id == "lse_1" + # The joiner abandoned its wait; it never raced an acquire of its own. + assert len(wire.acquire_calls) == 1 + + +async def test_a_joiner_waits_out_its_deadline_not_a_fresh_timeout(clock: VirtualClock) -> None: + # A check that already spent most of its budget acquiring and reserving + # must not get its whole timeout again for the extend join. + manager, store, wire = _make_manager(clock) + await _drawn_down_lease(store, clock) + arrived, release = wire.hold_extend() + wire.extend_responses.append({"lease": {"granted_total": 2000, "expires_at": clock() + 600}}) + + flight = asyncio.ensure_future(manager.maybe_extend("co_1", "ct_1")) + await arrived.wait() + + started = time.monotonic() + impatient = await manager.maybe_extend("co_1", "ct_1", 900, 30, started + 0.05) + assert impatient is None + assert time.monotonic() - started < 1 + assert len(wire.extend_calls) == 1 + + release.set() + await flight + + async def test_the_flight_cleanup_leaves_a_follow_up_registered(clock: VirtualClock) -> None: # A follow-up registers under the key of the flight it waited out, so that # flight's cleanup has to check identity before dropping the entry. From 7574cbd52c4f3537bcea49bd42b2af013faff3c1 Mon Sep 17 00:00:00 2001 From: Ben Papillon Date: Mon, 28 Sep 2026 14:30:49 -0700 Subject: [PATCH 5/5] check a joined extend against the requirement only --- src/schematic/leases/lease_manager.py | 9 ++++++++- tests/leases/test_lease_manager.py | 21 +++++++++++++++++++++ 2 files changed, 29 insertions(+), 1 deletion(-) diff --git a/src/schematic/leases/lease_manager.py b/src/schematic/leases/lease_manager.py index 4919b79..a06105b 100644 --- a/src/schematic/leases/lease_manager.py +++ b/src/schematic/leases/lease_manager.py @@ -367,8 +367,15 @@ async def _maybe_extend( # need, not ours: if a sibling's extend landed first it may # have skipped the wire call and left less than we need, so # hold its result to our own requirement before taking it. + # Only the requirement: a server that granted less than asked + # can leave the slot under the water mark, and re-extending + # for that would turn every water-mark joiner into a wire call. if additional_amount <= (inflight.requested_additional or 0.0): - if joined is None or not self._needs_extend(joined, resolved, required_credits): + if ( + joined is None + or required_credits is None + or joined.local_remaining_credits >= required_credits + ): return joined # It asked for less, or left us short. Go round again to # re-read the slot it just moved, so what we ask for next is diff --git a/tests/leases/test_lease_manager.py b/tests/leases/test_lease_manager.py index ddad932..dd96470 100644 --- a/tests/leases/test_lease_manager.py +++ b/tests/leases/test_lease_manager.py @@ -274,6 +274,27 @@ async def test_two_watermark_joiners_share_the_one_wire_call(clock: VirtualClock assert [entry.local_remaining_credits for entry in results if entry] == [1200] * 3 +async def test_water_mark_joiners_take_an_under_granted_extend(clock: VirtualClock) -> None: + # The server grants less than asked and the slot stays under the water + # mark. That is the server's answer for this round: the joiners must take + # it rather than each sending an extend of its own. + manager, store, wire = _make_manager(clock) + await _drawn_down_lease(store, clock) + arrived, release = wire.hold_extend() + wire.extend_responses.append({"lease": {"granted_total": 1050, "expires_at": clock() + 600}}) + + first = asyncio.ensure_future(manager.maybe_extend("co_1", "ct_1")) + await arrived.wait() + joiners = [asyncio.ensure_future(manager.maybe_extend("co_1", "ct_1")) for _ in range(3)] + await _settle() + + release.set() + results = await asyncio.gather(first, *joiners) + + assert len(wire.extend_calls) == 1 + assert [entry.local_remaining_credits for entry in results if entry] == [250] * 4 + + async def test_a_follow_up_does_not_inherit_another_callers_smaller_follow_up( clock: VirtualClock, ) -> None: