Loading .gitignore +4 −1 Original line number Diff line number Diff line Loading @@ -3,3 +3,6 @@ csst_dag.egg-info/ csst-msc-c9-50sqdeg-v3-plan/ build/ batch/ __pycache__/ *.pyc *.pyo csst_dag/v2/dag/__init__.py +4 −0 Original line number Diff line number Diff line Loading @@ -10,6 +10,10 @@ CSST_DAGS = { "csst-msc-l1-mbi", dispatcher=Dispatcher.dispatch_file, ), "csst-echo": FileDAG( "csst-echo", dispatcher=Dispatcher.dispatch_file, ), "csst-msc-l1-ast": FileDAG( "csst-msc-l1-ast", dispatcher=Dispatcher.dispatch_file, Loading csst_dag/v2/dag/base.py +4 −1 Original line number Diff line number Diff line Loading @@ -146,6 +146,7 @@ class BaseDAG(ABC): *, data: Optional[dict[str, Any]] = None, proc: Optional[dict[str, Any]] = None, docker_images: Optional[dict[str, Any]] = None, ) -> DagRunGroup: """ 统一触发入口:先解析三类依赖,再交给子类组装 payload。 Loading @@ -157,6 +158,7 @@ class BaseDAG(ABC): """ data = dict(data or {}) proc = dict(proc or {}) docker_images = dict(docker_images or {}) plan_query, data_query, upstream_query = self.build_dependency_queries( data=data, proc=proc Loading @@ -172,7 +174,7 @@ class BaseDAG(ABC): ), ) deps = self.filter_dependencies(deps) return self.build_trigger_payload(data=data, proc=proc, deps=deps) return self.build_trigger_payload(data=data, proc=proc, docker_images=docker_images, deps=deps) @abstractmethod def build_trigger_payload( Loading @@ -180,6 +182,7 @@ class BaseDAG(ABC): *, data: dict[str, Any], proc: dict[str, Any], docker_images: dict[str, Any], deps: DependencyBundle, ) -> DagRunGroup: """由子类实现:将输入与依赖解析结果组装成 DagRunGroup 对象。""" Loading csst_dag/v2/dag/brick.py +2 −0 Original line number Diff line number Diff line Loading @@ -80,12 +80,14 @@ class BrickDAG(BaseDAG): *, data: dict[str, Any], proc: dict[str, Any], docker_images: dict[str, Any], deps: DependencyBundle, ) -> DagRunGroup: dag_group = DagRunGroup( dag=self.dag_name, batch_id=data.get("batch_id"), priority=data.get("priority", "low"), docker_images=docker_images, ) if deps.data is None or len(deps.data) == 0: return dag_group Loading csst_dag/v2/dag/file.py +2 −0 Original line number Diff line number Diff line Loading @@ -151,12 +151,14 @@ class FileDAG(BaseDAG): *, data: dict[str, Any], proc: dict[str, Any], docker_images: dict[str, Any], deps: DependencyBundle, ) -> DagRunGroup: dag_group = DagRunGroup( dag=self.dag_name, batch_id=data.get("batch_id"), priority=data.get("priority", "low"), docker_images=docker_images, ) if deps.plan is None or deps.data is None or len(deps.plan) == 0 or len(deps.data) == 0: return dag_group Loading Loading
.gitignore +4 −1 Original line number Diff line number Diff line Loading @@ -3,3 +3,6 @@ csst_dag.egg-info/ csst-msc-c9-50sqdeg-v3-plan/ build/ batch/ __pycache__/ *.pyc *.pyo
csst_dag/v2/dag/__init__.py +4 −0 Original line number Diff line number Diff line Loading @@ -10,6 +10,10 @@ CSST_DAGS = { "csst-msc-l1-mbi", dispatcher=Dispatcher.dispatch_file, ), "csst-echo": FileDAG( "csst-echo", dispatcher=Dispatcher.dispatch_file, ), "csst-msc-l1-ast": FileDAG( "csst-msc-l1-ast", dispatcher=Dispatcher.dispatch_file, Loading
csst_dag/v2/dag/base.py +4 −1 Original line number Diff line number Diff line Loading @@ -146,6 +146,7 @@ class BaseDAG(ABC): *, data: Optional[dict[str, Any]] = None, proc: Optional[dict[str, Any]] = None, docker_images: Optional[dict[str, Any]] = None, ) -> DagRunGroup: """ 统一触发入口:先解析三类依赖,再交给子类组装 payload。 Loading @@ -157,6 +158,7 @@ class BaseDAG(ABC): """ data = dict(data or {}) proc = dict(proc or {}) docker_images = dict(docker_images or {}) plan_query, data_query, upstream_query = self.build_dependency_queries( data=data, proc=proc Loading @@ -172,7 +174,7 @@ class BaseDAG(ABC): ), ) deps = self.filter_dependencies(deps) return self.build_trigger_payload(data=data, proc=proc, deps=deps) return self.build_trigger_payload(data=data, proc=proc, docker_images=docker_images, deps=deps) @abstractmethod def build_trigger_payload( Loading @@ -180,6 +182,7 @@ class BaseDAG(ABC): *, data: dict[str, Any], proc: dict[str, Any], docker_images: dict[str, Any], deps: DependencyBundle, ) -> DagRunGroup: """由子类实现:将输入与依赖解析结果组装成 DagRunGroup 对象。""" Loading
csst_dag/v2/dag/brick.py +2 −0 Original line number Diff line number Diff line Loading @@ -80,12 +80,14 @@ class BrickDAG(BaseDAG): *, data: dict[str, Any], proc: dict[str, Any], docker_images: dict[str, Any], deps: DependencyBundle, ) -> DagRunGroup: dag_group = DagRunGroup( dag=self.dag_name, batch_id=data.get("batch_id"), priority=data.get("priority", "low"), docker_images=docker_images, ) if deps.data is None or len(deps.data) == 0: return dag_group Loading
csst_dag/v2/dag/file.py +2 −0 Original line number Diff line number Diff line Loading @@ -151,12 +151,14 @@ class FileDAG(BaseDAG): *, data: dict[str, Any], proc: dict[str, Any], docker_images: dict[str, Any], deps: DependencyBundle, ) -> DagRunGroup: dag_group = DagRunGroup( dag=self.dag_name, batch_id=data.get("batch_id"), priority=data.get("priority", "low"), docker_images=docker_images, ) if deps.plan is None or deps.data is None or len(deps.plan) == 0 or len(deps.data) == 0: return dag_group Loading