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).
//! A dedicated executor feed subscription pumps run on.
//!
//! Every other RPC family this listener serves (ingress, claims, versioned
//! reads, model-attempt reads, persona-memory-journal reads, query-audit
//! reads) runs its blocking work on `tokio::task::spawn_blocking`, releasing
//! its thread as soon as one call's synchronous work finishes. A feed
//! subscription's pump does not: it holds a blocking-pool thread for the
//! whole life of the stream. Sharing one pool between the two shapes of work
//! means a burst of long-lived subscriptions can starve every other family's
//! blocking calls of a thread — observed in production on 2026-09-11
//! (POLY-324): five projector families and Query held enough subscriptions
//! to leave the shared pool with no thread free for ingress, so an edge's
//! durable-receipt commit missed its budget and the message was lost, even
//! though the commit itself completed in milliseconds once it finally ran.
//!
//! This executor gives subscription work its own pool, so a listener sizes
//! the shared pool for the other families alone and a subscription burst
//! never touches it.

use std::io;

use crate::feed::MAX_CONCURRENT_SUBSCRIPTIONS;

/// Blocking-pool threads the feed executor reserves beyond the subscription
/// budget.
///
/// The admission gate (`crate::feed::tailer::Tailer::try_admit`) refuses a
/// subscription past budget rather than over-admitting it, so at most
/// [`MAX_CONCURRENT_SUBSCRIPTIONS`] pumps ever exist at once — a subscription
/// paused past `FeedStreamTuning::max_tail_idle` and reconnecting is refused
/// the same as any other arrival if the budget is already spent, never
/// admitted past it. This headroom is not covering an over-admission; it is
/// margin for the moment a finishing pump's thread has not yet returned to
/// the pool when a fresh dispatch needs one, so that dispatch runs
/// immediately instead of queuing briefly behind it.
pub const FEED_EXECUTOR_HEADROOM: usize = 16;

/// How many blocking-pool threads the feed executor runs.
///
/// One thread per subscription this listener admits, plus
/// [`FEED_EXECUTOR_HEADROOM`]. This is the whole reason the executor is
/// dedicated rather than shared: the listener's own blocking pool
/// (`polychrome-state`'s `main`) sizes for the *other* RPC families alone,
/// and never has to grow with the subscription budget.
pub const FEED_EXECUTOR_THREADS: usize = MAX_CONCURRENT_SUBSCRIPTIONS + FEED_EXECUTOR_HEADROOM;

/// The executor exceeds the subscription budget by real headroom.
///
/// Checked at compile time rather than by a test, for the same reason
/// `polychrome-state`'s own blocking-pool bound is: a headroom edited down to
/// zero is silent until the pause/resume churn it was margin for finally
/// exhausts the pool.
const _: () = {
    assert!(
        FEED_EXECUTOR_HEADROOM > 0,
        "zero headroom leaves no thread for a fresh dispatch to run on while a \
         finishing pump's thread has not yet returned to the pool"
    );
    assert!(
        FEED_EXECUTOR_THREADS > MAX_CONCURRENT_SUBSCRIPTIONS,
        "the executor must exceed the subscription budget, or the last admitted \
         subscription has no thread to pump on"
    );
};

/// Runs feed subscription work off the listener's shared blocking pool.
///
/// Built once per listener and shared by every [`crate::feed::FeedSvc`]
/// mounted over the same feed module — a feed module is served on more than
/// one router in production (the combined listener and its family-specific
/// fallback), and both must draw from the same pool, or the budget one
/// admits against is not the budget the other runs against.
///
/// # Dropping this from an async context
///
/// Always safe, deliberately: [`Drop`] shuts the owned runtime down in the
/// background rather than joining its threads synchronously. A bare
/// [`tokio::runtime::Runtime`] panics if its ordinary drop glue runs on a
/// thread another runtime is currently driving — exactly where the last
/// [`Arc`](std::sync::Arc) to this executor is dropped in production (inside
/// an async request handler) and in the loopback test harness (inside a
/// spawned task).
#[derive(Debug)]
pub struct FeedExecutor {
    // `Option` so `Drop` can take ownership of the runtime and hand it to
    // `shutdown_background`, which only accepts `self` by value.
    runtime: Option<tokio::runtime::Runtime>,
}

impl FeedExecutor {
    /// Starts an executor sized for [`FEED_EXECUTOR_THREADS`].
    ///
    /// # Errors
    ///
    /// Returns an error if the runtime's threads cannot be started.
    pub fn start_production() -> io::Result<Self> {
        Self::start(FEED_EXECUTOR_THREADS)
    }

    /// Starts an executor with `max_blocking` threads available to
    /// [`FeedExecutor::spawn_pump`].
    ///
    /// Never drives an async task of its own: this runtime exists only to
    /// own a blocking pool, so one implicit worker thread is enough, and
    /// that thread never runs anything (nothing here ever calls
    /// `block_on`). The blocking pool's threads are spawned lazily, on
    /// first use, independent of that worker.
    ///
    /// # Errors
    ///
    /// Returns an error if the runtime's threads cannot be started.
    pub fn start(max_blocking: usize) -> io::Result<Self> {
        tokio::runtime::Builder::new_current_thread()
            .max_blocking_threads(max_blocking)
            .thread_name("polychrome-feed-pump")
            .enable_all()
            .build()
            .map(|runtime| Self {
                runtime: Some(runtime),
            })
    }

    /// Runs `job` on this executor's own blocking pool.
    ///
    /// Never the caller's runtime, whatever that is — the whole point of a
    /// dedicated executor is that a caller on the listener's shared runtime
    /// pays nothing for this call beyond dispatching it.
    ///
    /// # Panics
    ///
    /// Panics if called after this executor was already dropped, which
    /// cannot happen through the public API: nothing hands back a
    /// `FeedExecutor` by value, only shared references to one still owned by
    /// an `Arc`.
    pub fn spawn_pump(&self, job: impl FnOnce() + Send + 'static) {
        self.runtime
            .as_ref()
            .expect("spawn_pump never runs after drop")
            .handle()
            .spawn_blocking(job);
    }
}

impl Drop for FeedExecutor {
    fn drop(&mut self) {
        if let Some(runtime) = self.runtime.take() {
            // Never `.shutdown_timeout` or the runtime's own default drop:
            // both block the calling thread until every blocking-pool
            // thread joins, which panics when the calling thread is itself
            // inside another runtime's async task — see this type's own
            // doc. `shutdown_background` returns immediately and lets the
            // pool's threads finish on their own.
            runtime.shutdown_background();
        }
    }
}

#[cfg(test)]
mod tests {
    use std::{
        sync::{
            Arc,
            atomic::{AtomicUsize, Ordering},
        },
        time::Duration,
    };

    use super::FeedExecutor;

    #[tokio::test]
    async fn dropping_the_executor_inside_another_runtimes_async_task_does_not_panic() {
        // The exact shape that panicked before `Drop` used
        // `shutdown_background`: build the executor, then drop the value
        // (not just a clone of a shared reference) from inside a `#[tokio::
        // test]`'s own async context — the same context a request handler,
        // or the loopback test harness's spawned task, drops the listener's
        // last `Arc<FeedExecutor>` in.
        let executor = FeedExecutor::start(2).expect("start");
        drop(executor);
    }

    #[test]
    fn spawned_work_runs_without_a_caller_runtime() {
        // No `#[tokio::test]` here on purpose: this proves the executor
        // needs no ambient runtime to run blocking work, which is exactly
        // what lets a listener build it once at startup, outside any
        // request's async context.
        let executor = FeedExecutor::start(4).expect("start");
        let ran = Arc::new(AtomicUsize::new(0));
        let done = Arc::clone(&ran);
        executor.spawn_pump(move || {
            done.fetch_add(1, Ordering::SeqCst);
        });
        let deadline = std::time::Instant::now() + Duration::from_secs(5);
        while ran.load(Ordering::SeqCst) == 0 {
            assert!(
                std::time::Instant::now() < deadline,
                "spawned job never ran"
            );
            std::thread::sleep(Duration::from_millis(5));
        }
    }

    #[test]
    fn many_spawned_jobs_run_concurrently_up_to_the_pool_size() {
        let executor = FeedExecutor::start(8).expect("start");
        let barrier = Arc::new(std::sync::Barrier::new(8));
        let mut handles = Vec::new();
        for _ in 0..8 {
            let (tx, rx) = std::sync::mpsc::channel();
            handles.push(rx);
            let barrier = Arc::clone(&barrier);
            executor.spawn_pump(move || {
                // Every one of the 8 jobs reaches the barrier: proves the
                // pool actually runs 8 blocking closures concurrently
                // rather than queueing them one at a time.
                barrier.wait();
                let _ = tx.send(());
            });
        }
        let deadline = std::time::Instant::now() + Duration::from_secs(5);
        for rx in handles {
            let remaining = deadline.saturating_duration_since(std::time::Instant::now());
            rx.recv_timeout(remaining)
                .expect("every job reaches the barrier and completes");
        }
    }
}