pub mod envelope;
use std::cell::{Cell, OnceCell};
use std::fmt;
use std::future::Future;
use std::sync::OnceLock;
use asupersync::runtime::Runtime;
use asupersync::runtime::RuntimeBuilder;
use asupersync::runtime::reactor::create_reactor;
use crate::crypto::{RandomDrawError, SECURITY_IDENTIFIER_BYTES, draw_security_identifier};
pub const SNAPSHOT_CLONE_IS_DETECTABLE: bool = false;
#[derive(Clone, Copy, PartialEq, Eq)]
pub struct ProcessGeneration {
pid: u32,
nonce: [u8; SECURITY_IDENTIFIER_BYTES],
generation: u64,
}
impl ProcessGeneration {
#[must_use]
pub const fn pid(&self) -> u32 {
self.pid
}
#[must_use]
pub const fn generation(&self) -> u64 {
self.generation
}
#[must_use]
pub const fn nonce(&self) -> &[u8; SECURITY_IDENTIFIER_BYTES] {
&self.nonce
}
#[must_use]
pub const fn observed(
pid: u32,
nonce: [u8; SECURITY_IDENTIFIER_BYTES],
generation: u64,
) -> Self {
Self {
pid,
nonce,
generation,
}
}
pub fn admit(&self, observed: Self) -> Result<(), ProcessGenerationError> {
if self.pid != observed.pid || self.nonce != observed.nonce {
return Err(ProcessGenerationError::ForkDetected {
installed_pid: self.pid,
observed_pid: observed.pid,
generation: self.generation,
});
}
if self.generation != observed.generation {
return Err(ProcessGenerationError::GenerationMismatch {
expected: self.generation,
observed: observed.generation,
});
}
Ok(())
}
}
impl fmt::Debug for ProcessGeneration {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("ProcessGeneration")
.field("pid", &self.pid)
.field("generation", &self.generation)
.field("nonce", &"<redacted>")
.finish()
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum ProcessGenerationError {
ForkDetected {
installed_pid: u32,
observed_pid: u32,
generation: u64,
},
GenerationMismatch {
expected: u64,
observed: u64,
},
NotInstalled,
EntropyUnavailable(RandomDrawError),
SnapshotCloneDeploymentUnsupported,
}
impl fmt::Display for ProcessGenerationError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::ForkDetected {
installed_pid,
observed_pid,
generation,
} => write!(
formatter,
"process generation {generation} was installed in pid {installed_pid} but is \
being used from pid {observed_pid}; inherited runtime, key, continuation, \
quota and supervisor state is not usable after fork — the child must exec or \
build a wholly new FastMCP instance"
),
Self::GenerationMismatch { expected, observed } => write!(
formatter,
"process generation mismatch: resource carries generation {observed} but the \
installed generation is {expected}"
),
Self::NotInstalled => formatter
.write_str("process-local state was used before ProcessGenerationGuard::install()"),
Self::EntropyUnavailable(source) => write!(
formatter,
"process generation nonce could not be drawn: {source}"
),
Self::SnapshotCloneDeploymentUnsupported => formatter.write_str(
"deployment permits live-memory snapshot/CRIU/VM/container cloning but neither \
disables ephemeral process-local protected state nor supplies a \
rollback-and-clone-resistant external epoch; process-local replay and nonce \
safety cannot be claimed under such cloning",
),
}
}
}
impl std::error::Error for ProcessGenerationError {}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct ExternalEpoch {
rollback_resistant: bool,
clone_resistant: bool,
}
impl ExternalEpoch {
#[must_use]
pub const fn new(rollback_resistant: bool, clone_resistant: bool) -> Self {
Self {
rollback_resistant,
clone_resistant,
}
}
#[must_use]
pub const fn is_conforming(&self) -> bool {
self.rollback_resistant && self.clone_resistant
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum SnapshotCloneStance {
NoLiveMemoryCloning,
LiveMemoryCloningPermitted {
ephemeral_protected_state_disabled: bool,
external_epoch: Option<ExternalEpoch>,
},
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct ProcessBoundToken {
minted_in: ProcessGeneration,
}
impl ProcessBoundToken {
#[must_use]
pub const fn minted_in(&self) -> ProcessGeneration {
self.minted_in
}
pub fn verify(&self) -> Result<(), ProcessGenerationError> {
let installed = ProcessGenerationGuard::installed()
.ok_or(ProcessGenerationError::NotInstalled)?
.generation;
let observed =
ProcessGeneration::observed(std::process::id(), installed.nonce, installed.generation);
installed.admit(observed)?;
self.minted_in.admit(observed)
}
}
pub struct ProcessGenerationGuard {
generation: ProcessGeneration,
}
static GUARD: OnceLock<ProcessGenerationGuard> = OnceLock::new();
impl ProcessGenerationGuard {
pub fn install() -> Result<&'static Self, ProcessGenerationError> {
if let Some(existing) = GUARD.get() {
existing.verify_current()?;
return Ok(existing);
}
let nonce =
draw_security_identifier().map_err(ProcessGenerationError::EntropyUnavailable)?;
let candidate = Self {
generation: ProcessGeneration {
pid: std::process::id(),
nonce: *nonce.as_bytes(),
generation: 0,
},
};
let installed = GUARD.get_or_init(|| candidate);
installed.verify_current()?;
Ok(installed)
}
#[must_use]
pub fn installed() -> Option<&'static Self> {
GUARD.get()
}
#[must_use]
pub const fn generation(&self) -> ProcessGeneration {
self.generation
}
#[must_use]
pub const fn token(&self) -> ProcessBoundToken {
ProcessBoundToken {
minted_in: self.generation,
}
}
pub fn verify_current(&self) -> Result<(), ProcessGenerationError> {
self.token().verify()
}
pub fn admit_snapshot_stance(
stance: SnapshotCloneStance,
) -> Result<(), ProcessGenerationError> {
match stance {
SnapshotCloneStance::NoLiveMemoryCloning => Ok(()),
SnapshotCloneStance::LiveMemoryCloningPermitted {
ephemeral_protected_state_disabled,
external_epoch,
} => {
if ephemeral_protected_state_disabled {
return Ok(());
}
match external_epoch {
Some(epoch) if epoch.is_conforming() => Ok(()),
_ => Err(ProcessGenerationError::SnapshotCloneDeploymentUnsupported),
}
}
}
}
}
const MAX_BLOCKING_THREADS: usize = 16;
thread_local! {
static RUNTIME: OnceCell<Runtime> = const { OnceCell::new() };
static BRIDGE_ACTIVE: Cell<bool> = const { Cell::new(false) };
static BRIDGE_NESTED_IN_TASK: Cell<bool> = const { Cell::new(false) };
static BLOCKING_LANE: Cell<bool> = const { Cell::new(false) };
}
#[must_use = "the lane declaration ends when the guard is dropped"]
pub struct BlockingLaneGuard {
previous: bool,
}
impl Drop for BlockingLaneGuard {
fn drop(&mut self) {
BLOCKING_LANE.with(|lane| lane.set(self.previous));
}
}
pub fn enter_blocking_lane() -> BlockingLaneGuard {
BlockingLaneGuard {
previous: BLOCKING_LANE.with(|lane| lane.replace(true)),
}
}
#[must_use]
pub fn bridge_would_starve_its_driver() -> bool {
BRIDGE_NESTED_IN_TASK.with(Cell::get) && !BLOCKING_LANE.with(Cell::get)
}
struct BridgeEntry {
previous_nested_in_task: bool,
}
impl BridgeEntry {
fn enter() -> Self {
BRIDGE_ACTIVE.with(|active| {
assert!(
!active.replace(true),
"nested fastmcp_core::runtime::block_on is not supported"
);
});
Self {
previous_nested_in_task: BRIDGE_NESTED_IN_TASK
.with(|nested| nested.replace(asupersync::Cx::is_active())),
}
}
}
impl Drop for BridgeEntry {
fn drop(&mut self) {
BRIDGE_NESTED_IN_TASK.with(|nested| nested.set(self.previous_nested_in_task));
BRIDGE_ACTIVE.with(|active| active.set(false));
}
}
pub fn block_on<F: Future>(future: F) -> F::Output {
ProcessGenerationGuard::install().unwrap_or_else(|error| {
panic!("refusing to drive a FastMCP runtime across a process generation: {error}")
});
let _entry = BridgeEntry::enter();
RUNTIME.with(|runtime| {
let runtime = runtime.get_or_init(|| {
let reactor = create_reactor().expect("failed to create platform I/O reactor");
RuntimeBuilder::current_thread()
.with_reactor(reactor)
.blocking_threads(0, MAX_BLOCKING_THREADS)
.build()
.expect("failed to build asupersync runtime")
});
runtime.block_on(future)
})
}
pub fn poll_on_cx<F: Future>(cx: &asupersync::Cx, future: F) -> F::Output {
use std::sync::Arc;
use std::task::{Context, Poll, Wake, Waker};
struct ThreadWake(std::thread::Thread);
impl Wake for ThreadWake {
fn wake(self: Arc<Self>) {
self.0.unpark();
}
fn wake_by_ref(self: &Arc<Self>) {
self.0.unpark();
}
}
let _entry = BridgeEntry::enter();
let _current = asupersync::Cx::set_current(Some(cx.clone()));
let waker = Waker::from(Arc::new(ThreadWake(std::thread::current())));
let mut task_cx = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
if let Poll::Ready(output) = future.as_mut().poll(&mut task_cx) {
return output;
}
std::thread::park_timeout(std::time::Duration::from_millis(1));
}
}
#[cfg(test)]
mod tests {
use super::{ProcessGenerationGuard, block_on};
#[test]
fn block_on_runs_async_blocks() {
let out = block_on(async { 1 + 1 });
assert_eq!(out, 2);
}
#[test]
fn block_on_can_be_called_multiple_times() {
let a = block_on(async { "a" });
let b = block_on(async { "b" });
assert_eq!(a, "a");
assert_eq!(b, "b");
}
#[test]
fn block_on_installs_guard_before_polling() {
const CHILD: &str = "FASTMCP_TEST_BRIDGE_FIRST_ENTRY";
if std::env::var_os(CHILD).is_none() {
let output = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"runtime::tests::block_on_installs_guard_before_polling",
"--nocapture",
"--test-threads=1",
])
.env(CHILD, "1")
.output()
.unwrap();
assert!(
output.status.success(),
"fresh-process bridge test failed:\n{}\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
return;
}
assert!(ProcessGenerationGuard::installed().is_none());
let generation = block_on(async {
let guard = ProcessGenerationGuard::installed()
.expect("the guard must exist before user code is polled");
guard.verify_current().unwrap();
guard.generation()
});
assert_eq!(generation.pid(), std::process::id());
assert_eq!(
ProcessGenerationGuard::install().unwrap().generation(),
generation
);
}
#[test]
fn nested_bridge_is_rejected_without_poisoning_outer_entry() {
block_on(async {
for _ in 0..2 {
let error = std::panic::catch_unwind(|| block_on(async { 1 }));
assert!(error.is_err(), "each nested entry must be rejected");
}
});
assert_eq!(block_on(async { 7 }), 7);
}
#[test]
fn poll_on_cx_drives_under_the_callers_cx_and_restores_the_thread() {
use asupersync::Cx;
let cx = Cx::for_testing();
assert!(Cx::current().is_none());
let driven = super::poll_on_cx(&cx, async {
asupersync::runtime::yield_now().await;
Cx::current().is_some()
});
assert!(
driven,
"the future is polled with the caller's cx installed"
);
assert!(
Cx::current().is_none(),
"the thread's previous cx is restored"
);
assert_eq!(super::poll_on_cx(&cx, async { 5 }), 5);
}
#[test]
fn poll_on_cx_entered_from_a_task_names_the_starved_driver() {
let runtime = asupersync::runtime::RuntimeBuilder::current_thread()
.build()
.unwrap();
let starved = runtime.block_on(async {
let cx = asupersync::Cx::current().expect("a runtime task has a cx");
super::poll_on_cx(&cx, async { super::bridge_would_starve_its_driver() })
});
assert!(starved, "a bridge that occupies its task's driver says so");
}
#[test]
fn poll_on_cx_entered_from_a_bare_thread_starves_no_driver() {
let cx = asupersync::Cx::for_testing();
let starved = super::poll_on_cx(&cx, async { super::bridge_would_starve_its_driver() });
assert!(!starved, "a bare thread occupies no task's driver");
}
#[test]
fn poll_on_cx_nested_in_a_bridge_is_rejected() {
block_on(async {
let cx = asupersync::Cx::for_testing();
let nested = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
super::poll_on_cx(&cx, async { 1 })
}));
assert!(nested.is_err(), "a second bridge on the thread is rejected");
});
let cx = asupersync::Cx::for_testing();
assert_eq!(super::poll_on_cx(&cx, async { 3 }), 3);
}
#[test]
fn bridge_entry_is_released_when_the_future_panics() {
let error = std::panic::catch_unwind(|| {
block_on(async { panic!("bridge test panic") });
});
assert!(error.is_err());
assert_eq!(block_on(async { 11 }), 11);
}
#[test]
fn guard_install_is_idempotent_across_threads() {
let expected = ProcessGenerationGuard::install().unwrap().generation();
let installers: Vec<_> = (0..8)
.map(|_| {
std::thread::spawn(|| {
let guard = ProcessGenerationGuard::install().unwrap();
guard.verify_current().unwrap();
guard.generation()
})
})
.collect();
for installer in installers {
assert_eq!(installer.join().unwrap(), expected);
}
}
}