mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
* refactor(rust): extract litellm-host-native as the shared Rust host driver Move service and hook dispatch out of host-http into a Driver that owns the machine and Rust handlers, returning at completion or a stream boundary and holding the demand reply until the consumer advances. Move the in-process runner onto the same driver. host-http now layers encoding, SSE, body polling and lifecycle observation over it. host-python keeps driving litellm-host directly Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(rust): interrupt the machine when the in-process stream consumer fails Restores the pre-refactor interruption path for StreamConsumer errors via Driver::fail and ports the generic run lifecycle tests into host-native. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(rust): separate the machine contract from coroutine execution * auth update * refactor(rust): use standard flow control for host requests * style(rust): keep host driver imports formatted * chores * mostly relocation * refactor(rust): separate interceptors from queued observers * refactor(rust): centralize legacy callback mappings and lifecycle * docs: define Python host boundaries and migration plan * refactor: enforce Python host and bridge boundaries * refactor(rust): separate operations from callback composition * refactor(rust): compose SDK policy through call hooks --------- Co-authored-by: Yujong Lee <yujong@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
170 lines
5.5 KiB
Rust
170 lines
5.5 KiB
Rust
use std::panic::{AssertUnwindSafe, catch_unwind};
|
|
|
|
use crate::panic_to_pyerr;
|
|
use pyo3::exceptions::{PyBaseException, PyRuntimeError};
|
|
use pyo3::gc::{PyTraverseError, PyVisit};
|
|
use pyo3::prelude::*;
|
|
|
|
pub type PythonLifecycle = for<'py> fn(Python<'py>) -> PyResult<Bound<'py, PyModule>>;
|
|
|
|
pub enum ExecutionStep {
|
|
Return(Py<PyAny>),
|
|
Await(Py<PyAny>),
|
|
/// The call streams: the caller gets a stream over this execution carrying this head,
|
|
/// and the execution stays suspended until the stream asks for a chunk.
|
|
Open(Py<PyAny>),
|
|
Yield(Py<PyAny>),
|
|
}
|
|
|
|
pub trait ExecutionBody: Send + Sync {
|
|
fn resume(&mut self, result: Option<PyResult<Py<PyAny>>>) -> PyResult<ExecutionStep>;
|
|
fn traverse(&self, visit: &PyVisit<'_>) -> Result<(), PyTraverseError>;
|
|
}
|
|
|
|
enum ExecutionState {
|
|
Created(Box<dyn ExecutionBody>),
|
|
Running,
|
|
Suspended(Box<dyn ExecutionBody>),
|
|
Closed,
|
|
}
|
|
|
|
#[pyclass]
|
|
pub struct Execution {
|
|
state: ExecutionState,
|
|
lifecycle: PythonLifecycle,
|
|
}
|
|
|
|
impl Execution {
|
|
pub fn new(body: impl ExecutionBody + 'static, lifecycle: PythonLifecycle) -> Self {
|
|
Self {
|
|
state: ExecutionState::Created(Box::new(body)),
|
|
lifecycle,
|
|
}
|
|
}
|
|
|
|
pub fn into_coroutine(self, py: Python<'_>) -> PyResult<Bound<'_, PyAny>> {
|
|
let binding = (self.lifecycle)(py)?;
|
|
let execution = Py::new(py, self)?;
|
|
binding.getattr("drive")?.call1((execution,))
|
|
}
|
|
|
|
pub(crate) fn into_sync_stream(
|
|
self,
|
|
py: Python<'_>,
|
|
head: Py<PyAny>,
|
|
) -> PyResult<Bound<'_, PyAny>> {
|
|
(self.lifecycle)(py)?
|
|
.getattr("SyncStream")?
|
|
.call1((Py::new(py, self)?, head))
|
|
}
|
|
|
|
/// An execution already started elsewhere and now waiting for its next input.
|
|
pub fn suspended(body: impl ExecutionBody + 'static, lifecycle: PythonLifecycle) -> Self {
|
|
Self {
|
|
state: ExecutionState::Suspended(Box::new(body)),
|
|
lifecycle,
|
|
}
|
|
}
|
|
|
|
fn advance(
|
|
slf: &Bound<'_, Self>,
|
|
py: Python<'_>,
|
|
result: Option<PyResult<Py<PyAny>>>,
|
|
) -> PyResult<Py<PyAny>> {
|
|
let mut body = {
|
|
let mut execution = slf.borrow_mut();
|
|
match (&execution.state, result.is_some()) {
|
|
(ExecutionState::Created(_), false) | (ExecutionState::Suspended(_), true) => {}
|
|
(ExecutionState::Running, _) => {
|
|
return Err(PyRuntimeError::new_err("execution is already running"));
|
|
}
|
|
(ExecutionState::Closed, _) => {
|
|
return Err(PyRuntimeError::new_err("execution is closed"));
|
|
}
|
|
_ => {
|
|
return Err(PyRuntimeError::new_err(
|
|
"execution requires start before resume and can only start once",
|
|
));
|
|
}
|
|
}
|
|
match std::mem::replace(&mut execution.state, ExecutionState::Running) {
|
|
ExecutionState::Created(body) | ExecutionState::Suspended(body) => body,
|
|
_ => unreachable!(),
|
|
}
|
|
};
|
|
let lifecycle = slf.borrow().lifecycle;
|
|
let outcome = catch_unwind(AssertUnwindSafe(|| {
|
|
let step = body.resume(result)?;
|
|
let (tag, value, suspended) = match step {
|
|
ExecutionStep::Await(value) => ("Await", value, true),
|
|
ExecutionStep::Open(head) => ("Open", head, true),
|
|
ExecutionStep::Yield(value) => ("Yield", value, true),
|
|
ExecutionStep::Return(value) => ("Complete", value, false),
|
|
};
|
|
let step = lifecycle(py)?.getattr(tag)?.call1((value,))?.unbind();
|
|
Ok((step, suspended))
|
|
}))
|
|
.map_err(panic_to_pyerr)
|
|
.and_then(|result| result);
|
|
match outcome {
|
|
Ok((step, true)) if matches!(slf.borrow().state, ExecutionState::Running) => {
|
|
slf.borrow_mut().state = ExecutionState::Suspended(body);
|
|
Ok(step)
|
|
}
|
|
outcome => {
|
|
slf.borrow_mut().state = ExecutionState::Closed;
|
|
drop(body);
|
|
outcome.and_then(|(step, suspended)| {
|
|
if suspended {
|
|
Err(PyRuntimeError::new_err(
|
|
"execution was closed while running",
|
|
))
|
|
} else {
|
|
Ok(step)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
#[pymethods]
|
|
impl Execution {
|
|
fn start(slf: &Bound<'_, Self>, py: Python<'_>) -> PyResult<Py<PyAny>> {
|
|
Self::advance(slf, py, None)
|
|
}
|
|
|
|
fn resume_value(
|
|
slf: &Bound<'_, Self>,
|
|
py: Python<'_>,
|
|
value: Py<PyAny>,
|
|
) -> PyResult<Py<PyAny>> {
|
|
Self::advance(slf, py, Some(Ok(value)))
|
|
}
|
|
|
|
fn resume_error(
|
|
slf: &Bound<'_, Self>,
|
|
py: Python<'_>,
|
|
error: Bound<'_, PyBaseException>,
|
|
) -> PyResult<Py<PyAny>> {
|
|
Self::advance(slf, py, Some(Err(PyErr::from_value(error.into_any()))))
|
|
}
|
|
|
|
fn close(slf: &Bound<'_, Self>) {
|
|
let state = std::mem::replace(&mut slf.borrow_mut().state, ExecutionState::Closed);
|
|
drop(state);
|
|
}
|
|
|
|
fn __traverse__(&self, visit: PyVisit<'_>) -> Result<(), PyTraverseError> {
|
|
match &self.state {
|
|
ExecutionState::Created(body) | ExecutionState::Suspended(body) => {
|
|
body.traverse(&visit)
|
|
}
|
|
_ => Ok(()),
|
|
}
|
|
}
|
|
|
|
fn __clear__(slf: &Bound<'_, Self>) {
|
|
Self::close(slf);
|
|
}
|
|
}
|