diff --git a/semantica/pipeline/pipeline_builder.py b/semantica/pipeline/pipeline_builder.py index d2044cbe..c64d5e64 100644 --- a/semantica/pipeline/pipeline_builder.py +++ b/semantica/pipeline/pipeline_builder.py @@ -272,7 +272,19 @@ class PipelineBuilder: step_name = step_config.get("name") step_type = step_config.get("type") if step_name and step_type: - self.add_step(step_name, step_type, **step_config.get("config", {})) + step = self.add_step( + step_name, step_type, **step_config.get("config", {}) + ) + step.dependencies = list( + step_config.get("dependencies", step.dependencies) + ) + step.delta_mode = step_config.get("delta_mode", step.delta_mode) + step.base_version_id = step_config.get( + "base_version_id", step.base_version_id + ) + step.target_version_id = step_config.get( + "target_version_id", step.target_version_id + ) # Set parallelism if specified if "parallelism" in pipeline_config: @@ -398,14 +410,29 @@ class PipelineSerializer: Returns: Serialized pipeline + + Notes: + Step handlers are runtime callables and are intentionally omitted from + the serialized representation. They must be rebound after deserialization. """ + reserved_config_keys = { + "handler", + "dependencies", + "delta_mode", + "base_version_id", + "target_version_id", + } pipeline_data = { "name": pipeline.name, "steps": [ { "name": step.name, "type": step.step_type, - "config": step.config, + "config": { + key: value + for key, value in step.config.items() + if key not in reserved_config_keys + }, "dependencies": step.dependencies, "delta_mode": getattr(step, "delta_mode", False), "base_version_id": getattr(step, "base_version_id", None), @@ -445,6 +472,18 @@ class PipelineSerializer: else: pipeline_data = serialized_pipeline + # Runtime handlers are process-local and cannot be reconstructed safely + # from serialized data. Copy before sanitizing so dict inputs are not mutated. + pipeline_data = dict(pipeline_data) + sanitized_steps = [] + for step_data in pipeline_data.get("steps", []): + sanitized_step = dict(step_data) + step_config = dict(sanitized_step.get("config", {})) + step_config.pop("handler", None) + sanitized_step["config"] = step_config + sanitized_steps.append(sanitized_step) + pipeline_data["steps"] = sanitized_steps + # Reconstruct pipeline builder = PipelineBuilder(**self.config) pipeline = builder.build_pipeline(pipeline_data, **options) diff --git a/tests/pipeline/test_pipeline_serializer.py b/tests/pipeline/test_pipeline_serializer.py new file mode 100644 index 00000000..fe616ddd --- /dev/null +++ b/tests/pipeline/test_pipeline_serializer.py @@ -0,0 +1,103 @@ +import copy +import json + +import pytest + +from semantica.pipeline.pipeline_builder import PipelineBuilder, PipelineSerializer + + +@pytest.mark.parametrize("serialization_format", ["dict", "json"]) +def test_roundtrip_preserves_dependencies_and_delta_metadata(serialization_format): + builder = PipelineBuilder() + builder.add_step("extract", "source") + builder.add_step( + "index", + "sink", + delta_mode=True, + base_version_id="v1", + target_version_id="v2", + ) + builder.connect_steps("extract", "index") + pipeline = builder.build("incremental-index") + + serializer = PipelineSerializer() + serialized = serializer.serialize_pipeline(pipeline, format=serialization_format) + restored = serializer.deserialize_pipeline(serialized) + + index_step = next(step for step in restored.steps if step.name == "index") + assert index_step.dependencies == ["extract"] + assert index_step.delta_mode is True + assert index_step.base_version_id == "v1" + assert index_step.target_version_id == "v2" + + +@pytest.mark.parametrize("serialization_format", ["dict", "json"]) +def test_serialization_omits_runtime_handlers(serialization_format): + def handler(data, **config): + return data + + builder = PipelineBuilder() + builder.add_step("extract", "source", handler=handler, batch_size=10) + pipeline = builder.build("handler-pipeline") + + serializer = PipelineSerializer() + serialized = serializer.serialize_pipeline(pipeline, format=serialization_format) + serialized_data = ( + json.loads(serialized) if isinstance(serialized, str) else serialized + ) + + assert serialized_data["steps"][0]["config"] == {"batch_size": 10} + + restored = serializer.deserialize_pipeline(serialized) + assert restored.steps[0].handler is None + assert restored.steps[0].config == {"batch_size": 10} + + +def test_deserialization_ignores_legacy_stringified_handler(): + serialized = json.dumps( + { + "name": "legacy-handler-pipeline", + "steps": [ + { + "name": "extract", + "type": "source", + "config": { + "handler": "", + "batch_size": 10, + }, + "dependencies": [], + } + ], + } + ) + + restored = PipelineSerializer().deserialize_pipeline(serialized) + + assert restored.steps[0].handler is None + assert restored.steps[0].config == {"batch_size": 10} + + +def test_deserialization_does_not_mutate_caller_owned_dict(): + payload = { + "name": "legacy-handler-pipeline", + "steps": [ + { + "name": "extract", + "type": "source", + "config": { + "handler": "", + "batch_size": 10, + }, + "dependencies": [], + } + ], + } + snapshot = copy.deepcopy(payload) + + restored = PipelineSerializer().deserialize_pipeline(payload) + + assert payload == snapshot + assert "handler" in payload["steps"][0]["config"] + assert payload is not snapshot + assert restored.steps[0].handler is None + assert restored.steps[0].config == {"batch_size": 10}