diff --git a/CHANGELOG.md b/CHANGELOG.md index 0b1bfca..8c8d390 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,31 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## [Unreleased] +### Added + +- **An Agent Skill for coding agents.** A model trained before 2026 has never + seen interlock, so an agent asked for a circuit breaker reaches for a + consecutive-failure counter and guesses at the API. + `skills/interlock-cb/SKILL.md` follows the open + [Agent Skills](https://agentskills.io/) format and installs with + `npx skills add bagowix/interlock` into Claude Code, Cursor, Codex, GitHub + Copilot, Gemini CLI and the other clients that read it. The skill is a + procedure: inventory the outbound calls, pick the integration per dependency, + size `Config` from observed traffic, roll out in `METRICS_ONLY`, map + rejections to `503 + Retry-After`, test with an injected clock, and migrate + from pybreaker, circuitbreaker, aiobreaker or purgatory. The reference + material stays in the docs, which the skill links. The README and the docs + landing page point to the install command. + +### Fixed + +- **The source distribution no longer carries the Hypothesis example + database.** The release job runs the test suite before `uv build`, and + Hypothesis leaves its `.hypothesis/` cache in the checkout. Git ignores it + through the nested `.gitignore` Hypothesis writes there, while hatchling + reads only the root one, so the 2.8.0 sdist shipped 37 opaque cache files. + The sdist target now excludes the directory explicitly. + ## [2.8.0] - 2026-09-01 ### Added diff --git a/README.md b/README.md index b72af1a..0ef04dc 100644 --- a/README.md +++ b/README.md @@ -256,6 +256,25 @@ the [`examples/`](https://github.com/bagowix/interlock/tree/main/examples) scripts or follow the [walkthrough](https://bagowix.github.io/interlock/demo/). +## Using with AI coding agents + +The repository ships an [Agent Skill](https://agentskills.io/) that walks a +coding agent through adding interlock to a codebase: an inventory of outbound +calls, the integration to use for each, threshold sizing, a shadow-mode +rollout, tests driven by a fake clock, and the migration from pybreaker or +circuitbreaker. It installs into Claude Code, Cursor, Codex, GitHub Copilot, +Gemini CLI and the other agents that read the open skills format: + +```bash +npx skills add bagowix/interlock +``` + +Then ask the agent to add circuit breakers to a service. Agents that read +documentation directly can use +[llms.txt](https://bagowix.github.io/interlock/llms.txt), the fully inlined +[llms-full.txt](https://bagowix.github.io/interlock/llms-full.txt) or +[Context7](https://context7.com/bagowix/interlock). + ## Contributing Bug reports and pull requests are welcome. See diff --git a/docs/guides/configuration.md b/docs/guides/configuration.md index 034d1d8..1e4fd03 100644 --- a/docs/guides/configuration.md +++ b/docs/guides/configuration.md @@ -59,9 +59,10 @@ Config(window_type=WindowType.TIME_BASED, window_size=30) A dependency that answers slowly but never errors will never trip a failure-rate breaker, yet it still exhausts your timeouts and threads. -Slow-call detection treats latency as a first-class failure signal. By default -`slow_call_rate_threshold=1.0` means slowness alone never trips the breaker -until you tune it down — safe to leave on while you observe. +Slow-call detection treats latency as a first-class failure signal. The default +`slow_call_rate_threshold=1.0` trips only when every call in the window is slow, +so latency is effectively off until you tune it down — safe to leave on while +you observe. ## Sharing config with a Registry diff --git a/docs/index.md b/docs/index.md index 0559d3f..d94465a 100644 --- a/docs/index.md +++ b/docs/index.md @@ -72,6 +72,24 @@ The same instance protects async callables, works as a (sync and async) context manager, and can be called directly via `breaker.call(fn, ...)` — start with [Getting started](getting-started.md). +## AI coding agents + +The repository ships an [Agent Skill](https://agentskills.io/) that walks a +coding agent through adding interlock to a codebase: an inventory of outbound +calls, the integration to use for each, threshold sizing, a shadow-mode +rollout, tests driven by a fake clock, and the migration from pybreaker or +circuitbreaker. It installs into Claude Code, Cursor, Codex, GitHub Copilot, +Gemini CLI and the other agents that read the open skills format: + +```bash +npx skills add bagowix/interlock +``` + +Then ask the agent to add circuit breakers to a service. Agents that read +documentation directly can use [llms.txt](llms.txt), the fully inlined +[llms-full.txt](llms-full.txt) or +[Context7](https://context7.com/bagowix/interlock). + ## Status interlock shipped a polished core first (state machine, windows, sync/async, diff --git a/docs/llms-full.txt b/docs/llms-full.txt index 8deb187..d9ffbb3 100644 --- a/docs/llms-full.txt +++ b/docs/llms-full.txt @@ -1203,8 +1203,9 @@ After migrating, expect these behavioural differences — all intended: Conversely, a single burst of failures below `minimum_number_of_calls` will *not* trip — the window has to fill first. - **Slow calls can trip too, but only when you opt in.** - `slow_call_rate_threshold` defaults to `1.0`, so latency alone never trips - until you tune it down ([configuration](guides/configuration.md#why-slow-calls-matter)). + `slow_call_rate_threshold` defaults to `1.0`, which trips only when every call + in the window is slow, so latency is effectively off until you tune it down + ([configuration](guides/configuration.md#why-slow-calls-matter)). - **Half-open is a budgeted probe round, not a single trial call.** Up to `permitted_calls_in_half_open` probes run (with a concurrency cap), and the breaker re-decides from their rate ([states](guides/states.md)). @@ -1543,9 +1544,10 @@ Config(window_type=WindowType.TIME_BASED, window_size=30) A dependency that answers slowly but never errors will never trip a failure-rate breaker, yet it still exhausts your timeouts and threads. -Slow-call detection treats latency as a first-class failure signal. By default -`slow_call_rate_threshold=1.0` means slowness alone never trips the breaker -until you tune it down — safe to leave on while you observe. +Slow-call detection treats latency as a first-class failure signal. The default +`slow_call_rate_threshold=1.0` trips only when every call in the window is slow, +so latency is effectively off until you tune it down — safe to leave on while +you observe. ## Sharing config with a Registry diff --git a/docs/llms.txt b/docs/llms.txt index 04baa21..96a8047 100644 --- a/docs/llms.txt +++ b/docs/llms.txt @@ -73,4 +73,5 @@ Key facts for answering questions about interlock: - [Full documentation](llms-full.txt): every page above inlined into one file. - Extras: `interlock-cb[fastapi]` (FastAPI 503 mapping), `interlock-cb[litestar]` (Litestar 503 mapping), `interlock-cb[httpx]` / `[httpx2]` (transports), `interlock-cb[aiohttp]` (client middleware), `interlock-cb[requests]` (session adapter), `interlock-cb[tenacity]` (retry glue), `interlock-cb[redis]` (shared distributed state), `interlock-cb[otel]` (OpenTelemetry metrics). +- Agent skill: `npx skills add bagowix/interlock` installs [skills/interlock-cb/SKILL.md](https://github.com/bagowix/interlock/blob/main/skills/interlock-cb/SKILL.md), a workflow for adding, sizing, rolling out and testing breakers with a coding agent. - Source and changelog: see the repository `CHANGELOG.md`. diff --git a/docs/migration.md b/docs/migration.md index 0edfe2a..f955833 100644 --- a/docs/migration.md +++ b/docs/migration.md @@ -367,8 +367,9 @@ After migrating, expect these behavioural differences — all intended: Conversely, a single burst of failures below `minimum_number_of_calls` will *not* trip — the window has to fill first. - **Slow calls can trip too, but only when you opt in.** - `slow_call_rate_threshold` defaults to `1.0`, so latency alone never trips - until you tune it down ([configuration](guides/configuration.md#why-slow-calls-matter)). + `slow_call_rate_threshold` defaults to `1.0`, which trips only when every call + in the window is slow, so latency is effectively off until you tune it down + ([configuration](guides/configuration.md#why-slow-calls-matter)). - **Half-open is a budgeted probe round, not a single trial call.** Up to `permitted_calls_in_half_open` probes run (with a concurrency cap), and the breaker re-decides from their rate ([states](guides/states.md)). diff --git a/pyproject.toml b/pyproject.toml index e886830..4e7ddec 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -54,6 +54,12 @@ default = true [tool.hatch.build.targets.wheel] packages = ["interlock"] +[tool.hatch.build.targets.sdist] +# The release job runs the tests before `uv build`, and Hypothesis leaves its +# example database in the checkout. Git ignores it through the nested +# .gitignore Hypothesis writes there; hatchling reads only the root one. +exclude = [".hypothesis/"] + [tool.ruff] src = ["interlock", "tests"] line-length = 100 @@ -190,6 +196,7 @@ plugins.MD029.style = "ordered" # consistent ordered list numbering plugins.MD033.enabled = false # allow inline HTML (useful for docs) plugins.MD046.enabled = false # pymdownx.tabbed indents fenced blocks; the style check misfires plugins.MD041.enabled = false # do not require an H1 at the top of every file +extensions.front-matter.enabled = true # SKILL.md carries YAML front matter (Agent Skills spec) [tool.pytest.ini_options] testpaths = ["tests"] diff --git a/skills/interlock-cb/SKILL.md b/skills/interlock-cb/SKILL.md new file mode 100644 index 0000000..9a411b3 --- /dev/null +++ b/skills/interlock-cb/SKILL.md @@ -0,0 +1,258 @@ +--- +name: interlock-cb +description: Add, size, roll out and test circuit breakers in Python 3.11+ code with interlock-cb (import interlock), or migrate to it from pybreaker, circuitbreaker, aiobreaker or purgatory. Covers per-host guards for httpx, httpx2, aiohttp and requests, decorators for SDK and database calls, threshold sizing, shadow-mode rollout, 503 + Retry-After mapping in FastAPI, Litestar, Flask or Django, retry and timeout composition, bulkheads and fallbacks. Use when a Python project needs a circuit breaker, fail-fast protection from a slow or failing dependency, an end to retry storms or cascading failures, or when the user mentions interlock. +license: MIT +metadata: + author: bagowix + docs: https://bagowix.github.io/interlock/ +--- + +# interlock-cb: add a circuit breaker the right way + +interlock-cb is a circuit breaker for Python 3.11+. One `CircuitBreaker` class guards sync and async callables as a decorator, a context manager or `breaker.call(fn, ...)`. It trips on the **failure rate over a sliding window** (count- or time-based) and can count **slow calls** as failures. The core has zero dependencies; HTTP clients, web frameworks, tenacity, Redis and OpenTelemetry plug in as optional extras. + +The library is from 2026, so your training data most likely predates it. Do not guess the API from pybreaker or circuitbreaker: every symbol you need is in this file, and the documentation links at the end cover the rest. + +## Ground rules + +- Package `interlock-cb`, import `interlock`. Python 3.11 or newer; if the project is older, say so and stop. pybreaker or circuitbreaker remain the right choice there. +- One breaker per dependency, keyed by host or by a named downstream. A failing host must never trip a healthy one. +- Thresholds come from observed traffic. Ship in shadow mode first (step 5), tune, then enforce. +- Every layer stays explicit: no hidden retries, no silent fallbacks, no swallowed `CircuitOpenError`. +- `Config` is a frozen dataclass validated at construction: to change a value, build a new one. +- Show the inventory and the plan to the user before editing code. + +## Workflow + +### 1. Inventory outbound calls + +Search the code for clients and for resilience code that already exists: + +- HTTP: `httpx.`, `httpx2.`, `aiohttp.ClientSession`, `requests.` (`Session`, `get`, `post`), `urllib3`. +- SDKs: `openai`, `anthropic`, `boto3`, `google.cloud`, `stripe`, `twilio`. +- Data and messaging: `redis`, `psycopg`, `asyncpg`, `sqlalchemy`, `pymongo`, `motor`, `aiokafka`, `pika`, `grpc`. +- Already present: `pybreaker`, `circuitbreaker`, `aiobreaker`, `purgatory`, `tenacity`, `backoff`, urllib3 `Retry(`, `max_retries=`, `timeout=`. + +Produce a table: dependency, call sites, client library, sync or async, existing timeout, existing retries. Each row gets one decision in step 2. + +### 2. Pick the guard per dependency + +| Client | Guard | Extra | +|---|---|---| +| httpx | `CircuitBreakerTransport(httpx.HTTPTransport())` or `AsyncCircuitBreakerTransport(httpx.AsyncHTTPTransport())` from `interlock.integrations.httpx`, passed as `transport=` to the client | `interlock-cb[httpx]` | +| httpx2 | the same two classes from `interlock.integrations.httpx2` | `interlock-cb[httpx2]` | +| aiohttp 3.12+ | `CircuitBreakerMiddleware()` from `interlock.integrations.aiohttp`, passed as `middlewares=(middleware,)`; call `await middleware.aclose()` at shutdown | `interlock-cb[aiohttp]` | +| requests | `CircuitBreakerAdapter()` from `interlock.integrations.requests`, mounted on the `Session` for both `https://` and `http://` | `interlock-cb[requests]` | +| OpenAI or Anthropic SDK | wrap the SDK's httpx transport (every endpoint of the host shares one breaker), or decorate the call with an SDK error classifier when you need slow-call timing of the whole operation; set `max_retries=0` on the SDK client so one layer owns retries | `interlock-cb[httpx]` | +| anything else | a named `CircuitBreaker` from a shared `Registry`: decorator, `call()` or context manager | none | + +HTTP integrations create one breaker per host lazily, so several hosts behind one client get separate breakers from a single transport, middleware or adapter, and its `registry` property lists them. Mount one `CircuitBreakerAdapter` instance for both URL prefixes. Responses are classified with the `HttpStatusClassifier` that each integration module exports (import it from the same module as the transport, middleware or adapter; `interlock` itself does not export it): `429, 500, 502, 503, 504` and transport exceptions count as failures, other `4xx` do not. Change the policy with `classifier=HttpStatusClassifier(failure_statuses={...}, excluded_exceptions=(...))`. A rejection raises a subclass of `CircuitOpenError` that is also a native error of the client (`CircuitOpenTransportError`, `CircuitOpenClientError`, `CircuitOpenRequestError`), so existing `except` blocks keep working. + +Every integration takes the same keyword arguments: `config`, `clock`, `initial_state`, `classifier`, `listener` and `name_resolver`, or a ready `registry` in place of the first five (combining both raises `ValueError`). A registry you supply replaces the integration's default classifier, so build it with `classifier=HttpStatusClassifier()` from that module; otherwise a returned `503` counts as a success and the breaker never trips on HTTP failures. Keep such a registry to HTTP clients: the classifier reads `status_code` off every result and raises on anything else. + +A service with several HTTP clients keeps one `Registry` per process and gives each client its own transport built on it, so a host reached from two clients is one breaker. With a registry you build yourself, also pass `unreachable_exceptions=(httpx.PoolTimeout,)`: the transport sets it only for a registry it creates, and without it an exhausted pool fails every half-open probe and the breaker stays open. Two httpx facts to check: a client given `transport=` ignores its own `limits=`, `verify=` and the proxy environment, so configure those on the inner transport; and the transport observes only the time to response headers, so a failure while reading a streamed body never reaches it. The factory, settings model, Prometheus listener, rollout stages and tests of a layout that has carried production traffic are in [references/production-httpx-service.md](references/production-httpx-service.md). + +Transport level, async httpx: + +```python +import httpx + +from interlock import Config, State +from interlock.integrations.httpx import AsyncCircuitBreakerTransport + +transport = AsyncCircuitBreakerTransport( + httpx.AsyncHTTPTransport(), + config=Config(failure_rate_threshold=0.5, minimum_number_of_calls=20), + initial_state=State.METRICS_ONLY, +) +client = httpx.AsyncClient(transport=transport, timeout=5.0) +``` + +Call level, any dependency: + +```python +from interlock import CircuitOpenError, Config, Registry, State + +registry = Registry( + config=Config(minimum_number_of_calls=20, slow_call_duration_threshold=2.0), + initial_state=State.METRICS_ONLY, +) +payments = registry.get('payments-gateway') + + +@payments +async def charge(amount: int) -> str: + return await gateway.charge(amount) + + +try: + await charge(100) +except CircuitOpenError as exc: + ... # exc.retry_after: seconds until the next probe, or None +``` + +The decorator keeps the wrapped signature for type checkers. The context manager (`with breaker:` / `async with breaker:`) sees only exceptions and duration, so use the decorator or `call()` when failure lives in a return value. + +### 3. Size the config + +| Field | Default | How to choose | +|---|---|---| +| `failure_rate_threshold` | `0.5` | Fraction of failed calls in the window that trips. Range `(0, 1]`. | +| `minimum_number_of_calls` | `10` | Calls needed before the rate is trusted. Raise it for busy dependencies (50), lower it for quiet ones (5). | +| `window_type`, `window_size` | `COUNT_BASED`, `100` | Last N calls, or last N seconds with `WindowType.TIME_BASED`. Time-based suits high throughput. | +| `slow_call_duration_threshold` | `60.0` | Seconds at or above which a call is slow. Set it to the client timeout or just above the p99 latency. | +| `slow_call_rate_threshold` | `1.0` | Fraction of slow calls that trips. `1.0` trips only when every call in the window is slow, so latency is effectively off until you tune it down; try `0.5` to `0.8` once shadow data exists. | +| `wait_duration_in_open` | `60.0` | Seconds to stay open before the first probe. | +| `wait_duration_backoff_multiplier`, `wait_duration_in_open_max` | `1.0`, `None` | Grow the wait after each failed probe round and cap it. The multiplier must stay `1.0` with a shared storage. | +| `permitted_calls_in_half_open`, `max_concurrent_probes` | `10`, `1` | Probe budget per half-open round, and how many probes run at once. `max_concurrent_probes` must stay within `[1, permitted_calls_in_half_open]`. | +| `auto_transition` | `False` | `True` lets a timer flip OPEN to HALF_OPEN as soon as the wait elapses, without waiting for a call; the first real call is still the first probe. | + +Validation raises `ValueError` at construction, so a bad value never reaches production. + +Slow calls are a second dimension. A call at or above `slow_call_duration_threshold` feeds `slow_call_rate`, and a slow success stays a success in `failure_rate`. Once the window holds `minimum_number_of_calls`, either rate reaching its threshold trips the breaker. + +Failure classification is separate from thresholds. The default counts every raised exception as a failure and every returned value as a success. Write a classifier when failure is a return value, or when business errors (a `404`, a validation error) must stay out of the rate: + +```python +class IgnoreNotFound: + def is_failure(self, *, result: object, exception: Exception | None) -> bool: + if isinstance(exception, NotFoundError): + return False + return exception is not None +``` + +Pass it as `classifier=`. For LLM SDKs, count `429, 500, 502, 503, 504, 529` and the SDK's connection and timeout errors; a `400` or `404` is the caller's bug. + +### 4. Compose timeouts and retries + +- Put the timeout inside the guarded call so the breaker records `CallTimeoutError`: `async with timeout(2.0):` from `interlock`, or the `sync_timeout(2.0)` decorator for blocking code. For HTTP transports the client timeout does the same job. +- Put retries outside the breaker and stop the moment the circuit opens. `retry_unless_open(*transient)` from `interlock.integrations.tenacity` is the tenacity predicate for that; `CircuitOpenError` is never transient. The one exception is background work that prefers waiting for the next probe over failing: there the predicate must include `CircuitOpenError` (`retry_if_exception_type((TimeoutError, CircuitOpenError))`) and the wait is `wait_probe(wait_exponential_jitter())`, which sleeps `retry_after` after a rejection and defers to the wrapped strategy otherwise. Pick one mode per call site. +- An existing tenacity decorator stays where it is: change its `retry=` to `retry_unless_open()` and keep its `stop`. With a transport, middleware or adapter, a retry decorator on the calling function already sits outside the breaker. +- Keep one retry layer. Disable SDK retries (`max_retries=0`) and urllib3 `Retry` when tenacity owns retrying. +- Retry predicates must not match the rejection. The typed rejections descend from the client library's error hierarchy, so a predicate on `httpx.TransportError`, `requests.exceptions.RequestException` or `aiohttp.ClientError` retries every rejection; key it on leaf types (`httpx.ConnectError`, `httpx.ReadTimeout`) or use `retry_unless_open`. +- When several concerns stack, use the pipeline. Order is explicit and the first strategy is the outermost: + +```python +from interlock import CircuitOpenError, Pipeline + +pipeline = ( + Pipeline.builder() + .fallback(lambda exc: [], on=(CircuitOpenError,)) # only for the listed errors + .retry(attempts=4) # interlock-cb[tenacity]; never retries an open circuit + .circuit_breaker(breaker) + .bulkhead(8) + .timeout(2.0) + .build() +) + + +@pipeline +async def fetch_picks(user: str) -> list[str]: + return await client.get_picks(user) +``` + +### 5. Roll out in shadow mode + +1. Deploy with `initial_state=State.METRICS_ONLY` on the breaker, `Registry` or transport, plus a listener: `LoggingEventListener()` from `interlock`, or `OTelEventListener()` from `interlock.integrations.otel` (`interlock-cb[otel]`). The breaker records real failure and slow-call rates and rejects nothing. +2. Read `breaker.snapshot()` over real traffic: it carries `total_calls`, `failed_calls`, `slow_calls`, `failure_rate` and `slow_call_rate`. `transport.registry.items()` lists the per-host breakers an HTTP integration has created so far. Adjust `Config`. +3. Deploy again with the default `initial_state=State.CLOSED`. The enforcing instance starts with a fresh window. + +Keep the mode and every threshold in deployment configuration, so the switch from shadow to enforcing is a config change and never a release; the settings model in the reference above does this. Leave observability exporters (metrics, tracing, logging backends) unguarded: a breaker there can degrade the process recursively. + +Map rejections at the web boundary so callers get `503` with `Retry-After`: + +- FastAPI (`interlock-cb[fastapi]`): `install_exception_handler(app)` once, then `breaker_dependency('orders-db', registry=registry)` behind `Depends` injects a breaker into a route. +- Litestar (`interlock-cb[litestar]`): `dependencies={'breaker': breaker_dependency('orders-db', registry=registry)}` and `exception_handlers={CircuitOpenError: circuit_open_handler}` on the app. +- Flask or Django: an error handler that returns `503` and sets `Retry-After` to `math.ceil(exc.retry_after)` when it is not `None`. The framework recipe in the docs has both. + +Share state through Redis (`interlock-cb[redis]`) only when many instances call the same downstream and should back off together under one global probe budget: `storage=RedisStorage(redis.Redis(...))` for sync code, `AsyncRedisStorage(redis.asyncio.Redis(...))` for async. A coordinated breaker serves only its storage's runtime. An unreachable Redis degrades the breaker to local state, never to a dead service. A coordinated breaker runs a background lane, so close it deterministically at shutdown with `breaker.close()` / `await breaker.aclose()`, or `registry.close_all()` / `await registry.aclose_all()`. Stay local when instances see different views of the dependency. + +### 6. Write tests without sleeping + +Time enters the breaker only through the injected `Clock`. Give the breaker a fake one and drive transitions explicitly. After `wait_duration_in_open` the next call is admitted as a probe and the state becomes `HALF_OPEN`; the round ends after `permitted_calls_in_half_open` probes, and the breaker closes when their failure and slow-call rates stay under the thresholds and re-opens otherwise. With the default of 10 probes, one successful probe leaves the breaker `HALF_OPEN`: + +```python +import pytest + +from interlock import CircuitBreaker, CircuitOpenError, Config, State + + +class FakeClock: + def __init__(self) -> None: + self.now = 0.0 + + def monotonic(self) -> float: + return self.now + + +def test__five_failures__opens_then_probes_after_wait() -> None: + clock = FakeClock() + breaker = CircuitBreaker( + name='svc', + config=Config(minimum_number_of_calls=5, wait_duration_in_open=30.0), + clock=clock, + ) + + def fail() -> None: + raise ConnectionError + + for _ in range(5): + with pytest.raises(ConnectionError): + breaker.call(fail) + assert breaker.state is State.OPEN + + with pytest.raises(CircuitOpenError): + breaker.call(fail) + + clock.now += 30.0 + assert breaker.call(lambda: 'ok') == 'ok' # the first probe is admitted + assert breaker.state is State.HALF_OPEN +``` + +Async callables get the same treatment through `await breaker.call(...)`. + +### 7. Migrate from pybreaker, circuitbreaker, aiobreaker or purgatory + +Those libraries trip on a streak of consecutive failures; interlock trips on a rate, so the numbers do not carry over. + +| Before | After | +|---|---| +| `fail_max=5` / `failure_threshold=5` | `Config(minimum_number_of_calls=5, failure_rate_threshold=0.8)` as a starting point, then tune in shadow mode | +| `reset_timeout` / `recovery_timeout` | `wait_duration_in_open` | +| pybreaker `exclude=[...]` (exceptions that do not count) | a `FailureClassifier` that returns `False` for them (step 3) | +| circuitbreaker `expected_exception=...` (the inverse: only these count) | a `FailureClassifier` that returns `True` only for them | +| listeners | an `EventListener` passed as `listener=` | +| `fallback_function` | `FallbackStrategy` or `.fallback(...)` in a pipeline, with an explicit `on=` | +| Redis-backed state | `RedisStorage` / `AsyncRedisStorage` | + +Keep the old breaker in place until shadow-mode data supports the new thresholds. The migration page has the full mapping per library. + +## Anti-patterns to refuse + +- One breaker shared by every host or dependency. +- Catching `CircuitOpenError` and retrying at once, or sleeping `2**n` seconds. Use `retry_unless_open`, or `wait_probe` when waiting is intended. +- Counting `4xx`, validation errors or cancellation as failures. +- Two retry layers on one call (SDK plus tenacity, urllib3 plus tenacity). +- A slow-call rate threshold without a duration threshold tied to real latency. +- `except Exception: return default` around a guarded call. Use a fallback with a listed exception type. +- `time.sleep` in tests. +- A shared storage together with `wait_duration_backoff_multiplier` above `1.0`, or a sync storage on an async call path. Both raise. +- Passing `registry=` together with `config`, `clock`, `initial_state`, `classifier` or `listener` to an integration. It raises; configure the registry instead. +- Handing an HTTP integration a registry built without `HttpStatusClassifier`: returned error statuses become successes. +- A retry predicate keyed on the client library's broad error base: it retries the typed rejection. + +## Report back + +End with the inventory table, the decision per dependency with the chosen thresholds and the reason, the rollout plan (shadow, observe, enforce), the tests added, and anything left out. + +## Documentation + +- Index: +- Everything inlined: +- Configuration: +- States and shadow mode: +- Retries: +- Pipeline: +- Integrations: +- Migration: +- API reference: diff --git a/skills/interlock-cb/references/production-httpx-service.md b/skills/interlock-cb/references/production-httpx-service.md new file mode 100644 index 0000000..d256a5f --- /dev/null +++ b/skills/interlock-cb/references/production-httpx-service.md @@ -0,0 +1,257 @@ +# Production layout for a service with several httpx clients + +This layout has carried production traffic: one process, several `httpx.AsyncClient` instances (a shared REST and JSON-RPC client plus SDK clients that accept an httpx client), every host guarded by one breaker registry, rolled out in shadow mode and observed through Prometheus. Names are placeholders; adapt them to the project's conventions. The snippets build on each other in order. + +## 1. Settings: thresholds live in deployment config + +Keep every `Config` field and the operating mode in the application settings, so the switch from shadow mode to enforcing is a configuration change and never a release: + +```python +import ssl +from typing import Literal + +import httpx +from interlock import Config, CoreEventListener, Registry, State, WindowType +from interlock.integrations.httpx import AsyncCircuitBreakerTransport, HttpStatusClassifier +from pydantic import BaseModel, Field + +BreakerMode = Literal[State.CLOSED, State.DISABLED, State.METRICS_ONLY] + + +class BreakerSettings(BaseModel): + mode: BreakerMode = State.METRICS_ONLY + failure_rate_threshold: float = Field(default=0.5, gt=0, le=1) + minimum_number_of_calls: int = Field(default=10, ge=1) + slow_call_duration_seconds: float = Field(default=5, gt=0) + slow_call_rate_threshold: float = Field(default=1, gt=0, le=1) + permitted_calls_in_half_open: int = Field(default=10, ge=1) + max_concurrent_probes: int = Field(default=1, ge=1) + wait_duration_in_open_seconds: float = Field(default=60, gt=0) + wait_duration_backoff_multiplier: float = Field(default=2, ge=1) + wait_duration_in_open_max_seconds: float = Field(default=300, gt=0) + window_type: WindowType = WindowType.TIME_BASED + window_size: int = Field(default=60, ge=1) + failure_statuses: frozenset[int] = frozenset({429, 500, 502, 503, 504}) + + def build_config(self) -> Config: + return Config( + failure_rate_threshold=self.failure_rate_threshold, + minimum_number_of_calls=self.minimum_number_of_calls, + slow_call_duration_threshold=self.slow_call_duration_seconds, + slow_call_rate_threshold=self.slow_call_rate_threshold, + permitted_calls_in_half_open=self.permitted_calls_in_half_open, + max_concurrent_probes=self.max_concurrent_probes, + wait_duration_in_open=self.wait_duration_in_open_seconds, + wait_duration_backoff_multiplier=self.wait_duration_backoff_multiplier, + wait_duration_in_open_max=self.wait_duration_in_open_max_seconds, + window_type=self.window_type, + window_size=self.window_size, + ) +``` + +Decisions behind it: + +- `mode` is a `Literal` over `State` values: one source of truth, no parallel enum. `FORCED_OPEN` is deliberately absent. As the initial state of every breaker, it turns a healthy process into one that rejects every outgoing call from startup, a silent failure where a crash would be honest. +- The default mode is shadow, so a fresh environment observes before it enforces. The half-open and open-wait fields are exposed from day one even though shadow mode ignores them: enabling `CLOSED` later touches only configuration. +- With the model above, an omitted deployment key silently uses the model default. A deployment that injects these values from a secret store resolved at process start fails on a missing key. Create the keys in every environment before the first deployment either way, so the values in production are the ones you tuned. +- A time-based window of 60 seconds suits uneven traffic; the slow-call duration sits near the client read timeout; a backoff multiplier of 2 capped at 300 seconds stops a dead dependency from being probed at full rate. +- `failure_statuses` can live in the model, but it is a classifier. Keep it out of the per-environment operational config. + +## 2. One registry per process, one transport per client + +```python +def breaker_name(request: httpx.Request) -> str: + # Normalise the host (strip an internal DNS suffix, lower-case it) so the + # breaker name equals the label the dashboards already use. + return request.url.host.removesuffix('.internal') + + +class BreakerTransportFactory: + def __init__(self, settings: BreakerSettings, listener: CoreEventListener) -> None: + # One registry per process: a host reached from several clients must be + # one breaker, otherwise minimum_number_of_calls fills at a fraction of + # the traffic and tripping is asymmetric between clients. + self._registry = Registry( + config=settings.build_config(), + initial_state=settings.mode, + classifier=HttpStatusClassifier(failure_statuses=settings.failure_statuses), + listener=listener, + # The transport sets this only for a registry it creates itself. A + # probe that never got a connection says nothing about the + # dependency; in CLOSED an exhausted pool stays a failure. + unreachable_exceptions=(httpx.PoolTimeout,), + ) + + @property + def registry(self) -> Registry: + return self._registry + + def create( + self, + transport: httpx.AsyncBaseTransport | None = None, + *, + limits: httpx.Limits | None = None, + verify: ssl.SSLContext | str | bool = True, + ) -> AsyncCircuitBreakerTransport: + # httpx applies limits, verify and proxy settings only to a transport it + # creates itself; a client given transport= drops them silently. limits + # and verify are configured here. A proxy is not: pass a transport built + # with proxy=... as the argument when the deployment needs one. + if transport is None: + transport = ( + httpx.AsyncHTTPTransport(verify=verify) + if limits is None + else httpx.AsyncHTTPTransport(verify=verify, limits=limits) + ) + + return AsyncCircuitBreakerTransport( + transport, + registry=self._registry, + name_resolver=breaker_name, + ) + + async def close(self) -> None: + # Last, after every client that uses a transport from this factory. + await self._registry.aclose_all() +``` + +Wiring: build the factory once at startup and hand each root client `transport=factory.create(...)`. An SDK client that accepts an httpx client gets one built the same way, with the SDK's own default `Limits` passed through (read the SDK constructor; the defaults often match httpx's, and relying on that is a trap). Shutdown order: SDK clients, the shared client, the factory. Use an exit stack, so the registry closes even when a client's close raises. + +Facts about httpx that bite, verified against httpx 0.28: + +- `allow_env_proxies = trust_env and transport is None`: a custom transport also disables `HTTP_PROXY` and `NO_PROXY` detection. Check the deployment manifests before the first rollout. +- Closing one client's transport leaves a shared registry alone; a transport closes only a registry it created. The owner closes the shared one, once. + +What to leave unwrapped: observability exporters (metrics, tracing, logging backends), where a breaker can degrade the process recursively; message brokers and database drivers, which are not HTTP and deserve a call-level decision of their own. + +Streaming: the transport observes the time to response headers. A failure while reading a streamed body (server-sent events, long downloads) is invisible to the transport breaker. Document that limitation in the runbook, or add a call-level breaker around the whole stream as a separate change. + +## 3. Rejections and retries + +`CircuitOpenTransportError` is an `httpx.TransportError`, so it lands in existing `except httpx.TransportError` blocks. Treat it as routine: log at WARNING with `retry_after`, answer `503` with `Retry-After` or whatever the edge policy is, and never log it as an exception. + +Retry predicates must not match the rejection. Key them on leaf types (`httpx.ConnectError`, `httpx.ReadTimeout`, a status-error predicate) or use `retry_unless_open`. A predicate on the base class (`httpx.TransportError`, `requests.exceptions.RequestException`, `aiohttp.ClientError`) retries every rejection, and the attempts burn in microseconds against a circuit that stays open. + +Retries sit above the transport, so the breaker sees every physical attempt. Its failure rate is per attempt, never per user operation. Say so on the dashboard and in the runbook. + +## 4. Metrics without the OpenTelemetry extra + +When the process exports Prometheus directly and runs no `MeterProvider`, a short listener does the job: + +```python +from collections.abc import Iterator + +from interlock import Outcome +from prometheus_client import Counter, Histogram +from prometheus_client.core import GaugeMetricFamily +from prometheus_client.registry import Collector + +# Buckets must reach past slow_call_duration_threshold; the prometheus_client +# defaults stop at 10 s and hide the p99 the threshold is tuned against. +BUCKETS = (0.05, 0.1, 0.25, 0.5, 1, 2, 5, 10, 30, 60, 120, 300) + +CALL_DURATION = Histogram( + 'circuit_breaker_call_duration_seconds', + 'Calls observed by circuit breakers.', + labelnames=('breaker', 'outcome'), + buckets=BUCKETS, +) +REJECTED = Counter( + 'circuit_breaker_rejected_total', + 'Calls rejected by an open circuit.', + labelnames=('breaker',), +) +STATE_CHANGES = Counter( + 'circuit_breaker_state_changes_total', + 'Circuit breaker transitions.', + labelnames=('breaker', 'from', 'to'), +) +RESETS = Counter('circuit_breaker_resets_total', 'Manual resets.', labelnames=('breaker',)) + + +class PrometheusBreakerListener(CoreEventListener): + def on_state_change(self, *, name: str, old: State, new: State) -> None: + STATE_CHANGES.labels(breaker=name, **{'from': str(old), 'to': str(new)}).inc() + + def on_call(self, *, name: str, outcome: Outcome, duration: float) -> None: + CALL_DURATION.labels(breaker=name, outcome=str(outcome)).observe(duration) + + def on_rejected(self, *, name: str) -> None: + REJECTED.labels(breaker=name).inc() + + def on_reset(self, *, name: str) -> None: + RESETS.labels(breaker=name).inc() + + +class BreakerStateCollector(Collector): + """Current state per breaker at scrape time. + + Transition counters answer "what happened"; during an incident the question + is "what is open right now", and reconstructing it from deltas takes time + nobody has. The registry knows; ask it. + """ + + def __init__(self, registry: Registry) -> None: + self._registry = registry + + def collect(self) -> Iterator[GaugeMetricFamily]: + gauge = GaugeMetricFamily( + 'circuit_breaker_state', + 'Current circuit breaker state, 1 for the active one.', + labels=('breaker', 'state'), + ) + + for name, breaker in self._registry.items(): + current = breaker.state + for state in State: + gauge.add_metric((name, str(state)), float(state is current)) + + yield gauge +``` + +Register the collector once per process with `prometheus_client.REGISTRY.register(...)`. Labels are the breaker name and the outcome, nothing else: no URL path, user or exception text. Dashboards per breaker: the share of `failure` plus `slow_failure`, the share of `slow_success` plus `slow_failure`, p95 and p99 duration, call volume, and `rejected_total`, which must stay at zero in shadow mode (alert on it). + +`DISABLED` stops only the sliding window; listener events keep flowing, so switching a breaker off does not blank the dashboards. That is intended: an empty panel reads as "no traffic". + +## 5. Rollout stages + +1. Pin the version (`interlock-cb[httpx]~=2.8`). +2. Settings, factory, listener, unit tests (section 6). +3. Wrap the shared client first: one change covers every REST and JSON-RPC client built on it. Prove on a test host that repeated `503` responses are recorded and nothing is rejected, and that request hooks, tracing, per-request timeouts and streaming still work. +4. SDK clients: pass the transport into the SDK constructor, or build the equivalent client by hand when the SDK's factory accepts none, keeping base URL, auth headers, timeouts, limits and TLS settings. Verify close on shutdown and the absence of pool leaks. +5. Observe at least one full business cycle. Treat low-traffic dependencies separately: their window may never reach `minimum_number_of_calls`. Record a baseline per host: volume, failure and slow-call rates, retry amplification, typical degradation windows. +6. A separate, reviewed change enables `CLOSED`, with a canary and a rollback path. + +## 6. Tests to write + +All run against `httpx.MockTransport`, no network: + +- Shadow mode records failures and rejects nothing: `snapshot().failed_calls` grows, no `CircuitOpenTransportError`. +- A custom failure status is recorded; a caller-side exception (`httpx.LocalProtocolError`) propagates and is not counted. +- Two hosts get independent breakers; two transports from the factory share one breaker per host (`first.registry is second.registry`). +- Closing one transport keeps the shared registry usable for the others. +- `limits` and `verify` reach the inner transport (`transport.wrapped`); httpx defaults stay when they are not given. +- Each mode maps to the expected initial state. +- A `PoolTimeout` on a half-open probe does not keep the breaker open. +- A rejection is not retried, while a `ConnectError` still is. +- `/metrics` exposes the series after a success, a retryable status and a transport error. + +The first one, as a pattern for the rest: + +```python +async def test__shadow_mode__records_failures_without_rejecting() -> None: + settings = BreakerSettings(minimum_number_of_calls=2, window_size=10) + factory = BreakerTransportFactory(settings, listener=PrometheusBreakerListener()) + transport = factory.create(httpx.MockTransport(lambda _request: httpx.Response(503))) + + async with httpx.AsyncClient(transport=transport) as client: + for _ in range(5): + response = await client.get('https://payments.internal/charge') + assert response.status_code == 503 + + breaker = transport.registry.get_existing('payments') + + assert breaker is not None + assert breaker.state is State.METRICS_ONLY + assert breaker.snapshot().failed_calls == 5 +```