refactor: add pipe module
This commit is contained in:
@@ -0,0 +1,83 @@
|
||||
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: List[str]) -> bool:
|
||||
for condition in conditions:
|
||||
if not self.resolve_value(condition):
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
PIPE_REGISTRY = PipeRegistry()
|
||||
Reference in New Issue
Block a user