refactor(rust): centralize bridge execution wrappers (#43871)

Co-authored-by: Yujong Lee <yujong@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
devin-ai-integration[bot] 2026-09-30 15:54:27 +00:00 • committed by GitHub
parent 9dda4d895f
commit b80052839e
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 34 additions and 33 deletions

View file

@ -1,5 +1,5 @@
use crate::cache::cache_error;
use crate::logger::run_sync_value;
use crate::execution::run_sync_value;
use litellm_cache_gcs::{DEFAULT_ENDPOINT, GcsConfig};
use litellm_cache_redis_semantic::RedisSemanticConfig;
use litellm_host_python::release_gil;

View file

@ -470,7 +470,7 @@ impl NativeResponseCache {
match self {
Self::Exact(_) | Self::QdrantSemantic(_) => {
let service = self.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move {
service
@ -495,7 +495,7 @@ impl NativeResponseCache {
match self {
Self::Exact(_) | Self::QdrantSemantic(_) => {
let service = self.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move { service.async_lookup(&request, now()).await },
cache_error,
@ -550,7 +550,7 @@ impl NativeResponseCache {
match self {
Self::Exact(_) | Self::QdrantSemantic(_) => {
let service = self.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move { service.async_store(&request, response, now()).await },
cache_error,
@ -619,7 +619,7 @@ impl NativeResponseCache {
match self {
Self::Exact(_) | Self::QdrantSemantic(_) => {
let service = self.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move { service.async_store_batch(entries, now()).await },
cache_error,

View file

@ -1,5 +1,5 @@
use crate::cache::cache_error;
use crate::logger::run_async;
use crate::execution::run_async;
use std::{collections::VecDeque, time::Duration};
use litellm_cache::Error;

View file

@ -144,7 +144,7 @@ impl NativeCacheHandle {
self.check_process()?;
let request = request(key, None)?;
let backend = self.backend.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move { backend.async_lookup(&request, super::request::now()).await },
cache_error,
@ -163,7 +163,7 @@ impl NativeCacheHandle {
let request = request(key, ttl)?;
let value: Value = from_py(value)?;
let backend = self.backend.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move {
backend
@ -188,7 +188,7 @@ impl NativeCacheHandle {
.map(|(key, value)| Ok((request(key, ttl)?, value)))
.collect::<PyResult<Vec<_>>>()?;
let backend = self.backend.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move {
backend
@ -202,19 +202,19 @@ impl NativeCacheHandle {
fn flush(&self, py: Python<'_>) -> PyResult<Py<PyAny>> {
self.check_process()?;
let backend = self.backend.clone();
crate::logger::run_sync(py, async move { backend.async_flush().await }, cache_error)
crate::execution::run_sync(py, async move { backend.async_flush().await }, cache_error)
}
fn async_flush<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
self.check_process()?;
let backend = self.backend.clone();
crate::logger::run_async(py, async move { backend.async_flush().await }, cache_error)
crate::execution::run_async(py, async move { backend.async_flush().await }, cache_error)
}
fn ping<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
self.check_process()?;
let storage = self.storage.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move {
match storage {
@ -229,7 +229,7 @@ impl NativeCacheHandle {
fn disconnect<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
self.check_process()?;
let storage = self.storage.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move {
match storage {
@ -244,7 +244,7 @@ impl NativeCacheHandle {
fn delete<'py>(&self, py: Python<'py>, keys: Vec<String>) -> PyResult<Bound<'py, PyAny>> {
self.check_process()?;
let storage = self.storage.clone();
crate::logger::run_async(
crate::execution::run_async(
py,
async move {
for key in keys {

View file

@ -1,4 +1,4 @@
use crate::logger::run_async;
use crate::execution::run_async;
use litellm_cache_response::PartialHits;
use litellm_host_python::{ExecutionStep, from_py, release_gil, to_py};
use pyo3::{

View file

@ -13,7 +13,7 @@ where
E: Send + 'static,
F: Future<Output = Result<T, E>> + Send + 'static,
{
litellm_host_python::run_sync(py, super::capture(py).instrument(future), map_error)
litellm_host_python::run_sync(py, crate::logger::capture(py).instrument(future), map_error)
}
pub(crate) fn run_async<T, E, F>(
@ -26,7 +26,7 @@ where
E: Send + 'static,
F: Future<Output = Result<T, E>> + Send + 'static,
{
litellm_host_python::run_async(py, super::capture(py).instrument(future), map_error)
litellm_host_python::run_async(py, crate::logger::capture(py).instrument(future), map_error)
}
pub(crate) fn run_sync_value<T, F>(py: Python<'_>, future: F) -> PyResult<T>
@ -34,7 +34,7 @@ where
T: Send + 'static,
F: Future<Output = PyResult<T>> + Send + 'static,
{
litellm_host_python::run_sync_value(py, super::capture(py).instrument(future))
litellm_host_python::run_sync_value(py, crate::logger::capture(py).instrument(future))
}
pub(crate) fn run_async_value<T, F>(py: Python<'_>, future: F) -> PyResult<Bound<'_, PyAny>>
@ -42,5 +42,5 @@ where
T: for<'py> IntoPyObject<'py> + Send + 'static,
F: Future<Output = PyResult<T>> + Send + 'static,
{
litellm_host_python::run_async_value(py, super::capture(py).instrument(future))
litellm_host_python::run_async_value(py, crate::logger::capture(py).instrument(future))
}

View file

@ -4,6 +4,7 @@ mod coercion;
mod credentials;
mod diagnostics;
mod errors;
mod execution;
mod http;
mod lifecycle;
mod logger;

View file

@ -1,7 +1,5 @@
mod execution;
mod machine;
pub(crate) use execution::{run_async, run_async_value, run_sync, run_sync_value};
pub(crate) use machine::LoggedMachine;
use litellm_host_python::Pythonized;

View file

@ -76,7 +76,7 @@ async fn traced_operation(_secret: &str) -> PyResult<()> {
#[pyfunction]
fn span_warning(py: Python<'_>) -> PyResult<Bound<'_, PyAny>> {
super::run_async_value(py, traced_operation("private-key-sentinel"))
crate::execution::run_async_value(py, traced_operation("private-key-sentinel"))
}
#[pyfunction]
@ -93,7 +93,7 @@ fn levels(py: Python<'_>) {
#[pyfunction]
fn asynchronous_warning(py: Python<'_>) -> PyResult<Bound<'_, PyAny>> {
super::run_async_value(py, async {
crate::execution::run_async_value(py, async {
tokio::task::yield_now().await;
litellm_tracing::warn!("async warning");
Ok(())
@ -102,7 +102,7 @@ fn asynchronous_warning(py: Python<'_>) -> PyResult<Bound<'_, PyAny>> {
#[pyfunction]
fn synchronous_warning(py: Python<'_>) -> PyResult<()> {
super::run_sync_value(py, async {
crate::execution::run_sync_value(py, async {
tokio::task::yield_now().await;
litellm_tracing::warn!("sync warning");
Ok(())
@ -111,7 +111,7 @@ fn synchronous_warning(py: Python<'_>) -> PyResult<()> {
#[pyfunction]
fn synchronous_failure(py: Python<'_>) -> PyResult<()> {
super::run_sync_value(py, async {
crate::execution::run_sync_value(py, async {
litellm_tracing::warn!("failure diagnostic");
Err(pyo3::exceptions::PyValueError::new_err("request failed"))
})

View file

@ -1,4 +1,4 @@
use crate::logger::{run_async, run_sync};
use crate::execution::{run_async, run_sync};
use litellm_core::audio_transcription::{
AudioTranscriptionRoute, Error, types::AudioTranscriptionRequest,
};

View file

@ -2,7 +2,7 @@ mod host;
use pyo3::types::{PyDict, PyTuple};
use crate::logger::{run_async, run_sync};
use crate::execution::{run_async, run_sync};
use litellm_core::chat_completions::{ChatCompletionsRoute, Error, types::ChatCompletionsRequest};
use litellm_llms_types::formats::chat_completions::ChatCompletionsResponse;
use pyo3::prelude::*;

View file

@ -142,7 +142,7 @@ impl ResponsesWebSocketConnection {
) -> PyResult<Bound<'py, PyAny>> {
let headers = marshal_headers(headers)?;
let timeout = optional_timeout(timeout_seconds);
crate::logger::run_async_value(py, async move {
crate::execution::run_async_value(py, async move {
let inner = RustResponsesWebSocketConnection::connect_url(&url, &headers, timeout)
.await
.map_err(route_error_to_pyerr)?;
@ -152,21 +152,21 @@ impl ResponsesWebSocketConnection {
fn send_text<'py>(&self, py: Python<'py>, text: String) -> PyResult<Bound<'py, PyAny>> {
let inner = self.inner.clone();
crate::logger::run_async_value(py, async move {
crate::execution::run_async_value(py, async move {
inner.send_text(text).await.map_err(route_error_to_pyerr)
})
}
fn recv_text<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
let inner = self.inner.clone();
crate::logger::run_async_value(py, async move {
crate::execution::run_async_value(py, async move {
inner.recv_text().await.map_err(route_error_to_pyerr)
})
}
fn close<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
let inner = self.inner.clone();
crate::logger::run_async_value(py, async move {
crate::execution::run_async_value(py, async move {
inner.close().await.map_err(route_error_to_pyerr)
})
}

View file

@ -1,4 +1,4 @@
use crate::logger::run_async;
use crate::execution::run_async;
use std::sync::Arc;
use std::{num::NonZero, thread::available_parallelism};

View file

@ -1,7 +1,7 @@
use std::{collections::BTreeMap, sync::Arc};
use litellm_core_utils::settings::{Lookup, ProcessEnvironment};
use litellm_host_python::{from_py, json_object_field, run_async_value, run_sync_value, to_py};
use litellm_host_python::{from_py, json_object_field, to_py};
use litellm_secrets::{
KeyManagementSettings, KeyManagementSystem, Secret, SecretManager, load_native_manager,
read_secret_from_python_manager,
@ -13,6 +13,8 @@ use pyo3::{
types::PyDict,
};
use crate::execution::{run_async_value, run_sync_value};
#[derive(Clone, PartialEq)]
struct Configuration {
system: KeyManagementSystem,