Skip to content

Programmable Rust hooks for sōzu — design proposal #1239

Description

@FlorentinDUBOIS

This issue is a self-contained design proposal for compile-time-linked Rust hooks in sozu-lib, in response to the original ask in #1178. It folds the maintainer reply, six design decisions, and findings from a multi-stream code review (correctness, project-guidelines conformance, simplification, security threat-hunt, and a second-model cross-check). Implementation does not start until this proposal is approved.

Implementation phasing, PR sequencing, risk register, items requiring resolution before implementation, and cross-references continue in the first comment below (the proposal is split across two posts because the full body exceeds GitHub's per-comment size limit).

Original ask: #1178.
Deployment target: kemeter/sozune — a Traefik-style higher-level reverse proxy built on top of sozu-lib. Stock sozu binary continues unchanged; Sozune ships the operator-facing middleware (basic auth, rate limit, headers, strip-prefix, redirects, compression, backend-timeout) by linking against sozu-lib and registering custom hooks.
Out of scope: WebAssembly / dynamic plugin loading. Pure Rust, compile-time-linked.


1. Context

1.1 The ask

@Shine-neko (issue #1178) wants a way to plug a function/callback at two stages of the pipeline: rewrite or modify requests (path, headers) before routing, and rewrite or modify responses (headers, status) before sending them back. @Wonshtrum (maintainer) replied that the closest existing primitives are kawa_h1::ParserCallbacks::on_request_headers / on_response_headers plus the URL-rewrite / header-injection / redirect / basic-auth feature stack that landed via PRs #1206-#1210, but these are static configuration, not operator-supplied logic. The maintainer explicitly invited a programmable surface.

1.2 The reference points

Three reference designs were studied to anchor sōzu's choices:

  • Pingora (Cloudflare): a single ProxyHttp trait with ~30 hooks, all async fn, monomorphised global per Service. Per-session type CTX associated type. Strong fit for tokio-based proxies; doesn't translate to sōzu's mio-only lib/.
  • HAProxy: three pluggability surfaces — native filter modules (C, sync function pointers, flt_ops struct with lifecycle + channel + HTTP/TCP hooks), Lua scripting (in-VM coroutines via core.register_action cheap and core.register_filter body-streaming), and SPOE (out-of-process NOTIFY/ACK protocol for verdict-only enrichment).
  • Sōzu's existing primitives: kawa_h1::ParserCallbacks::on_request_headers / on_response_headers (parser callbacks, internal), Router::route_from_request (resolves cluster + per-frontend policy), config-time per-frontend programmability (headers: Vec<Header>, redirect, rewrite_*, required_auth, answers."NNN", max_connections_per_ip).

1.3 The design decisions (operator-confirmed)

Six questions were asked, six answers locked:

Q Question Decision
Q1 Storage location — where the hook reference is held (b) Listener-owned registry, frontends/clusters reference hooks by ID
Q2 Sync vs async — what shape the hook signatures take (i) Pure sync + (ii) yield-on-pending mio integration. Both ship in v1, NOT deferred to v1.5
Q3 API surface — single trait or two-tier (b) Two-tier — HeaderHook cheap + always available, BodyHook opt-in
Q4 Operator deployment shape — distribution model (a) Pingora-style binary build — operator (Sozune) recompiles their own binary linking sozu-lib
Q5 ABI stability — facade layer or direct internal exposure Expose internal types directly — &mut HttpContext to hooks; Sozune accepts coupling
Q6 Header edit type — new HeaderEdit enum or reuse existing Header (a) Reuse Header — operators emit Vec<Header> from hooks, same downstream pipeline

Plus the per-stream state slot:

  • TypeId-keyed HashMap on HttpContext, lazy-allocated (Option<HashMap<TypeId, Box<dyn Any + Send + Sync>>>). Each hook's state type is its own slot — no collisions, no string IDs, no operator burden.

1.4 Hard constraints inherited from sōzu

  1. No async fn in lib/. mio + edge-triggered epoll, single-threaded per worker. Hooks cannot be async fn. The yield-on-pending path satisfies operator I/O needs without breaking this rule.
  2. No Arc<Mutex> inside the event loop. Slab-allocated session state. Hooks see &mut HttpContext per-session; cross-session shared state belongs in &self on the hook impl (e.g. Arc<RwLock<BanList>>).
  3. No panic on network-facing input. Hooks must convert errors into HookOutcome::Reject(status) or HookError. unwrap / expect / panic! forbidden in the protocol layer.
  4. Edge-triggered epoll discipline. When a hook queues bytes for write, it MUST call signal_pending_write on the readiness tracker — there is no wake-up for free.
  5. Buffer-pool-owned bodies. Body-stage hooks operate on per-session buffer-pool checkouts; chunk size cannot grow beyond the buffer slot.
  6. H2 stream ordering. Per-stream pending state must respect HTTP/2 framing — pending on stream 5 must NOT block stream 7 from making progress, but stream 5's response cannot go out before stream 5 finishes its hook chain.

2. Public API specification

2.1 Module layout

New module: lib/src/hooks/

lib/src/hooks/
├── mod.rs           # public trait + types + builder API
├── outcome.rs       # HookOutcome / WakeupSource / WakeupReason / HookError
├── registry.rs      # HookRegistry — per-listener storage
├── dispatch.rs      # Router-side dispatch (Phase B + C glue)
├── wakeup.rs        # mio Poll integration (Phase C)
└── builtins/        # Phase E — sōzu's own middlewares as hook impls
    ├── mod.rs
    ├── basic_auth.rs
    ├── header_edit.rs
    ├── rate_limit.rs
    ├── redirect.rs
    └── rewrite.rs

2.2 Trait surface

// lib/src/hooks/mod.rs
use std::any::Any;
use std::time::Duration;
use crate::protocol::http::editor::HttpContext;
use sozu_command::proto::command::Header;

/// Stage-1 hook — sees parsed request / response headers, can mutate
/// `HttpContext.headers_request` / `headers_response` (which feed the
/// existing `apply_request_rewrites_and_headers` /
/// `apply_response_header_edits` pipeline). Cheap path; always
/// available. Hook output is byte-validated in the dispatcher before
/// reaching the wire (see §2.7 step 6).
pub trait HeaderHook: Send + Sync + 'static {
    /// Operator-supplied identifier (no default). Used as the metric
    /// label suffix (`hooks.duration.<name>`) and log envelope tag —
    /// must be a stable identifier (`type_name::<Self>()` includes
    /// module paths and is documented as best-effort, so it is not
    /// the default).
    fn name(&self) -> &'static str;

    /// Called after `ParserCallbacks::on_request_headers` and after
    /// `Router::route_from_request` resolves the cluster.
    fn on_request_headers(
        &self,
        _ctx: &mut HttpContext,
        _wake: WakeupReason,
    ) -> HookResult {
        Ok(HookOutcome::Continue)
    }

    /// Called just before `kawa.prepare(...)` on the response side, in
    /// reverse scope order vs. request (cluster → frontend → listener).
    fn on_response_headers(
        &self,
        _ctx: &mut HttpContext,
        _wake: WakeupReason,
    ) -> HookResult {
        Ok(HookOutcome::Continue)
    }
}

/// Stage-2 hook — opt-in body-byte interception. Implementing this
/// trait also requires `HeaderHook` (the body stage runs after the
/// header stage of the same hook). Body hooks are gated behind a
/// listener-builder method `with_body_hook(name, Arc<dyn BodyHook>)`
/// distinct from `with_hook(...)` so the buffer-pool path is only
/// enabled for hooks that opt in.
pub trait BodyHook: HeaderHook {
    /// In-place body chunk mutation. Cannot grow the chunk beyond the
    /// buffer-pool slot. Hooks that shrink the chunk MUST keep
    /// downstream framing consistent — sōzu auto-recomputes
    /// `Content-Length` from the post-hook total on H1 (or marks the
    /// response as `Transfer-Encoding: chunked` if the original CL is
    /// authoritative on a different hop) before forwarding. H2:
    /// post-mutation DATA-frame size drives outbound flow-control
    /// debit. END_STREAM is never swapped under the operator's feet.
    /// Returns `HookOutcome::Continue` to forward, `Reject(status)`
    /// to abort BEFORE response bytes are emitted, `AbortStream { error }`
    /// to abort mid-response, `Pending(...)` to yield.
    fn on_request_body_chunk(
        &self,
        _ctx: &mut HttpContext,
        _chunk: &mut BodyChunk<'_>,
        _wake: WakeupReason,
    ) -> HookResult {
        Ok(HookOutcome::Continue)
    }

    fn on_response_body_chunk(
        &self,
        _ctx: &mut HttpContext,
        _chunk: &mut BodyChunk<'_>,
        _wake: WakeupReason,
    ) -> HookResult {
        Ok(HookOutcome::Continue)
    }
}

pub struct BodyChunk<'a> {
    pub bytes: &'a mut [u8],
    pub end_stream: bool,
}

pub type HookResult = Result<HookOutcome, HookError>;

2.3 Outcome enum

// lib/src/hooks/outcome.rs
use std::time::Duration;
use std::os::fd::RawFd;
use sozu_command::proto::command::Header;

/// Hook return value. Four terminal cases. Header edits are emitted
/// by mutating `HttpContext.headers_request` / `headers_response`
/// directly (Q6 reuses the existing `Header` plumbing); the
/// dispatcher byte-validates every emitted `Header` before it reaches
/// the wire (rejects CRLF / NUL / non-tchar in keys, C0 controls in
/// values per `parse_header_edit` in `command/src/config.rs:1397-1408`,
/// promoted to a shared crate location).
#[non_exhaustive]
pub enum HookOutcome {
    /// Proceed to the next hook in the chain (or to the next pipeline
    /// stage if this is the last hook).
    Continue,

    /// Short-circuit BEFORE response bytes have been emitted. Sōzu
    /// renders the matching default answer template via
    /// `set_default_answer` (`mux/answers.rs:193`); for 429 only,
    /// `set_default_answer_with_retry_after` (`answers.rs:207`)
    /// threading the cluster's effective retry_after via
    /// `SessionManager::effective_retry_after` (`server.rs:243`).
    /// `RejectStatus` is a newtype enforcing 4xx/5xx (and 301/302/308)
    /// at construction; out-of-range values surface as `HookError`
    /// at the dispatcher.
    Reject {
        status: RejectStatus,
        /// Optional cause carried into the per-hook error log; not
        /// rendered on the wire.
        message: Option<String>,
    },

    /// Abort an in-flight response stream. Used when a body hook
    /// detects a problem after response headers (or partial body)
    /// have already been emitted. H2: emit RST_STREAM with the
    /// operator's chosen H2 error code. H1: write-only shutdown
    /// (NEVER `Shutdown::Both` on TLS frontends — write-only
    /// shutdown avoids the TCP RST that truncates already-queued
    /// response bytes); the truncated body is the best-effort
    /// signal to the client.
    AbortStream {
        error: AbortStreamError,
    },

    /// Hook needs to wait for an external event. Sōzu's mio loop
    /// re-enters the hook (with the same `&mut HttpContext` and the
    /// hook's own state slot) when the wake source fires. `timeout`
    /// is capped per worker by `hooks.max_pending_timeout_seconds`
    /// (default 5 s, hard ceiling 60 s); concurrent `HookPending` per
    /// worker is capped at `hooks.max_concurrent_pending_per_worker`
    /// (default ≤ 32 768 to keep headroom inside the 65 536-token
    /// range — see §4.1 / §4.5).
    Pending {
        wake: WakeupSource,
        timeout: Duration,
        /// Same `RejectStatus` constraint as `Reject`.
        wake_timeout_status: RejectStatus,
    },
}

/// Newtype over `u16` enforcing 4xx/5xx (and 301/302/308) at
/// construction. `TryFrom<u16>` returns `InvalidRejectStatusError` for
/// out-of-range values; the dispatcher converts those into
/// `HookError { status: 500, message: "invalid Reject status: <N>" }`
/// and bumps `hooks.invalid_reject_status.<hook_name>`.
pub struct RejectStatus(u16);

#[non_exhaustive]
pub enum AbortStreamError {
    /// H2: emit RST_STREAM with this error code. H1: write-only shutdown.
    Internal,
    Cancel,
    EnhanceYourCalm,
    Custom(u32),
}

/// What the hook is waiting for. The mio loop registers the
/// appropriate event source on the worker's Poll.
#[non_exhaustive]
pub enum WakeupSource {
    /// File descriptor became readable (e.g. a TCP socket the hook
    /// opened to a sidecar Redis). Ownership transfers to sōzu on
    /// `Pending` return — `EPOLL_CTL_DEL` runs in
    /// `WakeupSourceHandle::Drop` before fd close, so a double-close
    /// cannot race a kernel fd-reuse. Operator MUST set `O_CLOEXEC`
    /// and MUST NOT register the same fd on any other Poll.
    FdReadable(OwnedFd),

    /// Channel message arrived. The caller must pre-register the
    /// channel via `HookRegistry::register_channel(token)` at
    /// listener-build time; this variant just stores the token. The
    /// registry monitors `Disconnected` on the channel handle — if
    /// the operator's background thread drops its `Sender`, the
    /// pending hook is re-entered with `WakeupReason::Timeout` (and
    /// `hooks.channel.disconnected.<hook_name>` increments) instead
    /// of leaking the wake-up entry until `timeout` elapses.
    ChannelReady(usize),

    /// Wall-clock delay (e.g. exponential backoff between retries).
    /// Integrates with the worker's existing Timer machinery.
    Delay(Duration),
}

/// Passed when a hook is invoked — `Initial` on first entry,
/// `Ready` / `Timeout` on re-entry from a previous `Pending`. Hooks
/// that issue multiple Pendings across phases use `source` to
/// identify which pending operation completed.
#[non_exhaustive]
pub enum WakeupReason {
    Initial,
    Ready { source: WakeupSourceId },
    Timeout { source: WakeupSourceId },
}

/// Opaque identifier for a `WakeupSource` returned by a `Pending`.
/// Generated at Pending-return time, plumbed through the registry,
/// surfaced on re-entry so hooks can disambiguate concurrent waits.
pub struct WakeupSourceId(NonZeroU64);

pub struct HookError {
    pub status: RejectStatus,
    pub message: String,
}

2.4 HttpContext extension

// lib/src/protocol/kawa_h1/editor.rs (existing struct extended)
use std::any::{Any, TypeId};
use std::collections::HashMap;

pub struct HttpContext {
    // ... all existing fields unchanged ...

    /// Per-session, per-hook state slot. Each hook stores its state
    /// keyed by its own `TypeId`, so two hooks with different state
    /// types never collide. `None` until the first hook stashes
    /// anything — stock sōzu deployments without hooks pay zero.
    pub(crate) hooks_state: Option<HashMap<TypeId, Box<dyn Any + Send + Sync>>>,
}

impl HttpContext {
    pub fn set_hook_state<T: Any + Send + Sync>(&mut self, state: T) {
        self.hooks_state
            .get_or_insert_with(HashMap::new)
            .insert(TypeId::of::<T>(), Box::new(state));
    }

    pub fn hook_state<T: Any + Send + Sync>(&self) -> Option<&T> {
        self.hooks_state
            .as_ref()?
            .get(&TypeId::of::<T>())?
            .downcast_ref::<T>()
    }

    pub fn hook_state_mut<T: Any + Send + Sync>(&mut self) -> Option<&mut T> {
        self.hooks_state
            .as_mut()?
            .get_mut(&TypeId::of::<T>())?
            .downcast_mut::<T>()
    }

    pub fn take_hook_state<T: Any + Send + Sync>(&mut self) -> Option<T> {
        self.hooks_state
            .as_mut()?
            .remove(&TypeId::of::<T>())?
            .downcast::<T>()
            .ok()
            .map(|b| *b)
    }
}

Construction site impact. HttpContext is constructed at five known sites today: mux/h2.rs:7422, mux/answers.rs:358, mux/mod.rs:494 (HttpContext::new), kawa_h1/editor.rs:833 + :862 (test fixtures), kawa_h1/mod.rs:314. The new hooks_state: Option<...> field defaults to None; the explicit struct literals at the first three sites still need their fields updated in Phase A.

Debug impl. HttpContext has #[derive(Debug)] at editor.rs:115 (struct definition at editor.rs:116). Box<dyn Any> is not Debug; the derive must be replaced with a hand-rolled impl Debug for HttpContext that lists every existing field by name and renders the new field as hooks_state: Option<...> (count only). A unit test in Phase A pins format!("{:?}", ctx) is panic-free with stashed Box<dyn Any> state — error-path log dumps ({:?} at e.g. mux/router.rs:~374) must keep working. The derivative crate is NOT a workspace dep; do not add it. No Clone / PartialEq impact.

State-discipline contract. Two hooks both stashing common types (Vec<u8>, String) collide under TypeId::of::<Vec<u8>>() — last writer wins, reader silently misclassifies. Concrete failure mode: auth hook stores user ID, logging hook overwrites with fingerprint, rate-limit hook reads ctx.hook_state::<Vec<u8>>() and gets the fingerprint — per-IP rate-limit silently degrades to per-fingerprint, allowing distributed credential stuffing under a single IP. Hook state MUST be a hook-private newtype:

struct MyAuthState(Vec<u8>);  // OK — distinct TypeId

lib/src/hooks/macros.rs ships a define_hook_state! macro that generates the newtype + Send/Sync bound + custom Debug (skips bytes). Documented as the only supported way to declare hook state. The registry logs a rate-limited warn! when two registered hooks call set_hook_state::<T> with the same T::TypeId.

Field-mutability matrix. Hooks see &mut HookMut (a sealed wrapper over HttpContext), NOT &mut HttpContext directly. The HookMut API gates writes on routing-frozen fields:

Field Hook access
cluster_id, backend_id read-only
method, authority, path (post-route) read-only
tls_*, keep_alive, protocol read-only
headers_request, headers_response read-write
hooks_state (via set_/hook_state* helpers) read-write
frontend_id, request_id, session_id read-only

Mutating authority after routing was the explicit operator-footgun risk; the read-only gate forces hooks to emit Header { key: "Host", val: ... } edits instead, which flow through downstream re-validation (see §2.7 step 7).

Backing store. §3 estimates 1–3 stashed states per request. For that size, Vec<(TypeId, Box<dyn Any + Send + Sync>)> with linear scan is cache-friendlier than HashMap (no hash, no drop overhead, smaller). The four convenience methods (set_/hook_state*) keep their signatures; backing store is an implementation detail.

2.5 Builder API

// lib/src/protocol/listener.rs (or HttpListener / HttpsListener equivalent)
use std::sync::Arc;
use crate::hooks::{HeaderHook, BodyHook, HookRegistry};

impl HttpListener {
    /// Register a header-stage hook by name. The name is the ID
    /// frontends and clusters reference in their `hook_ids: Vec<String>`.
    /// Must equal `hook.name()` — mismatch returns
    /// `HookRegisterError::NameMismatch`. Registering twice with the
    /// same name returns `HookRegisterError::Duplicate`. The registry
    /// also runs the state-discipline lint (§2.4) and emits a `warn!`
    /// when two hooks declare the same `TypeId` for their state slot.
    pub fn with_hook(
        &mut self,
        name: impl Into<String>,
        hook: Arc<dyn HeaderHook>,
    ) -> Result<&mut Self, HookRegisterError>;

    /// Register a body-stage hook. Implies header-stage registration
    /// at the same name (since `BodyHook: HeaderHook`).
    pub fn with_body_hook(
        &mut self,
        name: impl Into<String>,
        hook: Arc<dyn BodyHook>,
    ) -> Result<&mut Self, HookRegisterError>;
}

2.6 Proto schema additions

// command/src/command.proto
message RequestHttpFrontend {
    // ... existing fields ...
    repeated string hook_ids = 16;  // pick the next available tag at edit time
}

message Cluster {
    // ... existing fields ...
    repeated string hook_ids = 16;  // pick the next available tag at edit time
}

Validation: at config-load time AND at runtime AddHttpFrontend / AddCluster, every hook_id must resolve in the listener's registry. Unknown ID → reject the request with a typed HookConfigError carrying the offending ID, do not silently accept.

Cluster scope on a listener-owned registry — open question. Cluster is a global resource; AddCluster carries no listener context. Resolution of where the cluster's hook_ids get validated (worker-global registry vs. validation deferred to AddHttpFrontend time when a frontend pulls a cluster) is deferred to a dedicated planning session before Phase A starts. v1 keeps cluster hook scope per Q1 — only the validation site is to be decided. Tracked at §8.1 item 2.

Cargo-clean dependency. The workspace's prost map-field derive handling can leave stale generated code after a command.proto edit; run cargo clean -p sozu-command-lib && cargo build -p sozu-command-lib once after editing hook_ids tags before the next workspace build.

Cross-version safety. prost skips unknown fields silently → an older sōzu worker deserialising a config with hook_ids set silently treats the frontend as no-hook (degrade-open). For v1 this is the documented contract; future work may add a WorkerCapabilities advertisement so the master can refuse-application on mixed-version partial upgrades.

2.7 Dispatch order

The dispatch helpers in lib/src/hooks/dispatch.rs are called from Router::connect (mux/router.rs:90 → :141), NOT from inside Router::route_from_request (which would require widening that function's return type). Both H1 and H2 reach Router::connect through the unified mux dispatch path (HTTP/1.1 listeners flow through HttpStateMachine::Mux(MuxClear)), so a single hook insertion point covers both protocols.

Naming note. The "listener-post-route" chain fires AFTER route_from_request returns Ok — chosen for v1 because rewriting path/authority and then re-routing is a deeper change to Router than the wave can absorb. v2 may add a true pre-route stage that runs before route_from_request. The original issue #1178 ask for "rewrite path/headers before routing" is therefore deferred — operators can mutate via headers_request edits but those do not change the routing decision in v1.

Request side (fires AFTER existing config-time programmability):

1. ParserCallbacks::on_request_headers (sōzu invariants — Host validation, CL/TE reconciliation)
2. Router::route_from_request (config-time: redirect, auth, rewrite_*)
3. listener-post-route hook chain (listener.hook_ids on the listener config)
4. frontend hook chain (frontend.hook_ids)
5. cluster hook chain (cluster.hook_ids)
6. Hook output validation: every Header emitted by hooks 3-5 is byte-validated
   (rejects CRLF / NUL / non-tchar in keys, C0 in values) by the helper
   promoted from command/src/config.rs:1397-1408 (parse_header_edit) into a
   shared crate location. Invalid bytes → HookError { status: 500 }, bump
   hooks.invalid_header.<hook_name>.
7. Invariant re-validation: re-fire CL/TE consistency, Host single-value, no
   forbidden-byte check on the post-hook header set. Mandatory whenever the
   hook chain emitted any header edit. Failure → HookError { status: 500 }.
8. Backend connect

Trailer-frame coverage. pkawa::handle_trailer (lib/src/protocol/mux/pkawa.rs:1038) is plumbed through the same hook chain (steps 3-7) so trailer headers do not bypass listener-scoped policy. Concrete plumbing shape (a second invocation with a trailer-flag in WakeupReason, vs. a separate hook entry-point on HeaderHook) is settled at Phase B implementation time.

Response side (reverse):

1. Backend response received
2. cluster hook chain (response variant, reverse iteration)
3. frontend hook chain (response variant)
4. listener-post-route hook chain (response variant)
5. Hook output validation (same byte-validation gate as request side)
6. Invariant re-validation (response variant)
7. ParserCallbacks::on_response_headers + existing config-time response edits
8. kawa.prepare → wire bytes

Outcome routing:

  • First Continue accumulates header edits and proceeds. Header-edit ordering across scopes is deterministic last-writer-wins in dispatch order (listener → frontend → cluster on request, reversed on response).
  • First Reject(status) short-circuits via set_default_answer (mux/answers.rs:193) — or set_default_answer_with_retry_after (answers.rs:207) for 429 specifically, threading SessionManager::effective_retry_after (server.rs:243). Reject is valid only BEFORE response bytes have been emitted. A Reject on the response side after :status has been queued is converted to AbortStream { error: Internal } and a warn! is logged so the hook author sees the misuse.
  • First AbortStream { error } aborts an in-flight response. H2: emit RST_STREAM with the operator's chosen H2ErrorCode. H1: write-only shutdown (NEVER Shutdown::Both on TLS frontends — write-only shutdown avoids the TCP RST that truncates already-queued response bytes).
  • First Pending(...) parks the session via the wake-up registry (Phase C). Subject to per-worker caps (§4.5).

2.8 Panic-safety contract

&mut HookMut exposed to a hook is not naturally UnwindSafe (it wraps interior pointers, kawa storage). The dispatcher wraps every hook call in catch_unwind. On Err, the dispatcher MUST:

  1. Clear ctx.headers_request and ctx.headers_response post-route snapshots (the panicking hook may have left them half-populated).
  2. Drop ctx.hooks_state entirely (the panicking hook may have left dangling invariants in its own slot).
  3. Emit hooks.panicked.<hook_name> (rate-limited counter — operators see drift; the worker stays up).
  4. Route to set_default_answer with status 500.
  5. Mark the session for close after the default answer is written.

Phase A includes a regression test that panics from a HeaderHook mid-mutation and asserts no partial bytes reach the wire. Workspace Cargo.toml commits to panic = "unwind" so a future change is a deliberate decision the hook subsystem can react to. "Worker should not crash" is the load-bearing operator posture.


3. Per-stream state slot — semantics + cost

3.1 Lookup pattern

Each hook stashes state at TypeId::of::<MyHookPrivateState>() and reads it back at the same key. Two hooks with different nominal state types occupy different slots; common-type collisions (Vec<u8> × Vec<u8>) are an operator footgun (see §2.4 "State-discipline contract" — newtype mandatory; the registry warns on TypeId duplication at registration time).

References into per-stream kawa storage MUST NOT be stashed; only owned bytes (a Vec<u8> newtype, a String newtype, etc.). Buffer-pool slots are per-stream; aliasing across streams is undefined behaviour. The define_hook_state! macro in lib/src/hooks/macros.rs enforces the newtype + the Send/Sync bound at compile time.

3.2 Cost analysis

Steady-state per request: ~5-10 hook calls (3 scopes × 1-3 middlewares). Each hook access does 1 TypeId comparison + 1 HashMap lookup + 1 downcast = single-digit nanoseconds. Compared to TLS handshake (~ms) or backend connect (~ms), invisible.

Memory: each stashed state boxed once. Sozune-typical mix (rate-limit verdict, auth result, compression negotiation) = 3-5 boxed allocations per request, ~100 bytes. Negligible.

When hooks are unused (stock sōzu, no Sozune-style binary): hooks_state: None — one extra usize on HttpContext, zero allocations.

3.3 Why TypeId beats alternatives

Alternative Rejected because
Make HttpContext generic over Extra Refactor every call site. Sozune wants several middleware-specific state types, not one big struct.
Single untyped Box<dyn Any> slot Only one hook can use it. Two middlewares collide silently.
Pingora-style type CTX associated type Sōzu's hook trait is dyn-dispatched. Associated types don't survive type erasure.
Slab-keyed slot ID (operator manages a numeric ID per hook) Operator burden. TypeId map gets the same lookup cost without manual ID tracking; if profiling shows hot, swap the backing store.

4. Yield-on-pending wake-up subsystem (Phase C)

The most complex sub-system. Documented here in full because it touches the mio loop directly.

4.1 Token namespace

The worker's mio Poll currently allocates tokens for: front sockets (slab-keyed, low range), back sockets (slab + offset), command channel (fixed token), timer (fixed token), health-check probes (carved-out range [1<<24, 1<<24 + 1<<16), with the assertion at lib/src/health_check.rs:55-56 that session tokens stay capped well below 1<<24).

Hooks claim a fresh range:

const HOOKS_TOKEN_BASE: usize = 1 << 25;
const HOOKS_TOKEN_CAPACITY: usize = 1 << 16;  // 65 536 tokens

Allocation via a modulo allocator that skips in-flight tokens, refusing allocation rather than colliding. The wake-up registry caps active entries at min(HOOKS_TOKEN_CAPACITY / 2, hooks.max_concurrent_pending_per_worker) so a peer cannot exhaust the token range — refused allocations bump hooks.pending.refused_cap and the hook is rejected with HookError { status: 500 }. Verify before Phase A lands that max_connections * slab_entries_per_connection < (1 << 24) still holds with hooks adding pressure.

4.2 Registry

// lib/src/hooks/wakeup.rs
struct HookWakeupEntry {
    session_id: SessionId,
    stream_id: Option<StreamId>,    // Some on H2, None on H1 / TCP
    chain_position: ChainPosition,  // listener-post-route[i], frontend[j], cluster[k]
    direction: HookDirection,        // request | response
    hook: Arc<dyn HeaderHook>,
    /// The fd or channel handle registered on the Poll. Owned here
    /// so its lifetime is bounded by the registry; `Drop` runs
    /// `EPOLL_CTL_DEL` (for fds) or `disconnect` (for channels)
    /// before the source is closed, eliminating fd-reuse races.
    source: WakeupSourceHandle,
    /// Identifier surfaced on re-entry as
    /// `WakeupReason::Ready { source }` so hooks issuing multiple
    /// Pendings across phases can disambiguate.
    source_id: WakeupSourceId,
    deadline: Instant,
    wake_timeout_status: RejectStatus,
}

pub struct HookWakeupRegistry {
    entries: HashMap<Token, HookWakeupEntry>,
    /// Reverse index: per-(session_id, stream_id) so RST_STREAM and
    /// session-close sweeps run in O(active-pending-for-stream)
    /// without walking all entries.
    by_stream: HashMap<(SessionId, Option<StreamId>), SmallVec<[Token; 4]>>,
    next_token: Token,  // modulo allocator
    active_count: AtomicUsize,
    max_active: usize,
}

4.3 Re-entry path

When the mio loop returns an event with a token in the hooks range:

  1. Look up the HookWakeupEntry. If absent (already cancelled by RST_STREAM or session close), drop the event without bumping any counter.
  2. Resolve (session_id, stream_id) against the slab. If gone (session closed mid-pending), drop the entry + bump hooks.cancelled_pending.
  3. If alive: borrow &mut HttpContext, re-invoke the hook with WakeupReason::Ready { source: source_id } (or Timeout { source } from the timer path). The hook may return:
    • Continue — proceed to the next chain step (or pipeline stage).
    • Reject(status) — convert via set_default_answer (answers.rs:193) for non-429 codes, set_default_answer_with_retry_after for 429. If response bytes have already been emitted (response-side wake), convert to AbortStream { error: Internal } and warn!.
    • AbortStream { error } — emit RST_STREAM (H2) or write-only shutdown (H1).
    • Pending(...) — re-register on a new token; deadline accumulates against the original timeout-cap. A second-Pending after Timeout re-entry is rejected with wake_timeout_status to bound starvation.

Edge-triggered epoll discipline. On any path that queues bytes for write — header-edit application, default-answer write, AbortStream RST emission — the dispatcher MUST call signal_pending_write on the readiness tracker before returning to the loop, per §1.4 invariant 4. Edge-triggered epoll does not deliver a wake-up event for free when readiness is computed from the read path; forgetting this produces "stuck session" bugs that only manifest at exact byte boundaries.

4.4 Cancellation

Three cancellation paths:

  1. Session close. SessionManager::on_session_close calls HookWakeupRegistry::sweep_session(session_id), which uses the by_stream reverse index to find every entry for that session in O(active-pending-count-for-session) and drops them. Source handles drop → fds get unregistered from the Poll, channels disconnect their readiness signal. Counter hooks.cancelled_pending increments per swept entry.
  2. H2 RST_STREAM on a hook-pending stream. mux::h2::handle_rst_stream (or the equivalent pkawa path) calls HookWakeupRegistry::sweep_stream(session_id, stream_id) synchronously, BEFORE the H2 flood detector glitch counter sees the close. Late-firing wakes for that stream return "already cancelled" at step 1 of §4.3 and drop without bumping glitch_count — preserving the existing H2FloodDetector posture (the detector currently counts WINDOW_UPDATE on closed streams as a glitch, so race conditions between hook resume and stream close must be cleaned synchronously).
  3. Channel sender drop. When the operator's background thread panics and drops its mpsc::Sender, the registry detects Disconnected on the channel handle and re-enters the hook with WakeupReason::Timeout immediately (no need to wait for the configured timeout). The hook chooses to short-circuit or retry. hooks.channel.disconnected.<hook_name> increments.

4.5 Timeout enforcement

Every Pending carries timeout: Duration. The dispatcher caps this at hooks.max_pending_timeout_seconds (worker-config, default 5 s, hard ceiling 60 s) — operator-supplied timeouts above the cap are rejected with HookError { status: 500, message: "Pending timeout exceeds cap" } and hooks.pending.timeout_capped.<hook_name> increments. The worker's existing Timer machinery (the same one that drives front_timeout / back_timeout) registers a Timeout event at now + min(timeout, cap).

If the Timeout fires before a Ready wake, the hook is re-invoked with WakeupReason::Timeout { source }. If it returns Pending again, sōzu rejects with wake_timeout_status to bound starvation — the configured fallback. Per-worker concurrent-Pending count is also capped at hooks.max_concurrent_pending_per_worker (default ≤ 32 768) to keep headroom inside the 65 536-token range; exceeding the cap returns HookError { status: 500 } at registration without entering Pending state.

4.6 H2 stream ordering invariant

Per-stream pending state lives on the H2 connection's Stream slot in mux::stream. Hook-pending is a flag on Stream, NOT a new variant of StreamState. The actual StreamState enum at lib/src/protocol/mux/stream.rs:35 is a linkage state machine (Idle / Link / Linked(Token) / Unlinked / Recycle); RFC 9113 §5 protocol state lives elsewhere (front_received_end_of_stream, back_received_end_of_stream, per-stream readiness on the H2 Connection). Adding HookPending to StreamState would lose the backend Token carried by Linked(Token).

The shape:

// lib/src/protocol/mux/stream.rs (extended)
pub struct Stream {
    pub state: StreamState,
    // ... existing fields ...
    pub hook_pending: Option<HookPendingTag>,
}

pub struct HookPendingTag {
    pub token: Token,                      // index into HookWakeupRegistry
    pub direction: HookDirection,
    pub started_at: Instant,
}

Transitions: hook_pending = Some(_) on HookOutcome::Pending; cleared on hook re-entry returning anything other than Pending; cleared on RST_STREAM / session close. StreamState itself is untouched — Linked(backend_token) still carries its backend identity through a hook-pending window.

The connection-level scheduler skips streams with hook_pending.is_some() and rotates to ready streams. Response order is preserved by the existing FIFO bookkeeping that sequenced :status emission per RFC 9113 §5.

LIFECYCLE.md gets a new section: "Hook-pending streams and the response-order invariant" — written against the actual StreamState enum (linkage states), NOT the RFC-9113 abstract names.

4.6.1 Read-side flow control while hook-pending

While a request-side hook is parked (hook_pending.is_some()), the connection MUST NOT WINDOW_UPDATE the client for that stream and MUST NOT forward DATA payloads to a backend (which may not yet be connected — request hook chain fires after route_from_request but before backend connect, per §2.7). Without this gate, a hostile (or merely fast) client streaming the request body during a hook park is a memory-pressure DoS surface multiplied by concurrent connections.

Concretely:

  • mux::h2::handle_data_frame (lib/src/protocol/mux/h2.rs:4264) treats hook-pending streams as suspended for forwarding; received DATA accumulates against max_request_body_size and trips 413 on overflow. No WINDOW_UPDATE is emitted.
  • After resume (Continue clears hook_pending), the dispatcher emits the deferred WINDOW_UPDATE and resumes the read path.
  • After resume (Reject / AbortStream), the read path closes via the standard short-circuit; deferred bytes are discarded.
  • For response-side hook_pending, the read path (request body) is unaffected — only the write path is suspended.

4.7 Channel-fed background-thread pattern

The canonical I/O workaround for "I need to consult Redis":

  1. Operator's Sozune binary spawns a background thread holding the Redis connection.
  2. Background thread maintains an in-memory snapshot (Arc<RwLock<BanList>>).
  3. When a hook needs the snapshot fresh, it sends a request via a mpsc::Sender, then returns Pending { wake: WakeupSource::ChannelReady(channel_token), timeout: Duration::from_secs(1), wake_timeout_status: RejectStatus::SERVICE_UNAVAILABLE }.
  4. Background thread queries Redis, sends the verdict back via the channel.
  5. Hook re-enters with WakeupReason::Ready { source }, reads the verdict from ctx.hook_state::<RedisVerdict>(), emits Continue or Reject.

Sender-drop semantics. If the background thread panics and the Sender half drops, the registry detects Disconnected on the channel handle and re-enters the hook with WakeupReason::Timeout immediately — no waiting for timeout to elapse. The hook then chooses to short-circuit or retry. hooks.channel.disconnected.<hook_name> increments. Phase C ships an e2e cell "background-thread-panics-during-pending" to pin this contract.

Sozune's middleware crates would ship this pattern as a generic utility (SozuneAsyncResolver<T>) so individual middlewares don't reimplement it. Documented as an example in doc/programmable-hooks.md, NOT shipped inside sozu-lib.

4.8 What's NOT in scope of Phase C

  • Per-listener thread-pool for hook I/O (would break lib/ no-async by introducing a tokio runtime). Operators doing thread-pool I/O run their pool in their Sozune binary, send results via the channel pattern above.
  • Cross-worker coordination of pending state (each worker's wake-up registry is independent). A pending hook on worker 1 doesn't block worker 2; same as sōzu's existing per-worker isolation.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Labels

No labels
No labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions