diff --git a/docs/architecture.md b/docs/architecture.md index 93503fd..fe2c5bf 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -92,7 +92,8 @@ keys in `hosts.schemas.RESERVED_HOST_ENV_KEYS` are rejected. Two maintenance commands run as cron jobs from the same image: `hosts.janitor` reaps expired and orphaned hosts, `hosts.pool` keeps a -warm pool of pre-provisioned hosts (`POOL_SIZE`) to hide provider cold +warm pool of pre-provisioned hosts per provider (`POOL_SIZES`, with +`POOL_SIZE` as the default provider's target) to hide provider cold starts. ## Diagnostics diff --git a/docs/deploy.md b/docs/deploy.md index f9e8d9b..824a3a8 100644 --- a/docs/deploy.md +++ b/docs/deploy.md @@ -25,7 +25,8 @@ docker run --rm --env-file drukbox.env "$IMAGE" .venv/bin/python -m hosts.pool ``` The janitor reaps expired and orphaned hosts. The pool maintainer -pre-provisions warm hosts and only does anything when `POOL_SIZE > 0`. +pre-provisions warm hosts per provider and only does anything when at +least one provider has a warm target (`POOL_SIZES` / `POOL_SIZE`). Schedule both under your cron infrastructure (k8s `CronJob`, systemd timer) from the same image and env file. @@ -139,9 +140,10 @@ Core, optional: | `UVICORN_HOST` | `0.0.0.0` | API bind address. Set `127.0.0.1` to restrict to loopback. | | `PROVISIONING_GRACE_SECONDS` | `600` | Safety TTL on in-flight hosts so the janitor reaps row + VM if the client disconnects mid-provision. Must exceed the worst-case provision duration. | | `IDEMPOTENCY_KEY_TTL_HOURS` | `24` | Retention period for successful `Idempotency-Key` mappings. | -| `POOL_SIZE` | `0` | Warm hosts to keep ready. `0` disables pooling. | +| `POOL_SIZES` | `{}` | Warm hosts to keep ready per provider, as JSON (e.g. `{"exe": 2, "hetzner": 1}`). Overrides `POOL_SIZE` for the providers it names. | +| `POOL_SIZE` | `0` | Warm hosts to keep ready for the default provider. `0` disables its pool. | | `POOL_HOST_MAX_AGE_HOURS` | `4` | Max age before the janitor reaps an unclaimed pool host. | -| `POOL_MAX_CREATES_PER_TICK` | `2` | Upper bound on pool provisions per tick; caps over-provision blast radius when ticks overlap. | +| `POOL_MAX_CREATES_PER_TICK` | `2` | Upper bound on pool provisions per tick, across all providers; caps over-provision blast radius when ticks overlap. | Tailscale (required when `TAILSCALE_ENABLED=true`): diff --git a/src/core/settings.py b/src/core/settings.py index 1e57510..072b1f9 100644 --- a/src/core/settings.py +++ b/src/core/settings.py @@ -70,11 +70,22 @@ class Settings(BaseSettings): validation_alias="PROVISIONING_GRACE_SECONDS", description="Safety TTL on the host row while provisioning is in flight.", ) + pool_sizes: dict[str, Annotated[int, Field(ge=0)]] = Field( + default_factory=dict, + validation_alias="POOL_SIZES", + description=( + 'Pre-warmed host targets per provider, as JSON (e.g. {"exe": 2, "hetzner": 1}). ' + "Overrides POOL_SIZE for the providers it names." + ), + ) pool_size: int = Field( default=0, ge=0, validation_alias="POOL_SIZE", - description="Number of pre-warmed hosts to keep ready. 0 disables pooling.", + description=( + "Number of pre-warmed hosts to keep ready for the default provider. " + "0 disables its pool. A POOL_SIZES entry for that provider wins." + ), ) pool_host_max_age_hours: int = Field( default=4, @@ -86,9 +97,15 @@ class Settings(BaseSettings): default=2, ge=0, validation_alias="POOL_MAX_CREATES_PER_TICK", - description="Upper bound on pool-maintainer provisions per tick.", + description="Upper bound on pool-maintainer provisions per tick, across all providers.", ) + def get_pool_targets(self) -> dict[str, int]: + # POOL_SIZE seeds the default provider's target and POOL_SIZES + # overrides per provider; providers at zero drop out entirely. + targets = {self.default_host_provider: self.pool_size, **self.pool_sizes} + return {provider: target for provider, target in targets.items() if target > 0} + @lru_cache def get_settings() -> Settings: diff --git a/src/core/tests/test_settings.py b/src/core/tests/test_settings.py index de3e947..e404508 100644 --- a/src/core/tests/test_settings.py +++ b/src/core/tests/test_settings.py @@ -100,6 +100,49 @@ def test_numeric_settings_reject_negative_values(monkeypatch: pytest.MonkeyPatch _settings_with(monkeypatch, env) +def test_pool_size_seeds_the_default_providers_target(monkeypatch: pytest.MonkeyPatch) -> None: + env: dict[str, str | None] = { + **_base_env(), + "DEFAULT_HOST_PROVIDER": None, + "POOL_SIZE": "2", + "POOL_SIZES": None, + } + settings = _settings_with(monkeypatch, env) + assert settings.get_pool_targets() == {"exe": 2} + + +def test_pool_sizes_overrides_the_alias_for_the_same_provider( + monkeypatch: pytest.MonkeyPatch, +) -> None: + env: dict[str, str | None] = { + **_base_env(), + "DEFAULT_HOST_PROVIDER": None, + "POOL_SIZE": "5", + "POOL_SIZES": '{"exe": 2, "hetzner": 1}', + } + settings = _settings_with(monkeypatch, env) + assert settings.get_pool_targets() == {"exe": 2, "hetzner": 1} + + +def test_pool_targets_omit_zeroed_providers(monkeypatch: pytest.MonkeyPatch) -> None: + # An explicit zero in POOL_SIZES disables that provider's pool even when + # the POOL_SIZE alias would seed it. + env: dict[str, str | None] = { + **_base_env(), + "DEFAULT_HOST_PROVIDER": None, + "POOL_SIZE": "5", + "POOL_SIZES": '{"exe": 0}', + } + settings = _settings_with(monkeypatch, env) + assert settings.get_pool_targets() == {} + + +def test_pool_sizes_rejects_negative_targets(monkeypatch: pytest.MonkeyPatch) -> None: + env: dict[str, str | None] = {**_base_env(), "POOL_SIZES": '{"exe": -1}'} + with pytest.raises(ValueError, match="POOL_SIZES"): + _settings_with(monkeypatch, env) + + def test_load_test_env_overrides_ambient_values(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setenv("TAILSCALE_ENABLED", "false") conftest.load_test_env() diff --git a/src/hosts/pool.py b/src/hosts/pool.py index 112db00..b4efe45 100644 --- a/src/hosts/pool.py +++ b/src/hosts/pool.py @@ -1,5 +1,7 @@ import logging -from dataclasses import dataclass +import random +import uuid +from dataclasses import dataclass, field from datetime import timedelta from sqlalchemy import func, or_, select @@ -13,7 +15,7 @@ log = logging.getLogger(__name__) # Bail out after this many consecutive create failures (broker outage scenario) -# rather than logging the same failure pool_size times per tick. +# rather than logging the same failure once per deficit slot per tick. _MAX_CONSECUTIVE_CREATE_FAILURES = 3 @@ -21,53 +23,63 @@ class PoolMaintenanceSummary: created: int = 0 removed_excess: int = 0 - pool_size: int = 0 + targets: dict[str, int] = field(default_factory=dict) async def maintain_pool() -> PoolMaintenanceSummary: - """Top up / shrink the warm-host pool and reap failed pool members. + """Top up / shrink each provider's warm-host pool and reap failed pool members. - Each tick creates at most ``POOL_MAX_CREATES_PER_TICK`` hosts. Overlapping - ticks or accidental scheduler replicas can therefore over-provision by at - most that batch size, which the next tick sheds via the excess path. This - intentionally avoids any cross-tick locking — see CONTRIBUTING for context. + Each tick creates at most ``POOL_MAX_CREATES_PER_TICK`` hosts across all + providers. Overlapping ticks or accidental scheduler replicas can therefore + over-provision by at most that batch size, which the next tick sheds via the + excess path. This intentionally avoids any cross-tick locking — see + CONTRIBUTING for context. """ settings = get_settings() - if settings.pool_size <= 0: - return PoolMaintenanceSummary(pool_size=settings.pool_size) + targets = settings.get_pool_targets() + if not targets: + return PoolMaintenanceSummary() tailscale: Tailscale | None = Tailscale.from_settings() if settings.tailscale_enabled else None try: - return await _maintain(settings, tailscale) + return await _maintain(settings, targets, tailscale) finally: if tailscale is not None: await tailscale.aclose() -async def _maintain(settings: Settings, tailscale: Tailscale | None) -> PoolMaintenanceSummary: +async def _maintain( + settings: Settings, targets: dict[str, int], tailscale: Tailscale | None +) -> PoolMaintenanceSummary: now = utc_now() async with async_session_factory() as session: - # Count warm pool members that are still fresh and unclaimed. Only - # pool_member hosts count — demand-provisioned hosts also keep - # claimed_at NULL but must never be counted or shed. Exclude hosts whose - # expires_at has already passed — those are queued for janitor reaping. - current = await session.scalar( - select(func.count()) - .select_from(Host) + # Count warm pool members per provider that are still fresh and + # unclaimed. Only pool_member hosts count — demand-provisioned hosts + # also keep claimed_at NULL but must never be counted or shed. Exclude + # hosts whose expires_at has already passed — those are queued for + # janitor reaping. + counted = await session.execute( + select(Host.provider, func.count()) .where(Host.pool_member.is_(True)) .where(Host.claimed_at.is_(None)) .where(Host.status != HostStatus.ERROR.value) .where(or_(Host.expires_at.is_(None), Host.expires_at > now)) + .group_by(Host.provider) ) - current = current or 0 - - excess_ids: list = [] - if current > settings.pool_size: - shed = current - settings.pool_size - excess_ids = list( + current: dict[str, int] = {provider: count for provider, count in counted.all()} + + excess_ids: list[uuid.UUID] = [] + # Sweep the union of configured and present providers: one dropped from + # the targets sheds to zero instead of idling until its max-age TTL. + for provider in sorted(current.keys() | targets.keys()): + shed = current.get(provider, 0) - targets.get(provider, 0) + if shed <= 0: + continue + excess_ids.extend( ( await session.execute( select(Host.id) + .where(Host.provider == provider) .where(Host.pool_member.is_(True)) .where(Host.claimed_at.is_(None)) .where(Host.status != HostStatus.ERROR.value) @@ -78,12 +90,29 @@ async def _maintain(settings: Settings, tailscale: Tailscale | None) -> PoolMain ).scalars() ) - deficit = max(settings.pool_size - current, 0) - batch = min(deficit, settings.pool_max_creates_per_tick) + underfilled = [ + provider for provider, target in targets.items() if target > current.get(provider, 0) + ] + # Spread the per-tick budget round-robin across providers so one large + # deficit can't starve the others. The order is shuffled each tick: with + # no cross-tick state, a fixed order would starve the last provider + # whenever the cap is smaller than the number of underfilled providers — + # permanently so when an earlier provider's creates keep failing. + random.shuffle(underfilled) + deficits = {provider: targets[provider] - current.get(provider, 0) for provider in underfilled} + batch: list[str] = [] + while deficits and len(batch) < settings.pool_max_creates_per_tick: + for provider in list(deficits): + if len(batch) == settings.pool_max_creates_per_tick: + break + batch.append(provider) + deficits[provider] -= 1 + if not deficits[provider]: + del deficits[provider] created = 0 consecutive_failures = 0 - for _ in range(batch): + for provider in batch: try: async with async_session_factory() as session: service = HostService(session, settings=settings, tailscale=tailscale) @@ -92,13 +121,14 @@ async def _maintain(settings: Settings, tailscale: Tailscale | None) -> PoolMain env={}, image=None, expires_at=expires_at, + provider=provider, pool_member=True, ) created += 1 consecutive_failures = 0 except Exception: consecutive_failures += 1 - log.exception("pool: failed to create pool host") + log.exception("pool: failed to create pool host provider=%s", provider) if consecutive_failures >= _MAX_CONSECUTIVE_CREATE_FAILURES: log.error( "pool: aborting top-up after %d consecutive failures", @@ -118,16 +148,16 @@ async def _maintain(settings: Settings, tailscale: Tailscale | None) -> PoolMain if created or removed_excess: log.info( - "pool: maintained (created=%d, removed_excess=%d, target=%d)", + "pool: maintained (created=%d, removed_excess=%d, targets=%s)", created, removed_excess, - settings.pool_size, + targets, ) return PoolMaintenanceSummary( created=created, removed_excess=removed_excess, - pool_size=settings.pool_size, + targets=targets, ) diff --git a/src/hosts/service.py b/src/hosts/service.py index 8a2bdac..184331e 100644 --- a/src/hosts/service.py +++ b/src/hosts/service.py @@ -96,10 +96,14 @@ async def get_or_create_host( return existing host: Host | None = None - # Pool hosts are warmed with the default provider, so a pinned provider - # always provisions fresh. - if self.settings.pool_size > 0 and not env and image is None and not provider: - host = await self._try_claim_pool_host(expires_at=expires_at) + # Warm hosts are provider-specific, so the claim is scoped to the + # requested provider's pool. A request is pool-eligible only when it + # doesn't customize the host: default image and no env. + requested_provider = provider or self.settings.default_host_provider + if not env and image is None and self.settings.get_pool_targets().get(requested_provider): + host = await self._try_claim_pool_host( + provider=requested_provider, expires_at=expires_at + ) if host is None: host = await self.create_host( env=env, image=image, expires_at=expires_at, provider=provider @@ -119,7 +123,9 @@ async def get_or_create_host( return winner return host - async def _try_claim_pool_host(self, *, expires_at: datetime | None) -> Host | None: + async def _try_claim_pool_host( + self, *, provider: str, expires_at: datetime | None + ) -> Host | None: # Pick a candidate, then atomically claim it with UPDATE ... WHERE # claimed_at IS NULL ... RETURNING. The WHERE predicate is the actual # race guard — concurrent claimants resolve to a single winner per @@ -130,6 +136,7 @@ async def _try_claim_pool_host(self, *, expires_at: datetime | None) -> Host | N candidate_id = ( await self.session.execute( select(Host.id) + .where(Host.provider == provider) .where(Host.pool_member.is_(True)) .where(Host.claimed_at.is_(None)) .where(Host.status == HostStatus.ACTIVE.value) diff --git a/src/hosts/tests/test_api.py b/src/hosts/tests/test_api.py index 04496f7..7d90dc1 100644 --- a/src/hosts/tests/test_api.py +++ b/src/hosts/tests/test_api.py @@ -385,7 +385,7 @@ async def test_idempotency_loser_that_claimed_pool_host_is_returned_to_pool(monk async with async_session_factory() as session: service = HostService(session) - claimed = await service._try_claim_pool_host(expires_at=None) + claimed = await service._try_claim_pool_host(provider="exe", expires_at=None) assert claimed is not None assert claimed.claimed_at is not None await service._release_idempotency_loser(claimed) diff --git a/src/hosts/tests/test_pool.py b/src/hosts/tests/test_pool.py index 1085802..56387e0 100644 --- a/src/hosts/tests/test_pool.py +++ b/src/hosts/tests/test_pool.py @@ -3,6 +3,7 @@ from unittest.mock import AsyncMock import pytest +from sqlalchemy import func, select from uuid6 import uuid7 from core.database import async_session_factory @@ -15,8 +16,9 @@ @pytest.fixture def pooled_settings(monkeypatch): - """Force the pool to be enabled for this test.""" + """Force the pool to be enabled for this test, via the POOL_SIZE alias.""" monkeypatch.setenv("POOL_SIZE", "2") + monkeypatch.delenv("POOL_SIZES", raising=False) monkeypatch.setenv("POOL_HOST_MAX_AGE_HOURS", "4") # Raise the per-tick cap so tests can observe full deficit refills in a # single maintain_pool() call. Production defaults to a small cap so that @@ -27,9 +29,33 @@ def pooled_settings(monkeypatch): get_settings.cache_clear() +@pytest.fixture +def multi_pool_settings(monkeypatch): + """Warm targets for two providers ("exe" and the conftest "stub").""" + monkeypatch.setenv("POOL_SIZES", '{"exe": 2, "stub": 1}') + monkeypatch.delenv("POOL_SIZE", raising=False) + monkeypatch.setenv("POOL_HOST_MAX_AGE_HOURS", "4") + monkeypatch.setenv("POOL_MAX_CREATES_PER_TICK", "100") + get_settings.cache_clear() + yield get_settings() + get_settings.cache_clear() + + +async def _pool_counts_by_provider() -> dict[str, int]: + async with async_session_factory() as session: + counted = await session.execute( + select(Host.provider, func.count()) + .where(Host.pool_member.is_(True)) + .where(Host.claimed_at.is_(None)) + .group_by(Host.provider) + ) + return {provider: count for provider, count in counted.all()} + + async def _seed_pool_host( *, name: str, + provider: str = "exe", status: str = HostStatus.ACTIVE.value, claimed_at: datetime | None = None, expires_at: datetime | None = None, @@ -41,7 +67,7 @@ async def _seed_pool_host( id=uuid7(), name=name, status=status, - provider="exe", + provider=provider, image=ExeSettings().default_image, # pyright: ignore[reportCallIssue] env=env or {}, internal_ssh_host=f"{name}.example.ts.net", @@ -127,11 +153,12 @@ async def test_create_host_skips_pool_when_env_present(pooled_settings, monkeypa mocked_provision.assert_awaited_once() -async def test_create_host_skips_pool_when_provider_pinned( +async def test_create_host_skips_pool_when_provider_has_no_target( pooled_settings, monkeypatch, stub_provider ): - # Pool hosts are warmed with the default provider, so a pinned provider - # must provision fresh rather than hand back a default-provider member. + # POOL_SIZE targets only the default provider, so a request pinned to an + # untargeted provider must provision fresh rather than hand back a + # default-provider member. await _seed_pool_host(name="lb-pool-1") mocked_provision = AsyncMock() monkeypatch.setattr("hosts.service.HostService.provision", mocked_provision) @@ -145,6 +172,42 @@ async def test_create_host_skips_pool_when_provider_pinned( mocked_provision.assert_awaited_once() +async def test_create_host_claims_pool_member_of_the_pinned_provider( + multi_pool_settings, monkeypatch, stub_provider +): + await _seed_pool_host(name="lb-pool-exe") + stub_host = await _seed_pool_host(name="lb-pool-stub", provider="stub") + mocked_provision = AsyncMock() + monkeypatch.setattr("hosts.service.HostService.provision", mocked_provision) + + async with async_session_factory() as session: + service = HostService(session, settings=multi_pool_settings) + result = await service.get_or_create_host(env={}, image=None, provider="stub") + + assert result.id == stub_host.id + assert result.claimed_at is not None + mocked_provision.assert_not_awaited() + + +async def test_create_host_never_claims_another_providers_pool_member( + multi_pool_settings, monkeypatch, stub_provider +): + # Only an exe host is warm; a stub request has a pool target but must + # provision fresh rather than cross providers. + exe_host = await _seed_pool_host(name="lb-pool-exe") + mocked_provision = AsyncMock() + monkeypatch.setattr("hosts.service.HostService.provision", mocked_provision) + + async with async_session_factory() as session: + service = HostService(session, settings=multi_pool_settings) + result = await service.get_or_create_host(env={}, image=None, provider="stub") + + assert result.id != exe_host.id + assert result.provider == "stub" + assert result.claimed_at is None + mocked_provision.assert_awaited_once() + + async def test_create_host_skips_pool_when_image_override(pooled_settings, monkeypatch): await _seed_pool_host(name="lb-pool-1") mocked_provision = AsyncMock() @@ -244,13 +307,120 @@ async def test_maintain_pool_sheds_excess_when_pool_size_shrinks(pooled_settings assert summary.removed_excess == 3 assert summary.created == 0 async with async_session_factory() as session: - from sqlalchemy import select as _select - - result = await session.execute(_select(Host).where(Host.claimed_at.is_(None))) + result = await session.execute(select(Host).where(Host.claimed_at.is_(None))) remaining = list(result.scalars()) assert len(remaining) == 2 +async def test_maintain_pool_tops_up_each_provider_toward_its_target( + multi_pool_settings, monkeypatch, stub_provider +): + monkeypatch.setattr("hosts.service.HostService.provision", AsyncMock()) + + summary = await maintain_pool() + + assert summary.created == 3 + assert await _pool_counts_by_provider() == {"exe": 2, "stub": 1} + + +async def test_maintain_pool_sheds_excess_per_provider(monkeypatch, stub_provider): + # exe is over target while stub is exactly on target; only exe sheds, and + # only down to its own target. + monkeypatch.setenv("POOL_SIZES", '{"exe": 1, "stub": 2}') + monkeypatch.delenv("POOL_SIZE", raising=False) + get_settings.cache_clear() + try: + monkeypatch.setattr("networking.tailscale.Tailscale.release_device", AsyncMock()) + monkeypatch.setattr("providers.exe.provider.ExeProvider.delete_vm", AsyncMock()) + for i in range(3): + await _seed_pool_host(name=f"lb-exe-{i}") + for i in range(2): + await _seed_pool_host(name=f"lb-stub-{i}", provider="stub") + + summary = await maintain_pool() + + assert summary.removed_excess == 2 + assert summary.created == 0 + assert await _pool_counts_by_provider() == {"exe": 1, "stub": 2} + finally: + get_settings.cache_clear() + + +async def test_maintain_pool_sheds_provider_dropped_from_targets(monkeypatch, stub_provider): + # stub carries no target at all: its warm hosts shed to zero instead of + # idling until the max-age TTL, while exe tops up toward its own target. + monkeypatch.setenv("POOL_SIZES", '{"exe": 1}') + monkeypatch.delenv("POOL_SIZE", raising=False) + get_settings.cache_clear() + try: + monkeypatch.setattr("hosts.service.HostService.provision", AsyncMock()) + monkeypatch.setattr("networking.tailscale.Tailscale.release_device", AsyncMock()) + stub_host = await _seed_pool_host(name="lb-stub-0", provider="stub") + + summary = await maintain_pool() + + assert summary.removed_excess == 1 + assert summary.created == 1 + assert stub_provider.deleted == [stub_host.name] + assert await _pool_counts_by_provider() == {"exe": 1} + finally: + get_settings.cache_clear() + + +async def test_maintain_pool_caps_total_creates_across_providers(monkeypatch, stub_provider): + monkeypatch.setenv("POOL_SIZES", '{"exe": 5, "stub": 5}') + monkeypatch.delenv("POOL_SIZE", raising=False) + monkeypatch.setenv("POOL_MAX_CREATES_PER_TICK", "3") + get_settings.cache_clear() + try: + provision = AsyncMock() + monkeypatch.setattr("hosts.service.HostService.provision", provision) + + summary = await maintain_pool() + + assert summary.created == 3 + assert provision.await_count == 3 + # The budget round-robins across providers instead of filling one + # first: two providers under a cap of 3 split 2/1, in an order that + # is shuffled per tick. + counts = await _pool_counts_by_provider() + assert set(counts) == {"exe", "stub"} + assert sorted(counts.values()) == [1, 2] + finally: + get_settings.cache_clear() + + +async def test_maintain_pool_reaches_other_providers_when_one_keeps_failing( + monkeypatch, stub_provider +): + # With a per-tick cap of 1 and exe's creates permanently failing, a fixed + # top-up order would retry exe every tick and never warm stub. The + # per-tick shuffle must reach stub across ticks (the odds of 50 shuffles + # all leading with exe are 2**-50). + monkeypatch.setenv("POOL_SIZES", '{"exe": 2, "stub": 1}') + monkeypatch.delenv("POOL_SIZE", raising=False) + monkeypatch.setenv("POOL_MAX_CREATES_PER_TICK", "1") + get_settings.cache_clear() + try: + monkeypatch.setattr("hosts.service.HostService.provision", AsyncMock()) + real_create_host = HostService.create_host + + async def create_host_with_exe_down(self, **kwargs): + if kwargs.get("provider") == "exe": + raise RuntimeError("exe provider down") + return await real_create_host(self, **kwargs) + + monkeypatch.setattr("hosts.service.HostService.create_host", create_host_with_exe_down) + + for _ in range(50): + await maintain_pool() + if await _pool_counts_by_provider() == {"stub": 1}: + break + assert await _pool_counts_by_provider() == {"stub": 1} + finally: + get_settings.cache_clear() + + async def test_maintain_pool_excludes_expired_hosts_from_count(pooled_settings, monkeypatch): # An unclaimed host past its expires_at is queued for janitor reaping and # must not count toward the pool target. Maintainer should top up to N.