fix(pipeline): wire registered step handlers (#1215)

* fix(pipeline): wire registered step handlers

Resolve handlers registered by step type, keep explicit handlers authoritative, and prevent builder control fields from leaking into runtime kwargs.

Refs #1214

* fix(pipeline): preserve dependencies on deserialize

* fix(pipeline): dispatch falsy handlers via identity check

---------

Co-authored-by: 江俊杰 <jiangjunjie.37@jd.com>
This commit is contained in:
cxzg007
2026-08-27 13:11:42 +05:00
committed by GitHub
co-authored by 江俊杰
parent 9cec305a75
commit c9c777993b
4 changed files with 99 additions and 8 deletions
+4 -4
View File
@@ -222,10 +222,10 @@ builder.register_step_handler("ner_extract", run_ner)
builder.register_step_handler("triplet_extract", run_triplets)
builder.register_step_handler("kg_merge", merge_into_graph)
builder.add_step("ingest", "file_ingest", handler=ingest_stix_bundles, path="./stix_bundles/")
builder.add_step("ner", "ner_extract", handler=run_ner, confidence_threshold=0.75)
builder.add_step("triplets", "triplet_extract", handler=run_triplets, include_temporal=True)
builder.add_step("store", "kg_merge", handler=merge_into_graph, output_path="./cti_output/")
builder.add_step("ingest", "file_ingest", path="./stix_bundles/")
builder.add_step("ner", "ner_extract", confidence_threshold=0.75)
builder.add_step("triplets", "triplet_extract", include_temporal=True)
builder.add_step("store", "kg_merge", output_path="./cti_output/")
# ingest feeds both ner and triplets in parallel
builder.connect_steps("ingest", "ner")
+1 -1
View File
@@ -379,7 +379,7 @@ class ExecutionEngine:
data = delta_result
if step.handler:
if step.handler is not None:
return step.handler(data, **step.config, **options)
else:
return data
+6 -2
View File
@@ -131,13 +131,17 @@ class PipelineBuilder:
delta_mode = config.pop("delta_mode", False)
base_version_id = config.pop("base_version_id", None)
target_version_id = config.pop("target_version_id", None)
dependencies = config.pop("dependencies", [])
handler = config.pop("handler", None)
if handler is None:
handler = self.step_registry.get(step_type)
step = PipelineStep(
name=step_name,
step_type=step_type,
config=config,
dependencies=config.get("dependencies", []),
handler=config.get("handler"),
dependencies=dependencies,
handler=handler,
delta_mode = delta_mode,
base_version_id=base_version_id,
target_version_id=target_version_id,
+88 -1
View File
@@ -66,6 +66,93 @@ class TestPipelineModule(unittest.TestCase):
self.assertEqual(result.output, 12) # (5 + 1) * 2 = 12
self.assertEqual(pipeline.steps[0].status, StepStatus.COMPLETED)
def test_registered_step_handler_executes(self):
"""A handler registered by step type should execute."""
def increment(data):
return data + 1
builder = PipelineBuilder()
builder.register_step_handler("math", increment)
builder.add_step("increment", "math")
result = ExecutionEngine().execute_pipeline(
builder.build("registered"), data=1
)
self.assertTrue(result.success)
self.assertEqual(result.output, 2)
def test_explicit_handler_does_not_receive_control_fields(self):
"""Builder-only fields should not be passed to strict handlers."""
def source(data):
return data
def increment(data, amount):
return data + amount
builder = PipelineBuilder()
builder.add_step("source", "source", handler=source)
step = builder.add_step(
"increment",
"math",
handler=increment,
dependencies=["source"],
amount=2,
)
result = ExecutionEngine().execute_pipeline(
builder.build("strict"), data=1
)
self.assertTrue(result.success)
self.assertEqual(result.output, 3)
self.assertEqual(step.config, {"amount": 2})
def test_explicit_handler_overrides_registered_handler(self):
"""An explicit step handler should take precedence over the registry."""
builder = PipelineBuilder()
builder.register_step_handler("math", lambda data: data + 100)
builder.add_step("increment", "math", handler=lambda data: data + 1)
result = ExecutionEngine().execute_pipeline(
builder.build("override"), data=1
)
self.assertTrue(result.success)
self.assertEqual(result.output, 2)
def test_step_without_handler_passes_input_through(self):
"""A step with no explicit or registered handler should be a no-op."""
builder = PipelineBuilder()
builder.add_step("passthrough", "unregistered")
result = ExecutionEngine().execute_pipeline(
builder.build("handlerless"), data={"value": 1}
)
self.assertTrue(result.success)
self.assertEqual(result.output, {"value": 1})
def test_falsy_explicit_handler_is_invoked(self):
"""A handler whose __bool__ is False must still be dispatched."""
class FalseyHandler:
def __bool__(self):
return False
def __call__(self, data):
return {"explicit": data}
builder = PipelineBuilder()
builder.add_step("step1", "mytype", handler=FalseyHandler())
result = ExecutionEngine().execute_pipeline(
builder.build("falsy"), data=1
)
self.assertTrue(result.success)
self.assertEqual(result.output, {"explicit": 1})
def test_execution_engine_failure(self):
"""Test pipeline failure handling."""
def failing_handler(data, **kwargs):
@@ -112,8 +199,8 @@ class TestPipelineModule(unittest.TestCase):
_ = semantica.pipeline
from semantica.pipeline import PipelineBuilder, PipelineValidator
from semantica.deduplication import DuplicateDetector
from semantica.pipeline import PipelineBuilder, PipelineValidator
builder = PipelineBuilder()
builder.add_step("step1", "dummy")