mirror of
https://github.com/run-llama/workflows-py.git
synced 2026-08-24 20:01:34 -04:00
Add idle workflow tracking via internal event (#261)
This commit is contained in:
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user