polyc-state-connect 2026.10.2

State plane transport adapter: capability-specific Connect clients and server-trait glue mapping the generated wire types onto the polyc-state kernel — typed outcomes, per-call admission, and the conformance surface the authenticated shell proves itself against (docs/proposals/separated-planes.md).
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
//! Runs synchronous State authority calls on Tokio's blocking pool
//! (POLY-372).
//!
//! A durable authority call takes a synchronous mutex and can wait on an
//! `fsync` while it holds it. Run that call on a runtime worker thread and
//! one slow commit holds the worker for the whole disk round trip. The State
//! pod's runtime has one worker per core, so a few concurrent commits starve
//! every request, feed pump, and probe sharing the runtime. The blocking
//! pool exists for exactly this shape of work.
//!
//! The request span does not cross [`tokio::task::spawn_blocking`] on its
//! own. Instrumenting the returned `JoinHandle` enters the span only while a
//! worker polls the handle; the closure itself runs on a blocking thread with
//! no span entered. Both helpers here enter the request span inside the
//! closure instead, so an event the authority emits is attributed to the call
//! that caused it.
//!
//! # One family's share of the pool
//!
//! The pool is one per runtime, so every family the helpers serve draws from
//! a per-family semaphore that admits at most [`MAX_FAMILY_BLOCKING_SLOTS`]
//! closures at once. A family's authority serializes on one mutex: the first
//! slot does the work and every further slot only parks a thread behind that
//! mutex. The cap stops one family whose mutex is queued behind a slow disk
//! from parking the whole pool with it. `polychrome-state`'s
//! `MAX_BLOCKING_THREADS` doc carries the arithmetic this constant appears
//! in.
//!
//! Families that serialize on the *same* mutex share one semaphore rather
//! than one each — see [`slot_family`]. The partition journal, the durable
//! commit feed, and the projection catalog all wait on the journal's one
//! mutation mutex, so their calls draw one shared share: anything wider would
//! let callers parked on that single lock occupy more of the pool than the
//! bound intends, which is the starvation shape this limiter exists to
//! remove. Verified replay is the opposite case on the same authority: it
//! releases that mutex before its expensive pass, so it queues on its own
//! [`REPLAY_BLOCKING_SLOTS`] share and a long replay cannot hold commit
//! slots.
//!
//! A call that finds its family's slots taken waits on the semaphore
//! asynchronously — it holds no pool thread while it waits — until its own
//! deadline ends the wait. The deadline is the tighter of the transport's
//! remaining deadline and the budget the call declared, which admission
//! computes once and hands the helper as `wait`.
//!
//! # Ordering and the failure model
//!
//! A call that is waiting for a slot has queued nothing. Its refusal is
//! [`StateError::Unavailable`] with [`OutageReach::NoDurableEffect`] for
//! mutations and reads alike, because no closure ever ran: the caller may
//! retry without risking a duplicate durable effect. A call that took a slot
//! behaves exactly as a call did before the limiter existed — only a lost
//! join handle keeps the weaker claim, `PossiblyApplied` for a mutation whose
//! closure may have committed before its outcome was lost, `NoDurableEffect`
//! for a read.
//!
//! The semaphore permit lives inside the blocking closure, so a slot is held
//! until the authority call ends rather than until the caller stops waiting.
//! A caller that drops its request mid-commit cannot free the slot early and
//! let a second call overlap durable work that is still running. A panic
//! inside the closure unwinds the permit with it, and a request dropped while
//! still waiting for a slot never held one — both leave the family able to
//! serve the next call.

use std::{
    collections::BTreeMap,
    sync::{Arc, Mutex, MutexGuard, OnceLock, PoisonError},
    time::Duration,
};

use connectrpc::ConnectError;
use polyc_state::{
    error::{OutageReach, StateError},
    feed,
    id::OperationFamily,
    journal, projection,
};
use prometheus::{IntCounterVec, IntGaugeVec, register_int_counter_vec, register_int_gauge_vec};
use tokio::sync::{OwnedSemaphorePermit, Semaphore};

use crate::error::WireOutcome;

/// Blocking-pool threads one operation family may hold at once.
///
/// A pooled family's authority serializes on one synchronous mutex, so one
/// slot does the work and the rest only queue behind it. Eight lets a burst
/// park a handful of threads on that mutex; calls past it wait on the
/// semaphore asynchronously, holding no thread, until their deadline.
/// `polychrome-state`'s `MAX_BLOCKING_THREADS` doc shows the sum this
/// constant feeds.
pub const MAX_FAMILY_BLOCKING_SLOTS: usize = 8;

/// Blocking-pool threads the verified-replay share admits at once.
///
/// Smaller than [`MAX_FAMILY_BLOCKING_SLOTS`] on purpose: a replay releases
/// the journal's mutation lock before its expensive pass, so it queues on
/// its own share rather than the journal gate's — but that pass can hold a
/// pool thread for a whole partition's history, the longest residency any
/// State closure has. Four lets a handful of partitions verify at once while
/// bounding how much of the pool that residency can occupy. The per-partition
/// stripe already serializes two replays of one partition, so the cap bounds
/// concurrency across partitions, not per partition.
pub const REPLAY_BLOCKING_SLOTS: usize = 4;

/// Runs one mutating authority call on the blocking pool under `span`.
///
/// `wait` is the longest the call may queue for a family slot before it is
/// refused; [`crate::admission::slot_wait`] computes it at admission. A
/// refusal carries [`OutageReach::NoDurableEffect`]: nothing was queued. A
/// lost join handle maps to [`OutageReach::PossiblyApplied`]: the closure may
/// have committed before its outcome was lost.
///
/// # Errors
///
/// Returns the authority's typed refusal unchanged, and
/// [`StateError::Unavailable`] — `NoDurableEffect` when no slot freed inside
/// `wait`, `PossiblyApplied` when the pool task's join resolves to an error.
pub(crate) async fn mutation<R, E>(
    span: tracing::Span,
    family: OperationFamily,
    wait: Duration,
    operation: impl FnOnce() -> Result<R, E> + Send + 'static,
) -> Result<R, ConnectError>
where
    R: Send + 'static,
    E: WireOutcome + Send + 'static,
{
    run(span, family, wait, OutageReach::PossiblyApplied, operation).await
}

/// Runs one read-only authority call on the blocking pool under `span`.
///
/// `wait` is the longest the call may queue for a family slot before it is
/// refused; [`crate::admission::slot_wait`] computes it at admission. Both a
/// refusal and a lost join handle map to [`OutageReach::NoDurableEffect`]: a
/// read holds no durable effect whether it queued, ran, or was lost.
///
/// # Errors
///
/// Returns the authority's typed refusal unchanged, and
/// [`StateError::Unavailable`] with `NoDurableEffect` when no slot freed
/// inside `wait` or the pool task's join resolves to an error.
pub(crate) async fn read<R, E>(
    span: tracing::Span,
    family: OperationFamily,
    wait: Duration,
    operation: impl FnOnce() -> Result<R, E> + Send + 'static,
) -> Result<R, ConnectError>
where
    R: Send + 'static,
    E: WireOutcome + Send + 'static,
{
    run(span, family, wait, OutageReach::NoDurableEffect, operation).await
}

/// The per-family semaphores the helpers draw from.
///
/// Keyed by [`OperationFamily`], not by service or route: one family's
/// authority is reachable through every router that mounts it, and the family
/// name is the budget those routes share. A semaphore is created on first use
/// and never closed, so a registered family keeps its slots for the life of
/// the process. The map lock is held for one lookup and never across a wait.
fn slots() -> MutexGuard<'static, BTreeMap<OperationFamily, Arc<Semaphore>>> {
    static SLOTS: OnceLock<Mutex<BTreeMap<OperationFamily, Arc<Semaphore>>>> = OnceLock::new();
    SLOTS
        .get_or_init(|| Mutex::new(BTreeMap::new()))
        .lock()
        .unwrap_or_else(PoisonError::into_inner)
}

/// The semaphore `family`'s calls queue on.
///
/// Usually the family itself. The commit-feed and projection-catalog
/// families are the exception: every call those services serve serializes on
/// the journal's single mutation mutex — the feed takes it through
/// `begin_feed_mutation`, the catalog takes it in every method — so the
/// three families draw one share keyed by the journal's name. Three
/// independent shares would let callers parked on that one lock hold
/// `3 * MAX_FAMILY_BLOCKING_SLOTS` pool threads at once: the same
/// starvation shape as letting subscriptions onto this pool. `state.versioned`
/// is the precedent: fifteen services already alias one family there because
/// they serialize on one `DurableVersionedState` mutex.
fn slot_family(family: &OperationFamily) -> OperationFamily {
    if *family == feed::family() || *family == projection::family() {
        journal::family()
    } else {
        family.clone()
    }
}

/// The slots one keyed share admits.
///
/// Every share gets [`MAX_FAMILY_BLOCKING_SLOTS`] except the verified-replay
/// family, which gets [`REPLAY_BLOCKING_SLOTS`] — see that constant for why
/// it is smaller.
fn slot_count(key: &OperationFamily) -> usize {
    if *key == journal::replay_family() {
        REPLAY_BLOCKING_SLOTS
    } else {
        MAX_FAMILY_BLOCKING_SLOTS
    }
}

/// Returns the semaphore `family`'s calls queue on, creating it on first use.
///
/// The map key is [`slot_family`]'s, not the caller's: families sharing one
/// mutex share one semaphore. Metrics and typed refusals keep the caller's
/// own family name, so a refusal still names the service that was refused.
fn family_slots(family: &OperationFamily) -> Arc<Semaphore> {
    let key = slot_family(family);
    slots()
        .entry(key.clone())
        .or_insert_with(|| Arc::new(Semaphore::new(slot_count(&key))))
        .clone()
}

/// Blocking-pool slots a family is holding, labeled `family`.
///
/// The gauge is the pool's active count per family: a family at
/// [`MAX_FAMILY_BLOCKING_SLOTS`] is where a parked mutex starts refusing new
/// calls rather than parking more threads.
fn slots_in_use() -> &'static IntGaugeVec {
    static METRIC: OnceLock<IntGaugeVec> = OnceLock::new();
    METRIC.get_or_init(|| {
        register_int_gauge_vec!(
            "polychrome_state_blocking_pool_slots_in_use",
            "Blocking-pool slots each State operation family is holding right now.",
            &["family"]
        )
        .expect("the blocking-pool slots gauge registers once into the default registry")
    })
}

/// Slot waits that ended in refusal, labeled `family`.
///
/// A refusal means the family's slots stayed taken until the call's own
/// deadline ended the wait — the signal the re-enable window reads for a
/// family whose pool share is saturated.
fn slot_refusals() -> &'static IntCounterVec {
    static METRIC: OnceLock<IntCounterVec> = OnceLock::new();
    METRIC.get_or_init(|| {
        register_int_counter_vec!(
            "polychrome_state_blocking_pool_slot_refusals_total",
            "State calls refused because no blocking-pool slot freed inside the call's deadline, \
             by family.",
            &["family"]
        )
        .expect("the slot-refusals counter registers once into the default registry")
    })
}

/// One family's slot: the semaphore permit and the in-use count live and die
/// together.
///
/// Moved into the blocking closure so the slot is held for the authority
/// call's whole run, not the caller's patience. A panic inside the closure
/// unwinds it; the counter then shows the family drawing one slot fewer.
struct HeldSlot {
    _permit: OwnedSemaphorePermit,
    family: OperationFamily,
}

impl HeldSlot {
    fn new(permit: OwnedSemaphorePermit, family: &OperationFamily) -> Self {
        slots_in_use().with_label_values(&[family.as_str()]).inc();
        Self {
            _permit: permit,
            family: family.clone(),
        }
    }
}

impl Drop for HeldSlot {
    fn drop(&mut self) {
        slots_in_use()
            .with_label_values(&[self.family.as_str()])
            .dec();
    }
}

/// Refuses a call that could not take a family slot inside `wait`.
///
/// Nothing was queued, so the reach is [`OutageReach::NoDurableEffect`] for
/// mutations and reads alike. The warning names the family and the bound it
/// waited against, and never the request's contents.
fn refuse<E: WireOutcome>(family: OperationFamily, wait: Duration) -> ConnectError {
    slot_refusals().with_label_values(&[family.as_str()]).inc();
    tracing::warn!(
        family = %family,
        wait_ms = wait.as_millis(),
        "a State call found no blocking-pool slot inside its deadline"
    );
    E::from(StateError::Unavailable {
        family,
        reach: OutageReach::NoDurableEffect,
    })
    .to_connect()
}

/// Dispatches `operation` onto the blocking pool with `span` entered on the
/// blocking thread, after the call takes a slot on its family's semaphore.
///
/// `reach` is the claim a lost join handle may make: `PossiblyApplied` for a
/// mutation, `NoDurableEffect` for a read. A refusal before queueing claims
/// `NoDurableEffect` regardless — no closure ever ran.
async fn run<R, E>(
    span: tracing::Span,
    family: OperationFamily,
    wait: Duration,
    reach: OutageReach,
    operation: impl FnOnce() -> Result<R, E> + Send + 'static,
) -> Result<R, ConnectError>
where
    R: Send + 'static,
    E: WireOutcome + Send + 'static,
{
    let slots = family_slots(&family);
    // The registry never closes a family's semaphore; treating a close as
    // the same refusal keeps a hypothetical one from panicking a worker.
    let Ok(Ok(permit)) = tokio::time::timeout(wait, slots.acquire_owned()).await else {
        return Err(refuse::<E>(family, wait));
    };
    let slot = HeldSlot::new(permit, &family);
    tokio::task::spawn_blocking(move || {
        let _slot = slot;
        span.in_scope(operation)
    })
    .await
    .map_err(|_| {
        E::from(StateError::Unavailable {
            family: family.clone(),
            reach,
        })
        .to_connect()
    })?
    .map_err(|error| error.to_connect())
}

#[cfg(test)]
mod tests {
    use std::collections::BTreeSet;

    use polyc_state::{
        burn, claims, feed, id::OperationFamily, immutable, ingress, journal, model_attempt,
        observation, passkey, persona_memory, projection, query_audit, spend, tasks, usage,
        versioned, wallet,
    };

    use super::{slot_count, slot_family};

    /// The semaphore-key census behind `MAX_BLOCKING_THREADS`'s arithmetic
    /// in `crates/state-service/src/main.rs`: every family a state-connect
    /// service can name at a `blocking::*` call site, mapped through
    /// [`slot_family`], must collapse to the nine keys that fit the pool.
    ///
    /// Mutation: set [`super::REPLAY_BLOCKING_SLOTS`] to 64 — the sum
    /// outgrows the pool and this fails.
    #[test]
    fn the_family_slot_shares_fit_the_pool() {
        // The aliases are the point: burn, passkey, spend, tasks, usage,
        // wallet, model-attempt, and persona-memory store calls all resolve
        // to `state.versioned` because they serialize on the one
        // `DurableVersionedState` mutex, the way journal, feed, and
        // projection calls resolve to `state.partition-journal`.
        let caller_families = [
            versioned::family(),
            burn::family(),
            passkey::family(),
            spend::family(),
            tasks::family(),
            usage::family(),
            wallet::family(),
            model_attempt::family(),
            persona_memory::store::family(),
            persona_memory::journal::family(),
            claims::family(),
            ingress::family(),
            immutable::family(),
            observation::family(),
            query_audit::family(),
            journal::family(),
            feed::family(),
            projection::family(),
            journal::replay_family(),
        ];
        let keys: BTreeSet<_> = caller_families.iter().map(slot_family).collect();
        let names: BTreeSet<_> = keys.iter().map(OperationFamily::as_str).collect();
        assert_eq!(
            names,
            BTreeSet::from([
                "state.claims",
                "state.immutable_object",
                "state.ingress",
                "state.observation",
                "state.partition-journal",
                "state.partition-journal-replay",
                "state.persona-memory-journal",
                "state.query_audit",
                "state.versioned",
            ]),
            "the distinct semaphore keys the services draw"
        );
        let ceilings: usize = keys.iter().map(slot_count).sum();
        // The extra draw is the one-time startup conformance task; 128 is
        // `MAX_BLOCKING_THREADS` in `crates/state-service/src/main.rs`, so
        // `ceilings + 1 <= 128` reads here as `ceilings < 128`.
        assert!(
            ceilings < 128,
            "{ceilings} slot ceilings plus the startup task must fit the pool"
        );
    }
}