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={}))