72 lines
2.3 KiB
Python
72 lines
2.3 KiB
Python
from typing import Any, Dict, Optional, Callable, List, TypeVar, Generic
|
|
from sqlalchemy.orm import Session
|
|
from app.core.auth import CurrentUser
|
|
|
|
from app.serviceInput.base import ServiceInput
|
|
from app.serviceOutput.base import ServiceOutput
|
|
|
|
InputT = TypeVar("InputT")
|
|
OutputT = TypeVar("OutputT")
|
|
|
|
class Pipeline(Generic[InputT, OutputT]):
|
|
"""
|
|
Pipeline base que unifica entrada, procesamiento y salida.
|
|
"""
|
|
|
|
def __init__(self, db: Session, user: CurrentUser, operation: str):
|
|
self.db = db
|
|
self.user = user
|
|
self.operation = operation
|
|
self._steps: List[Callable] = []
|
|
self._input: Optional[ServiceInput] = None
|
|
self._output: Optional[ServiceOutput] = None
|
|
self._input_data: Any = None
|
|
self._schema = None
|
|
self._permission: Optional[str] = None
|
|
|
|
def set_input_data(self, data: Any) -> 'Pipeline':
|
|
"""Define los datos de entrada"""
|
|
self._input_data = data
|
|
return self
|
|
|
|
def with_schema(self, schema: Any) -> 'Pipeline':
|
|
"""Define el schema de validación"""
|
|
self._schema = schema
|
|
return self
|
|
|
|
def with_permission(self, permission: Optional[str]) -> 'Pipeline':
|
|
"""Define el permiso requerido"""
|
|
self._permission = permission
|
|
return self
|
|
|
|
def add_step(self, name: str, func: Callable) -> 'Pipeline':
|
|
"""Agrega un paso al pipeline"""
|
|
self._steps.append({"name": name, "func": func})
|
|
return self
|
|
|
|
async def execute(self, context: Dict = None) -> Dict:
|
|
"""
|
|
Ejecuta el pipeline completo:
|
|
1. Validación de entrada
|
|
2. Verificación de permisos
|
|
3. Ejecución de pasos
|
|
4. Retorna resultado
|
|
"""
|
|
if context is None:
|
|
context = {}
|
|
|
|
context["db"] = self.db
|
|
context["user"] = self.user
|
|
context["operation"] = self.operation
|
|
context["input_data"] = self._input_data
|
|
context["schema"] = self._schema
|
|
context["permission"] = self._permission
|
|
|
|
result = context
|
|
|
|
for step in self._steps:
|
|
result = await step["func"](result)
|
|
if result.get("error"):
|
|
return {"success": False, "error": result["error"]}
|
|
|
|
return {"success": True, "data": result.get("output")} |