_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}