commit d8144e249edcf74bbd567106288773fa5016976d Author: asepharyana Date: Fri Aug 28 16:05:04 2026 +0700 Add corex library with resource-aware execution primitives - 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. diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..e3063a9 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,86 @@ +name: CI + +on: + push: + branches: [main] + pull_request: + branches: [main] + +env: + CARGO_TERM_COLOR: always + +jobs: + fmt: + name: Rustfmt + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + with: + components: rustfmt + - run: cargo fmt --all --check + + clippy: + name: Clippy + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + with: + components: clippy + - name: Clippy (all features) + run: cargo clippy --workspace --all-features -- -D warnings + - name: Clippy (no default features) + run: cargo clippy --workspace --no-default-features -- -D warnings + + test-matrix: + name: Test (${{ matrix.name }}) + runs-on: ubuntu-latest + strategy: + fail-fast: false + matrix: + include: + - name: default features + flags: "" + - name: all features + flags: "--all-features" + - name: io only + flags: "--no-default-features --features io" + - name: compute only + flags: "--no-default-features --features compute" + - name: bg only + flags: "--no-default-features --features bg" + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - name: Build + run: cargo build --workspace ${{ matrix.flags }} + - name: Test + run: cargo test --workspace ${{ matrix.flags }} + + example: + name: Run example + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - run: cargo run --example main --features full + + docs: + name: Docs check + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - run: RUSTDOCFLAGS="-D warnings" cargo doc --workspace --all-features --no-deps + + package: + name: Cargo package dry-run + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - name: Package corex-core + run: cargo package -p corex-core --no-verify + - name: Package corex + run: cargo package -p corex --no-verify diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml new file mode 100644 index 0000000..c8f1b66 --- /dev/null +++ b/.github/workflows/publish.yml @@ -0,0 +1,23 @@ +name: Publish to crates.io + +on: + push: + tags: ["v*.*.*"] + +jobs: + publish: + name: Publish + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - name: Publish corex-core + run: cargo publish -p corex-core --allow-dirty --no-verify + env: + CARGO_REGISTRY_TOKEN: ${{ secrets.CARGO_REGISTRY_TOKEN }} + - name: Wait for crates.io indexing + run: sleep 30 + - name: Publish corex + run: cargo publish -p corex --allow-dirty --no-verify + env: + CARGO_REGISTRY_TOKEN: ${{ secrets.CARGO_REGISTRY_TOKEN }} diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..7cafc37 --- /dev/null +++ b/.gitignore @@ -0,0 +1,15 @@ +# Cargo build artifacts +/target +/crates/*/target + +# corex is a library workspace: the lockfile is not committed, per Rust's +# convention for libraries (binaries/applications should commit theirs). +# Cargo regenerates it locally and in CI on every build. +Cargo.lock + +# Editor / OS noise +.DS_Store +Thumbs.db +*.swp +.idea/ +.vscode/ diff --git a/Cargo.toml b/Cargo.toml new file mode 100644 index 0000000..bc001f2 --- /dev/null +++ b/Cargo.toml @@ -0,0 +1,62 @@ +[workspace] +resolver = "2" +members = ["crates/corex-core"] + +[workspace.package] +version = "0.1.0" +edition = "2021" +rust-version = "1.75" +license = "MIT OR Apache-2.0" +repository = "https://github.com/username/corex" +homepage = "https://github.com/username/corex" +documentation = "https://docs.rs/corex" +authors = ["The corex Authors"] + +[workspace.dependencies] +tokio = { version = "1.53", features = ["full"] } +rayon = "1.12" +num_cpus = "1.17" +tracing = "0.1" +corex-core = { path = "crates/corex-core", version = "0.1.0", default-features = false } + +[package] +name = "corex" +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 = "Resource-aware abstractions for async I/O, heavy compute, and background queue management." +readme = "README.md" +keywords = ["async", "concurrency", "rayon", "tokio", "resource-management"] +categories = ["asynchronous", "concurrency", "rust-patterns"] + +[lib] +name = "corex" +path = "src/lib.rs" + +[features] +default = [] +io = ["corex-core/io"] +compute = ["corex-core/compute"] +bg = ["corex-core/bg"] +full = ["io", "compute", "bg"] + +[dependencies] +corex-core = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true } +tracing-subscriber = "0.3" + +[[example]] +name = "main" +path = "examples/main.rs" +required-features = ["full"] + +[package.metadata.docs.rs] +all-features = true +rustdoc-args = ["--cfg", "docsrs"] diff --git a/LICENSE-APACHE b/LICENSE-APACHE new file mode 100644 index 0000000..0da389e --- /dev/null +++ b/LICENSE-APACHE @@ -0,0 +1,201 @@ + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + APPENDIX: How to apply the Apache License to your work. + + To apply the Apache License to your work, attach the following + boilerplate notice, with the fields enclosed by brackets "[]" + replaced with your own identifying information. (Don't include + the brackets!) The text should be enclosed in the appropriate + comment syntax for the file format. We also recommend that a + file or class name and description of purpose be included on the + same "printed page" as the copyright notice for easier + identification within third-party archives. + + Copyright 2026 The corex Authors + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. diff --git a/LICENSE-MIT b/LICENSE-MIT new file mode 100644 index 0000000..687a34c --- /dev/null +++ b/LICENSE-MIT @@ -0,0 +1,21 @@ +MIT License + +Copyright (c) 2026 The corex Authors + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/README.md b/README.md new file mode 100644 index 0000000..76cb037 --- /dev/null +++ b/README.md @@ -0,0 +1,81 @@ +# corex + +[![Crates.io](https://img.shields.io/crates/v/corex.svg)](https://crates.io/crates/corex) +[![Documentation](https://docs.rs/corex/badge.svg)](https://docs.rs/corex) +[![License](https://img.shields.io/badge/license-MIT%20OR%20Apache--2.0-blue.svg)](LICENSE-MIT) + +Resource-aware execution primitives for Rust: async I/O, heavy compute, and background queue management, sized automatically from the host's logical core count and exposed through a single, lazily-initialized engine context. + +## Resource Sizing + +Given $N$ logical cores (via `num_cpus::get()`): + +| Subsystem | Sizing Formula | Default on 8 cores | Backing Primitive | +| :--- | :--- | :--- | :--- | +| **Async I/O** | $N$ | 8 | Ambient `tokio::spawn` + `tracing` span | +| **Compute** | $\max(1, N - 1)$ | 7 | Sized `rayon::ThreadPool` + `catch_unwind` | +| **Background Queue** | $\max(2, \lfloor N / 2 \rfloor)$ | 4 | `tokio::sync::Semaphore` + `tokio::spawn` | + +## Features + +- **`io`**: enables `corex::spawn_io`, instrumented async task spawning. +- **`compute`**: enables `corex::compute`, panic-isolated execution on a sized Rayon pool. +- **`bg`**: enables `corex::spawn_bg`, semaphore-bounded background tasks. +- **`full`**: enables all three subsystems. + +Zero features enabled by default (`default = []`), so you only pull in the dependencies your application actually uses. + +## Quick Start + +Add to your `Cargo.toml`: + +```toml +[dependencies] +corex = { version = "0.1", features = ["full"] } +``` + +Use the entry points directly: + +```rust +#[tokio::main] +async fn main() { + // Optional explicit bootstrap: logs or validates resource sizing upfront. + // Omit it and the first call to any primitive below will initialize it lazily. + let ctx = corex::init(); + println!( + "io_threads={} compute_threads={} bg_concurrency={}", + ctx.io_threads, ctx.compute_threads, ctx.bg_concurrency + ); + + // 1. Async I/O (instrumented with tracing) + let io = corex::spawn_io(async { + // ... network / disk work ... + 42 + }); + + // 2. Heavy Compute (isolated from worker panics) + let sum = corex::compute(|| (1..=1_000_000u64).sum::())?; + + // 3. Background Queue (concurrency-bounded) + let bg = corex::spawn_bg(async { + // ... deferred cleanup / telemetry ... + }).await; + + let _ = (io.await, bg.await); +} +``` + +## Running the Example + +```bash +cargo run --example main --features full +``` + +## License + +Licensed under either of: + +- Apache License, Version 2.0 ([LICENSE-APACHE](LICENSE-APACHE) or ) +- MIT license ([LICENSE-MIT](LICENSE-MIT) or ) + +at your option. diff --git a/crates/corex-core/Cargo.toml b/crates/corex-core/Cargo.toml new file mode 100644 index 0000000..27556a1 --- /dev/null +++ b/crates/corex-core/Cargo.toml @@ -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"] diff --git a/crates/corex-core/README.md b/crates/corex-core/README.md new file mode 100644 index 0000000..2851b47 --- /dev/null +++ b/crates/corex-core/README.md @@ -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. diff --git a/crates/corex-core/src/bg.rs b/crates/corex-core/src/bg.rs new file mode 100644 index 0000000..cb12053 --- /dev/null +++ b/crates/corex-core/src/bg.rs @@ -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(future: F) -> tokio::task::JoinHandle +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); + } +} diff --git a/crates/corex-core/src/compute.rs b/crates/corex-core/src/compute.rs new file mode 100644 index 0000000..e853473 --- /dev/null +++ b/crates/corex-core/src/compute.rs @@ -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: F) -> Result +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) -> String { + if let Some(message) = payload.downcast_ref::<&str>() { + (*message).to_string() + } else if let Some(message) = payload.downcast_ref::() { + 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 = compute(|| panic!("boom")); + assert!(matches!(panicked, Err(CorexError::ComputePanic(_)))); + + let recovered = compute(|| 1 + 1); + assert_eq!(recovered.unwrap(), 2); + } +} diff --git a/crates/corex-core/src/context.rs b/crates/corex-core/src/context.rs new file mode 100644 index 0000000..b6402c8 --- /dev/null +++ b/crates/corex-core/src/context.rs @@ -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 = 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); + } +} diff --git a/crates/corex-core/src/error.rs b/crates/corex-core/src/error.rs new file mode 100644 index 0000000..43f29f4 --- /dev/null +++ b/crates/corex-core/src/error.rs @@ -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 {} diff --git a/crates/corex-core/src/io.rs b/crates/corex-core/src/io.rs new file mode 100644 index 0000000..d684bad --- /dev/null +++ b/crates/corex-core/src/io.rs @@ -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(future: F) -> tokio::task::JoinHandle +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); + } +} diff --git a/crates/corex-core/src/lib.rs b/crates/corex-core/src/lib.rs new file mode 100644 index 0000000..5d2f32a --- /dev/null +++ b/crates/corex-core/src/lib.rs @@ -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; diff --git a/examples/main.rs b/examples/main.rs new file mode 100644 index 0000000..b524e0d --- /dev/null +++ b/examples/main.rs @@ -0,0 +1,41 @@ +//! Demonstrates `corex`'s three execution primitives end to end, including +//! automatic recovery from a panicking compute closure. + +#[tokio::main] +async fn main() { + tracing_subscriber::fmt::init(); + + let ctx = corex::init(); + println!( + "engine context: io_threads={} compute_threads={} bg_concurrency={}", + ctx.io_threads, ctx.compute_threads, ctx.bg_concurrency + ); + + let io_handle = corex::spawn_io(async { + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + 42u64 + }); + + let sum_result = corex::compute(|| (1..=1_000u64).sum::()); + + let panic_result: Result = corex::compute(|| { + panic!("intentional panic to demonstrate isolation"); + }); + + let recovery_result = corex::compute(|| 2u64 + 2u64); + + let bg_handle = corex::spawn_bg(async { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + "bg-task-done" + }) + .await; + + let io_value = io_handle.await.expect("io task panicked"); + let bg_value = bg_handle.await.expect("bg task panicked"); + + println!("spawn_io result: {io_value}"); + println!("compute sum result: {sum_result:?}"); + println!("compute panic-isolation result: {panic_result:?}"); + println!("compute pool still usable after panic: {recovery_result:?}"); + println!("spawn_bg result: {bg_value}"); +} diff --git a/src/lib.rs b/src/lib.rs new file mode 100644 index 0000000..c667e28 --- /dev/null +++ b/src/lib.rs @@ -0,0 +1,39 @@ +//! # corex +//! +//! Resource-aware abstractions for async I/O, heavy compute, and +//! background queue management, built on a single lazily-initialized +//! global engine context. +//! +//! Call [`init`] once at startup (or simply let the first call to any +//! entry point below trigger it lazily) and then use whichever of the +//! feature-gated entry points your workload needs: +//! +//! - [`spawn_io`] (feature `io`) — spawn an async I/O task, tracing-instrumented. +//! - [`compute`] (feature `compute`) — run CPU-bound work on a sized Rayon pool, panic-isolated. +//! - [`spawn_bg`] (feature `bg`) — spawn a background task under bounded concurrency. +//! +//! Enable the `full` feature to pull in all three at once. + +pub use corex_core::context::{context, EngineContext}; + +#[cfg(feature = "compute")] +pub use corex_core::error::CorexError; + +#[cfg(feature = "io")] +pub use corex_core::io::spawn_io; + +#[cfg(feature = "compute")] +pub use corex_core::compute::compute; + +#[cfg(feature = "bg")] +pub use corex_core::bg::spawn_bg; + +/// Bootstraps the global corex [`EngineContext`]. +/// +/// See [`corex_core::context::init`] for full semantics: this is safe to +/// call any number of times, from any thread, and is equivalent to letting +/// the first call to [`spawn_io`], [`compute`], or [`spawn_bg`] trigger +/// initialization implicitly. +pub fn init() -> &'static EngineContext { + corex_core::context::init() +}