diff --git a/reme4/components/job/base_job.py b/reme4/components/job/base_job.py index 0d2432cf..860849f7 100644 --- a/reme4/components/job/base_job.py +++ b/reme4/components/job/base_job.py @@ -29,33 +29,36 @@ class BaseJob(BaseComponent): self.description = description self.parameters = parameters or {} self.step_configs = steps or [] + # Resolved at start: (cls, params) pairs. Steps are re-instantiated per call so they stay + # stateless across runs and concurrent invocations don't share mutable step state. self.step_specs: list[tuple[type["BaseStep"], dict]] = [] async def _start(self) -> None: - """Resolve step configs into (cls, params) pairs; defer instantiation to __call__.""" if self.app_context is None: raise RuntimeError(f"app_context must be provided for job '{self.name}'") - for raw in self.step_configs: - config = raw if isinstance(raw, ComponentConfig) else ComponentConfig(**raw) - if not config.backend: - raise ValueError("Step is missing the required 'backend' field") - step_cls = R.get(ComponentEnum.STEP, config.backend) - if not step_cls: - raise ValueError(f"Unregistered backend '{config.backend}' of type '{ComponentEnum.STEP}'") - params = config.model_dump() - params["app_context"] = self.app_context - self.step_specs.append((step_cls, params)) + self.step_specs = [self._resolve_step(raw) for raw in self.step_configs] async def _close(self) -> None: - """Release all step specs.""" self.step_specs.clear() + def _resolve_step(self, raw: ComponentConfig | dict) -> tuple[type["BaseStep"], dict]: + """Validate a step config and look up its class via the registry.""" + config = raw if isinstance(raw, ComponentConfig) else ComponentConfig(**raw) + if not config.backend: + raise ValueError("Step is missing the required 'backend' field") + step_cls = R.get(ComponentEnum.STEP, config.backend) + if not step_cls: + raise ValueError(f"Unregistered backend '{config.backend}' of type '{ComponentEnum.STEP}'") + params = config.model_dump() + params["app_context"] = self.app_context + return step_cls, params + def _build_steps(self) -> list["BaseStep"]: - """Instantiate fresh step instances from stored specs.""" + # dict(params) copies kwargs so steps cannot mutate the shared spec. return [step_cls(**dict(params)) for step_cls, params in self.step_specs] async def __call__(self, **kwargs) -> Response: - """Execute all steps in order and return the final response.""" + """Run all steps in order, capturing any failure into the response.""" context = RuntimeContext(**kwargs) try: for step in self._build_steps():