Add corex library with resource-aware execution primitives
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (default features) (push) Canceled after 0s
CI / Test (all features) (push) Canceled after 0s
CI / Test (bg only) (push) Canceled after 0s
CI / Test (compute only) (push) Canceled after 0s
CI / Test (io only) (push) Canceled after 0s
CI / Run example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (push) Canceled after 0s

- Implement CI workflow for formatting, linting, testing, and documentation checks.
- Create publish workflow for automated publishing to crates.io.
- Add .gitignore to exclude build artifacts and editor files.
- Define Cargo.toml for corex and corex-core with dependencies and metadata.
- Add README.md files for corex and corex-core with usage instructions and licensing.
- Implement core execution primitives: spawn_io, compute, and spawn_bg with panic isolation.
- Establish global engine context for resource management based on logical CPU cores.
- Introduce error handling for compute panics.
This commit is contained in:
asepharyana
2026-08-28 16:05:04 +07:00
commit d8144e249e
17 changed files with 898 additions and 0 deletions
+31
View File
@@ -0,0 +1,31 @@
[package]
name = "corex-core"
version.workspace = true
edition.workspace = true
rust-version.workspace = true
license.workspace = true
repository.workspace = true
homepage.workspace = true
documentation.workspace = true
authors.workspace = true
description = "Core allocation logic and execution primitives for corex."
readme = "README.md"
keywords = ["async", "concurrency", "rayon", "tokio", "resource-management"]
categories = ["asynchronous", "concurrency", "rust-patterns"]
[features]
default = []
io = ["dep:tokio"]
compute = ["dep:rayon"]
bg = ["dep:tokio"]
full = ["io", "compute", "bg"]
[dependencies]
tokio = { workspace = true, optional = true }
rayon = { workspace = true, optional = true }
num_cpus = { workspace = true }
tracing = { workspace = true }
[package.metadata.docs.rs]
all-features = true
rustdoc-args = ["--cfg", "docsrs"]
+9
View File
@@ -0,0 +1,9 @@
# corex-core
Core allocation logic, global context initialization, and execution primitives for [`corex`](https://crates.io/crates/corex).
Applications should depend on the `corex` facade crate rather than this crate directly.
## License
Licensed under either of Apache License, Version 2.0 or MIT license at your option.
+53
View File
@@ -0,0 +1,53 @@
//! Bounded-concurrency background task execution.
use tracing::Instrument;
use crate::context::context;
/// Spawns `future` as a background task once a concurrency permit is
/// available, returning its [`tokio::task::JoinHandle`].
///
/// At most [`crate::context::EngineContext::bg_concurrency`] background
/// tasks run at any one time; awaiting `spawn_bg` blocks the caller until a
/// slot frees up, which is what provides the bound. A task's permit is held
/// for the task's full lifetime and released automatically when it
/// completes.
///
/// Panic isolation is provided by Tokio itself: a panicking background task
/// cannot crash the runtime or any sibling task, and is surfaced to the
/// caller as `Err(JoinError)` when the returned handle is awaited, exactly
/// as with [`crate::io::spawn_io`].
///
/// # Panics
///
/// Panics if called outside the context of a running Tokio runtime.
pub async fn spawn_bg<F>(future: F) -> tokio::task::JoinHandle<F::Output>
where
F: std::future::Future + Send + 'static,
F::Output: Send + 'static,
{
let permit = context()
.bg_semaphore
.acquire()
.await
.expect("corex: bg semaphore closed unexpectedly");
let span = tracing::info_span!("corex_bg_task");
tokio::spawn(
async move {
let _permit = permit;
future.await
}
.instrument(span),
)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn spawn_bg_roundtrips_a_value() {
let handle = spawn_bg(async { 9u32 }).await;
assert_eq!(handle.await.unwrap(), 9);
}
}
+60
View File
@@ -0,0 +1,60 @@
//! Panic-isolated heavy compute on a sized [`rayon::ThreadPool`].
use std::panic::{catch_unwind, AssertUnwindSafe};
use crate::context::context;
use crate::error::CorexError;
/// Runs `f` on the global compute thread pool and returns its result.
///
/// If `f` panics, the panic is caught and converted into
/// [`CorexError::ComputePanic`] instead of unwinding across the pool
/// boundary or poisoning the pool; subsequent calls to [`compute`] continue
/// to work normally.
///
/// # Panic-safety caveat
///
/// `f` is wrapped in [`AssertUnwindSafe`] so that closures capturing
/// ordinary references or non-[`UnwindSafe`](std::panic::UnwindSafe) state
/// can be submitted without a compile error. This is sound with respect to
/// the compute pool itself, since a panicking closure's stack (and any
/// locals it holds) is discarded entirely rather than observed afterward.
/// It does not, however, guarantee exception-safety of state the closure
/// captured by mutable reference: if `f` panics partway through mutating a
/// captured `&mut T`, that `T` may be left in an inconsistent state from
/// the caller's perspective.
pub fn compute<F, R>(f: F) -> Result<R, CorexError>
where
F: FnOnce() -> R + Send,
R: Send,
{
let wrapped = AssertUnwindSafe(f);
context()
.compute_pool
.install(move || catch_unwind(wrapped))
.map_err(|payload| CorexError::ComputePanic(panic_payload_to_string(payload)))
}
fn panic_payload_to_string(payload: Box<dyn std::any::Any + Send>) -> String {
if let Some(message) = payload.downcast_ref::<&str>() {
(*message).to_string()
} else if let Some(message) = payload.downcast_ref::<String>() {
message.clone()
} else {
"compute closure panicked with a non-string payload".to_string()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn compute_panic_is_isolated_and_pool_survives() {
let panicked: Result<u32, CorexError> = compute(|| panic!("boom"));
assert!(matches!(panicked, Err(CorexError::ComputePanic(_))));
let recovered = compute(|| 1 + 1);
assert_eq!(recovered.unwrap(), 2);
}
}
+98
View File
@@ -0,0 +1,98 @@
//! Global engine context providing resource-aware execution primitives.
//!
//! The context is initialized lazily on first access via [`context`], or
//! explicitly via [`init`]. Resource sizing is derived from the number of
//! logical CPU cores reported by [`num_cpus::get`].
use std::sync::OnceLock;
static CONTEXT: OnceLock<EngineContext> = OnceLock::new();
/// The global, lazily-initialized engine context.
///
/// Holds the computed thread and concurrency counts for each corex
/// subsystem, along with the resource pools those counts were used to
/// build. The context lives for the lifetime of the process once
/// initialized: it is stored in a `'static` [`OnceLock`] and is never
/// dropped.
pub struct EngineContext {
/// Logical CPU core count used for sizing async I/O scheduling.
pub io_threads: usize,
/// Number of worker threads allocated to the compute [`rayon::ThreadPool`].
pub compute_threads: usize,
/// Maximum number of background tasks permitted to run concurrently.
pub bg_concurrency: usize,
#[cfg(feature = "compute")]
pub(crate) compute_pool: rayon::ThreadPool,
#[cfg(feature = "bg")]
pub(crate) bg_semaphore: tokio::sync::Semaphore,
}
impl EngineContext {
/// Builds a new [`EngineContext`] sized from the current machine's
/// logical core count.
///
/// # Panics
///
/// Panics if the compute thread pool cannot be constructed. This only
/// happens under an unrecoverable environment failure, such as the
/// operating system refusing to spawn any new thread.
fn build() -> Self {
let logical_cores = num_cpus::get();
let io_threads = logical_cores;
let compute_threads = logical_cores.saturating_sub(1).max(1);
let bg_concurrency = (logical_cores / 2).max(2);
#[cfg(feature = "compute")]
let compute_pool = rayon::ThreadPoolBuilder::new()
.num_threads(compute_threads)
.thread_name(|index| format!("corex-compute-{index}"))
.build()
.expect("corex: failed to build rayon compute thread pool");
#[cfg(feature = "bg")]
let bg_semaphore = tokio::sync::Semaphore::new(bg_concurrency);
Self {
io_threads,
compute_threads,
bg_concurrency,
#[cfg(feature = "compute")]
compute_pool,
#[cfg(feature = "bg")]
bg_semaphore,
}
}
}
/// Returns the global [`EngineContext`], initializing it on first access.
///
/// Safe to call from any thread at any time; initialization happens
/// exactly once regardless of how many callers race to trigger it.
pub fn context() -> &'static EngineContext {
CONTEXT.get_or_init(EngineContext::build)
}
/// Bootstraps the global [`EngineContext`].
///
/// Behaviorally identical to [`context`]; provided as the explicit,
/// discoverable entry point applications call at startup to force
/// initialization eagerly, for example so resource sizing can be logged
/// before any workload runs. Calling it more than once, or never calling
/// it at all before using [`context`], is equally correct.
pub fn init() -> &'static EngineContext {
CONTEXT.get_or_init(EngineContext::build)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn context_numbers_are_sane() {
let ctx = context();
assert!(ctx.io_threads >= 1);
assert!(ctx.compute_threads >= 1);
assert!(ctx.bg_concurrency >= 2);
}
}
+28
View File
@@ -0,0 +1,28 @@
//! Shared error type for corex execution primitives.
/// Errors surfaced by corex's panic-isolated execution primitives.
///
/// Marked `#[non_exhaustive]` so new variants can be added without a
/// breaking change; downstream `match` expressions should include a
/// wildcard arm.
#[non_exhaustive]
#[derive(Debug)]
pub enum CorexError {
/// A closure submitted to [`crate::compute::compute`] panicked.
///
/// The contained string is a best-effort rendering of the panic
/// payload; the compute thread pool itself remains usable afterward.
ComputePanic(String),
}
impl std::fmt::Display for CorexError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::ComputePanic(message) => {
write!(f, "compute closure panicked: {message}")
}
}
}
}
impl std::error::Error for CorexError {}
+31
View File
@@ -0,0 +1,31 @@
//! Async I/O task spawning, instrumented with [`tracing`].
use tracing::Instrument;
/// Spawns `future` onto the ambient Tokio runtime, wrapped in a
/// `corex_io_task` tracing span.
///
/// # Panics
///
/// Panics if called outside the context of a running Tokio runtime; corex
/// does not construct or own a runtime of its own, it schedules onto
/// whichever runtime the caller is already inside.
pub fn spawn_io<F>(future: F) -> tokio::task::JoinHandle<F::Output>
where
F: std::future::Future + Send + 'static,
F::Output: Send + 'static,
{
let span = tracing::info_span!("corex_io_task");
tokio::spawn(future.instrument(span))
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn spawn_io_roundtrips_a_value() {
let handle = spawn_io(async { 7u32 });
assert_eq!(handle.await.unwrap(), 7);
}
}
+19
View File
@@ -0,0 +1,19 @@
//! Core allocation logic, global context initialization, and execution
//! primitives for [corex](https://docs.rs/corex).
//!
//! This crate is not typically consumed directly; applications should
//! depend on the `corex` facade crate instead, which re-exports the pieces
//! of this crate behind feature flags.
pub mod context;
#[cfg(feature = "compute")]
pub mod error;
#[cfg(feature = "io")]
pub mod io;
#[cfg(feature = "compute")]
pub mod compute;
#[cfg(feature = "bg")]
pub mod bg;