Skip to content

Commit e23ed30

Browse files
committed
refactor(#55): separate adapter provider from execution context
- Add AdapterProvider and internal ResolvedAdapterSet seam - Centralize adapter config precedence and runtime adapter construction in provider - Write resolved adapter configs back into RunConfig for compatibility - Route PipelineExecutor run plans through explicit registry project adapter config - Pass resolved adapter sets through Pipeline, PipelineRunner, and ExecutionContextBuilder - Keep runtime adapter construction inside retry-wrapped runner operations - Preserve direct Pipeline.run/run_async fallback compatibility - Add provider, retry-scope, Ray policy, explicit-none, and custom-adapter coverage Tests: .venv/bin/python -m pytest tests/ -q --tb=short (469 passed, 1 warning) Ruff: .venv/bin/python -m ruff check <touched files> Final code review: Standards + Spec clean
1 parent 2007927 commit e23ed30

10 files changed

Lines changed: 533 additions & 178 deletions

File tree

Lines changed: 164 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,164 @@
1+
"""Adapter resolution and construction for pipeline execution."""
2+
3+
from __future__ import annotations
4+
5+
from dataclasses import dataclass
6+
from typing import Any
7+
8+
from ..cfg.pipeline.run import RunConfig, WithAdapterConfig
9+
from ..utils.adapter import AdapterManager
10+
11+
12+
@dataclass(frozen=True)
13+
class ResolvedAdapterSet:
14+
"""Resolved adapter configs and runtime adapter instances for a run."""
15+
16+
with_adapter_cfg: WithAdapterConfig
17+
pipeline_adapter_cfg: Any
18+
project_adapter_cfg: Any
19+
runtime_adapters: list[Any]
20+
21+
22+
class AdapterProvider:
23+
"""Resolve adapter config precedence and construct runtime adapters."""
24+
25+
def __init__(self, adapter_manager: AdapterManager | None = None) -> None:
26+
self._adapter_manager = adapter_manager or AdapterManager()
27+
28+
def resolve(
29+
self,
30+
run_config: RunConfig,
31+
pipeline_config: Any,
32+
project_adapter_base: Any = None,
33+
*,
34+
construct_runtime: bool = True,
35+
) -> ResolvedAdapterSet:
36+
"""Resolve adapter configs into ``run_config`` and optionally create adapters."""
37+
explicit_overrides = set(run_config.explicit_overrides or [])
38+
39+
pipeline_adapter_cfg = self._resolve_pipeline_adapter_config(
40+
run_config,
41+
pipeline_config,
42+
explicit_overrides,
43+
)
44+
project_adapter_cfg = self._resolve_project_adapter_config(
45+
run_config,
46+
project_adapter_base,
47+
explicit_overrides,
48+
)
49+
with_adapter_cfg = run_config.with_adapter or WithAdapterConfig()
50+
run_config.with_adapter = with_adapter_cfg
51+
52+
runtime_adapters = (
53+
self._create_runtime_adapters(
54+
run_config,
55+
with_adapter_cfg,
56+
pipeline_adapter_cfg,
57+
project_adapter_cfg,
58+
)
59+
if construct_runtime
60+
else []
61+
)
62+
63+
return ResolvedAdapterSet(
64+
with_adapter_cfg=with_adapter_cfg,
65+
pipeline_adapter_cfg=pipeline_adapter_cfg,
66+
project_adapter_cfg=project_adapter_cfg,
67+
runtime_adapters=runtime_adapters,
68+
)
69+
70+
def construct_runtime_adapters(
71+
self,
72+
run_config: RunConfig,
73+
adapter_set: ResolvedAdapterSet,
74+
) -> ResolvedAdapterSet:
75+
"""Construct runtime adapters for an already-resolved adapter set."""
76+
return ResolvedAdapterSet(
77+
with_adapter_cfg=adapter_set.with_adapter_cfg,
78+
pipeline_adapter_cfg=adapter_set.pipeline_adapter_cfg,
79+
project_adapter_cfg=adapter_set.project_adapter_cfg,
80+
runtime_adapters=self._create_runtime_adapters(
81+
run_config,
82+
adapter_set.with_adapter_cfg,
83+
adapter_set.pipeline_adapter_cfg,
84+
adapter_set.project_adapter_cfg,
85+
),
86+
)
87+
88+
def _create_runtime_adapters(
89+
self,
90+
run_config: RunConfig,
91+
with_adapter_cfg: WithAdapterConfig,
92+
pipeline_adapter_cfg: Any,
93+
project_adapter_cfg: Any,
94+
) -> list[Any]:
95+
adapters = self._adapter_manager.create_adapters(
96+
with_adapter_cfg,
97+
pipeline_adapter_cfg,
98+
project_adapter_cfg,
99+
)
100+
if run_config.adapter:
101+
adapters.extend(run_config.adapter.values())
102+
return adapters
103+
104+
def _resolve_pipeline_adapter_config(
105+
self,
106+
run_config: RunConfig,
107+
pipeline_config: Any,
108+
explicit_overrides: set[str],
109+
) -> Any:
110+
from ..cfg.pipeline.adapter import AdapterConfig as PipelineAdapterConfig
111+
112+
pipeline_adapter_base = getattr(pipeline_config, "adapter", None)
113+
if pipeline_adapter_base is not None and not (
114+
"pipeline_adapter_cfg" in explicit_overrides
115+
and run_config.pipeline_adapter_cfg is None
116+
):
117+
run_config.pipeline_adapter_cfg = (
118+
self._adapter_manager.resolve_pipeline_adapter_config(
119+
run_config.pipeline_adapter_cfg,
120+
pipeline_adapter_base,
121+
)
122+
)
123+
124+
if run_config.pipeline_adapter_cfg is None:
125+
run_config.pipeline_adapter_cfg = PipelineAdapterConfig()
126+
return run_config.pipeline_adapter_cfg
127+
128+
def _resolve_project_adapter_config(
129+
self,
130+
run_config: RunConfig,
131+
project_adapter_base: Any,
132+
explicit_overrides: set[str],
133+
) -> Any:
134+
from ..cfg.project.adapter import AdapterConfig as ProjectAdapterConfig
135+
136+
if not (
137+
"project_adapter_cfg" in explicit_overrides
138+
and run_config.project_adapter_cfg is None
139+
):
140+
run_config.project_adapter_cfg = (
141+
self._adapter_manager.resolve_project_adapter_config(
142+
run_config.project_adapter_cfg,
143+
project_adapter_base,
144+
)
145+
)
146+
147+
if run_config.project_adapter_cfg is None:
148+
run_config.project_adapter_cfg = ProjectAdapterConfig()
149+
return run_config.project_adapter_cfg
150+
151+
152+
def resolve_run_config_adapter_configs(
153+
run_config: RunConfig,
154+
pipeline_config: Any,
155+
project_adapter_base: Any = None,
156+
) -> RunConfig:
157+
"""Compatibility helper that resolves adapter configs without exposing adapters."""
158+
AdapterProvider().resolve(
159+
run_config,
160+
pipeline_config,
161+
project_adapter_base,
162+
construct_runtime=False,
163+
)
164+
return run_config

src/flowerpower/pipeline/execution_context.py

Lines changed: 11 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,9 @@
99
from loguru import logger
1010

1111
from ..cfg.pipeline.run import ExecutorConfig, RunConfig
12-
from ..utils.adapter import AdapterManager
12+
from .adapter_provider import ResolvedAdapterSet, resolve_run_config_adapter_configs
13+
14+
__all__ = ["ExecutionContextBuilder", "resolve_run_config_adapter_configs"]
1315

1416

1517
class ExecutionContextBuilder:
@@ -19,25 +21,27 @@ def __init__(
1921
self,
2022
*,
2123
executor_factory: Any,
22-
adapter_manager: Any,
23-
pipeline_config: Any,
24-
project_context: Any,
24+
adapter_manager: Any = None,
25+
pipeline_config: Any = None,
26+
project_context: Any = None,
2527
) -> None:
2628
self._executor_factory = executor_factory
2729
self._adapter_manager = adapter_manager
2830
self._pipeline_config = pipeline_config
2931
self._project_context = project_context
3032

3133
def build(
32-
self, run_config: RunConfig
34+
self,
35+
run_config: RunConfig,
36+
adapter_set: ResolvedAdapterSet,
3337
) -> tuple[h_executors.BaseExecutor, Callable | None, list]:
3438
"""Create executor, shutdown function, and adapters for a pipeline run."""
3539
executor_cfg = run_config.executor or ExecutorConfig()
3640
executor, cleanup_fn = self._create_executor(
3741
executor_cfg,
38-
run_config.project_adapter_cfg,
42+
adapter_set.project_adapter_cfg,
3943
)
40-
adapters = self._create_adapters(run_config)
44+
adapters = list(adapter_set.runtime_adapters)
4145
logger.debug(
4246
"Execution context created. executor={executor} adapters={adapters}",
4347
executor=executor_cfg.type,
@@ -62,27 +66,6 @@ def _create_executor(
6266
cleanup_fn = ray_module.shutdown if should_shutdown else None
6367
return executor, cleanup_fn
6468

65-
def _create_adapters(self, run_config: RunConfig) -> list:
66-
with_adapter_cfg = run_config.with_adapter
67-
68-
pipeline_adapter_cfg = run_config.pipeline_adapter_cfg
69-
if pipeline_adapter_cfg is None:
70-
from ..cfg.pipeline.adapter import AdapterConfig as PipelineAdapterConfig
71-
pipeline_adapter_cfg = PipelineAdapterConfig()
72-
73-
project_adapter_cfg = run_config.project_adapter_cfg
74-
if project_adapter_cfg is None:
75-
from ..cfg.project.adapter import AdapterConfig as ProjectAdapterConfig
76-
project_adapter_cfg = ProjectAdapterConfig()
77-
78-
adapters = self._adapter_manager.create_adapters(
79-
with_adapter_cfg, pipeline_adapter_cfg, project_adapter_cfg
80-
)
81-
82-
if run_config.adapter:
83-
adapters.extend(run_config.adapter.values())
84-
85-
return adapters
8669

8770
@staticmethod
8871
def _get_optional_ray():
@@ -98,34 +81,3 @@ def _get_optional_ray():
9881
return None
9982

10083

101-
def resolve_run_config_adapter_configs(
102-
run_config: RunConfig,
103-
pipeline_config: Any,
104-
project_adapter_base: Any = None,
105-
) -> RunConfig:
106-
"""Merge pipeline and project adapter defaults into a resolved RunConfig.
107-
108-
This resolution must happen before runtime object construction so that the
109-
execution-context builder can consume the resolved RunConfig values without
110-
re-deciding precedence against pipeline defaults.
111-
"""
112-
explicit_overrides = set(run_config.explicit_overrides or [])
113-
manager = AdapterManager()
114-
if pipeline_config is not None:
115-
pipeline_adapter = getattr(pipeline_config, "adapter", None)
116-
if pipeline_adapter is not None and not (
117-
"pipeline_adapter_cfg" in explicit_overrides
118-
and run_config.pipeline_adapter_cfg is None
119-
):
120-
run_config.pipeline_adapter_cfg = manager.resolve_pipeline_adapter_config(
121-
run_config.pipeline_adapter_cfg, pipeline_adapter
122-
)
123-
124-
if not (
125-
"project_adapter_cfg" in explicit_overrides
126-
and run_config.project_adapter_cfg is None
127-
):
128-
run_config.project_adapter_cfg = manager.resolve_project_adapter_config(
129-
run_config.project_adapter_cfg, project_adapter_base
130-
)
131-
return run_config

src/flowerpower/pipeline/executor.py

Lines changed: 13 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@
1111
)
1212
from ..utils.logging import setup_logging
1313
from ..utils.security import validate_pipeline_name
14-
from .execution_context import resolve_run_config_adapter_configs
14+
from .adapter_provider import AdapterProvider, ResolvedAdapterSet
1515

1616
if TYPE_CHECKING:
1717
from .config_manager import PipelineConfigManager
@@ -26,6 +26,7 @@ class PipelineRunPlan:
2626
pipeline_config: Any
2727
run_config: RunConfig
2828
pipeline: Any
29+
adapter_set: ResolvedAdapterSet
2930

3031

3132
class PipelineExecutor:
@@ -91,7 +92,10 @@ def run(self, name: str, run_config: RunConfig | None = None, **kwargs) -> dict[
9192
"""
9293
plan = self._build_run_plan(name, run_config, **kwargs)
9394
self._apply_run_logging(plan)
94-
return plan.pipeline._run_resolved(run_config=plan.run_config)
95+
return plan.pipeline._run_resolved(
96+
run_config=plan.run_config,
97+
adapter_set=plan.adapter_set,
98+
)
9599

96100
async def run_async(
97101
self, name: str, run_config: RunConfig | None = None, **kwargs
@@ -110,7 +114,10 @@ async def run_async(
110114
"""
111115
plan = self._build_run_plan(name, run_config, **kwargs)
112116
self._apply_run_logging(plan)
113-
return await plan.pipeline._run_resolved_async(run_config=plan.run_config)
117+
return await plan.pipeline._run_resolved_async(
118+
run_config=plan.run_config,
119+
adapter_set=plan.adapter_set,
120+
)
114121

115122
def _build_run_plan(
116123
self,
@@ -129,12 +136,11 @@ def _build_run_plan(
129136
if kwargs:
130137
run_config = merge_run_config_with_kwargs(run_config, kwargs)
131138

132-
# Fold pipeline and project adapter defaults into the resolved RunConfig,
133-
# so runtime object construction consumes the resolved values only.
134-
run_config = resolve_run_config_adapter_configs(
139+
adapter_set = AdapterProvider().resolve(
135140
run_config,
136141
pipeline_config,
137142
self._project_adapter_base(),
143+
construct_runtime=False,
138144
)
139145

140146
# Guard against non-clearable fields that were left unset.
@@ -152,6 +158,7 @@ def _build_run_plan(
152158
pipeline_config=pipeline_config,
153159
run_config=run_config,
154160
pipeline=pipeline,
161+
adapter_set=adapter_set,
155162
)
156163

157164
@staticmethod

0 commit comments

Comments
 (0)