Service Client & Change-Feed Watcher¶
hyperdjango.serviceclient provides two reusable building blocks for programs
that call an internal JSON-over-HTTP service and follow its change feed:
ServiceClient, a retrying JSON transport, and ChangeFeedWatcher, an
ordered, self-healing feed consumer.
The module is stdlib-only (urllib, socket, ssl, a minimal RFC 6455
WebSocket client) and depends on no application. It is the shared foundation
for outbound SDKs: instead of each SDK hand-rolling a retry loop, a bearer
header, an mTLS context, and a feed watcher, it wraps ServiceClient and
ChangeFeedWatcher and parameterizes the specifics (base URL, paths, auth
header, error text).
When to use it¶
Use ServiceClient whenever a program makes JSON HTTP calls to a service and
wants bounded, correct retries and typed errors. Use ChangeFeedWatcher when
that service exposes a change feed and the program must consume it without
losing events across disconnects. One watcher speaks three delivery models —
durable ledger, ephemeral, and catchup — and runs whichever the hub
advertises in its hello frame (see Delivery models).
ServiceClient¶
from hyperdjango.serviceclient import ServiceClient, RetryPolicy
client = ServiceClient(
"https://svc.internal:8960",
token="hsk_...", # optional credential
token_header="Authorization", # or "X-API-Key", etc.
token_scheme="Bearer", # "" sends the raw token (API-key style)
timeout=5.0,
retry=RetryPolicy(max_attempts=3, base_backoff=0.1, max_backoff=10.0),
ca_file="ca.crt", # optional mTLS identity
client_cert_file="client.crt",
client_key_file="client.key",
)
data = client.request("GET", "/v1/things", params={"limit": "50"})
client.request("POST", "/v1/things", json_body={"name": "x"})
request(method, path, *, json_body=None, params=None, idempotent=None)
returns the parsed JSON body (or None when the response is empty). A 2xx
whose body is not decodable JSON raises ResponseError — a 200 carrying a
captive-portal or proxy error page is a contract violation, not data.
Raw-status requests¶
request_raw(method, path, ...) -> (status, headers, body) runs the exact same
retry, backoff, TLS, and no-redirect machinery as request, but returns the
HTTP status instead of mapping a non-2xx to an exception. Use it when a status
is a first-class result rather than an error — a conditional 304, or a 404 /
409 a caller wants to branch on:
status, headers, body = client.request_raw(
"GET", "/v1/secrets/x", params={"known_version": "7"}
)
if status == 304:
... # not modified — serve the cached copy
body is the parsed JSON when the body decodes, else None. A transport
failure still retries per the policy and raises ServiceUnavailable when
exhausted, and a body over the size cap still raises ResponseError; only
definitive HTTP statuses are returned rather than raised.
Redirects and size caps¶
- Redirects are never followed. A JSON API has no valid
3xxanswer, and following one would re-send the credential to the redirect target (possibly a different host). A3xxsurfaces as aRequestErrorfromrequestand as the returned status fromrequest_raw. - Every response body is capped at
max_response_bytes(default 32 MiB) — success and error alike. An over-cap2xxbody raisesResponseError; an over-cap error body degrades to an emptydetailand the status still maps to its typed error, so a hostile or misconfigured server cannot balloon memory through either path. A transport failure while reading a body — an error body, or a2xxbody cut short by a mid-flight reset (anIncompleteReadfrom a truncated chunked stream, aBadStatusLinefrom a torn connection) — is caught like any other transport failure: it retries per the policy and surfaces typed (aServiceErrorsubtype), never as a bareOSErrororIncompleteRead.
Retry contract¶
Retries apply only to idempotent requests. idempotent defaults to True
for GET/HEAD/OPTIONS and False otherwise; pass idempotent=True
explicitly for a safe-to-repeat POST (for example one carrying a dedupe key).
- HTTP status responses are definitive — a
4xxor5xxis never retried. The server saw the request and answered; repeating it is the caller's decision, not the transport's. - Transport failures (connection refused, timeout, reset, or a body read
cut short mid-flight —
URLError,TimeoutError,OSError, or anhttp.client.HTTPExceptionsuch asIncompleteRead/BadStatusLine) retry up to the policy, then raiseServiceUnavailable.
RetryPolicy backoff before retry n (0-based) is
min(base_backoff * 2**n, max_backoff) plus uniform jitter in
[0, base_backoff). The jitter breaks up synchronized retry storms; the cap
bounds the worst-case wait. base_backoff and max_backoff must be
non-negative — 0 is allowed and means retry immediately (no wait), a negative
value is rejected at construction.
Local-resource exhaustion is retried out-of-band. A connect-time
EADDRNOTAVAIL / EADDRINUSE (ephemeral-port starvation behind a busy NAT, or
a flood of short-lived connections leaving sockets in TIME_WAIT) is a purely
local condition — the request never left the host — so it is not the server's
fault and consumes no RetryPolicy attempt. It is instead retried against a
wall-clock deadline (~30s) on a short backoff, regardless of idempotency, since
nothing was sent; these conditions self-heal within seconds as ports leave
TIME_WAIT. This keeps a chatty client (or the full parallel test suite) from
failing an otherwise-fine request on a transient local hiccup. Once the deadline
passes it falls back to raising ServiceUnavailable like any other transport
failure.
mTLS¶
Pass ca_file (to pin the server) and optionally client_cert_file /
client_key_file (to present a client certificate) with an https base URL.
build_ssl_context(ca_file, client_cert_file, client_key_file) is exposed for
callers that want to build the context directly. With a client certificate the
token is optional — the server may authenticate the certificate's subject. A
client certificate is honored independently of ca_file: presenting a client
identity does not require also pinning a CA (without a CA the default trust
store verifies the server).
Error hierarchy¶
| Exception | Raised for |
|---|---|
ServiceError |
base type; carries .status and .detail |
AuthError |
401 (bad/revoked credential) or 403 (missing grant/scope) |
RequestError |
any other 4xx, or a refused 3xx redirect — the request itself is wrong |
ServerError |
5xx |
ResponseError |
a 2xx body that is not decodable JSON, or a body over the size cap |
ServiceUnavailable |
transport failure after the retry policy is exhausted |
.detail carries the server's message when the error body was JSON with a
detail field. Applications typically alias or subclass these (for example
SecretNotFound(RequestError)).
classify_status(status, detail="") exposes the status→type mapping as a plain
function (401/403 → AuthError; a 3xx or any other 4xx → RequestError;
5xx → ServerError). An SDK that layers its own meanings on a few statuses (a
404 that means "not found", a 409 that means "conflict") special-cases those
and delegates every other status here, instead of re-deriving the base taxonomy.
Env-driven construction¶
ServiceClient never reads the process environment itself. Client programs
that want env-driven construction use the helper, which centralizes the
variable shape behind a prefix:
from hyperdjango.serviceclient import service_client_from_env
client = service_client_from_env("HYPERSECRET") # any override kwarg wins
reads HYPERSECRET_URL, HYPERSECRET_TOKEN, HYPERSECRET_CA_FILE,
HYPERSECRET_CLIENT_CERT, and HYPERSECRET_CLIENT_KEY.
An SDK subclass cannot go through service_client_from_env (that builds a
plain ServiceClient). It instead reads the same variable shape with
service_client_env_kwargs(prefix) -> dict and splats the result into its own
constructor:
from hyperdjango.serviceclient import service_client_env_kwargs
class MySdk(ServiceClient):
@classmethod
def from_env(cls, prefix, **app_specific):
return cls(**service_client_env_kwargs(prefix), **app_specific)
The returned dict holds base_url, token, ca_file, client_cert_file, and
client_key_file. service_client_from_env builds on it, so the {PREFIX}_*
variable shape lives in exactly one place.
ChangeFeedWatcher¶
from hyperdjango.serviceclient import ChangeFeedWatcher
def on_event(event):
refresh(event)
def on_reset(response):
full_resync_from_producer()
watcher = ChangeFeedWatcher(
client,
replay_path="/v1/events",
ws_path="/ws/feed", # None for poll-only
on_event=on_event,
on_reset=on_reset,
cursor=last_seen_cursor,
limit=500,
poll_interval=30.0,
).start()
...
watcher.stop()
Delivery models¶
The hub advertises a delivery mode in its hello frame and the watcher runs
the matching state machine — one class, one internal mode switch. The frames
are:
- subscribe (client → hub, on every connect):
{"type":"subscribe","prefixes":[...],"client_id":<str|null>,"last_seq":<int|null>,"cursor":<int|null>,"epoch":<str|null>} - hello (hub → client):
{"type":"hello","mode":"ephemeral"|"catchup"|"ledger","seq":<int>,"cursor":<int>,"resync":<bool>,"epoch":<str|null>} - event (hub → client, ephemeral/catchup):
{"type":"event","subject":...,"kind":...,"seq":N,"metadata":{...}} - wake (hub → client, ledger):
{"type":"wake","cursor":N}— a hint only
| Model | Configure | On (re)connect | Delivery |
|---|---|---|---|
| ledger | replay_path (or replay) |
pull replay to catch up on anything missed while down | replay pages are the source of truth; a WS frame only wakes it |
| ephemeral | ws_path, no replay_path |
on_reset (invalidate + lazily re-fetch) — every time |
each event frame → on_event; no per-client server state |
| catchup | ws_path, no replay_path |
send (client_id, last_seq); hub replays only missed events |
each event frame → on_event, advancing last_seq |
In catchup, the hub retains a per-client buffer keyed by client_id (stable
for the watcher's lifetime — pass one to survive a process restart, else a
per-watcher uuid is generated). On reconnect the watcher resumes from the
retained last_seq, so a brief disconnect replays only the events missed in the
gap — the server is more persistent than the client. The hub stamps a random
epoch at startup and advertises it in the hello; the watcher echoes it in the
next subscribe. If the hub evicted past last_seq (its ring overran or the
client_id is unknown) or the epoch changed (the hub restarted, so its
in-memory sequence is a fresh, unrelated space) the watcher does a full
on_reset instead of resuming — so a stale last_seq that happens to fall
inside a new incarnation's sequence range can never be mistaken for a valid
resume. Delivery is never partial or duplicated.
Negotiation. The watcher's configured intent proposes a model — a durable
replay_path means ledger, a ws-only feed means in-frame delivery — and the
hub's advertised mode decides. A mismatch always falls back to a
frame-delivery resync, which is safe: a hub advertising ledger to a watcher
with no replay endpoint is served as ephemeral, and a hub advertising an
in-frame model to a ledger-configured watcher quiesces the replay drain and
takes the hub's frames. A hub that sends no hello at all is treated as ledger
whenever a replay_path is configured, so a pure wake-hint hub needs no
handshake.
The rest of this section describes the ledger model in depth; it is the durable tier and the one with the strongest ordering guarantee.
The ordering principle: wake hint + pull replay¶
The correctness of the watcher rests on one invariant:
The replay endpoint is the single ordered source of truth. Live WebSocket traffic is only a wake-up hint. The cursor advances only through contiguous replay pages.
The watcher never delivers a live WebSocket payload and never advances the
cursor from one. A live frame — of any content, in any order — only nudges the
watcher to pull replay pages. Delivery and cursor movement happen exclusively
through replay(after=cursor), drained page by page in ledger order, with the
cursor advancing per delivered event.
This makes ordering and at-least-once delivery hold by construction under sharded fan-out, concurrent publishers, replay/live interleave, and dropped wake-ups: an out-of-order, duplicated, or missing wake cannot reorder, dupe, or lose an event, because the wake is never the source of the data.
A periodic poll tick (poll_interval, default 30s) pulls replay even when
no wake arrives, so the loop self-heals when every wake is lost — a live
WebSocket is an optimization for latency, never a requirement for correctness.
Wake targets and re-drain¶
A wake frame may advertise the cursor the server has reached (the field named by
wake_cursor_field, default cursor). The watcher records the highest such
hint as a liveness target only — it is never delivered and never advances
the cursor. When a drain finishes with the cursor still below that target, the
server's replay ceiling lagged the wake (typically an unrelated long transaction
holding the visibility horizon back), so the watcher re-drains on a short,
bounded backoff (a capped exponential from ~0.05s up to max(3s, poll_interval)
— 30s at the default poll_interval) until the target is reached, instead of
sleeping a full poll_interval. The floor is restored only when a fresh wake
raises the target: a bogus or mis-mapped hint that can never be reached — even
while unrelated events keep the cursor inching forward — decays to the poll
cadence rather than pinning the loop at the floor. This closes an
intermittent latency gap where an event is committed and announced but not yet
visible to replay for a few polls: it is delivered the instant replay reveals
it, in order, exactly once — because the re-drain still pulls every event from
the replay pages, never from the hint.
Reset¶
When the client's cursor falls below the server's retention floor, the missed
events are gone and cannot be replayed. The server flags the replay response
with reset=true; the watcher invokes on_reset(response) (so the app can
resync from the producer), advances the cursor past the trimmed gap, and
continues delivering from the floor forward. on_reset fires once per cursor
position: a degenerate reset that never advances the cursor (an empty page with
no or a non-integer cursor_field) is not re-fired every poll — it fires once
for that stuck position and resets counts it once — while a later reset at a
new cursor value fires again.
Reconnect backoff¶
An idle wake connection is held open by client keepalive pings (every
ws_ping_interval, default 20s), so a quiet hub does not trip a read timeout
and churn reconnects. A ping requires an answer, though: if several consecutive
ping intervals pass with zero inbound bytes — a black-holed peer (NAT drop,
power-off, no RST) that a live hub would never resemble — the connection is
declared dead and the reconnect logic takes over, rather than pinging into the
void for hours while wake latency silently degrades to poll_interval.
Announced frames are capped at ws_max_frame_bytes (default 8 MiB) — a larger
frame is rejected before its payload is read — and the upgrade handshake
response header block is itself capped (64 KiB) so a peer that never sends the
terminator cannot grow the read buffer without bound. recv_json handles one
whole text/binary frame per message: it does not reassemble fragmented
messages and treats continuation frames as unsupported. The wake loop does not
use recv_json — it consumes frames tolerantly, treating any complete message
as a wake and reading the cursor only as a best-effort hint: a non-JSON payload
or a fragmented (continuation-split) message is still a wake, reassembled and
drained, never a reconnect trigger, since the wake's content is never the source
of delivered data. A _WebSocketConnection is single-thread-at-a-time for
send/recv — the
watcher uses one reader thread and only stop() touches the socket concurrently
(to close it and unblock the reader).
The wake WebSocket reconnects with exponential backoff and jitter. Backoff
resets to base after a healthy session — one that delivered a wake or stayed
connected past stable_period (default 30s) — not on mere connect success: a
hub that accepts then immediately drops keeps backing off instead of being
hammered.
Callback safety and lifecycle¶
on_event and on_reset exceptions are contained (they never kill the loop)
and counted on the watcher (event_callback_errors, reset_callback_errors).
No drain-side error kills the feed thread either — a ServiceError (unreachable
replay, a non-JSON or oversized page) is retried on the next tick with the
cursor untouched, and any other unexpected error is counted at drain_errors.
start() launches the background threads; stop(timeout=5.0) sets the stop
flag, wakes the drain loop, closes an idle wake connection so its blocked
receive returns at once, then joins both threads up to timeout. The threads
are daemons: if the drain thread is mid retry-backoff sleep inside an
in-flight replay call when stop() is called, that sleep is not interruptible,
so stop() may return before the thread has fully unwound. No new work starts
once the flag is set, the thread carries no external state, and being a daemon
it cannot outlive the process — but callers wanting a hard guarantee that the
thread has exited should pass a timeout comfortably above the retry backoff
ceiling. The ledger cursor is readable at watcher.cursor, and the last
in-frame seq delivered (catchup's resume key, ephemeral's high-water mark) at
watcher.last_seq. A ws-only ephemeral/catchup watcher has no replay drain
thread — start()/stop() manage only the WebSocket thread.
Response shape¶
The watcher expects a replay response shaped like:
The field names (events, cursor, reset, and each event's id), the wake
frame's cursor field, and the query parameter names (after, limit) are all
constructor parameters (events_field, cursor_field, reset_field,
event_id_field, wake_cursor_field, after_param, limit_param,
extra_params), so a service with different names is expressed without changing
the watcher.
A replay page must be a JSON object; a top-level non-object (for example a bare
JSON array) is rejected as a bad page — counted at drain_errors, never
delivered, and the cursor is left untouched for the next tick.