Skip to main content

_lib/
lib.rs

1// Location: ./bindings/python/src/lib.rs
2// Copyright 2025
3// SPDX-License-Identifier: Apache-2.0
4// Authors: Ted Habeck
5//
6// cpex-python — PyO3 native extension module `cpex._lib`.
7//
8// Module registration and tokio runtime initialization.
9// The runtime is initialized once with a multi-thread builder that honours
10// the `CPEX_PY_WORKER_THREADS` environment variable, mirroring the
11// `CPEX_FFI_WORKER_THREADS` knob in cpex-ffi.
12
13use pyo3::prelude::*;
14
15mod builtins;
16mod conversions;
17mod error;
18mod manager;
19mod result;
20
21use manager::PyPluginManager;
22use result::PyPipelineResult;
23
24/// Name of the env var operators set to bound worker threads.
25const ENV_WORKER_THREADS: &str = "CPEX_PY_WORKER_THREADS";
26
27/// Parse `CPEX_PY_WORKER_THREADS`. Returns `Some(n)` for valid positive
28/// integers, `None` otherwise (falls back to tokio default `num_cpus`).
29fn worker_threads_from_env() -> Option<usize> {
30    let raw = std::env::var(ENV_WORKER_THREADS).ok()?;
31    match raw.parse::<usize>() {
32        Ok(n) if n > 0 => {
33            tracing::info!(
34                "cpex-python: runtime using {} worker threads (from {})",
35                n,
36                ENV_WORKER_THREADS,
37            );
38            Some(n)
39        },
40        _ => {
41            tracing::warn!(
42                "cpex-python: {}={:?} is not a positive integer; using num_cpus default",
43                ENV_WORKER_THREADS,
44                raw,
45            );
46            None
47        },
48    }
49}
50
51#[pymodule]
52fn _lib(m: &Bound<'_, PyModule>) -> PyResult<()> {
53    // Initialize the pyo3-async-runtimes tokio runtime with a multi-thread
54    // builder so async methods are dispatched onto a real thread pool rather
55    // than a single-threaded executor. This must run before any `future_into_py`
56    // call — doing it here at module import time is the correct hook.
57    //
58    // This is a separate runtime from cpex-ffi's SHARED_RUNTIME; the
59    // shared-budget philosophy is mirrored, not the runtime instance.
60    let mut builder = tokio::runtime::Builder::new_multi_thread();
61    builder.enable_all();
62    if let Some(n) = worker_threads_from_env() {
63        builder.worker_threads(n);
64    }
65    pyo3_async_runtimes::tokio::init(builder);
66
67    m.add_class::<PyPluginManager>()?;
68    m.add_class::<PyPipelineResult>()?;
69
70    Ok(())
71}