Skip to content
Open
29 changes: 29 additions & 0 deletions docs/releases/3.2.2.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,35 @@ Reported by [@Dolecor](https://github.com/Dolecor).

[#643](https://github.com/fgmacedo/python-statemachine/issues/643).

### Events due at the same moment kept their order only by chance

Both event queues are priority queues ordered by each trigger's due time, and triggers due at
the same time came out in no particular order. A clock coarse enough to stamp two events alike
could run them out of the order they were sent or raised: twenty events raised in one callback,
all stamped with the same time, ran as `0, 2, 6, 14, 19, 18, ...`. Each trigger now also holds
the order it was created in, and triggers due at the same moment run in that order, on both
engines and both queues.

### Cancelling a delayed event could disorder the others and strand its caller

`cancel_event()` removed a trigger by filtering the list behind the queue's heap, and a filtered
heap is not always a heap. The triggers left could then run out of due order: of ten delayed
events due from 10 to 90 ms, cancelling the first left the one due at 50 ms running before the
one due at 40 ms. The trigger's future went with it, so on the async engine a coroutine
awaiting the cancelled event's `send()` was never answered.

A cancelled trigger now stays where it is, marked as cancelled, and is dropped when it reaches
the head of the queue, without the processing loop waiting for its due time. A coroutine
awaiting it gets `None` as soon as it is cancelled.

### Delays measured on the wall clock

Every trigger was stamped with `time.time()` plus its delay, and the processing loop held a
trigger until `time.time()` reached its stamp. A step of the wall clock, from a time sync or a
change of the system time, moved every due time with it: stepped back, it held even an event
sent with no delay until the clock caught up; stepped forward, it ran delayed events early.
Delays are now measured on `time.monotonic()`, which no setting of the wall clock moves.

### `delay` and `internal` dropped from an explicit `Event`

`Event(dark.to(lit), delay=50)` rebuilt the event without its `delay`, so
Expand Down
2 changes: 2 additions & 0 deletions docs/statechart.md
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,8 @@ sm.send("timeout", delay=5000, send_id="my_timeout")
sm.cancel_event("my_timeout")
```

On the async engine, a coroutine awaiting the `send()` of a cancelled event gets `None`.

```{note}
The delay is **blocking** in the sync engine — the processing loop sleeps until the
delay elapses, holding the calling thread. In the async engine, delays are scheduled
Expand Down
250 changes: 118 additions & 132 deletions statemachine/engines/async_.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import contextvars
from collections.abc import Callable
from itertools import chain
from time import time
from time import monotonic
from typing import TYPE_CHECKING

from ..event_data import EventData
Expand Down Expand Up @@ -113,7 +113,7 @@ async def _conditions_match(self, transition: "Transition", trigger_data: Trigge
transition.cond.key, *args, on_error=on_error, **kwargs
)

async def _first_transition_that_matches( # type: ignore[override]
async def _first_transition_that_matches(
self,
state: State,
trigger_data: TriggerData,
Expand All @@ -130,7 +130,7 @@ async def _first_transition_that_matches( # type: ignore[override]
return transition
return None

async def _select_transitions( # type: ignore[override]
async def _select_transitions(
self, trigger_data: TriggerData, predicate: Callable
) -> "OrderedSet[Transition]":
enabled_transitions: "OrderedSet[Transition]" = OrderedSet()
Expand All @@ -147,7 +147,7 @@ async def _select_transitions( # type: ignore[override]
async def select_eventless_transitions(self, trigger_data: TriggerData):
return await self._select_transitions(trigger_data, lambda t, _e: t.is_eventless)

async def select_transitions(self, trigger_data: TriggerData) -> "OrderedSet[Transition]": # type: ignore[override]
async def select_transitions(self, trigger_data: TriggerData) -> "OrderedSet[Transition]":
return await self._select_transitions(trigger_data, lambda t, e: t.match(e))

async def _execute_transition_content(
Expand All @@ -172,32 +172,30 @@ async def _execute_transition_content(

return result

async def _exit_states( # type: ignore[override]
async def _exit_states(
self, enabled_transitions: "list[Transition]", trigger_data: TriggerData
) -> "OrderedSet[State]":
ordered_states, result = self._prepare_exit_states(enabled_transitions)
on_error = self._on_error_handler()

for info in ordered_states:
# Cancel invocations for this state before executing exit handlers.
if info.state is not None: # pragma: no branch
self._invoke_manager.cancel_for_state(info.state)
self._invoke_manager.cancel_for_state(info.state)

args, kwargs = await self._get_args_kwargs(
info.transition, trigger_data, source=info.state
)

if info.state is not None: # pragma: no branch
self._debug("%s Exiting state: %s", self._log_id, info.state)
await self.sm._callbacks.async_call(
info.state.exit.key, *args, on_error=on_error, **kwargs
)
self._debug("%s Exiting state: %s", self._log_id, info.state)
await self.sm._callbacks.async_call(
info.state.exit.key, *args, on_error=on_error, **kwargs
)

self._remove_state_from_configuration(info.state)

return result

async def _enter_states( # noqa: C901
async def _enter_states(
self,
enabled_transitions: "list[Transition]",
trigger_data: TriggerData,
Expand Down Expand Up @@ -329,19 +327,6 @@ async def microstep(self, transitions: "list[Transition]", trigger_data: Trigger

# --- Engine loop ---

async def _run_microstep(self, enabled_transitions, trigger_data): # pragma: no cover
"""Run a microstep for internal/eventless transitions with error handling.

Note: microstep() handles its own errors internally, so this try/except
is a safety net that is not expected to be reached in normal operation.
"""
try:
await self.microstep(list(enabled_transitions), trigger_data)
except InvalidDefinition:
raise
except Exception as e:
self._handle_error(e, trigger_data)

async def activate_initial_state(self, **kwargs):
"""Activate the initial state.

Expand All @@ -354,14 +339,12 @@ async def activate_initial_state(self, **kwargs):
"""
return await self.processing_loop()

async def processing_loop( # noqa: C901
self, caller_future: "asyncio.Future[object] | None" = None
):
"""Process event triggers with the 3-phase macrostep architecture.
async def processing_loop(self, caller_future: "asyncio.Future[object] | None" = None):
"""Process event triggers, a macrostep at a time.

Phase 1: Eventless transitions + internal queue until quiescence.
Phase 2: Remaining internal events (safety net for invoke-generated events).
Phase 3: External events.
Each macrostep takes eventless transitions and internal events until neither enables a
transition, spawns the invocations of the states it entered, and then takes the
external events that are due.

When ``caller_future`` is provided, the caller can ``await`` it to
receive its own event's result — even if another coroutine holds the
Expand All @@ -381,110 +364,14 @@ async def processing_loop( # noqa: C901
took_events = True
while took_events and self.running:
self.clear_cache()
took_events = False
macrostep_done = False

# Phase 1: eventless transitions and internal events
while not macrostep_done:
self._microstep_count = 0
self._debug(
"%s Macrostep %d: eventless/internal queue",
self._log_id,
self._macrostep_count,
)

self.clear_cache()
internal_event = TriggerData(self.sm, event=None) # null object for eventless
enabled_transitions = await self.select_eventless_transitions(internal_event)
if not enabled_transitions:
if self.internal_queue.is_empty():
macrostep_done = True
else:
internal_event = self.internal_queue.pop()
enabled_transitions = await self.select_transitions(internal_event)
if enabled_transitions:
self._debug(
"%s Enabled transitions: %s", self._log_id, enabled_transitions
)
took_events = True
await self._run_microstep(enabled_transitions, internal_event)
took_events = await self._run_macrostep()

# Spawn invoke handlers for states entered during this macrostep.
await self._invoke_manager.spawn_pending_async()
self._check_root_final_state()

# Phase 2: remaining internal events
while not self.internal_queue.is_empty(): # pragma: no cover
internal_event = self.internal_queue.pop()
enabled_transitions = await self.select_transitions(internal_event)
if enabled_transitions:
await self._run_microstep(enabled_transitions, internal_event)

# Phase 3: external events
self._debug("%s Macrostep %d: external queue", self._log_id, self._macrostep_count)
while not self.external_queue.is_empty():
self.clear_cache()
took_events = True
external_event = self.external_queue.pop()
current_time = time()
if external_event.execution_time > current_time:
self.put(external_event, _delayed=True)
await asyncio.sleep(self.sm._loop_sleep_in_ms)
# Break to Phase 1 so internal events and eventless
# transitions can be processed while we wait.
break

self._macrostep_count += 1
self._microstep_count = 0
self._debug(
"%s macrostep %d: event=%s",
self._log_id,
self._macrostep_count,
external_event.event,
)

# Handle lazy initial state activation.
# Break out of phase 3 so the outer loop restarts from phase 1
# (eventless/internal), ensuring internal events queued during
# initial entry are processed before any external events.
if external_event.event == "__initial__":
transitions = self._initial_transitions(external_event)
await self._enter_states(
transitions, external_event, OrderedSet(), OrderedSet()
)
break

# Finalize + autoforward for active invocations
self._invoke_manager.handle_external_event(external_event)

event_future = external_event.future
try:
enabled_transitions = await self.select_transitions(external_event)
self._debug(
"%s Enabled transitions: %s", self._log_id, enabled_transitions
)
if enabled_transitions:
result = await self.microstep(
list(enabled_transitions), external_event
)
self._resolve_future(event_future, result)
if first_result is self._sentinel:
first_result = result
else:
if not self.sm.allow_event_without_transition:
tna = TransitionNotAllowed(
external_event.event, self.sm.configuration
)
self._reject_future(event_future, tna)
self._reject_pending_futures(tna)
raise tna
# Event allowed but no transition — resolve with None
self._resolve_future(event_future, None)
except Exception as exc:
self._reject_future(event_future, exc)
self._reject_pending_futures(exc)
self.clear()
raise
took_external_events, first_result = await self._take_external_events(first_result)
took_events = took_events or took_external_events

except Exception as exc:
if caller_future is not None:
Expand All @@ -510,6 +397,105 @@ async def processing_loop( # noqa: C901
return await caller_future
return result

async def _run_macrostep(self) -> bool:
"""Take eventless transitions and internal events until neither enables a transition.

Returns whether any transition was taken. The internal queue is empty on return.
"""
took_events = False
while True:
self._microstep_count = 0
self._debug(
"%s Macrostep %d: eventless/internal queue",
self._log_id,
self._macrostep_count,
)

self.clear_cache()
internal_event = TriggerData(self.sm, event=None) # null object for eventless
enabled_transitions = await self.select_eventless_transitions(internal_event)
if not enabled_transitions:
if self.internal_queue.is_empty():
return took_events
internal_event = self.internal_queue.pop()
enabled_transitions = await self.select_transitions(internal_event)
if enabled_transitions:
self._debug("%s Enabled transitions: %s", self._log_id, enabled_transitions)
took_events = True
await self.microstep(list(enabled_transitions), internal_event)

async def _take_external_events(self, first_result):
"""Take external events in due order until the queue is empty or holds one not yet due.

Returns whether any event was taken from the queue, and ``first_result`` updated with
the result of the first event that enabled a transition.
"""
self._debug("%s Macrostep %d: external queue", self._log_id, self._macrostep_count)
took_events = False
while not self.external_queue.is_empty():
self.clear_cache()
took_events = True
external_event = self.external_queue.pop()
current_time = monotonic()
if external_event.execution_time > current_time:
self.put(external_event, _delayed=True)
await asyncio.sleep(self.sm._loop_sleep_in_ms)
# Break to the macrostep so internal events and eventless
# transitions can be processed while we wait.
break

self._macrostep_count += 1
self._microstep_count = 0
self._debug(
"%s macrostep %d: event=%s",
self._log_id,
self._macrostep_count,
external_event.event,
)

# Handle lazy initial state activation.
# Break to the macrostep so internal events queued during
# initial entry are processed before any external events.
if external_event.event == "__initial__":
transitions = self._initial_transitions(external_event)
await self._enter_states(transitions, external_event, OrderedSet(), OrderedSet())
break

result = await self._take_external_event(external_event)
if first_result is self._sentinel:
first_result = result
return took_events, first_result

async def _take_external_event(self, external_event: TriggerData):
"""Run one external event and settle its future.

Returns the event's microstep result, or the sentinel when it enabled no transition.
"""
# Finalize + autoforward for active invocations
self._invoke_manager.handle_external_event(external_event)

event_future = external_event.future
try:
enabled_transitions = await self.select_transitions(external_event)
self._debug("%s Enabled transitions: %s", self._log_id, enabled_transitions)
if enabled_transitions:
result = await self.microstep(list(enabled_transitions), external_event)
self._resolve_future(event_future, result)
return result
if not self.sm.allow_event_without_transition:
tna = TransitionNotAllowed(external_event.event, self.sm.configuration)
self._reject_future(event_future, tna)
self._reject_pending_futures(tna)
raise tna
# Event allowed but no transition — resolve with None
self._resolve_future(event_future, None)
return self._sentinel
except Exception as exc:
self._reject_future(event_future, exc)
self._reject_pending_futures(exc)
self.clear()
raise

async def enabled_events(self, *args, **kwargs):
sm = self.sm
enabled = {}
Expand Down
Loading