Add idle workflow tracking via internal event (#261)

This commit is contained in:
Adrian Lyjak
2026-01-09 15:07:53 -05:00
committed by GitHub
parent 894445a355
commit 3b043b8801
7 changed files with 377 additions and 0 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"llama-index-workflows": patch
---
Track when workflows are idle (waiting on external input)
@@ -277,6 +277,22 @@ class InternalDispatchEvent(Event):
pass
class WorkflowIdleEvent(InternalDispatchEvent):
"""Emitted when workflow transitions to idle (waiting on external input).
A workflow is idle when:
1. The workflow is running (hasn't completed/failed/cancelled)
2. All steps have no pending events in their queues
3. All steps have no workers currently executing
4. At least one step has an active waiter (from ctx.wait_for_event())
This event is intentionally minimal - no metadata beyond the event type.
Resumption from idle is signaled by StepStateChanged with StepState.RUNNING.
"""
pass
class StepState(Enum):
# is enqueued, but no capacity yet available to run
PREPARING = "preparing"
@@ -22,6 +22,7 @@ from workflows.events import (
StepState,
StepStateChanged,
StopEvent,
WorkflowIdleEvent,
)
from workflows.runtime.types.commands import (
CommandCompleteRun,
@@ -367,6 +368,25 @@ def rewind_in_progress(
return state, commands
def _check_idle_state(state: BrokerState) -> bool:
"""Returns True if workflow is idle (waiting only on external events).
A workflow is idle when:
1. The workflow is running (hasn't completed/failed/cancelled)
2. All steps have no pending events in their queues
3. All steps have no workers currently executing
4. At least one step has an active waiter (from ctx.wait_for_event())
"""
if not state.is_running:
return False
for worker_state in state.workers.values():
if worker_state.queue or worker_state.in_progress:
return False
return any(ws.collected_waiters for ws in state.workers.values())
def _process_step_result_tick(
tick: TickStepResult[R], init: BrokerState, now_seconds: float
) -> tuple[BrokerState, list[WorkflowCommand]]:
@@ -542,6 +562,14 @@ def _process_step_result_tick(
event, tick.step_name, worker_state, now_seconds
)
commands.extend(subcommands)
# Check for idle transition at end of processing
was_idle = _check_idle_state(init)
now_idle = _check_idle_state(state)
if now_idle and not was_idle:
commands.append(CommandPublishEvent(WorkflowIdleEvent()))
return state, commands
@@ -37,6 +37,7 @@ class PersistentHandler(BaseModel):
started_at: datetime | None = None
updated_at: datetime | None = None
completed_at: datetime | None = None
idle_since: datetime | None = None
ctx: dict[str, Any] = {}
@field_validator("result", mode="before")
@@ -32,6 +32,7 @@ from workflows.events import (
StepState,
StepStateChanged,
StopEvent,
WorkflowIdleEvent,
)
from workflows.handler import WorkflowHandler
from workflows.protocol import (
@@ -1428,6 +1429,7 @@ class _WorkflowHandler:
_workflow_store: AbstractWorkflowStore
_persistence_backoff: list[float]
_on_finish: Callable[[], Awaitable[None]] | None = None
idle_since: datetime | None = None
def _as_persistent(self) -> PersistentHandler:
"""Persist the current handler state immediately to the workflow store."""
@@ -1445,6 +1447,7 @@ class _WorkflowHandler:
started_at=self.started_at,
updated_at=self.updated_at,
completed_at=self.completed_at,
idle_since=self.idle_since,
ctx=self.run_handler.ctx.to_dict() if self.run_handler.ctx else {},
)
return persistent
@@ -1570,6 +1573,15 @@ class _WorkflowHandler:
await self.checkpoint()
self._on_finish = on_finish
async for event in self.run_handler.stream_events(expose_internal=True):
# Track idle state transitions
if isinstance(event, WorkflowIdleEvent):
self.idle_since = datetime.now(timezone.utc)
elif (
isinstance(event, StepStateChanged)
and event.step_state == StepState.RUNNING
):
self.idle_since = None # Resumed from idle
if ( # Watch for a specific internal event that signals the step is complete
isinstance(event, StepStateChanged)
and event.step_state == StepState.NOT_RUNNING
@@ -28,6 +28,7 @@ from workflows.events import (
StartEvent,
StepStateChanged,
StopEvent,
WorkflowIdleEvent,
)
from workflows.retry_policy import ConstantDelayRetryPolicy, RetryPolicy
from workflows.runtime.control_loop import control_loop
@@ -796,3 +797,98 @@ async def test_control_loop_retry_exhaustion_respects_total_time(
# Attempts should increment: 1, 2, 3, 4
expected_attempts = list(range(1, len(policy.observed_attempts) + 1))
assert policy.observed_attempts == expected_attempts
@pytest.mark.asyncio
async def test_control_loop_emits_idle_event_when_waiting(
test_plugin: MockRuntimePlugin,
) -> None:
"""
Test that WorkflowIdleEvent is emitted when workflow becomes idle waiting for event.
A workflow is idle when it has active waiters but no pending work. This test
validates that the idle event is published to the stream when this state is reached.
"""
class AwaitedEvent(Event):
value: str
class IdleTrackingWorkflow(Workflow):
@step
async def waiter(self, ev: StartEvent, ctx: Context) -> StopEvent:
awaited = await ctx.wait_for_event(
AwaitedEvent,
waiter_event=InputRequiredEvent(),
)
return StopEvent(result=f"received_{awaited.value}")
wf = IdleTrackingWorkflow(timeout=2.0)
task = asyncio.create_task(
run_control_loop(
workflow=wf,
start_event=StartEvent(),
plugin=test_plugin,
)
)
# Collect events until we see the WorkflowIdleEvent
idle_event_found = False
input_required_found = False
events_before_idle: list[Event] = []
while True:
ev = await test_plugin.get_stream_event(timeout=1.0)
if isinstance(ev, WorkflowIdleEvent):
idle_event_found = True
break
if isinstance(ev, InputRequiredEvent):
input_required_found = True
events_before_idle.append(ev)
if isinstance(ev, StopEvent):
break
# WorkflowIdleEvent should be emitted after the workflow enters wait state
assert idle_event_found, (
"WorkflowIdleEvent should be emitted when workflow is waiting for external event"
)
assert input_required_found, "InputRequiredEvent should be emitted before idle"
# Now send the awaited event to complete the workflow
await test_plugin.send_event(TickAddEvent(event=AwaitedEvent(value="test")))
result = await asyncio.wait_for(task, timeout=1.0)
assert isinstance(result, StopEvent)
assert result.result == "received_test"
@pytest.mark.asyncio
async def test_control_loop_idle_event_not_emitted_on_completion(
test_plugin: MockRuntimePlugin,
) -> None:
"""
Test that WorkflowIdleEvent is NOT emitted when workflow completes normally.
Even if a workflow has waiters, if it completes (StopEvent), it should not
emit an idle event because the workflow is no longer running.
"""
result = await run_control_loop(
workflow=SimpleWorkflow(timeout=1.0),
start_event=StartEvent(),
plugin=test_plugin,
)
# Verify the workflow completed
assert isinstance(result, StopEvent)
assert result.result == "processed_42"
# Drain and verify no WorkflowIdleEvent was emitted
all_events: list[Event] = []
while test_plugin.has_stream_events():
ev = await test_plugin.get_stream_event(timeout=0.1)
all_events.append(ev)
idle_events = [e for e in all_events if isinstance(e, WorkflowIdleEvent)]
assert len(idle_events) == 0, (
"WorkflowIdleEvent should not be emitted when workflow completes normally"
)
@@ -22,10 +22,12 @@ from workflows.events import (
StepState,
StepStateChanged,
StopEvent,
WorkflowIdleEvent,
)
from workflows.retry_policy import ConstantDelayRetryPolicy
from workflows.runtime.control_loop import (
_add_or_enqueue_event,
_check_idle_state,
_process_add_event_tick,
_process_cancel_run_tick,
_process_publish_event_tick,
@@ -613,3 +615,220 @@ def test_step_worker_failed_retry_preserves_first_attempt_at(
assert len(queue_cmds) == 1
assert queue_cmds[0].attempts == 3 # incremented from 2
assert queue_cmds[0].first_attempt_at == original_first_attempt_at # preserved!
# =============================================================================
# Idle Workflow Tracking Tests
# =============================================================================
def test_check_idle_state_not_running(base_state: BrokerState) -> None:
"""A workflow that is not running is not idle."""
base_state.is_running = False
assert _check_idle_state(base_state) is False
def test_check_idle_state_has_queued_events(base_state: BrokerState) -> None:
"""A workflow with queued events is not idle."""
base_state.workers["test_step"].queue.append(
EventAttempt(event=MyTestEvent(value=1))
)
assert _check_idle_state(base_state) is False
def test_check_idle_state_has_in_progress(base_state: BrokerState) -> None:
"""A workflow with in-progress workers is not idle."""
add_worker(base_state, MyTestEvent(value=1))
assert _check_idle_state(base_state) is False
def test_check_idle_state_no_waiters(base_state: BrokerState) -> None:
"""A workflow with no waiters is not idle (even with empty queues)."""
# State is running, no queue, no in_progress, but no waiters either
assert _check_idle_state(base_state) is False
def test_check_idle_state_is_idle_with_waiter(base_state: BrokerState) -> None:
"""A running workflow with only waiters and no work is idle."""
waiter = StepWorkerWaiter(
waiter_id="w1",
event=MyTestEvent(value=1),
waiting_for_event=OtherEvent,
requirements={},
has_requirements=False,
resolved_event=None,
)
base_state.workers["test_step"].collected_waiters.append(waiter)
assert _check_idle_state(base_state) is True
def test_idle_event_emitted_on_transition_to_idle(base_state: BrokerState) -> None:
"""WorkflowIdleEvent is emitted when workflow transitions to idle."""
event = MyTestEvent(value=42)
add_worker(base_state, event)
# Add a waiter so the workflow can become idle
waiter = StepWorkerWaiter(
waiter_id="w1",
event=event,
waiting_for_event=OtherEvent,
requirements={},
has_requirements=False,
resolved_event=None,
)
base_state.workers["test_step"].collected_waiters.append(waiter)
# Process result that completes the worker but leaves waiter active
tick: TickStepResult[Any] = TickStepResult(
step_name="test_step",
worker_id=0,
event=event,
result=[
AddWaiter(
waiter_id="w1",
waiter_event=None,
requirements={},
timeout=None,
event_type=OtherEvent,
)
],
)
new_state, commands = _process_step_result_tick(tick, base_state, now_seconds=110.0)
# Should have WorkflowIdleEvent as the last command
idle_commands = [
c
for c in commands
if isinstance(c, CommandPublishEvent) and isinstance(c.event, WorkflowIdleEvent)
]
assert len(idle_commands) == 1
assert _check_idle_state(new_state) is True
def test_check_idle_state_multi_step_not_idle_if_one_has_work(
base_state: BrokerState,
) -> None:
"""With multiple steps, not idle if any step has work."""
# Add a second step
other_step_cfg = StepConfig(
accepted_events=[OtherEvent],
event_name="ev",
return_types=[StopEvent, type(None)],
context_parameter="ctx",
retry_policy=None,
num_workers=1,
resources=[],
)
base_state.config.steps["other_step"] = InternalStepConfig(
accepted_events=[OtherEvent], retry_policy=None, num_workers=1
)
base_state.workers["other_step"] = InternalStepWorkerState(
queue=[],
config=other_step_cfg,
in_progress=[],
collected_events={},
collected_waiters=[],
)
# Add waiter to test_step (which alone would make it idle)
waiter = StepWorkerWaiter(
waiter_id="w1",
event=MyTestEvent(value=1),
waiting_for_event=OtherEvent,
requirements={},
has_requirements=False,
resolved_event=None,
)
base_state.workers["test_step"].collected_waiters.append(waiter)
# Without work in other_step, workflow is idle
assert _check_idle_state(base_state) is True
# Add in_progress work to other_step - now not idle
base_state.workers["other_step"].in_progress.append(
InProgressState(
event=OtherEvent(data="test"),
worker_id=0,
shared_state=StepWorkerState(
step_name="other_step",
collected_events={},
collected_waiters=[],
),
attempts=0,
first_attempt_at=100.0,
)
)
assert _check_idle_state(base_state) is False
# Or with queued work
base_state.workers["other_step"].in_progress = []
base_state.workers["other_step"].queue.append(
EventAttempt(event=OtherEvent(data="queued"))
)
assert _check_idle_state(base_state) is False
def test_no_idle_event_when_work_remains(base_state: BrokerState) -> None:
"""WorkflowIdleEvent is not emitted if there's still work to do."""
event = MyTestEvent(value=42)
add_worker(base_state, event)
# Queue another event so work remains after processing
base_state.workers["test_step"].queue.append(
EventAttempt(event=MyTestEvent(value=99))
)
tick: TickStepResult[Any] = TickStepResult(
step_name="test_step",
worker_id=0,
event=event,
result=[StepWorkerResult(result=None)], # Completes but queue has more
)
_, commands = _process_step_result_tick(tick, base_state, now_seconds=110.0)
idle_commands = [
c
for c in commands
if isinstance(c, CommandPublishEvent) and isinstance(c.event, WorkflowIdleEvent)
]
assert len(idle_commands) == 0
def test_no_idle_event_when_workflow_completes(base_state: BrokerState) -> None:
"""WorkflowIdleEvent is not emitted when workflow completes (StopEvent)."""
event = MyTestEvent(value=42)
add_worker(base_state, event)
# Add a waiter
waiter = StepWorkerWaiter(
waiter_id="w1",
event=event,
waiting_for_event=OtherEvent,
requirements={},
has_requirements=False,
resolved_event=None,
)
base_state.workers["test_step"].collected_waiters.append(waiter)
# Complete the workflow with StopEvent
tick: TickStepResult[Any] = TickStepResult(
step_name="test_step",
worker_id=0,
event=event,
result=[StepWorkerResult(result=StopEvent(result="done"))],
)
new_state, commands = _process_step_result_tick(tick, base_state, now_seconds=110.0)
# Workflow is no longer running
assert new_state.is_running is False
# No idle event should be emitted
idle_commands = [
c
for c in commands
if isinstance(c, CommandPublishEvent) and isinstance(c.event, WorkflowIdleEvent)
]
assert len(idle_commands) == 0