"""Api01: standalone client for 127.0.0.1. Generated by mimicplus (template backend).

Learned from 15 API calls (2 noise flows filtered);
mode: auth. Zero dependencies (stdlib only); no proxy or mitmweb needed at runtime.
Type-checks under mypy --strict. Spec v4; secret audit: clean.

    from api01_client import Api01
    api = Api01()
    python this_file.py endpoints | selfcheck | call <method> key=value ...

Generated by mimicplus. Use only on data you are allowed to access.
"""

# ruff: noqa
# fmt: off
from __future__ import annotations

from typing import Any, Dict, List, Optional, Tuple, TypedDict, cast
from urllib.parse import quote as _quote


def _q(s: str) -> str:
    return _quote(s, safe="")


# ---- BEGIN MIMICPLUS RUNTIME ----
import binascii as _binascii
import contextlib as _contextlib
import gzip as _gzip
import ipaddress as _ipaddress
import json as _json
import logging as _logging
import mimetypes as _mimetypes
import os as _os
import random as _random
import socket as _socket
import sys as _sys
import threading as _threading
import time as _time
import uuid as _uuid
import zlib as _zlib
from collections.abc import Callable, Iterable, Iterator, Mapping, Sequence
from dataclasses import dataclass as _dataclass
from dataclasses import field as _field
from email.utils import parsedate_to_datetime as _parsedate
from typing import Any, ClassVar
from urllib import error as _uerror
from urllib import request as _urequest
from urllib.parse import urlencode as _urlencode
from urllib.parse import urlsplit as _urlsplit
from urllib.robotparser import RobotFileParser as _RobotFileParser

MIMICPLUS_RUNTIME_VERSION = "0.4.0"
_RETRY_STATUSES = frozenset({408, 425, 429, 500, 502, 503, 504})
_IDEMPOTENT = frozenset({"GET", "HEAD", "OPTIONS", "PUT", "DELETE", "TRACE"})
_log = _logging.getLogger("mimicplus.client")

Headers = dict[str, str]
QueryItems = list[tuple[str, Any]]
JSONSchema = dict[str, Any]


def _event(level: int, name: str, **fields: Any) -> None:
    if _log.isEnabledFor(level):
        _log.log(level, name, extra={"fields": fields})


# ---------------------------------------------------------------- errors ---
class MimicError(RuntimeError):
    """Base class of every error a generated client raises."""


class TransportError(MimicError):
    """The request never produced an HTTP response (DNS, TCP, TLS, reset...)."""


class RequestTimeout(TransportError):
    """An attempt exceeded ``timeout`` or the whole call exceeded ``deadline``."""


class HTTPError(MimicError):
    """The server answered with a 4xx/5xx status."""

    def __init__(
        self, status: int, url: str, body: bytes = b"", headers: Mapping[str, str] | None = None
    ) -> None:
        self.status, self.url, self.body = status, url, body
        self.headers: Headers = dict(headers or {})
        super().__init__(f"HTTP {status} for {url.split('?')[0]}: {body[:300]!r}")


class ClientHTTPError(HTTPError):
    """4xx (the request is wrong or unauthorized)."""


class RateLimited(ClientHTTPError):
    """429 Too Many Requests after all retries."""


class ServerHTTPError(HTTPError):
    """5xx (the server failed)."""


class RobotsDisallowed(MimicError):
    """robots.txt forbids the URL for this user agent."""


class SSRFBlocked(MimicError):
    """The URL resolves to a private/loopback/link-local address that was not allowed."""


class CircuitOpenError(MimicError):
    """Too many consecutive failures for this host; calls are short-circuited for a while."""


class AuthFlowError(MimicError):
    """A learned login/refresh chain could not produce a token."""


def _http_error(status: int, url: str, body: bytes, headers: Mapping[str, str]) -> HTTPError:
    if status == 429:
        return RateLimited(status, url, body, headers)
    if 400 <= status < 500:
        return ClientHTTPError(status, url, body, headers)
    if status >= 500:
        return ServerHTTPError(status, url, body, headers)
    return HTTPError(status, url, body, headers)


# ---------------------------------------------------------------- schemas ---
def _jtype(v: Any) -> str:
    if v is None:
        return "null"
    if isinstance(v, bool):
        return "boolean"
    if isinstance(v, int):
        return "integer"
    if isinstance(v, float):
        return "number"
    if isinstance(v, list):
        return "array"
    if isinstance(v, dict):
        return "object"
    return "string"


def infer_schema(value: Any, _max_items: int = 50, _depth: int = 0) -> JSONSchema:
    """Infer a small JSON-schema subset (type/properties/required/items) from a value."""
    t = _jtype(value)
    s: JSONSchema = {"type": [t]}
    if _depth > 24:
        return s
    if t == "object":
        s["properties"] = {str(k): infer_schema(v, _max_items, _depth + 1) for k, v in value.items()}
        s["required"] = sorted(str(k) for k in value)
    elif t == "array":
        item: JSONSchema | None = None
        for v in value[:_max_items]:
            item = merge_schema(item, infer_schema(v, _max_items, _depth + 1))
        if item is not None:
            s["items"] = item
    return s


def merge_schema(a: JSONSchema | None, b: JSONSchema | None) -> JSONSchema | None:
    """Union of two inferred schemas (fields missing on either side become optional)."""
    if a is None:
        return b
    if b is None:
        return a
    types = sorted(set(a.get("type", [])) | set(b.get("type", [])))
    if "integer" in types and "number" in types:
        types.remove("integer")
    out: JSONSchema = {"type": types}
    if "properties" in a or "properties" in b:
        pa: dict[str, JSONSchema] = a.get("properties", {})
        pb: dict[str, JSONSchema] = b.get("properties", {})
        out["properties"] = {k: merge_schema(pa.get(k), pb.get(k)) for k in sorted(set(pa) | set(pb))}
        ra = set(a.get("required", [])) if "properties" in a else set(pb)
        rb = set(b.get("required", [])) if "properties" in b else set(pa)
        out["required"] = sorted(ra & rb)
    if "items" in a or "items" in b:
        out["items"] = merge_schema(a.get("items"), b.get("items"))
    return out


def diff_schema(old: JSONSchema | None, new: JSONSchema | None, path: str = "$") -> list[dict[str, Any]]:
    """How ``new`` drifted from ``old``: list of {path, kind, detail, breaking}."""
    out: list[dict[str, Any]] = []
    if old is None or new is None:
        return out

    def norm(ts: Sequence[str]) -> set[str]:
        return {"number" if t == "integer" else t for t in ts if t != "null"}

    ot, nt = norm(old.get("type", [])), norm(new.get("type", []))
    if ot and nt:  # null-only observations carry no type information
        if not (ot & nt):
            out.append(
                {
                    "path": path,
                    "kind": "type-changed",
                    "detail": f"{sorted(old['type'])} -> {sorted(new['type'])}",
                    "breaking": True,
                }
            )
            return out
        if not nt <= ot:
            out.append(
                {
                    "path": path,
                    "kind": "type-widened",
                    "detail": f"{sorted(old['type'])} -> {sorted(new['type'])}",
                    "breaking": False,
                }
            )
    if "properties" in old and "properties" in new:
        po: dict[str, JSONSchema] = old["properties"]
        pn: dict[str, JSONSchema] = new["properties"]
        req = set(old.get("required", []))
        for k in sorted(po):
            if k not in pn:
                out.append(
                    {
                        "path": f"{path}.{k}",
                        "kind": "removed",
                        "detail": "field missing from live response",
                        "breaking": k in req,
                    }
                )
            else:
                out.extend(diff_schema(po[k], pn[k], f"{path}.{k}"))
        out.extend(
            {"path": f"{path}.{k}", "kind": "added", "detail": "new field", "breaking": False}
            for k in sorted(set(pn) - set(po))
        )
    if "items" in old and "items" in new:
        out.extend(diff_schema(old["items"], new["items"], f"{path}[]"))
    return out


def get_path(value: Any, dotted: str | Sequence[str] | None) -> Any:
    """Walk a dotted path ("a.b.0.c"); integer parts index lists. Missing => None."""
    if not dotted:
        return value
    parts = list(dotted) if isinstance(dotted, (list, tuple)) else str(dotted).split(".")
    for part in parts:
        if isinstance(value, dict):
            value = value.get(part)
        elif isinstance(value, list) and str(part).lstrip("-").isdigit():
            i = int(part)
            value = value[i] if -len(value) <= i < len(value) else None
        else:
            return None
    return value


def set_path(obj: dict[str, Any], path: str | Sequence[str], value: Any) -> dict[str, Any]:
    """Set a nested field, creating dicts as needed: set_path(b, ["filters", "country"], "FR")."""
    parts = list(path) if isinstance(path, (list, tuple)) else str(path).split(".")
    cur = obj
    for part in parts[:-1]:
        nxt = cur.get(part)
        if not isinstance(nxt, dict):
            nxt = {}
            cur[part] = nxt
        cur = nxt
    cur[parts[-1]] = value
    return obj


# ------------------------------------------------------------- pagination ---
_NEXT_KEYS = (
    "next",
    "next_page",
    "nextPage",
    "next_url",
    "nextUrl",
    "nextLink",
    "next_link",
    "@odata.nextLink",
    "odata.nextLink",
    "next_href",
    "nextPageUrl",
    "next_page_url",
)
_CURSOR_KEYS = (
    "next_cursor",
    "nextCursor",
    "cursor",
    "next_page_token",
    "nextPageToken",
    "continuation",
    "continuationToken",
    "continuation_token",
    "iterationNextToken",
    "nextToken",
    "next_token",
    "scroll_id",
    "_scroll_id",
    "after",
    "endCursor",
)
_WRAPPERS = (
    "",
    "links",
    "_links",
    "paging",
    "pagination",
    "meta",
    "page_metadata",
    "pageInfo",
    "page_info",
    "metadata",
    "response_metadata",
    "cursor",
)
_TOTAL_KEYS = (
    "total",
    "totalCount",
    "total_count",
    "hitCount",
    "totalResults",
    "total_results",
    "totalNoticeCount",
    "count",
    "numFound",
    "totalElements",
    "total_hits",
)


def _link_header_next(headers: Mapping[str, str] | None) -> str | None:
    v = {k.lower(): x for k, x in (headers or {}).items()}.get("link") or ""
    for part in v.split(","):
        seg = part.split(";")
        if len(seg) > 1 and any(p.strip().replace(" ", "") in ('rel="next"', "rel=next") for p in seg[1:]):
            return seg[0].strip().strip("<>")
    return None


def find_next_link(response: Any, headers: Mapping[str, str] | None = None) -> tuple[str | None, str | None]:
    """(dotted_path, url) of a 'next page' URL in a response (or the HTTP Link header)."""
    if isinstance(response, dict):
        for w in _WRAPPERS:
            node = get_path(response, w) if w else response
            if not isinstance(node, dict):
                continue
            for k in _NEXT_KEYS:
                v = node.get(k)
                key = k
                if isinstance(v, dict):
                    key = k + (".href" if v.get("href") else ".url")
                    v = v.get("href") or v.get("url")
                if isinstance(v, str) and (v.startswith(("http://", "https://", "/")) or "?" in v):
                    return (f"{w}.{key}" if w else key), v
    link = _link_header_next(headers)
    return ("<Link>", link) if link else (None, None)


def find_cursor(response: Any) -> tuple[str | None, Any]:
    """(dotted_path, value) of an opaque next-page cursor in a response."""
    if not isinstance(response, dict):
        return None, None
    for w in _WRAPPERS:
        node = get_path(response, w) if w else response
        if not isinstance(node, dict):
            continue
        for k in _CURSOR_KEYS:
            v = node.get(k)
            if (
                isinstance(v, (str, int))
                and not isinstance(v, bool)
                and str(v)
                and not str(v).startswith(("http://", "https://"))
            ):
                return (f"{w}.{k}" if w else k), v
    return None, None


def find_total(response: Any) -> tuple[str | None, int | None]:
    if isinstance(response, dict):
        for w in ("", "meta", "page_metadata", "pagination", "paging", "data", "result"):
            node = get_path(response, w) if w else response
            if isinstance(node, dict):
                for k in _TOTAL_KEYS:
                    v = node.get(k)
                    if isinstance(v, (int, str)) and not isinstance(v, bool) and str(v).isdigit():
                        return (f"{w}.{k}" if w else k), int(v)
    return None, None


# ------------------------------------------------------------- politeness ---
class RateLimiter:
    """Minimum interval between requests, shared by every client instance per host."""

    _locks: ClassVar[dict[str, _threading.Lock]] = {}
    _last: ClassVar[dict[str, float]] = {}
    _guard: ClassVar[_threading.Lock] = _threading.Lock()

    def __init__(self, host: str, min_interval: float = 1.0) -> None:
        self.host, self.min_interval = host, float(min_interval)
        with RateLimiter._guard:
            RateLimiter._locks.setdefault(host, _threading.Lock())

    def reserve(self) -> float:
        """Claim the next slot; returns how long the caller must sleep before sending."""
        with RateLimiter._locks[self.host]:
            now = _time.monotonic()
            slot = max(now, RateLimiter._last.get(self.host, 0.0) + self.min_interval)
            RateLimiter._last[self.host] = slot
            return slot - now

    def wait(self) -> None:
        delay = self.reserve()
        if delay > 0:
            _time.sleep(delay)


class CircuitBreaker:
    """Per-host breaker: ``threshold`` consecutive failures open it for ``reset_after`` seconds,
    then one trial call is let through (half-open); success closes it, failure re-opens it.
    """

    _registry: ClassVar[dict[str, CircuitBreaker]] = {}
    _guard: ClassVar[_threading.Lock] = _threading.Lock()

    def __init__(self, host: str, threshold: int = 10, reset_after: float = 30.0) -> None:
        self.host, self.threshold, self.reset_after = host, max(1, threshold), reset_after
        self.failures, self.opened_at, self.state = 0, 0.0, "closed"
        self._lock = _threading.Lock()

    @classmethod
    def for_host(cls, host: str, threshold: int = 10, reset_after: float = 30.0) -> CircuitBreaker:
        with cls._guard:
            br = cls._registry.get(host)
            if br is None:
                br = cls._registry[host] = cls(host, threshold, reset_after)
            return br

    @classmethod
    def reset_all(cls) -> None:
        with cls._guard:
            cls._registry.clear()

    def before(self) -> None:
        with self._lock:
            if self.state == "open":
                if _time.monotonic() - self.opened_at < self.reset_after:
                    raise CircuitOpenError(
                        f"circuit open for {self.host} after {self.failures} consecutive failures; "
                        f"retry in {self.reset_after - (_time.monotonic() - self.opened_at):.0f}s"
                    )
                self.state = "half-open"
                _event(_logging.INFO, "circuit_half_open", host=self.host)

    def success(self) -> None:
        with self._lock:
            if self.state != "closed":
                _event(_logging.INFO, "circuit_closed", host=self.host)
            self.failures, self.state = 0, "closed"

    def failure(self) -> None:
        with self._lock:
            self.failures += 1
            if self.state == "half-open" or self.failures >= self.threshold:
                if self.state != "open":
                    _event(_logging.WARNING, "circuit_open", host=self.host, failures=self.failures)
                self.state, self.opened_at = "open", _time.monotonic()


# ---------------------------------------------------------------- SSRF ---
_CGNAT = _ipaddress.ip_network("100.64.0.0/10")


def is_private_address(ip: str) -> bool:
    """True for loopback, private, link-local (incl. cloud metadata), CGNAT, multicast, reserved."""
    try:
        addr = _ipaddress.ip_address(ip.split("%")[0])
    except ValueError:
        return True
    if isinstance(addr, _ipaddress.IPv6Address) and addr.ipv4_mapped is not None:
        addr = addr.ipv4_mapped
    return bool(
        addr.is_private
        or addr.is_loopback
        or addr.is_link_local
        or addr.is_multicast
        or addr.is_reserved
        or addr.is_unspecified
        or (isinstance(addr, _ipaddress.IPv4Address) and addr in _CGNAT)
    )


class SSRFGuard:
    """Refuse URLs that resolve to internal addresses unless explicitly allowed.

    ``trusted`` hosts (the ones you configured: BASE_URL/HOSTS/``base_url=``/``allow_hosts=``)
    may be private, e.g. a client learned against ``localhost``. Everything else, including
    URLs that come from *responses* (next-page links, redirects), must resolve to public
    addresses. Resolution results are cached for ``ttl`` seconds.
    """

    def __init__(self, trusted: Sequence[str] = (), allow_private: bool = False, ttl: float = 60.0) -> None:
        self.trusted = {h.lower() for h in trusted if h}
        self.allow_private, self.ttl = allow_private, ttl
        self._cache: dict[str, tuple[float, bool]] = {}

    def check(self, url: str) -> None:
        p = _urlsplit(url)
        if p.scheme not in ("http", "https"):
            raise SSRFBlocked(f"scheme {p.scheme!r} not allowed (http/https only): {url[:80]}")
        host = (p.hostname or "").lower()
        if not host:
            raise SSRFBlocked(f"no host in URL {url[:80]!r}")
        if self.allow_private or host in self.trusted or (p.netloc.lower() in self.trusted):
            return
        hit = self._cache.get(host)
        if hit is None or _time.monotonic() - hit[0] > self.ttl:
            try:
                infos = _socket.getaddrinfo(
                    host, p.port or (443 if p.scheme == "https" else 80), proto=_socket.IPPROTO_TCP
                )
                private = any(is_private_address(str(i[4][0])) for i in infos) or not infos
            except (OSError, UnicodeError):
                private = False  # unresolvable: the request itself will fail with a TransportError
            hit = (_time.monotonic(), private)
            self._cache[host] = hit
        if hit[1]:
            _event(_logging.WARNING, "ssrf_blocked", host=host)
            raise SSRFBlocked(
                f"{host} resolves to a private/internal address; pass allow_hosts=[{host!r}] or "
                f"allow_private=True (or set MIMICPLUS_ALLOW_PRIVATE=1) if this is intended"
            )


# ------------------------------------------------------------- transport ---
@_dataclass
class PreparedRequest:
    method: str
    url: str
    headers: Headers
    body: bytes | None
    timeout: float
    retryable: bool = False
    own_host: bool = True
    meta: dict[str, Any] = _field(default_factory=dict)


@_dataclass
class RawResponse:
    status: int
    headers: Headers
    body: bytes
    url: str = ""
    set_cookies: list[str] = _field(default_factory=list)  # every Set-Cookie (headers dict keeps one)


Transport = Callable[[PreparedRequest], RawResponse]


class _GuardedRedirects(_urequest.HTTPRedirectHandler):
    """Re-run the SSRF guard on every redirect hop and never forward credentials cross-host."""

    def __init__(self, guard: SSRFGuard, strip: Sequence[str]) -> None:
        super().__init__()
        self.guard, self.strip = guard, {s.lower() for s in strip}

    def redirect_request(
        self, req: _urequest.Request, fp: Any, code: int, msg: str, headers: Any, newurl: str
    ) -> _urequest.Request | None:
        self.guard.check(newurl)
        new = super().redirect_request(req, fp, code, msg, headers, newurl)
        if new is not None and _urlsplit(newurl).netloc != _urlsplit(req.full_url).netloc:
            for h in list(new.headers):
                if h.lower() in self.strip or h.lower() in ("authorization", "cookie"):
                    del new.headers[h]
        return new


class UrllibTransport:
    """Default stdlib transport (urllib) with guarded redirects."""

    def __init__(self, guard: SSRFGuard, secret_headers: Sequence[str] = ()) -> None:
        self.opener = _urequest.build_opener(_GuardedRedirects(guard, secret_headers))

    def __call__(self, req: PreparedRequest) -> RawResponse:
        r = _urequest.Request(req.url, data=req.body, headers=req.headers, method=req.method)
        try:
            with self.opener.open(r, timeout=req.timeout) as resp:
                return RawResponse(
                    resp.status,
                    dict(resp.headers),
                    resp.read(),
                    resp.geturl(),
                    resp.headers.get_all("set-cookie") or [],
                )
        except _uerror.HTTPError as e:
            ehs: Any = e.headers  # may be None for synthetic errors
            cookies = ehs.get_all("set-cookie") if ehs is not None else None
            return RawResponse(e.code, dict(e.headers or {}), e.read() or b"", req.url, cookies or [])
        except TimeoutError as e:
            raise RequestTimeout(
                f"{req.method} {req.url.split('?')[0]} timed out after {req.timeout}s"
            ) from e
        except _uerror.URLError as e:
            if isinstance(e.reason, TimeoutError):
                raise RequestTimeout(f"{req.method} {req.url.split('?')[0]} timed out") from e
            if isinstance(e.reason, MimicError):
                raise e.reason from e
            raise TransportError(f"{req.method} {req.url.split('?')[0]} failed: {e.reason}") from e
        except (ConnectionError, OSError) as e:
            raise TransportError(f"{req.method} {req.url.split('?')[0]} failed: {e}") from e


ROBOTS_POLICIES = ("deny", "allow", "warn")


class Robots:
    """robots.txt per RFC 9309: 4xx => allow all.

    5xx / network failure ("unreachable") follows a policy: "deny" (RFC default: assume
    complete disallow), "allow", or "warn" (allow and log a warning). Set it per client
    (robots_unreachable=...) or globally with $MIMICPLUS_ROBOTS_UNREACHABLE.
    """

    _cache: ClassVar[dict[str, _RobotFileParser | None]] = {}
    UNREACHABLE: ClassVar[None] = None

    @classmethod
    def fetch(cls, root: str, user_agent: str, timeout: float = 15) -> _RobotFileParser | None:
        if root in cls._cache:
            return cls._cache[root]
        rp: _RobotFileParser | None = _RobotFileParser()
        try:
            req = _urequest.Request(
                root + "/robots.txt", headers={"User-Agent": user_agent, "Accept-Encoding": "gzip"}
            )
            with _urequest.urlopen(req, timeout=timeout) as r:
                raw = _decompress(r.read(), dict(r.headers))
                assert rp is not None  # noqa: S101 - narrows the Optional for mypy
                rp.parse(raw.decode("utf-8", "replace").splitlines())
        except _uerror.HTTPError as e:
            if 400 <= e.code < 500 and rp is not None:
                rp.parse(["User-agent: *", "Allow: /"])  # RFC 9309: 4xx => no restrictions
            else:
                rp = cls.UNREACHABLE
        except (OSError, ValueError):
            rp = cls.UNREACHABLE
        cls._cache[root] = rp
        return rp

    @classmethod
    def allowed(cls, url: str, user_agent: str, timeout: float = 15, unreachable: str = "deny") -> bool:
        p = _urlsplit(url)
        rp = cls.fetch(f"{p.scheme}://{p.netloc}", user_agent, timeout)
        if rp is None:
            if unreachable == "warn":
                _log.warning(
                    "robots.txt for %s unreachable; proceeding (robots_unreachable='warn')",
                    p.netloc,
                    extra={"fields": {"host": p.netloc}},
                )
            return unreachable in ("allow", "warn")
        return rp.can_fetch(user_agent, url)


# ------------------------------------------------------------ pagination ---
class _Pager:
    """Pagination state machine shared by the sync and async drivers.

    ``plan()`` says what to call next: ``("call", kwargs)`` / ``("get", url)`` / ``("stop", None)``;
    ``feed(resp, items)`` advances the state from a page.
    """

    def __init__(
        self, rc: Mapping[str, Any], kw: Mapping[str, Any], max_pages: int, max_items: int | None
    ) -> None:
        self.rc: dict[str, Any] = dict(rc)
        self.style: str | None = self.rc.get("style")
        self.kw: dict[str, Any] = dict(kw)
        self.max_pages, self.max_items = max_pages, max_items
        self.page_size = self.kw.get(self.rc["size_param"]) if self.rc.get("size_param") else None
        if self.style in ("page", "offset") and self.rc.get("param") not in self.kw:
            self.kw[self.rc["param"]] = int(self.rc.get("start", 1 if self.style == "page" else 0))
        self.pages, self.n_items, self.seen = 0, 0, set[str]()
        self.next_url: str | None = None
        self.done = False

    def plan(self) -> tuple[str, Any]:
        if self.done or self.pages >= self.max_pages:
            return "stop", None
        if self.next_url:
            return "get", self.next_url
        return "call", dict(self.kw)

    def take(self, items: list[Any]) -> list[Any]:
        """Items of this page to yield (respects max_items; detects repeated pages)."""
        self.pages += 1
        sig = _json.dumps(items[:3], sort_keys=True, default=str)
        if not items or sig in self.seen:
            self.done = True
            return []
        self.seen.add(sig)
        if self.max_items is not None:
            items = items[: max(0, self.max_items - self.n_items)]
        self.n_items += len(items)
        if self.max_items is not None and self.n_items >= self.max_items:
            self.done = True
        return items

    def feed(
        self, resp: Any, items: list[Any], last_url: str | None, base_url: str, headers: Mapping[str, str]
    ) -> None:
        if self.done:
            return
        rc, style = self.rc, self.style
        total = get_path(resp, rc["total_path"]) if rc.get("total_path") else find_total(resp)[1]
        reached = total is not None and str(total).isdigit() and self.n_items >= int(total)
        ps = self.page_size
        short = ps is not None and bool(ps) and len(items) < int(ps)
        if style in ("page", "offset"):
            if short or reached:
                self.done = True
            else:
                step = 1 if style == "page" else len(items)
                self.kw[rc["param"]] = int(self.kw[rc["param"]]) + step
            return
        if style == "cursor":
            cur = get_path(resp, rc["cursor_path"]) if rc.get("cursor_path") else find_cursor(resp)[1]
            if cur in (None, "", False):
                self.done = True
            else:
                self.kw[rc["param"]] = cur
            return
        nxt: Any = None
        if rc.get("next_path") and rc["next_path"] != "<Link>":
            nxt = get_path(resp, rc["next_path"])
            if isinstance(nxt, dict):
                nxt = nxt.get("href") or nxt.get("url")
        if not nxt:
            nxt = find_next_link(resp, headers)[1]
        if not nxt:
            if style is None:  # auto mode: fall back to a response cursor
                cpath, cur = find_cursor(resp)
                cparam = rc.get("param") or (cpath.split(".")[-1] if cpath else None)
                if cur not in (None, "") and cparam:
                    self.style, self.rc["param"] = "cursor", cparam
                    self.kw[cparam] = cur
                    return
            self.done = True
            return
        nxt = str(nxt)
        if nxt.startswith("/"):
            p = _urlsplit(last_url or base_url)
            nxt = f"{p.scheme}://{p.netloc}{nxt}"
        elif nxt.startswith("?"):
            nxt = (last_url or base_url).split("?")[0] + nxt
        self.style = style or "next"
        self.next_url = nxt


# ---------------------------------------------------------------- client ---
# ------------------------------------------------------- cookies / CSRF ---
def parse_set_cookie(header: str) -> tuple[str, str, bool] | None:
    """``(name, value, expired)`` of one ``Set-Cookie`` header value, or None if malformed."""
    first, *attrs = header.split(";")
    name, sep, value = first.strip().partition("=")
    if not sep or not name.strip():
        return None
    expired = not value.strip()
    for a in attrs:
        k, _, v = a.strip().partition("=")
        if k.lower() == "max-age":
            with _contextlib.suppress(ValueError):
                expired = expired or int(v) <= 0
        elif k.lower() == "expires":
            with _contextlib.suppress(TypeError, ValueError, IndexError):
                expired = expired or _parsedate(v).timestamp() < _time.time()
    return name.strip(), value.strip().strip('"'), expired


def merge_cookie_header(existing: str | None, jar: Mapping[str, str]) -> str:
    """``existing`` Cookie header with ``jar`` entries added/overriding by name."""
    pairs: dict[str, str] = {}
    for part in (existing or "").split(";"):
        k, sep, v = part.strip().partition("=")
        if sep and k:
            pairs[k] = v
    pairs.update(jar)
    return "; ".join(f"{k}={v}" for k, v in pairs.items())


_CSRF_HTML = (
    r"<meta[^>]*name=[\"']{name}[\"'][^>]*content=[\"']([^\"']+)[\"']",
    r"<meta[^>]*content=[\"']([^\"']+)[\"'][^>]*name=[\"']{name}[\"']",
    r"<input[^>]*name=[\"']{name}[\"'][^>]*value=[\"']([^\"']+)[\"']",
    r"<input[^>]*value=[\"']([^\"']+)[\"'][^>]*name=[\"']{name}[\"']",
)


def csrf_from_html(html: str, name: str) -> str | None:
    """The anti-CSRF token named ``name`` in a page's ``<meta>`` / hidden ``<input>``."""
    import re as _re

    for rx in _CSRF_HTML:
        m = _re.search(rx.replace("{name}", _re.escape(name)), html, _re.I)
        if m:
            return m.group(1)
    return None


# ------------------------------------------------------ Server-Sent Events ---
def parse_sse_lines(lines: Iterable[str]) -> Iterator[dict[str, Any]]:
    """Dispatch ``text/event-stream`` lines into ``{"event", "id", "data", "retry"}`` dicts.

    ``data`` is JSON-decoded when it parses, else the raw string (WHATWG: lines joined by ``\\n``).
    """
    ev: dict[str, Any] = {}
    data: list[str] = []
    for raw in _chain_blank(lines):
        line = raw.rstrip("\r\n")
        if not line:
            if data or ev:
                text = "\n".join(data)
                try:
                    value: Any = _json.loads(text) if text else None
                except ValueError:
                    value = text
                yield {
                    "event": ev.get("event") or "message",
                    "id": ev.get("id"),
                    "retry": ev.get("retry"),
                    "data": value,
                }
            ev, data = {}, []
            continue
        if line.startswith(":"):
            continue
        field, _, value_s = line.partition(":")
        value_s = value_s[1:] if value_s.startswith(" ") else value_s
        if field == "data":
            data.append(value_s)
        elif field in ("event", "id"):
            ev[field] = value_s
        elif field == "retry" and value_s.isdigit():
            ev["retry"] = int(value_s)


def _chain_blank(lines: Iterable[str]) -> Iterator[str]:
    yield from lines
    yield ""


# ------------------------------------------------------------- WebSocket ---
_WS_GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"


class WebSocketError(MimicError):
    """WebSocket handshake or protocol failure."""


class WebSocket:
    """Minimal RFC 6455 client (stdlib sockets): text/binary frames, fragmentation, ping/pong, close.

    Use as an iterator (``for msg in ws``) or with ``recv()``/``recv_json()``; ``send``/``send_json``.
    """

    def __init__(self, sock: Any, subprotocol: str | None = None, buffered: bytes = b"") -> None:
        self.sock, self.subprotocol, self._buf, self.closed = sock, subprotocol, buffered, False

    @classmethod
    def connect(
        cls,
        url: str,
        headers: Mapping[str, str] | None = None,
        subprotocols: Sequence[str] = (),
        timeout: float = 30.0,
    ) -> WebSocket:
        import base64 as _b64
        import hashlib as _hashlib
        import ssl as _ssl

        u = _urlsplit(url)
        if u.scheme not in ("ws", "wss"):
            raise WebSocketError(f"not a ws:// or wss:// URL: {url}")
        port = u.port or (443 if u.scheme == "wss" else 80)
        sock: Any = _socket.create_connection((u.hostname or "", port), timeout=timeout)
        if u.scheme == "wss":
            sock = _ssl.create_default_context().wrap_socket(sock, server_hostname=u.hostname)
        key = _b64.b64encode(_os.urandom(16)).decode()
        target = (u.path or "/") + (f"?{u.query}" if u.query else "")
        host = u.hostname if u.port in (None, 80, 443) else f"{u.hostname}:{u.port}"
        lines = [
            f"GET {target} HTTP/1.1",
            f"Host: {host}",
            "Upgrade: websocket",
            "Connection: Upgrade",
            f"Sec-WebSocket-Key: {key}",
            "Sec-WebSocket-Version: 13",
        ]
        if subprotocols:
            lines.append(f"Sec-WebSocket-Protocol: {', '.join(subprotocols)}")
        skip = {"host", "upgrade", "connection", "content-type", "content-length", "accept-encoding"}
        lines += [f"{k}: {v}" for k, v in (headers or {}).items() if k.lower() not in skip]
        sock.sendall(("\r\n".join(lines) + "\r\n\r\n").encode())
        resp = b""
        while b"\r\n\r\n" not in resp:
            chunk = sock.recv(4096)
            if not chunk:
                raise WebSocketError(f"connection closed during handshake with {url.split('?')[0]}")
            resp += chunk
            if len(resp) > 65536:
                raise WebSocketError("handshake response too large")
        head, _, rest = resp.partition(b"\r\n\r\n")
        status_line, *hlines = head.decode("latin-1").split("\r\n")
        hdrs = {k.strip().lower(): v.strip() for k, _, v in (h.partition(":") for h in hlines)}
        parts = status_line.split(" ", 2)
        if len(parts) < 2 or parts[1] != "101":
            sock.close()
            code = int(parts[1]) if len(parts) > 1 and parts[1].isdigit() else 0
            raise WebSocketError(f"WebSocket upgrade refused: {status_line} ({code})")
        want = _b64.b64encode(_hashlib.sha1((key + _WS_GUID).encode()).digest()).decode()  # noqa: S324
        if hdrs.get("sec-websocket-accept") != want:
            sock.close()
            raise WebSocketError("bad Sec-WebSocket-Accept (not a WebSocket server?)")
        return cls(sock, hdrs.get("sec-websocket-protocol") or None, rest)

    # ---- framing -----------------------------------------------------------
    def _exact(self, n: int) -> bytes:
        while len(self._buf) < n:
            chunk = self.sock.recv(max(4096, n - len(self._buf)))
            if not chunk:
                raise WebSocketError("connection closed")
            self._buf += chunk
        out, self._buf = self._buf[:n], self._buf[n:]
        return out

    def _send_frame(self, opcode: int, payload: bytes) -> None:
        import struct as _struct

        n = len(payload)
        head = bytes([0x80 | opcode])
        if n < 126:
            head += bytes([0x80 | n])
        elif n < 65536:
            head += bytes([0x80 | 126]) + _struct.pack(">H", n)
        else:
            head += bytes([0x80 | 127]) + _struct.pack(">Q", n)
        mask = _os.urandom(4)
        self.sock.sendall(head + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload)))

    def _recv_frame(self) -> tuple[bool, int, bytes]:
        import struct as _struct

        b1, b2 = self._exact(2)
        n = b2 & 0x7F
        if n == 126:
            n = _struct.unpack(">H", self._exact(2))[0]
        elif n == 127:
            n = _struct.unpack(">Q", self._exact(8))[0]
        if n > 64 * 1024 * 1024:
            raise WebSocketError(f"frame too large ({n} bytes)")
        mask = self._exact(4) if b2 & 0x80 else b""
        data = self._exact(n)
        if mask:
            data = bytes(c ^ mask[i % 4] for i, c in enumerate(data))
        return bool(b1 & 0x80), b1 & 0x0F, data

    def send(self, data: str | bytes) -> None:
        self._send_frame(1 if isinstance(data, str) else 2, data.encode() if isinstance(data, str) else data)

    def send_json(self, value: Any) -> None:
        self.send(_json.dumps(value))

    def recv(self) -> str | bytes | None:
        """Next text (str) or binary (bytes) message; None once the connection is closed."""
        if self.closed:
            return None
        parts: list[bytes] = []
        first_op = 0
        while True:
            try:
                fin, op, data = self._recv_frame()
            except (WebSocketError, OSError):
                self.closed = True
                return None
            if op == 9:  # ping
                self._send_frame(10, data)
                continue
            if op == 10:
                continue
            if op == 8:
                self.closed = True
                with _contextlib.suppress(OSError):
                    self._send_frame(8, data[:2])
                return None
            if op in (1, 2):
                first_op = op
            parts.append(data)
            if fin:
                raw = b"".join(parts)
                return raw.decode("utf-8", "replace") if first_op == 1 else raw

    def recv_json(self) -> Any:
        msg = self.recv()
        if msg is None:
            return None
        try:
            return _json.loads(msg)
        except ValueError:
            return msg

    def __iter__(self) -> Iterator[Any]:
        while True:
            msg = self.recv_json()
            if msg is None:
                return
            yield msg

    def close(self) -> None:
        if not self.closed:
            self.closed = True
            with _contextlib.suppress(OSError):
                self._send_frame(8, b"\x03\xe8")
        with _contextlib.suppress(OSError):
            self.sock.close()

    def __enter__(self) -> WebSocket:
        return self

    def __exit__(self, *exc: object) -> None:
        self.close()


# --------------------------------------------------------------- gRPC-web ---
def pb_varint(n: int) -> bytes:
    n &= (1 << 64) - 1
    out = bytearray()
    while True:
        c, n = n & 0x7F, n >> 7
        out.append(c | (0x80 if n else 0))
        if not n:
            return bytes(out)


def pb_encode(fields: Sequence[tuple[int, Any]] | Mapping[int, Any]) -> bytes:
    """Protobuf wire encoding of ``{number: value}``: int/bool varint, float double, str/bytes, dict."""
    import struct as _struct

    items = list(fields.items()) if isinstance(fields, Mapping) else list(fields)
    out = bytearray()
    for num, v in items:
        num = int(num)
        if isinstance(v, list):
            for x in v:
                out += pb_encode([(num, x)])
        elif isinstance(v, (bool, int)):
            out += pb_varint(num << 3) + pb_varint(int(v))
        elif isinstance(v, float):
            out += pb_varint(num << 3 | 1) + _struct.pack("<d", v)
        elif v is not None:
            raw = pb_encode({int(k): x for k, x in v.items()}) if isinstance(v, dict) else v
            raw = raw.encode() if isinstance(raw, str) else bytes(raw)
            out += pb_varint(num << 3 | 2) + pb_varint(len(raw)) + raw
    return bytes(out)


def _pb_read_varint(b: bytes, i: int) -> tuple[int, int]:
    val, shift = 0, 0
    while True:
        if i >= len(b) or shift > 63:
            raise ValueError("truncated varint")
        c = b[i]
        i += 1
        val |= (c & 0x7F) << shift
        shift += 7
        if not c & 0x80:
            return val, i


def pb_decode(b: bytes, depth: int = 0) -> list[tuple[int, str, Any]]:
    """``[(number, wire, value)]`` of a protobuf message (schema-less); ValueError if malformed."""
    import base64 as _b64
    import struct as _struct

    out: list[tuple[int, str, Any]] = []
    i = 0
    while i < len(b):
        key, i = _pb_read_varint(b, i)
        num, wt = key >> 3, key & 7
        if num < 1 or wt not in (0, 1, 2, 5):
            raise ValueError(f"bad field key {key}")
        if wt == 0:
            v, i = _pb_read_varint(b, i)
            out.append((num, "varint", v))
        elif wt in (1, 5):
            size = 8 if wt == 1 else 4
            if i + size > len(b):
                raise ValueError("truncated fixed field")
            fmt = "<d" if wt == 1 else "<f"
            out.append((num, "fixed64" if wt == 1 else "fixed32", _struct.unpack(fmt, b[i : i + size])[0]))
            i += size
        else:
            n, i = _pb_read_varint(b, i)
            if i + n > len(b):
                raise ValueError("truncated length-delimited field")
            chunk = b[i : i + n]
            i += n
            value: Any
            try:
                text = chunk.decode("utf-8")
                value = text if text.isprintable() else None
            except UnicodeDecodeError:
                value = None
            if value is None and depth < 4 and chunk:
                try:
                    value = {str(k): x for k, _, x in pb_decode(chunk, depth + 1)}
                except ValueError:
                    value = None
            out.append((num, "bytes", value if value is not None else _b64.b64encode(chunk).decode()))
    return out


def grpc_web_frames(body: bytes) -> tuple[list[bytes], dict[str, str]]:
    """``(messages, trailers)`` of a gRPC-web (binary) body."""
    import struct as _struct

    msgs: list[bytes] = []
    trailers: dict[str, str] = {}
    i = 0
    while i + 5 <= len(body):
        flag, n = body[i], _struct.unpack(">I", body[i + 1 : i + 5])[0]
        chunk = body[i + 5 : i + 5 + n]
        if len(chunk) < n:
            break
        if flag & 0x80:
            for line in chunk.decode("utf-8", "replace").split("\r\n"):
                k, _, v = line.partition(":")
                if k.strip():
                    trailers[k.strip().lower()] = v.strip()
        else:
            msgs.append(chunk)
        i += 5 + n
    return msgs, trailers


class GrpcError(MimicError):
    """A gRPC call answered with a non-zero ``grpc-status``."""

    def __init__(self, msg: str, status: int = 2) -> None:
        super().__init__(msg)
        self.grpc_status = status


class Client:
    """Base class for generated clients. Subclasses set HOST/BASE_URL/... and add methods."""

    HOST: ClassVar[str | None] = None
    BASE_URL: ClassVar[str | None] = None
    MODE: ClassVar[str] = "public"  # "public" (no auth) or "auth"
    DEFAULT_HEADERS: ClassVar[dict[str, str]] = {}
    AUTH_HEADERS: ClassVar[list[str]] = []  # header *names* needed in auth mode (values loaded at runtime)
    AUTH_QUERY_PARAMS: ClassVar[list[str]] = []
    HOSTS: ClassVar[dict[str, str]] = {}  # multi-host clients: key -> base URL (auth goes only to these)
    ROBOTS_UNREACHABLE: ClassVar[str] = "deny"  # robots.txt 5xx/unreachable: "deny" | "allow" | "warn"
    PAGINATION: ClassVar[dict[str, dict[str, Any]]] = {}  # name -> pagination recipe (see paginate())
    PROBES: ClassVar[dict[str, dict[str, Any]]] = {}  # name -> recorded probe for selfcheck()
    ITEMS: ClassVar[dict[str, list[str | None]]] = {}  # name -> [items_path, id_key]
    SAFE: ClassVar[dict[str, bool]] = {}  # name -> read-only (=> retryable)
    IDEMPOTENT: ClassVar[dict[str, bool]] = {}  # name -> safe to retry although it writes
    QUERY_PARAM: ClassVar[dict[str, str]] = {}  # name -> free-text query kwarg
    AUTH_FLOW: ClassVar[dict[str, Any] | None] = None  # learned login -> token -> use chain
    SESSION: ClassVar[dict[str, Any] | None] = None  # learned cookie session + CSRF recipe
    CHANNELS: ClassVar[dict[str, dict[str, Any]]] = {}  # name -> learned SSE / WebSocket channel
    USER_AGENT: ClassVar[str] = f"mimicplus/{MIMICPLUS_RUNTIME_VERSION}"

    def __init__(
        self,
        base_url: str | None = None,
        headers: Mapping[str, str] | None = None,
        auth: Mapping[str, str] | None = None,
        user_agent: str | None = None,
        min_interval: float = 1.0,
        max_retries: int = 4,
        backoff: float = 0.5,
        max_backoff: float = 60.0,
        timeout: float = 30.0,
        respect_robots: bool | None = None,
        verbose: bool = False,
        robots_unreachable: str | None = None,
        bases: Mapping[str, str] | None = None,
        *,
        deadline: float | None = 300.0,
        transport: Transport | None = None,
        allow_private: bool | None = None,
        allow_hosts: Sequence[str] = (),
        circuit_threshold: int = 10,
        circuit_reset: float = 30.0,
        idempotency_keys: bool = False,
        cookies: Mapping[str, str] | None = None,
        use_cookies: bool = True,
    ) -> None:
        self.base_url = (base_url or self.BASE_URL or f"https://{self.HOST}").rstrip("/")
        self.bases = {k: v.rstrip("/") for k, v in dict(self.HOSTS).items()}
        self.bases.update({k: v.rstrip("/") for k, v in (bases or {}).items()})
        self.user_agent = user_agent or _os.environ.get("MIMICPLUS_USER_AGENT") or self.USER_AGENT
        self.headers: Headers = {k.lower(): v for k, v in dict(self.DEFAULT_HEADERS).items()}
        self.headers.update({k.lower(): v for k, v in (headers or {}).items()})
        self.auth: Headers = {
            k.lower(): v for k, v in (auth if auth is not None else self._auth_from_env()).items()
        }
        self.auth_params: Headers = {}
        for k in list(self.auth):
            if k.startswith("?"):
                self.auth_params[k[1:]] = self.auth.pop(k)
        self.cookies: Headers = dict(cookies or {})  # in-memory jar for the client's own hosts
        self.use_cookies = use_cookies
        self._csrf: str | None = None
        self._busy = False  # re-entrancy guard while fetching a CSRF token
        login_ep = (self.SESSION or {}).get("login_endpoint")
        if (
            self.MODE == "auth"
            and not (self.auth or self.auth_params or self.cookies)
            and not (self.AUTH_FLOW or login_ep)
        ):
            _log.warning(
                "%s was learned from authenticated traffic but no auth is loaded; pass auth={...}, "
                "use .from_har()/.from_curl(), or set %s",
                type(self).__name__,
                self._env_name(),
            )
        self.min_interval = min_interval
        self.limiter = RateLimiter(_urlsplit(self.base_url).netloc, min_interval)
        self._limiters = {self.limiter.host: self.limiter}
        pol = robots_unreachable or _os.environ.get("MIMICPLUS_ROBOTS_UNREACHABLE") or self.ROBOTS_UNREACHABLE
        if pol not in ROBOTS_POLICIES:
            raise ValueError(f"robots_unreachable must be one of {ROBOTS_POLICIES}, not {pol!r}")
        self.robots_unreachable = pol
        self.max_retries, self.backoff, self.max_backoff = max_retries, backoff, max_backoff
        self.timeout, self.deadline = timeout, deadline
        self.respect_robots = (self.MODE == "public") if respect_robots is None else respect_robots
        self.verbose = verbose
        if allow_private is None:
            allow_private = _os.environ.get("MIMICPLUS_ALLOW_PRIVATE", "") in ("1", "true", "yes")
        trusted = [
            _urlsplit(self.base_url).hostname or "",
            *(_urlsplit(b).hostname or "" for b in self.bases.values()),
            *allow_hosts,
        ]
        self.guard = SSRFGuard(trusted, allow_private=allow_private)
        self.transport: Transport = transport or UrllibTransport(self.guard, self.AUTH_HEADERS)
        self.circuit_threshold, self.circuit_reset = circuit_threshold, circuit_reset
        self.idempotency_keys = idempotency_keys
        self.last_status: int | None = None
        self.last_headers: Headers = {}
        self.last_url: str | None = None
        self._token_expiry: float | None = None
        self._refresh_token: str | None = None
        self._refreshing = False

    # ---- auth loading (never stored in source) ----------------------------
    @classmethod
    def _env_name(cls) -> str:
        return f"{cls.__name__.upper()}_AUTH"

    @classmethod
    def _auth_from_env(cls) -> dict[str, str]:
        """JSON {header: value} from $<CLASS>_AUTH (inline JSON or a path) or $MIMICPLUS_AUTH."""
        raw = _os.environ.get(cls._env_name()) or _os.environ.get("MIMICPLUS_AUTH") or ""
        if not raw:
            return {}
        if _os.path.exists(raw):
            with open(raw, encoding="utf-8") as f:
                raw = f.read()
        data = _json.loads(raw)
        if not isinstance(data, dict):
            raise ValueError(f"${cls._env_name()} must be a JSON object of header -> value")
        return {str(k): str(v) for k, v in data.items()}

    @classmethod
    def from_har(cls, path: str, **kw: Any) -> Any:
        """Load auth headers for HOST from the newest matching request in a HAR export."""
        with open(path, encoding="utf-8") as f:
            entries = _json.load(f).get("log", {}).get("entries", [])
        for e in reversed(entries):
            req = e.get("request", {})
            if _urlsplit(req.get("url", "")).hostname == cls.HOST:
                hs = {h["name"].lower(): h["value"] for h in req.get("headers", [])}
                auth = {k: hs[k] for k in cls.AUTH_HEADERS if k in hs}
                if auth or cls.MODE == "public":
                    return cls(auth=auth, **kw)
        raise MimicError(f"no request to {cls.HOST} carrying {cls.AUTH_HEADERS} in {path}")

    @classmethod
    def from_curl(cls, text: str, **kw: Any) -> Any:
        """Load auth headers from a 'Copy as cURL' paste."""
        import shlex

        toks = shlex.split(text.replace("\\\n", " "))
        auth: dict[str, str] = {}
        for i, t in enumerate(toks[:-1]):
            if t in ("-H", "--header"):
                k, _, v = toks[i + 1].partition(":")
                if k.strip().lower() in cls.AUTH_HEADERS:
                    auth[k.strip().lower()] = v.strip()
            elif t in ("-b", "--cookie") and "cookie" in cls.AUTH_HEADERS:
                auth["cookie"] = toks[i + 1]
        return cls(auth=auth, **kw)

    # ---- learned auth flow (login -> token -> use, refresh) ---------------
    def _install_token(self, resp: Any) -> None:
        flow = self.AUTH_FLOW or {}
        token = get_path(resp, flow.get("token_path"))
        if not isinstance(token, (str, int)) or not str(token):
            raise AuthFlowError(f"no token at {flow.get('token_path')!r} in the login response")
        value = str(flow.get("token_format") or "{token}").replace("{token}", str(token))
        if flow.get("token_in") == "query":
            self.auth_params[str(flow["token_name"])] = value
        else:
            self.auth[str(flow["token_name"]).lower()] = value
        exp = get_path(resp, flow.get("expires_in_path")) if flow.get("expires_in_path") else None
        self._token_expiry = (
            (_time.monotonic() + max(0.0, float(exp) - 30.0))
            if isinstance(exp, (int, float)) and not isinstance(exp, bool)
            else None
        )
        rt = get_path(resp, flow.get("refresh_token_path")) if flow.get("refresh_token_path") else None
        if isinstance(rt, str) and rt:
            self._refresh_token = rt
        _event(
            _logging.INFO,
            "auth_token_installed",
            expires=self._token_expiry is not None,
            refreshable=self._refresh_token is not None,
        )

    def authenticate(self, **credentials: Any) -> Any:
        """Run the learned login endpoint with ``credentials`` and install the issued token.

        Credentials are passed straight through and never stored; only the token is kept,
        in memory. Returns the login response.
        """
        flow = self.AUTH_FLOW
        if not flow:
            raise AuthFlowError(f"{type(self).__name__} has no learned auth flow")
        resp = getattr(self, flow["login_endpoint"])(**credentials)
        self._install_token(resp)
        return resp

    def _can_refresh(self) -> bool:
        flow = self.AUTH_FLOW or {}
        return bool(
            flow.get("refresh_endpoint")
            and flow.get("refresh_param")
            and self._refresh_token
            and not self._refreshing
        )

    def refresh_auth(self) -> bool:
        """Use the stored refresh token with the learned refresh endpoint. True on success."""
        if not self._can_refresh():
            return False
        flow = self.AUTH_FLOW or {}
        self._refreshing = True
        try:
            resp = getattr(self, flow["refresh_endpoint"])(**{flow["refresh_param"]: self._refresh_token})
            self._install_token(resp)
            return True
        finally:
            self._refreshing = False

    # ---- request pipeline ---------------------------------------------------
    def _u(self, key: str, path: str) -> str:
        """URL for ``path`` on learned host ``key`` (multi-host clients)."""
        return self.bases.get(key, self.base_url) + path

    def _auth_hosts(self) -> set[str]:
        return {_urlsplit(self.base_url).netloc} | {_urlsplit(b).netloc for b in self.bases.values()}

    def _limiter_for(self, url: str) -> RateLimiter:
        host = _urlsplit(url).netloc
        lim = self._limiters.get(host)
        if lim is None:
            lim = self._limiters[host] = RateLimiter(host, self.min_interval)
        return lim

    def _breaker_for(self, url: str) -> CircuitBreaker:
        return CircuitBreaker.for_host(_urlsplit(url).netloc, self.circuit_threshold, self.circuit_reset)

    def _prepare(
        self,
        method: str,
        path: str,
        params: Mapping[str, Any] | Sequence[tuple[str, Any]] | None,
        json: Any,
        data: Any,
        headers: Mapping[str, str] | None,
        safe: bool | None,
        files: Mapping[str, Any] | None,
        multipart: bool,
    ) -> PreparedRequest:
        method = method.upper()
        url = path if path.startswith(("http://", "https://")) else self.base_url + path
        self.guard.check(url)
        own_host = _urlsplit(url).netloc in self._auth_hosts()
        q: QueryItems = list(params.items() if isinstance(params, Mapping) else (params or []))
        q = [(k, v) for k, v in q if v is not None] + (list(self.auth_params.items()) if own_host else [])
        if q:
            url += ("&" if "?" in url else "?") + _urlencode(q, doseq=True)
        hdrs = dict(self.headers)
        if own_host:
            hdrs.update(self.auth)
            if self.cookies and self.use_cookies:
                hdrs["cookie"] = merge_cookie_header(hdrs.get("cookie"), self.cookies)
            csrf = self._csrf_cfg()
            unsafe = not safe if safe is not None else method not in ("GET", "HEAD", "OPTIONS")
            if csrf and unsafe and self._csrf:
                hdrs[str(csrf["header"])] = self._csrf
        hdrs["user-agent"] = self.user_agent
        hdrs["accept-encoding"] = "gzip, deflate"
        body: bytes | None = None
        if files or multipart:
            body, ctype = encode_multipart(data if isinstance(data, (dict, list)) else None, files)
            hdrs["content-type"] = ctype
        elif json is not None:
            body = _json.dumps(json).encode()
            hdrs.setdefault("content-type", "application/json")
        elif data is not None:
            raw = _urlencode(data) if isinstance(data, (dict, list)) else data
            body = raw.encode() if isinstance(raw, str) else raw
            hdrs.setdefault("content-type", "application/x-www-form-urlencoded")
        else:
            hdrs.pop("content-type", None)
        extra_h = {k.lower(): v for k, v in (headers or {}).items()}
        if (files or multipart) and "multipart" in extra_h.get("content-type", ""):
            extra_h.pop("content-type")  # keep our boundary
        hdrs.update(extra_h)
        if self.idempotency_keys and method in ("POST", "PATCH") and not safe:
            hdrs.setdefault("idempotency-key", _uuid.uuid4().hex)
        if self.respect_robots and not Robots.allowed(
            url, self.user_agent, min(self.timeout, 15), self.robots_unreachable
        ):
            _event(_logging.WARNING, "robots_disallowed", url=url.split("?")[0])
            raise RobotsDisallowed(
                f"robots.txt disallows {url} for {self.user_agent!r} "
                f"(pass respect_robots=False only if you have permission)"
            )
        # Idempotency-aware retries: reads, idempotent verbs, or writes carrying an idempotency key.
        retryable = bool(safe) or (safe is None and method in _IDEMPOTENT) or "idempotency-key" in hdrs
        return PreparedRequest(method, url, hdrs, body, self.timeout, retryable, own_host)

    def _backoff_delay(self, attempt: int, retry_after: float | None) -> float:
        delay = (
            retry_after
            if retry_after is not None
            else _random.uniform(0, min(self.max_backoff, self.backoff * (2**attempt)))
        )
        return min(delay, self.max_backoff)

    def _on_transport_error(
        self,
        prep: PreparedRequest,
        exc: TransportError,
        attempt: int,
        started: float,
        breaker: CircuitBreaker,
    ) -> float:
        """Delay before retrying a transport failure, or re-raise it."""
        breaker.failure()
        # A failed write may still have reached the server: only retry if that is harmless.
        if prep.retryable and attempt < self.max_retries:
            delay = self._backoff_delay(attempt, None)
            if self.deadline is None or (_time.monotonic() - started + delay) < self.deadline:
                _event(
                    _logging.INFO if not self.verbose else _logging.WARNING,
                    "retry",
                    url=prep.url.split("?")[0],
                    attempt=attempt + 1,
                    delay=round(delay, 2),
                    reason=type(exc).__name__,
                )
                return delay
        raise exc

    def _on_response(
        self,
        prep: PreparedRequest,
        resp: RawResponse,
        attempt: int,
        started: float,
        breaker: CircuitBreaker,
        raw: bool,
    ) -> tuple[str, Any]:
        """Decide what to do with a response: ("ok", value) | ("retry", delay) | ("reauth", None)."""
        payload = _decompress(resp.body, resp.headers)
        status, rh = resp.status, resp.headers
        self.last_status, self.last_headers, self.last_url = status, rh, prep.url
        if prep.own_host and self.use_cookies and resp.set_cookies:
            self._store_cookies(resp.set_cookies)
        _event(
            _logging.DEBUG,
            "response",
            method=prep.method,
            url=prep.url.split("?")[0],
            status=status,
            attempt=attempt,
            ms=round((_time.monotonic() - started) * 1000),
        )
        if status >= 500:
            breaker.failure()
        else:
            breaker.success()
        if status == 401 and prep.own_host and self._can_refresh() and not prep.meta.get("reauthed"):
            return "reauth", None
        if (
            status in (403, 419)
            and prep.own_host
            and self._csrf_cfg()
            and not prep.meta.get("csrf_retried")
            and str(self._csrf_cfg().get("header")) in prep.headers
            and not self._busy
        ):
            return "csrf", None
        if status in _RETRY_STATUSES and prep.retryable and attempt < self.max_retries:
            delay = self._backoff_delay(attempt, _retry_after(rh))
            if self.deadline is None or (_time.monotonic() - started + delay) < self.deadline:
                _event(
                    _logging.INFO if not self.verbose else _logging.WARNING,
                    "retry",
                    url=prep.url.split("?")[0],
                    attempt=attempt + 1,
                    delay=round(delay, 2),
                    status=status,
                )
                return "retry", delay
        if status >= 400:
            raise _http_error(status, prep.url, payload, rh)
        return "ok", (payload if raw else _parse_body(payload, rh))

    def _attempt_timeout(self, prep: PreparedRequest, started: float) -> float:
        if self.deadline is None:
            return self.timeout
        left = self.deadline - (_time.monotonic() - started)
        if left <= 0:
            raise RequestTimeout(f"{prep.method} {prep.url.split('?')[0]} exceeded deadline {self.deadline}s")
        return min(self.timeout, left)

    def request(
        self,
        method: str,
        path: str,
        params: Mapping[str, Any] | Sequence[tuple[str, Any]] | None = None,
        json: Any = None,
        data: Any = None,
        headers: Mapping[str, str] | None = None,
        safe: bool | None = None,
        raw: bool = False,
        files: Mapping[str, Any] | None = None,
        multipart: bool = False,
        _reauthed: bool = False,
    ) -> Any:
        """Make a request; returns parsed JSON (or text/bytes). Raises HTTPError on 4xx/5xx.

        json=  -> application/json body;  data= -> form-encoded (or raw str/bytes);
        files= -> real multipart/form-data: {field: bytes | "path/to/file" | (filename, bytes[, ctype])}
                  (any ``data`` dict is sent as the text parts); multipart=True forces multipart.
        safe=  -> True marks a read-only call (retried); None infers it from the HTTP method.
        Auth headers/params are only attached for the client's own hosts, never to foreign URLs.
        """
        if self._token_expiry is not None and _time.monotonic() >= self._token_expiry and self._can_refresh():
            self.refresh_auth()
        if self._csrf_needed(method, safe):
            self.fetch_csrf()
        prep = self._prepare(method, path, params, json, data, headers, safe, files, multipart)
        prep.meta["reauthed"] = _reauthed
        breaker = self._breaker_for(prep.url)
        started, attempt = _time.monotonic(), 0
        while True:
            breaker.before()
            delay = self._limiter_for(prep.url).reserve()
            if delay > 0:
                _time.sleep(delay)
            prep.timeout = self._attempt_timeout(prep, started)
            try:
                resp = self.transport(prep)
            except TransportError as e:
                _time.sleep(self._on_transport_error(prep, e, attempt, started, breaker))
                attempt += 1
                continue
            action, value = self._on_response(prep, resp, attempt, started, breaker, raw)
            if action == "ok":
                return value
            if action == "reauth" and self.refresh_auth():
                return self.request(
                    method, path, params, json, data, headers, safe, raw, files, multipart, True
                )
            if action == "reauth":
                raise _http_error(401, prep.url, resp.body, resp.headers)
            if action == "csrf":  # stale/missing CSRF token: reload it once and retry
                self._csrf = None
                self.fetch_csrf()
                prep = self._prepare(method, path, params, json, data, headers, safe, files, multipart)
                prep.meta.update(reauthed=_reauthed, csrf_retried=True)
                continue
            _time.sleep(value)
            attempt += 1

    # ---- cookie session + CSRF (learned SESSION) ----------------------------
    def _store_cookies(self, set_cookies: Sequence[str]) -> None:
        for h in set_cookies:
            parsed = parse_set_cookie(h)
            if parsed is None:
                continue
            name, value, expired = parsed
            if expired:
                self.cookies.pop(name, None)
            else:
                self.cookies[name] = value
                if self._csrf_cfg().get("source") == "cookie" and name == self._csrf_cfg().get("name"):
                    self._csrf = value

    def _csrf_cfg(self) -> dict[str, Any]:
        return dict((self.SESSION or {}).get("csrf") or {})

    def _csrf_needed(self, method: str, safe: bool | None) -> bool:
        cfg = self._csrf_cfg()
        unsafe = not safe if safe is not None else method.upper() not in ("GET", "HEAD", "OPTIONS")
        return (
            bool(cfg) and unsafe and self._csrf is None and cfg.get("source") != "cookie" and not self._busy
        )

    def _csrf_from(self, resp: Any) -> None:
        cfg = self._csrf_cfg()
        src, name = cfg.get("source"), str(cfg.get("name") or "")
        tok: Any = None
        if src == "body":
            tok = get_path(resp, name)
        elif src == "header":
            tok = self.last_headers.get(name) or {k.lower(): v for k, v in self.last_headers.items()}.get(
                name
            )
        elif src == "html":
            text = resp.decode("utf-8", "replace") if isinstance(resp, bytes) else str(resp)
            tok = csrf_from_html(text, name)
        if not isinstance(tok, (str, int)) or not str(tok):
            raise AuthFlowError(f"no CSRF token ({src} {name!r}) in the response")
        self._csrf = str(tok)
        _event(_logging.INFO, "csrf_token_installed", source=src)

    def fetch_csrf(self) -> str | None:
        """(Re)load the anti-CSRF token the learned way (endpoint body/header or page meta tag)."""
        cfg = self._csrf_cfg()
        if not cfg or cfg.get("source") == "cookie":
            return self._csrf
        self._busy = True
        try:
            if cfg.get("endpoint"):
                resp = getattr(self, str(cfg["endpoint"]))()
            else:
                resp = self.request("GET", str(cfg.get("page_path") or "/"), safe=True, raw=True)
            self._csrf_from(resp)
        finally:
            self._busy = False
        return self._csrf

    def login(self, **credentials: Any) -> Any:
        """Log in with the learned flow: bearer auth flow (token) or cookie session (Set-Cookie).

        Credentials are passed straight through and never stored.
        """
        if self.AUTH_FLOW:
            return self.authenticate(**credentials)
        ep = (self.SESSION or {}).get("login_endpoint")
        if not ep:
            raise AuthFlowError(f"{type(self).__name__} has no learned login flow")
        resp = getattr(self, ep)(**credentials)
        _event(_logging.INFO, "session_login", cookies=sorted(self.cookies))
        return resp

    def logout(self) -> Any:
        """Call the learned logout endpoint (if any) and forget cookies + CSRF token."""
        ep = (self.SESSION or {}).get("logout_endpoint")
        resp = getattr(self, ep)() if ep else None
        self.cookies.clear()
        self._csrf = None
        return resp

    def save_session(self, path: str) -> None:
        """Persist cookies + CSRF token to ``path`` (mode 0600) to resume later without logging in."""
        fd = _os.open(path, _os.O_WRONLY | _os.O_CREAT | _os.O_TRUNC, 0o600)
        with _os.fdopen(fd, "w", encoding="utf-8") as f:
            _json.dump({"cookies": self.cookies, "csrf": self._csrf}, f)

    def load_session(self, path: str) -> None:
        with open(path, encoding="utf-8") as f:
            data = _json.load(f)
        self.cookies.update({str(k): str(v) for k, v in (data.get("cookies") or {}).items()})
        self._csrf = data.get("csrf") or self._csrf

    # ---- streaming: SSE, WebSocket, GraphQL subscriptions -------------------
    def stream(
        self,
        path: str,
        params: Mapping[str, Any] | Sequence[tuple[str, Any]] | None = None,
        method: str = "GET",
        json: Any = None,
        headers: Mapping[str, str] | None = None,
        max_events: int | None = None,
        event: str | None = None,
    ) -> Iterator[dict[str, Any]]:
        """Server-Sent Events from ``path``: yields ``{"event", "id", "retry", "data"}`` as they arrive.

        ``event=`` keeps only that event type; ``max_events`` stops after N events.
        """
        hdrs = {"accept": "text/event-stream", "cache-control": "no-cache", **dict(headers or {})}
        prep = self._prepare(method, path, params, json, None, hdrs, True, None, False)
        prep.timeout = self.timeout
        n = 0
        for ev in parse_sse_lines(self._stream_lines(prep)):
            if event and ev["event"] != event:
                continue
            yield ev
            n += 1
            if max_events is not None and n >= max_events:
                return

    def _stream_lines(self, prep: PreparedRequest) -> Iterator[str]:
        if isinstance(self.transport, UrllibTransport):
            r = _urequest.Request(prep.url, data=prep.body, headers=prep.headers, method=prep.method)
            try:
                resp = self.transport.opener.open(r, timeout=prep.timeout)
            except _uerror.HTTPError as e:
                raise _http_error(e.code, prep.url, e.read() or b"", dict(e.headers or {})) from e
            except (_uerror.URLError, OSError) as e:
                raise TransportError(f"{prep.method} {prep.url.split('?')[0]} failed: {e}") from e
            with resp:
                self.last_status, self.last_headers = resp.status, dict(resp.headers)
                if prep.own_host:
                    self._store_cookies(resp.headers.get_all("set-cookie") or [])
                for raw in resp:
                    yield raw.decode("utf-8", "replace")
            return
        rr = self.transport(prep)  # injected transports (tests, mocks): buffered body
        self.last_status, self.last_headers = rr.status, rr.headers
        if rr.status >= 400:
            raise _http_error(rr.status, prep.url, rr.body, rr.headers)
        yield from _decompress(rr.body, rr.headers).decode("utf-8", "replace").splitlines()

    def ws_url(self, path: str, params: Mapping[str, Any] | Sequence[tuple[str, Any]] | None = None) -> str:
        url = path if path.startswith(("ws://", "wss://", "http://", "https://")) else self.base_url + path
        url = "ws" + url[4:] if url.startswith("http") else url
        q = list(params.items() if isinstance(params, Mapping) else (params or []))
        q = [(k, v) for k, v in q if v is not None]
        return url + (("&" if "?" in url else "?") + _urlencode(q, doseq=True) if q else "")

    def ws_connect(
        self,
        path: str,
        params: Mapping[str, Any] | Sequence[tuple[str, Any]] | None = None,
        subprotocols: Sequence[str] = (),
        headers: Mapping[str, str] | None = None,
    ) -> WebSocket:
        """Open a WebSocket on ``path`` with this client's auth headers + cookies (own hosts only)."""
        url = self.ws_url(path, params)
        http_url = "http" + url[2:]
        self.guard.check(http_url)
        own = _urlsplit(http_url).netloc in self._auth_hosts()
        hdrs = dict(self.headers)
        if own:
            hdrs.update(self.auth)
            if self.cookies:
                hdrs["cookie"] = merge_cookie_header(hdrs.get("cookie"), self.cookies)
            if self.auth_params:
                url += ("&" if "?" in url else "?") + _urlencode(list(self.auth_params.items()))
        hdrs["user-agent"] = self.user_agent
        hdrs["origin"] = self.base_url
        hdrs.update({k.lower(): v for k, v in (headers or {}).items()})
        _event(_logging.DEBUG, "ws_connect", url=url.split("?")[0])
        return WebSocket.connect(url, hdrs, subprotocols, self.timeout)

    def listen(
        self,
        path: str,
        messages: Sequence[Any] = (),
        params: Mapping[str, Any] | Sequence[tuple[str, Any]] | None = None,
        subprotocols: Sequence[str] = (),
        max_messages: int | None = None,
        until: Callable[[Any], bool] | None = None,
    ) -> Iterator[Any]:
        """Connect, send ``messages`` (JSON-encoded unless str/bytes), yield every received message."""
        with self.ws_connect(path, params, subprotocols) as ws:
            for m in messages:
                ws.send(m) if isinstance(m, (str, bytes)) else ws.send_json(m)
            for n, msg in enumerate(ws, 1):
                yield msg
                if (max_messages is not None and n >= max_messages) or (until is not None and until(msg)):
                    return

    def graphql_subscribe(
        self,
        path: str,
        query: str,
        variables: Mapping[str, Any] | None = None,
        operation_name: str | None = None,
        protocol: str = "graphql-transport-ws",
        max_events: int | None = None,
        connection_params: Mapping[str, Any] | None = None,
    ) -> Iterator[Any]:
        """A GraphQL subscription over WebSocket; yields each event's ``data``.

        Speaks ``graphql-transport-ws`` (graphql-ws library) and the legacy ``graphql-ws``
        (subscriptions-transport-ws) protocol. GraphQL ``errors`` raise :class:`MimicError`.
        """
        legacy = protocol == "graphql-ws"
        with self.ws_connect(path, subprotocols=[protocol]) as ws:
            ws.send_json({"type": "connection_init", "payload": dict(connection_params or {})})
            payload = {"query": query, "variables": dict(variables or {}), "operationName": operation_name}
            ws.send_json({"id": "1", "type": "start" if legacy else "subscribe", "payload": payload})
            n = 0
            for msg in ws:
                t = msg.get("type") if isinstance(msg, dict) else None
                if t == "ping":
                    ws.send_json({"type": "pong"})
                elif t in ("next", "data"):
                    body = msg.get("payload") or {}
                    if body.get("errors"):
                        raise MimicError(f"GraphQL subscription error: {body['errors']}")
                    yield body.get("data")
                    n += 1
                    if max_events is not None and n >= max_events:
                        return
                elif t in ("error", "connection_error"):
                    raise MimicError(f"GraphQL subscription error: {msg.get('payload')}")
                elif t == "complete":
                    return

    def _grpc_body(self, fields: Mapping[int, Any] | Sequence[tuple[int, Any]]) -> bytes:
        import struct as _struct

        msg = pb_encode(fields)
        return b"\x00" + _struct.pack(">I", len(msg)) + msg

    def _grpc_result(self, raw: Any) -> dict[str, Any]:
        msgs, trailers = grpc_web_frames(raw if isinstance(raw, bytes) else bytes(str(raw), "latin-1"))
        trailers.update({k.lower(): v for k, v in self.last_headers.items() if k.lower().startswith("grpc-")})
        status = int(trailers.get("grpc-status", "0") or 0)
        if status:
            raise GrpcError(f"gRPC status {status}: {trailers.get('grpc-message', '')}", status)
        out: dict[str, Any] = {}
        for m in msgs[:1]:
            for num, _wire, v in pb_decode(m):
                key = f"field_{num}"
                if key in out:
                    prev = out[key]
                    out[key] = [*prev, v] if isinstance(prev, list) else [prev, v]
                else:
                    out[key] = v
        return out

    def grpc_call(
        self,
        path: str,
        fields: Mapping[int, Any] | Sequence[tuple[int, Any]],
        content_type: str = "application/grpc-web+proto",
        safe: bool | None = None,
    ) -> dict[str, Any]:
        """Unary gRPC-web call: encode ``{field_number: value}``, decode the reply as ``{"field_N": v}``."""
        raw = self.request(
            "POST",
            path,
            data=self._grpc_body(fields),
            headers={"content-type": content_type, "x-grpc-web": "1", "accept": content_type},
            safe=safe,
            raw=True,
        )
        return self._grpc_result(raw)

    # ---- helpers for generated methods -------------------------------------
    def items(self, name: str, response: Any) -> list[Any]:
        """The result list inside a response of endpoint ``name`` (per ITEMS), else []."""
        spec = self.ITEMS.get(name)
        if not spec:
            return response if isinstance(response, list) else []
        got = get_path(response, spec[0])
        return got if isinstance(got, list) else []

    def item_id(self, name: str, item: Any) -> str:
        spec = self.ITEMS.get(name) or [None, None]
        if spec[1] and isinstance(item, dict) and item.get(spec[1]) is not None:
            return str(item[spec[1]])
        return _json.dumps(item, sort_keys=True, default=str)

    def _pager(
        self,
        name: str | Callable[..., Any],
        kw: Mapping[str, Any],
        max_pages: int,
        max_items: int | None,
        recipe: Mapping[str, Any] | None,
    ) -> tuple[Callable[..., Any], str, _Pager]:
        fn: Callable[..., Any] = name if callable(name) else getattr(self, name)
        key = getattr(fn, "__name__", str(name)) if callable(name) else name
        return fn, key, _Pager(recipe or self.PAGINATION.get(key) or {}, kw, max_pages, max_items)

    def _page_items(self, pager: _Pager, key: str, resp: Any) -> list[Any]:
        ip = pager.rc.get("items_path")
        items = get_path(resp, ip) if ip else self.items(key, resp)
        return items if isinstance(items, list) else []

    def paginate(
        self,
        name: str | Callable[..., Any],
        *args: Any,
        max_pages: int = 10,
        max_items: int | None = None,
        recipe: Mapping[str, Any] | None = None,
        **kw: Any,
    ) -> Iterator[Any]:
        """Yield result items across pages of endpoint ``name`` (a method name or a callable).

        The recipe comes from PAGINATION[name] (learned from traffic) or ``recipe=``; without one,
        next-links / cursors in the response (and the HTTP Link header) are auto-detected.
        Recipe keys: style = page | offset | cursor | next;  param (python kwarg to vary);
        start; size_param; cursor_path; next_path; total_path; items_path.
        Stops on an empty page, a short page, a repeated page, total reached, or max_pages.
        """
        fn, key, pager = self._pager(name, kw, max_pages, max_items, recipe)
        while True:
            what, arg = pager.plan()
            if what == "stop":
                return
            resp = fn(*args, **arg) if what == "call" else self.request("GET", arg, safe=True)
            items = self._page_items(pager, key, resp)
            yield from pager.take(items)
            pager.next_url = None
            pager.feed(resp, items, self.last_url, self.base_url, self.last_headers)

    # ---- drift detection ----------------------------------------------------
    def _probe_kwargs(self, probe: Mapping[str, Any]) -> dict[str, Any]:
        ct = probe.get("content_type") or ""
        body = probe.get("body") or None
        kw: dict[str, Any] = {"params": probe.get("query") or None, "safe": True}
        if body is not None:
            kw["json" if "json" in ct else "data"] = _json.loads(body) if "json" in ct else body
        return kw

    def _probe_result(self, name: str, probe: Mapping[str, Any], live: Any) -> dict[str, Any]:
        changes = diff_schema(probe.get("schema"), infer_schema(live))
        breaking = [c for c in changes if c["breaking"]]
        return {
            "status": self.last_status,
            "ok": not breaking,
            "breaking": breaking,
            "info": [c for c in changes if not c["breaking"]],
            "items": len(self.items(name, live)) if name in self.ITEMS else None,
        }

    def selfcheck(self, names: Sequence[str] | None = None) -> dict[str, Any]:
        """Replay every recorded read-only probe; compare live response schema to the recorded one.

        Returns {"ok": bool, "checked": n, "endpoints": {name: {...}}}. ``ok`` is False if any
        probe failed or showed *breaking* drift (removed required field / changed type).
        """
        report: dict[str, Any] = {"ok": True, "checked": 0, "client": type(self).__name__, "endpoints": {}}
        for name, probe in self.PROBES.items():
            if names and name not in names:
                continue
            try:
                live = self.request(probe["method"], probe["path"], **self._probe_kwargs(probe))
                r = self._probe_result(name, probe, live)
            except MimicError as e:
                r = {
                    "status": getattr(e, "status", None),
                    "ok": False,
                    "breaking": [],
                    "info": [],
                    "items": None,
                    "error": str(e),
                }
            report["endpoints"][name] = r
            report["checked"] += 1
            report["ok"] = report["ok"] and r["ok"]
        return report


def encode_multipart(
    fields: Mapping[str, Any] | Sequence[tuple[str, Any]] | None = None,
    files: Mapping[str, Any] | None = None,
    boundary: str | None = None,
) -> tuple[bytes, str]:
    """RFC 7578 multipart/form-data body. Returns (body_bytes, content_type)."""
    boundary = boundary or ("mimicplus" + _binascii.hexlify(_os.urandom(12)).decode())
    out: list[bytes] = []
    items = list(fields.items() if isinstance(fields, Mapping) else (fields or []))
    for k, v in items:
        if v is None:
            continue
        out.append(f'--{boundary}\r\nContent-Disposition: form-data; name="{_mp_quote(k)}"\r\n\r\n'.encode())
        out.append(v if isinstance(v, bytes) else str(v).encode("utf-8"))
        out.append(b"\r\n")
    for k, v in (files or {}).items():
        if v is None:
            continue
        ctype: str | None = None
        if isinstance(v, (tuple, list)):
            fname, content = str(v[0]), v[1]
            ctype = v[2] if len(v) > 2 else None
        elif isinstance(v, str) and _os.path.isfile(v):
            fname = _os.path.basename(v)
            with open(v, "rb") as fh:
                content = fh.read()
        else:
            fname, content = k, v
        if isinstance(content, str):
            content = content.encode("utf-8")
        ctype = ctype or _mimetypes.guess_type(fname)[0] or "application/octet-stream"
        out.append(
            (
                f'--{boundary}\r\nContent-Disposition: form-data; name="{_mp_quote(k)}"; '
                f'filename="{_mp_quote(fname)}"\r\nContent-Type: {ctype}\r\n\r\n'
            ).encode()
        )
        out.append(bytes(content))
        out.append(b"\r\n")
    out.append(f"--{boundary}--\r\n".encode())
    return b"".join(out), f"multipart/form-data; boundary={boundary}"


def _mp_quote(s: object) -> str:
    return str(s).replace('"', "%22").replace("\r", "%0D").replace("\n", "%0A")


def decode_multipart(body: bytes, content_type: str) -> list[dict[str, Any]]:
    """Parse a multipart/form-data body -> list of {name, filename, content_type, value(bytes)}."""
    import email.message
    import email.parser
    import email.policy

    msg = email.parser.BytesParser(policy=email.policy.HTTP).parsebytes(
        b"Content-Type: " + content_type.encode() + b"\r\n\r\n" + body
    )
    parts: list[dict[str, Any]] = []
    if not msg.is_multipart():
        return parts
    for part in msg.iter_parts():
        if not isinstance(part, email.message.EmailMessage):
            continue
        parts.append(
            {
                "name": part.get_param("name", header="content-disposition"),
                "filename": part.get_filename(),
                "content_type": part.get_content_type() if part.get("content-type") else None,
                "value": part.get_payload(decode=True) or b"",
            }
        )
    return parts


def _retry_after(headers: Mapping[str, str]) -> float | None:
    v = {k.lower(): x for k, x in headers.items()}.get("retry-after")
    if not v:
        return None
    try:
        return max(0.0, float(v))
    except ValueError:
        try:
            return max(0.0, _parsedate(v).timestamp() - _time.time())
        except (TypeError, ValueError):
            return None


def _decompress(payload: bytes, headers: Mapping[str, str]) -> bytes:
    enc = {k.lower(): v for k, v in headers.items()}.get("content-encoding", "").lower()
    try:
        if enc == "gzip":
            return _gzip.decompress(payload)
        if enc == "deflate":
            try:
                return _zlib.decompress(payload)
            except _zlib.error:
                return _zlib.decompress(payload, -_zlib.MAX_WBITS)
    except (OSError, EOFError, _zlib.error):
        pass
    return payload


def _parse_body(payload: bytes, headers: Mapping[str, str]) -> Any:
    if not payload:
        return None
    try:
        return _json.loads(payload.decode("utf-8"))
    except (UnicodeDecodeError, ValueError):
        try:
            return payload.decode("utf-8")
        except UnicodeDecodeError:
            return payload


def _cli(cls: type[Client], argv: Sequence[str] | None = None) -> int:
    """Tiny CLI for a generated client: endpoints | selfcheck | call <method> k=v ..."""
    import argparse

    p = argparse.ArgumentParser(prog=_sys.argv[0], description=cls.__doc__)
    p.add_argument("--log", default=_os.environ.get("MIMICPLUS_LOG", "WARNING"), help="log level")
    sub = p.add_subparsers(dest="cmd", required=True)
    sub.add_parser("endpoints", help="list methods")
    sc = sub.add_parser("selfcheck", help="replay probes and detect schema drift")
    sc.add_argument("--json", action="store_true")
    c = sub.add_parser("call", help="call a method: call search_items q=shoes rows=10")
    c.add_argument("method")
    c.add_argument("args", nargs="*")
    c.add_argument("--items", action="store_true", help="print only the result items (JSONL)")
    a = p.parse_args(argv)
    _logging.basicConfig(level=str(a.log).upper(), format="[%(name)s] %(levelname)s: %(message)s")
    if a.cmd == "endpoints":
        for name in sorted(cls.SAFE):
            fn = getattr(cls, name)
            doc = (fn.__doc__ or "").strip().splitlines()
            print(f"{name}{'' if cls.SAFE[name] else '  [writes]'}  - {doc[0] if doc else ''}")
        return 0
    client = cls()
    if a.cmd == "selfcheck":
        rep = client.selfcheck()
        if a.json:
            print(_json.dumps(rep, indent=2))
        else:
            _print_selfcheck(rep)
        return 0 if rep["ok"] else 1
    kw: dict[str, Any] = {}
    for s in a.args:
        k, _, v = s.partition("=")
        try:
            kw[k] = _json.loads(v)
        except ValueError:
            kw[k] = v
    out = getattr(client, a.method)(**kw)
    if a.items:
        for it in client.items(a.method, out):
            print(_json.dumps(it, default=str))
    else:
        print(_json.dumps(out, indent=2, default=str))
    return 0


def _print_selfcheck(rep: Mapping[str, Any]) -> None:
    print(f"selfcheck {rep['client']}: {'OK' if rep['ok'] else 'DRIFT/FAIL'} ({rep['checked']} probes)")
    for name, r in rep["endpoints"].items():
        mark = "ok  " if r["ok"] else "FAIL"
        extra = f" items={r['items']}" if r.get("items") is not None else ""
        print(
            f"  [{mark}] {name}: HTTP {r['status']}{extra} breaking={len(r['breaking'])} "
            f"info={len(r['info'])}" + (f" error={r['error'][:120]}" if r.get("error") else "")
        )
        for ch in r["breaking"][:10]:
            print(f"         ! {ch['kind']} {ch['path']} {ch['detail']}")
        for ch in r["info"][:5]:
            print(f"         ~ {ch['kind']} {ch['path']} {ch['detail']}")


# ---- END MIMICPLUS RUNTIME ----


# ---- BEGIN API ----
LoginResponse = TypedDict("LoginResponse", {
    "access_token": str,
    "refresh_token": str,
    "token_type": str,
    "expires_in": int,
}, total=False)

GetMeResponse = TypedDict("GetMeResponse", {
    "id": int,
    "username": str,
    "email": str,
    "plan": str,
}, total=False)

ListProductsResponseItemsItem = TypedDict("ListProductsResponseItemsItem", {
    "id": int,
    "in_stock": bool,
    "name": str,
    "price": float,
    "tags": List[str],
}, total=False)

ListProductsResponse = TypedDict("ListProductsResponse", {
    "items": List[ListProductsResponseItemsItem],
    "next_page": Optional[int],
    "page": int,
    "per_page": int,
    "total": int,
}, total=False)

GetProductResponse = TypedDict("GetProductResponse", {
    "id": int,
    "in_stock": bool,
    "name": str,
    "price": float,
    "tags": List[str],
}, total=False)

CreateItemResponseItem = TypedDict("CreateItemResponseItem", {
    "product_id": int,
    "qty": int,
}, total=False)

CreateItemResponse = TypedDict("CreateItemResponse", {
    "cart_size": int,
    "item": CreateItemResponseItem,
}, total=False)

SearchProductsResponseDataSearchProducts = TypedDict("SearchProductsResponseDataSearchProducts", {
    "nodes": List[Dict[str, Any]],
    "totalCount": int,
}, total=False)

SearchProductsResponseData = TypedDict("SearchProductsResponseData", {
    "searchProducts": SearchProductsResponseDataSearchProducts,
}, total=False)

SearchProductsResponse = TypedDict("SearchProductsResponse", {
    "data": SearchProductsResponseData,
}, total=False)

GetOrderResponseDataOrder = TypedDict("GetOrderResponseDataOrder", {
    "id": str,
    "status": str,
    "total": float,
    "items": List[int],
}, total=False)

GetOrderResponseData = TypedDict("GetOrderResponseData", {
    "order": GetOrderResponseDataOrder,
}, total=False)

GetOrderResponse = TypedDict("GetOrderResponse", {
    "data": GetOrderResponseData,
}, total=False)

AddReviewResponseDataAddReview = TypedDict("AddReviewResponseDataAddReview", {
    "id": str,
    "rating": int,
    "ok": bool,
}, total=False)

AddReviewResponseData = TypedDict("AddReviewResponseData", {
    "addReview": AddReviewResponseDataAddReview,
}, total=False)

AddReviewResponse = TypedDict("AddReviewResponse", {
    "data": AddReviewResponseData,
}, total=False)

RefreshResponse = TypedDict("RefreshResponse", {
    "access_token": str,
    "refresh_token": str,
    "token_type": str,
    "expires_in": int,
}, total=False)


class Api01(Client):
    """Client for 127.0.0.1 (needs auth headers ['authorization'] loaded at runtime (env API01_AUTH, .from_har(), .from_curl(), or .authenticate(**credentials) via the learned login flow); never hardcoded)."""

    HOST = '127.0.0.1'
    BASE_URL = 'http://127.0.0.1:8765'
    MODE = 'auth'
    DEFAULT_HEADERS = {'accept': 'application/json'}
    AUTH_HEADERS = ['authorization']
    AUTH_QUERY_PARAMS = []

    def login(
            self,
            username: Optional[str] = None,
            password: Optional[str] = None,
            **extra: Any) -> LoginResponse:
        """POST /api/auth/login -> 200 (auth login step: use .authenticate(**credentials); WRITES: not auto-retried; seen 1x in capture).

        Args:
            username: credential body field 'username'; never stored, omitted when None
            password: credential body field 'password'; never stored, omitted when None
            **extra: additional body fields
        """
        params: List[Tuple[str, Any]] = []
        body: Dict[str, Any] = {
            'username': username,
            'password': password,
        }
        body = {k: v for k, v in body.items() if not (v is None and k in ['username', 'password'])}
        body.update(extra)
        return cast('LoginResponse', self.request('POST', '/api/auth/login', params=params, json=body, headers={'content-type': 'application/json'}, safe=False))

    def get_me(
            self,
            **extra: Any) -> GetMeResponse:
        """GET /api/me -> 200 (read-only, retried on 429/5xx; seen 1x in capture).
        """
        params: List[Tuple[str, Any]] = []
        params += list(extra.items())
        return cast('GetMeResponse', self.request('GET', '/api/me', params=params, safe=True))

    def list_products(
            self,
            page: int = 1,
            per_page: int = 2,
            **extra: Any) -> ListProductsResponse:
        """GET /api/products -> 200 (read-only, retried on 429/5xx; results in `items` keyed by `id`; paginated (page): use .paginate('list_products', ...); seen 3x in capture).

        Args:
            page: query field 'page' (varied in capture)
            per_page: query field 'per_page'
            **extra: additional query params
        """
        params: List[Tuple[str, Any]] = [('page', page), ('per_page', per_page)]
        params += list(extra.items())
        return cast('ListProductsResponse', self.request('GET', '/api/products', params=params, safe=True))

    def get_product(
            self,
            product_id: int,
            **extra: Any) -> GetProductResponse:
        """GET /api/products/{product_id} -> 200 (read-only, retried on 429/5xx; seen 4x in capture).

        Args:
            product_id: path segment (e.g. '102')
            **extra: additional query params
        """
        params: List[Tuple[str, Any]] = []
        params += list(extra.items())
        return cast('GetProductResponse', self.request('GET', f'/api/products/{_q(str(product_id))}', params=params, safe=True))

    def create_item(
            self,
            product_id: int = 104,
            qty: int = 1,
            **extra: Any) -> CreateItemResponse:
        """POST /api/cart/items -> 201 (WRITES: not auto-retried; seen 1x in capture).

        Args:
            product_id: body field 'product_id'
            qty: body field 'qty'
            **extra: additional body fields
        """
        params: List[Tuple[str, Any]] = []
        body: Dict[str, Any] = {
            'product_id': product_id,
            'qty': qty,
        }
        body.update(extra)
        return cast('CreateItemResponse', self.request('POST', '/api/cart/items', params=params, json=body, headers={'content-type': 'application/json'}, safe=False))

    def search_products(
            self,
            q: str,
            first: int = 5,
            **extra: Any) -> SearchProductsResponse:
        """POST /graphql -> 200 (GraphQL query SearchProducts; read-only, retried on 429/5xx; results in `data.searchProducts.nodes` keyed by `id`; seen 2x in capture).

        Args:
            q: free-text query (sent as body field 'q')
            first: body field 'variables.first'
            **extra: additional body fields
        """
        params: List[Tuple[str, Any]] = []
        body: Dict[str, Any] = {
            'operationName': 'SearchProducts',
            'query': ('query SearchProducts($q: String!, $first: Int) { searchProducts(q: $q, first: $first) { '
     'totalCount nodes { id name price } } }'),
            'variables': {'q': 'lamp', 'first': 5},
        }
        set_path(body, ['variables', 'q'], q)
        set_path(body, ['variables', 'first'], first)
        body.update(extra)
        return cast('SearchProductsResponse', self.request('POST', '/graphql', params=params, json=body, headers={'content-type': 'application/json'}, safe=True))

    def get_order(
            self,
            id: str = 'o-1',
            **extra: Any) -> GetOrderResponse:
        """POST /graphql -> 200 (GraphQL query GetOrder; read-only, retried on 429/5xx; seen 1x in capture).

        Args:
            id: body field 'variables.id'
            **extra: additional body fields
        """
        params: List[Tuple[str, Any]] = []
        body: Dict[str, Any] = {
            'operationName': 'GetOrder',
            'query': 'query GetOrder($id: ID!) { order(id: $id) { id status total items } }',
            'variables': {'id': 'o-1'},
        }
        set_path(body, ['variables', 'id'], id)
        body.update(extra)
        return cast('GetOrderResponse', self.request('POST', '/graphql', params=params, json=body, headers={'content-type': 'application/json'}, safe=True))

    def add_review(
            self,
            text: str,
            product_id: str = '101',
            rating: int = 5,
            **extra: Any) -> AddReviewResponse:
        """POST /graphql -> 200 (GraphQL mutation AddReview; WRITES: not auto-retried; seen 1x in capture).

        Args:
            text: free-text query (sent as body field 'text')
            product_id: body field 'variables.productId'
            rating: body field 'variables.rating'
            **extra: additional body fields
        """
        params: List[Tuple[str, Any]] = []
        body: Dict[str, Any] = {
            'operationName': 'AddReview',
            'query': ('mutation AddReview($productId: ID!, $rating: Int!, $text: String) { addReview(productId: '
     '$productId, rating: $rating, text: $text) { id rating ok } }'),
            'variables': {'productId': '101', 'rating': 5, 'text': 'great'},
        }
        set_path(body, ['variables', 'productId'], product_id)
        set_path(body, ['variables', 'rating'], rating)
        set_path(body, ['variables', 'text'], text)
        body.update(extra)
        return cast('AddReviewResponse', self.request('POST', '/graphql', params=params, json=body, headers={'content-type': 'application/json'}, safe=False))

    def refresh(
            self,
            refresh_token: Optional[str] = None,
            **extra: Any) -> RefreshResponse:
        """POST /api/auth/refresh -> 200 (auth refresh step (called automatically); WRITES: not auto-retried; seen 1x in capture).

        Args:
            refresh_token: credential body field 'refresh_token'; never stored, omitted when None
            **extra: additional body fields
        """
        params: List[Tuple[str, Any]] = []
        body: Dict[str, Any] = {
            'refresh_token': refresh_token,
        }
        body = {k: v for k, v in body.items() if not (v is None and k in ['refresh_token'])}
        body.update(extra)
        return cast('RefreshResponse', self.request('POST', '/api/auth/refresh', params=params, json=body, headers={'content-type': 'application/json'}, safe=False))

# ---- END API ----


# ---- BEGIN METADATA (selfcheck/watch data; regenerate with `mimicplus gen`) ----
Api01.ITEMS = {'list_products': ['items', 'id'], 'search_products': ['data.searchProducts.nodes', 'id']}
Api01.SAFE = {'login': False,
 'get_me': True,
 'list_products': True,
 'get_product': True,
 'create_item': False,
 'search_products': True,
 'get_order': True,
 'add_review': False,
 'refresh': False}
Api01.QUERY_PARAM = {'search_products': 'q', 'add_review': 'text'}
Api01.PAGINATION = {'list_products': {'items_path': 'items',
                   'total_path': 'total',
                   'size_param': 'per_page',
                   'style': 'page',
                   'param': 'page',
                   'start': 1}}
Api01.AUTH_FLOW = {'login_endpoint': 'login',
 'token_path': 'access_token',
 'token_in': 'header',
 'token_name': 'authorization',
 'token_format': 'Bearer {token}',
 'expires_in_path': 'expires_in',
 'refresh_endpoint': 'refresh',
 'refresh_token_path': 'refresh_token',
 'refresh_param': 'refresh_token',
 'uses': ['get_me',
          'list_products',
          'get_product',
          'create_item',
          'search_products',
          'get_order',
          'add_review']}
Api01.PROBES = {'get_me': {'method': 'GET',
            'path': '/api/me',
            'query': [],
            'body': None,
            'content_type': '',
            'status': 200,
            'schema': {'type': ['object'],
                       'properties': {'id': {'type': ['integer']},
                                      'username': {'type': ['string']},
                                      'email': {'type': ['string']},
                                      'plan': {'type': ['string']}},
                       'required': ['email', 'id', 'plan', 'username']}},
 'list_products': {'method': 'GET',
                   'path': '/api/products',
                   'query': [['page', '3'], ['per_page', '2']],
                   'body': None,
                   'content_type': '',
                   'status': 200,
                   'schema': {'type': ['object'],
                              'properties': {'items': {'type': ['array'],
                                                       'items': {'type': ['object'],
                                                                 'properties': {'id': {'type': ['integer']},
                                                                                'in_stock': {'type': ['boolean']},
                                                                                'name': {'type': ['string']},
                                                                                'price': {'type': ['number']},
                                                                                'tags': {'type': ['array'],
                                                                                         'items': {'type': ['string']}}},
                                                                 'required': ['id',
                                                                              'in_stock',
                                                                              'name',
                                                                              'price',
                                                                              'tags']}},
                                             'next_page': {'type': ['integer', 'null']},
                                             'page': {'type': ['integer']},
                                             'per_page': {'type': ['integer']},
                                             'total': {'type': ['integer']}},
                              'required': ['items', 'next_page', 'page', 'per_page', 'total']}},
 'get_product': {'method': 'GET',
                 'path': '/api/products/102',
                 'query': [],
                 'body': None,
                 'content_type': '',
                 'status': 200,
                 'schema': {'type': ['object'],
                            'properties': {'id': {'type': ['integer']},
                                           'in_stock': {'type': ['boolean']},
                                           'name': {'type': ['string']},
                                           'price': {'type': ['number']},
                                           'tags': {'type': ['array'],
                                                    'items': {'type': ['string']}}},
                            'required': ['id', 'in_stock', 'name', 'price', 'tags']}},
 'search_products': {'method': 'POST',
                     'path': '/graphql',
                     'query': [],
                     'body': '{"operationName": "SearchProducts", "query": "query '
                             'SearchProducts($q: String!, $first: Int) { searchProducts(q: $q, '
                             'first: $first) { totalCount nodes { id name price } } }", '
                             '"variables": {"q": "lamp", "first": 5}}',
                     'content_type': 'application/json',
                     'status': 200,
                     'schema': {'type': ['object'],
                                'properties': {'data': {'type': ['object'],
                                                        'properties': {'searchProducts': {'type': ['object'],
                                                                                          'properties': {'nodes': {'type': ['array'],
                                                                                                                   'items': {'type': ['object'],
                                                                                                                             'properties': {'id': {'type': ['string']},
                                                                                                                                            'name': {'type': ['string']},
                                                                                                                                            'price': {'type': ['number']}},
                                                                                                                             'required': ['id',
                                                                                                                                          'name',
                                                                                                                                          'price']}},
                                                                                                         'totalCount': {'type': ['integer']}},
                                                                                          'required': ['nodes',
                                                                                                       'totalCount']}},
                                                        'required': ['searchProducts']}},
                                'required': ['data']}},
 'get_order': {'method': 'POST',
               'path': '/graphql',
               'query': [],
               'body': '{"operationName": "GetOrder", "query": "query GetOrder($id: ID!) { '
                       'order(id: $id) { id status total items } }", "variables": {"id": "o-1"}}',
               'content_type': 'application/json',
               'status': 200,
               'schema': {'type': ['object'],
                          'properties': {'data': {'type': ['object'],
                                                  'properties': {'order': {'type': ['object'],
                                                                           'properties': {'id': {'type': ['string']},
                                                                                          'status': {'type': ['string']},
                                                                                          'total': {'type': ['number']},
                                                                                          'items': {'type': ['array'],
                                                                                                    'items': {'type': ['integer']}}},
                                                                           'required': ['id',
                                                                                        'items',
                                                                                        'status',
                                                                                        'total']}},
                                                  'required': ['order']}},
                          'required': ['data']}}}
# ---- END METADATA ----

CLIENT = Api01

if __name__ == "__main__":
    raise SystemExit(_cli(Api01))
