Skip to content

just_dna_enricher.net

just_dna_enricher.net

Shared HTTP-politeness primitives for the network tier.

Extracted from gnomad.py when a second and third rate-limited service arrived (NCBI eutils, the PMC ID converter, Europe PMC, OLS4, HGNC). Nothing here knows about any particular API — it is the pacing, batching and ordering discipline every client in this package has to obey, in one place so the rule cannot drift between them.

Two of these look trivial and are not:

  • PacingGate takes its clock and its sleep as parameters, so a test can prove a six-second interval is honoured without a suite that really sleeps six seconds per request.
  • dedupe preserves first-occurrence order rather than going through a set, because the order requests are made in decides the order rows are emitted in, and emitted order is part of artifact.digest (Principle 7).

PacingGate dataclass

PacingGate(
    interval: float,
    clock: Callable[[], float] = time.monotonic,
    sleeper: Callable[[float], None] = time.sleep,
    last: float | None = None,
    spent: int = 0,
    _lock: Lock = threading.Lock(),
)

Enforce a minimum interval between requests, on an injectable clock.

The clock and the sleep are parameters rather than direct time calls so a test can prove the gate honours the interval without a suite that really sleeps six seconds per request. Monotonic by default, so a wall-clock adjustment mid-run cannot collapse the interval to zero.

One gate is safe to share across threads, and it had to become so (S15). LookupClients' docstring tells callers to hold a client and reuse it — precisely because a fresh one per question would discard this state — so a server running its blocking work through a thread pool shares one gate by following our own advice. The unsynchronized version read last, slept, then wrote it, so two threads could both find the interval elapsed, both skip the sleep, and turn a published 3/s budget into 6/s. That budget is a courtesy someone else enforces by blocking the operator's IP, so "single-threaded callers only" was not a contract worth keeping unstated or worth keeping.

The lock covers the bookkeeping, not the sleep. Each caller reserves the next free slot and then waits for it alone, so N callers get N slots spaced interval apart rather than serializing behind one lock held across a sleep — same guarantee, and no thread is blocked by another's wait. Behaviour on a single thread is unchanged.

spent is what the gate admitted, and the one number a host cannot otherwise get (S95). A proxy metering egress per upstream had to charge by the shape of a request — an upper bound, because nothing downstream reported the calls actually made. Every egressing client waits on its gate once per attempt, inside its retry loop (gnomad._post, eutils._request), so one increment is one upstream attempt: a 429 retried three times counts three, and a snapshot hit that never reached the gate counts nothing. Monotonic, bumped under the same lock as the slot, and never reset — a reader that wants a rate takes two readings.

attempt_floor

attempt_floor(default: int)

Bases: stop_base

stop_after_attempt(default), resolved per call so a deployment can raise it (RM42).

Drop-in for the stop_after_attempt(n) it replaces, and deliberately only for a bare one: a composed policy (stop_after_attempt(3) | stop_after_delay(60)) means both, and raising one term silently changes something whose author meant the conjunction. None of this tier's policies is composed today; the rule matters the day one is.

Deliberately no count in this prose. It read "the nine policies" in two places while the tree carried twelve, and nobody noticed because the guard was a floor (len(found) >= 9) walking seven of the nine modules that own one -- so three new policies and two whole unwalked modules were both invisible. @registry-completeness, the same shape as RM96's _ALL_MODELS hole: a number in prose is a registry nothing iterates. test_gated_snapshots.py now discovers the modules instead of listing them and asserts an equality over what it walked, which is a claim that cannot go stale.

Source code in enricher/src/just_dna_enricher/net.py
def __init__(self, default: int) -> None:
    self.default = default

StreamedFile dataclass

StreamedFile(
    path: Path,
    sha256: str,
    etag: str | None = None,
    last_modified: str | None = None,
)

What one bulk download established: where it landed, its digest, and the source's own labels.

The digest is returned rather than logged. Four of the eleven copies computed a sha256 while streaming and then only wrote it to the log, so a caller that needed to record the provenance of the bytes it had just fetched had to hash the file again (@dont-discard-computed). Two of them later grew a tuple[Path, str] return for exactly that reason, one lane at a time.

etag and last_modified are the source's, None when it sends neither. They are how a lane can later ask has this file changed without downloading it again — MANE, CIViC and PubMind record them today, and the rest get them for free rather than growing their own copy later.

batched

batched(items: list[T], size: int) -> Iterator[list[T]]

Split into batches of at most size, preserving order (so emission stays deterministic).

Source code in enricher/src/just_dna_enricher/net.py
def batched[T](items: list[T], size: int) -> Iterator[list[T]]:
    """Split into batches of at most `size`, preserving order (so emission stays deterministic)."""
    iterator = iter(items)
    while batch := list(itertools.islice(iterator, size)):
        yield batch

dedupe

dedupe(items: Iterable[T]) -> list[T]

First-occurrence-order de-duplication (Principle 7: never set iteration for emitted order).

Source code in enricher/src/just_dna_enricher/net.py
def dedupe[T](items: Iterable[T]) -> list[T]:
    """First-occurrence-order de-duplication (Principle 7: never `set` iteration for emitted order)."""
    seen: set[T] = set()
    out: list[T] = []
    for item in items:
        if item not in seen:
            seen.add(item)
            out.append(item)
    return out

retry_attempts

retry_attempts(default: int) -> int

How many attempts this client may make: its own default, raised to the configured floor.

A floor, never a flat setting. The per-client numbers are deliberate — gnomAD and eutils sit at 4 because their budgets are the tightest — so a single value that set every client would flatten tuning that was chosen on purpose, while one that raises preserves it. Below the default it is a no-op rather than a way to make a client give up sooner; there is no deployment that wants less persistence than an author at a terminal, and allowing it would turn one variable into a footgun.

Why a knob exists at all: three attempts is right for the audience the CLI was written for — a person who would rather see a failure in ten seconds than wait out a flapping upstream. It is wrong for the other shape the 0.5 tiering created, a server running enrich() inside an unattended publish, where giving up on a transient 502 does not cost ten seconds, it costs the publisher a whole re-upload of a module the server had already accepted and validated. Two callers wanting opposite things from one constant is the definition of a knob.

Safe to raise because every gated client with a published rate budget paces before it retries: an extra attempt spends a slot of that budget rather than bursting past it. The AlphaGenome Atlas client's gate has a zero interval, because the Atlas publishes no budget and none was found up to 21 calls/s (RM307); it still counts every attempt.

Source code in enricher/src/just_dna_enricher/net.py
def retry_attempts(default: int) -> int:
    """How many attempts this client may make: its own default, **raised** to the configured floor.

    **A floor, never a flat setting.** The per-client numbers are deliberate — gnomAD and eutils sit at
    4 because their budgets are the tightest — so a single value that *set* every client would flatten
    tuning that was chosen on purpose, while one that *raises* preserves it. Below the default it is a
    no-op rather than a way to make a client give up sooner; there is no deployment that wants less
    persistence than an author at a terminal, and allowing it would turn one variable into a footgun.

    Why a knob exists at all: three attempts is right for the audience the CLI was written for — a
    person who would rather see a failure in ten seconds than wait out a flapping upstream. It is wrong
    for the other shape the 0.5 tiering created, a **server** running `enrich()` inside an unattended
    publish, where giving up on a transient 502 does not cost ten seconds, it costs the publisher a
    whole re-upload of a module the server had already accepted and validated. Two callers wanting
    opposite things from one constant is the definition of a knob.

    Safe to raise because every gated client with a published rate budget **paces before it retries**:
    an extra attempt spends a slot of that budget rather than bursting past it. The AlphaGenome Atlas
    client's gate has a zero interval, because the Atlas publishes no budget and none was found up to
    21 calls/s (RM307); it still counts every attempt.
    """
    global _env_loaded, _file_attempts
    if not _env_loaded:
        # Imported here rather than at module scope: `locations` is a leaf and `net` is a leaf, and
        # making one import the other for one call would couple them permanently. This is the guarded
        # exception the house rule allows for exactly this shape.
        from just_dna_enricher.locations import env_value

        _file_attempts = env_value(RETRY_ATTEMPTS_ENV)
        _env_loaded = True
    stated = os.environ.get(RETRY_ATTEMPTS_ENV, _file_attempts)
    raw = (stated or "").strip()
    if not raw:
        return default
    try:
        configured = int(raw)
    except ValueError:
        logger.warning(
            "%s=%r is not an integer; using this client's own %d attempt(s).",
            RETRY_ATTEMPTS_ENV,
            raw,
            default,
        )
        return default
    return max(default, configured)

stream_to_file

stream_to_file(
    dest: Path,
    url: str,
    *,
    error_cls: type[Exception],
    what: str,
    timeout: float | None = None,
    remedy: str = "",
) -> StreamedFile

Stream url to dest atomically, hashing as it goes; retry the transport, translate the rest.

The one body every download_* in this package calls. Four properties, and each is a defect that reached a user before it was one:

Atomic. The bytes go to <dest>.part and are renamed only once the stream finished, so a failed fetch leaves the directory as it found it rather than truncating a good file already there (@a-failed-fetch-is-not-a-no-op). The partial is removed on failure — a .part left behind is the one residue a re-run would have to reason about.

Retried, but only where retrying is honest. httpx.TransportError covers a connection cut mid-body — RemoteProtocolError subclasses it, which is the failure that motivated this — and a second attempt genuinely fixes it. A status error is not retried: a 404 from a mistyped release tag is the same 404 four times over, and the caller wants it now rather than after three backoffs. attempt_floor reads $JUST_DNA_HTTP_RETRY_ATTEMPTS at retry time like every other policy here, so a deployment can raise the floor without touching code (RM42).

Translated. httpx's exceptions do not leave this function. A caller of this package may not be made to depend on the transport library's exception tree to know that a fetch failed, and the lane adapters catch their own builder's type — so a leak is not merely untidy, it is a lane that cannot report built=False (@client-exception-contract).

Restarted from zero on each attempt. The hasher and the output file are created inside the attempt, not outside it: a retry that appended to a partial body would produce a file whose digest is real and whose contents are nonsense, which no footer check and no raise_for_status would catch.

what names the thing being fetched for the message ("the ClinVar VCF"); remedy is an optional sentence saying what the caller can do instead, for the lanes that accept a local file.

Source code in enricher/src/just_dna_enricher/net.py
def stream_to_file(
    dest: "Path",
    url: str,
    *,
    error_cls: type[Exception],
    what: str,
    timeout: float | None = None,
    remedy: str = "",
) -> StreamedFile:
    """Stream `url` to `dest` atomically, hashing as it goes; retry the transport, translate the rest.

    The one body every `download_*` in this package calls. Four properties, and each is a defect that
    reached a user before it was one:

    **Atomic.** The bytes go to `<dest>.part` and are renamed only once the stream finished, so a
    failed fetch leaves the directory as it found it rather than truncating a good file already there
    (`@a-failed-fetch-is-not-a-no-op`). The partial is removed on failure — a `.part` left behind is
    the one residue a re-run would have to reason about.

    **Retried, but only where retrying is honest.** `httpx.TransportError` covers a connection cut
    mid-body — `RemoteProtocolError` subclasses it, which is the failure that motivated this — and a
    second attempt genuinely fixes it. A **status** error is not retried: a 404 from a mistyped
    release tag is the same 404 four times over, and the caller wants it now rather than after three
    backoffs. `attempt_floor` reads `$JUST_DNA_HTTP_RETRY_ATTEMPTS` at retry time like every other
    policy here, so a deployment can raise the floor without touching code (RM42).

    **Translated.** `httpx`'s exceptions do not leave this function. A caller of this package may not
    be made to depend on the transport library's exception tree to know that a fetch failed, and the
    lane adapters catch their own builder's type — so a leak is not merely untidy, it is a lane that
    cannot report `built=False` (`@client-exception-contract`).

    **Restarted from zero on each attempt.** The hasher and the output file are created inside the
    attempt, not outside it: a retry that appended to a partial body would produce a file whose
    digest is real and whose contents are nonsense, which no footer check and no `raise_for_status`
    would catch.

    `what` names the thing being fetched for the message ("the ClinVar VCF"); `remedy` is an optional
    sentence saying what the caller can do instead, for the lanes that accept a local file.
    """
    dest = Path(dest)
    dest.parent.mkdir(parents=True, exist_ok=True)
    tmp = dest.with_name(dest.name + ".part")

    @retry(
        stop=attempt_floor(3),
        wait=wait_exponential_jitter(initial=2.0, max=30.0),
        retry=retry_if_exception_type(httpx.TransportError),
        reraise=True,
    )
    def _attempt() -> StreamedFile:
        hasher = hashlib.sha256()
        with httpx.stream("GET", url, follow_redirects=True, timeout=timeout) as response:
            response.raise_for_status()
            etag = response.headers.get("ETag")
            last_modified = response.headers.get("Last-Modified")
            with tmp.open("wb") as handle:
                for chunk in response.iter_bytes():
                    handle.write(chunk)
                    hasher.update(chunk)
        return StreamedFile(tmp, hasher.hexdigest(), etag, last_modified)

    logger.info("Downloading %s from %s ...", what, url)
    try:
        streamed = _attempt()
    except httpx.HTTPError as exc:
        tmp.unlink(missing_ok=True)
        message = f"could not download {what} from {url}: {exc}"
        raise error_cls(f"{message}. {remedy}" if remedy else message) from exc
    except BaseException:
        # Not only the transport: a disk that fills or a mount that turns read-only raises `OSError`
        # from `handle.write`, and the partial was left behind on exactly that leg while the
        # docstring promised none. The type is the caller's (the lane adapters catch `OSError`),
        # so it is not translated — only the residue is removed.
        tmp.unlink(missing_ok=True)
        raise
    tmp.replace(dest)
    logger.info("Downloaded %s (sha256 %s)", dest, streamed.sha256)
    return dataclasses.replace(streamed, path=dest)