mirror of
https://github.com/gorules/zen.git
synced 2026-10-05 08:02:28 +00:00
feat: py bindings refactoring (#317)
* feat: py bindings refactoring * unset version in pyproject.toml
This commit is contained in:
@@ -1,46 +1,72 @@
|
||||
use std::future::Future;
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::content::PyZenDecisionContentJson;
|
||||
use anyhow::anyhow;
|
||||
use pyo3::{Py, PyAny, PyObject, Python};
|
||||
|
||||
use either::Either;
|
||||
use pyo3::{IntoPyObjectExt, Py, PyAny, PyResult, Python};
|
||||
use pyo3_async_runtimes::TaskLocals;
|
||||
use zen_engine::loader::{DecisionLoader, LoaderError, LoaderResponse};
|
||||
use zen_engine::model::DecisionContent;
|
||||
|
||||
#[derive(Default)]
|
||||
pub(crate) struct PyDecisionLoader(Option<Py<PyAny>>);
|
||||
|
||||
impl From<PyObject> for PyDecisionLoader {
|
||||
fn from(value: PyObject) -> Self {
|
||||
Self(Some(value))
|
||||
}
|
||||
pub(crate) struct PyDecisionLoader {
|
||||
callback: Option<Py<PyAny>>,
|
||||
task_locals: Option<TaskLocals>,
|
||||
}
|
||||
|
||||
impl From<Option<Py<PyAny>>> for PyDecisionLoader {
|
||||
fn from(value: Option<Py<PyAny>>) -> Self {
|
||||
Self(value)
|
||||
impl PyDecisionLoader {
|
||||
pub fn new(callback: Option<Py<PyAny>>, task_locals: Option<TaskLocals>) -> Self {
|
||||
Self {
|
||||
callback,
|
||||
task_locals,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl PyDecisionLoader {
|
||||
fn load_element(&self, key: &str) -> Result<Arc<DecisionContent>, anyhow::Error> {
|
||||
let Some(object) = &self.0 else {
|
||||
async fn load_element(&self, key: &str) -> Result<Arc<DecisionContent>, anyhow::Error> {
|
||||
let Some(callable) = &self.callback else {
|
||||
return Err(anyhow!("Loader is not defined"));
|
||||
};
|
||||
|
||||
let content = Python::with_gil(|py| {
|
||||
let result = object.call1(py, (key,))?;
|
||||
result.extract::<String>(py)
|
||||
})?;
|
||||
let maybe_result: PyResult<_> = Python::with_gil(|py| {
|
||||
let result = callable.call1(py, (key,))?;
|
||||
let is_coroutine = result.getattr(py, "__await__").is_ok();
|
||||
if !is_coroutine {
|
||||
return Ok(Either::Left(
|
||||
result.extract::<PyZenDecisionContentJson>(py)?,
|
||||
));
|
||||
}
|
||||
|
||||
Ok(serde_json::from_str::<DecisionContent>(&content)?.into())
|
||||
let Some(task_locals) = &self.task_locals else {
|
||||
Err(anyhow!("Task locals are required in async context"))?
|
||||
};
|
||||
|
||||
let result_future = pyo3_async_runtimes::into_future_with_locals(
|
||||
task_locals,
|
||||
result.into_bound_py_any(py)?,
|
||||
)?;
|
||||
|
||||
Ok(Either::Right(result_future))
|
||||
});
|
||||
|
||||
match maybe_result? {
|
||||
Either::Left(result) => Ok(result.0 .0),
|
||||
Either::Right(future) => {
|
||||
let result = future.await?;
|
||||
let content =
|
||||
Python::with_gil(|py| result.extract::<PyZenDecisionContentJson>(py))?;
|
||||
Ok(content.0 .0)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl DecisionLoader for PyDecisionLoader {
|
||||
fn load<'a>(&'a self, key: &'a str) -> impl Future<Output = LoaderResponse> + 'a {
|
||||
async move {
|
||||
self.load_element(key).map_err(|e| {
|
||||
self.load_element(key).await.map_err(|e| {
|
||||
LoaderError::Internal {
|
||||
source: e,
|
||||
key: key.to_string(),
|
||||
|
||||
Reference in New Issue
Block a user