85 lines
2.6 KiB
Python
85 lines
2.6 KiB
Python
import importlib
|
|
import re
|
|
from pathlib import Path
|
|
from typing import Any, Generic, List, Optional, Type, TypeVar
|
|
|
|
__all__ = ("Pipe", "PipeRegistry", "PipeParser", "PipeRunner", "PIPE_REGISTRY")
|
|
|
|
I = TypeVar("I")
|
|
O = TypeVar("O")
|
|
|
|
|
|
class Pipe(Generic[I, O]):
|
|
type: str
|
|
|
|
def apply(self, value: I, basepath: Path) -> Optional[O]:
|
|
raise NotImplementedError
|
|
|
|
|
|
class PipeRegistry:
|
|
def __init__(self):
|
|
self._pipes = {}
|
|
|
|
def register(self, pipe_type: Type[Pipe]):
|
|
self._pipes[pipe_type.type] = pipe_type
|
|
|
|
def register_module(self, module_name: str):
|
|
module = importlib.import_module(module_name)
|
|
modulepath = Path(module.__file__).parent
|
|
for item in modulepath.iterdir():
|
|
if item.suffix == ".py" and not item.name.startswith("__"):
|
|
sub_module = importlib.import_module(f"{module_name}.{item.stem}")
|
|
for v in sub_module.__dict__.values():
|
|
if isinstance(v, type):
|
|
for baseclass in v.__bases__:
|
|
if baseclass == Pipe:
|
|
self.register(v)
|
|
|
|
def get(self, pipe_type: str) -> Type[Pipe]:
|
|
return self._pipes[pipe_type]
|
|
|
|
|
|
class PipeParser:
|
|
_PIPE_PATTERN = re.compile("^(\\w+)\\((.*)\\)")
|
|
|
|
def parse(self, pipeline: str) -> List[Pipe]:
|
|
chain = [item.strip() for item in pipeline.split("|")]
|
|
result = []
|
|
for part in chain:
|
|
matched = self._PIPE_PATTERN.match(part)
|
|
if matched:
|
|
pipe_name = matched.group(1)
|
|
pipe_args = [item.strip() for item in matched.group(2).split(",")]
|
|
else:
|
|
pipe_name = part.strip()
|
|
pipe_args = []
|
|
pipe_type = PIPE_REGISTRY.get(pipe_name)
|
|
result.append(pipe_type(*pipe_args))
|
|
return result
|
|
|
|
|
|
class PipeRunner:
|
|
def __init__(self, basepath: Path, parser: Optional[PipeParser] = None):
|
|
self._basepath = basepath
|
|
self._parser = parser or PipeParser()
|
|
|
|
def run(self, pipes: List[Pipe]) -> Any:
|
|
result = None
|
|
for pipe in pipes:
|
|
result = pipe.apply(result, self._basepath)
|
|
return result
|
|
|
|
def resolve_value(self, pipeline: str) -> Any:
|
|
pipe = self._parser.parse(pipeline)
|
|
return self.run(pipe)
|
|
|
|
def check_conditions(self, conditions: Optional[List[str]]) -> bool:
|
|
if conditions:
|
|
for condition in conditions:
|
|
if not self.resolve_value(condition):
|
|
return False
|
|
return True
|
|
|
|
|
|
PIPE_REGISTRY = PipeRegistry()
|