From b2f1bef1d0284c11d139809ce4621fddcbffa5da Mon Sep 17 00:00:00 2001 From: ogabrielluiz Date: Mon, 4 May 2026 17:49:09 -0300 Subject: [PATCH] feat(execution): route Loop subgraph through coordinator (seam #4) --- src/lfx/src/lfx/base/flow_controls/loop_utils.py | 7 ++++++- .../unit/components/flow_controls/test_loop_events.py | 4 ++-- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/src/lfx/src/lfx/base/flow_controls/loop_utils.py b/src/lfx/src/lfx/base/flow_controls/loop_utils.py index 575ededb97..cb2b6ba160 100644 --- a/src/lfx/src/lfx/base/flow_controls/loop_utils.py +++ b/src/lfx/src/lfx/base/flow_controls/loop_utils.py @@ -265,8 +265,13 @@ async def execute_loop_body( # Execute subgraph and collect results # Pass event_manager so UI receives events from subgraph execution + from lfx.execution import StepResult, get_default_coordinator + results = [] - async for result in iteration_subgraph.async_start(event_manager=event_manager): + async for step in get_default_coordinator().run(iteration_subgraph, inputs=[], event_manager=event_manager): + if not isinstance(step, StepResult): + continue + result = step.payload results.append(result) # Stop all on error (as per design decision) if hasattr(result, "valid") and not result.valid: diff --git a/src/lfx/tests/unit/components/flow_controls/test_loop_events.py b/src/lfx/tests/unit/components/flow_controls/test_loop_events.py index 5431770a95..ea71e6fef3 100644 --- a/src/lfx/tests/unit/components/flow_controls/test_loop_events.py +++ b/src/lfx/tests/unit/components/flow_controls/test_loop_events.py @@ -82,7 +82,7 @@ class TestEventManagerPropagation: received_event_manager = None # Create a mock subgraph that captures the event_manager - async def mock_async_start(event_manager=None): + async def mock_async_start(event_manager=None, **kwargs): # noqa: ARG001 nonlocal received_event_manager received_event_manager = event_manager yield MagicMock(valid=True, result_dict=MagicMock(outputs={})) @@ -119,7 +119,7 @@ class TestEventManagerPropagation: mock_event_manager = MagicMock() event_manager_calls = [] - async def mock_async_start(event_manager=None): + async def mock_async_start(event_manager=None, **kwargs): # noqa: ARG001 event_manager_calls.append(event_manager) yield MagicMock(valid=True, result_dict=MagicMock(outputs={}))