drop and restore idle workflows in workflow server (#266)

This commit is contained in:
Adrian Lyjak
2026-01-12 16:52:22 -05:00
committed by GitHub
parent f96faa2d03
commit 2ff316d0c9
24 changed files with 1622 additions and 168 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"llama-index-workflows": minor
---
Updates workflow server with functionality to drop and restore idle workflow handlers that are waiting on external input.
+1 -1
View File
@@ -29,7 +29,7 @@ repos:
name: ty
language: system
entry: uv run ty check packages
exclude: ^examples/|^docs/
pass_filenames: false
- repo: https://github.com/pre-commit/mirrors-mypy
rev: v1.15.0
hooks:
@@ -366,6 +366,22 @@ class WorkflowIdleEvent(InternalDispatchEvent):
pass
class UnhandledEvent(InternalDispatchEvent):
"""Emitted when an incoming event is not handled by any step or waiter.
This helps callers understand when an external event is ignored and whether
the workflow is idle after processing the event.
"""
event_type: str = Field(description="Class name of the unhandled event.")
qualified_name: str = Field(description="Fully qualified name of the event type.")
step_name: str | None = Field(
default=None,
description="Target step name if the event was addressed to a step.",
)
idle: bool = Field(description="Whether the workflow is idle after processing.")
class StepState(Enum):
# is enqueued, but no capacity yet available to run
PREPARING = "preparing"
@@ -9,7 +9,13 @@ from typing import AsyncGenerator, Callable
from workflows.decorators import P, R
from workflows.events import Event, StopEvent
from workflows.runtime.types.plugin import Plugin, SnapshottableRuntime, WorkflowRuntime
from workflows.runtime.types.plugin import (
ControlLoopFunction,
Plugin,
RegisteredWorkflow,
SnapshottableRuntime,
WorkflowRuntime,
)
from workflows.runtime.types.step_function import StepWorkerFunction
from workflows.runtime.types.ticks import WorkflowTick
from workflows.workflow import Workflow
@@ -19,10 +25,10 @@ class BasicRuntime:
def register(
self,
workflow: Workflow,
workflow_function: Callable[P, R],
steps: dict[str, StepWorkerFunction[R]],
) -> None:
return
workflow_function: ControlLoopFunction,
steps: dict[str, StepWorkerFunction],
) -> None | RegisteredWorkflow:
return None
def new_runtime(self, run_id: str) -> WorkflowRuntime:
snapshottable: SnapshottableRuntime = AsyncioWorkflowRuntime(run_id)
@@ -34,6 +34,13 @@ class HandlersListResponse(BaseModel):
class HealthResponse(BaseModel):
status: Literal["healthy"]
loaded_workflows: int = Field(
description="Number of workflow handlers currently loaded in memory"
)
active_workflows: int = Field(
description="Number of workflow handlers that are active (not idle)"
)
idle_workflows: int = Field(description="Number of workflow handlers that are idle")
class WorkflowsListResponse(BaseModel):
@@ -23,6 +23,7 @@ from workflows.events import (
StepState,
StepStateChanged,
StopEvent,
UnhandledEvent,
WorkflowCancelledEvent,
WorkflowFailedEvent,
WorkflowIdleEvent,
@@ -537,7 +538,7 @@ def _process_step_result_tick(
new_waiter = StepWorkerWaiter(
waiter_id=result.waiter_id,
event=this_execution.event,
waiting_for_event=result.event_type,
waiting_for_event=result.event_type, # ty: ignore[invalid-argument-type] - ty choking here, with result.event_type resolved as "object"
requirements=result.requirements,
has_requirements=bool(len(result.requirements)),
resolved_event=None,
@@ -662,11 +663,13 @@ def _process_add_event_tick(
state = init.deepcopy()
# iterate through the steps, and add to steps work queue if it's accepted.
commands: list[WorkflowCommand] = []
handled = False
if isinstance(tick.event, StartEvent):
state.is_running = True
for step_name, step_config in state.config.steps.items():
is_accepted = type(tick.event) in step_config.accepted_events
if is_accepted and (tick.step_name is None or tick.step_name == step_name):
handled = True
subcommands = _add_or_enqueue_event(
EventAttempt(
event=tick.event,
@@ -690,6 +693,7 @@ def _process_add_event_tick(
for k, v in wait_condition.requirements.items()
)
if is_match:
handled = True
wait_condition.resolved_event = tick.event
subcommands = _add_or_enqueue_event(
EventAttempt(event=wait_condition.event),
@@ -698,6 +702,22 @@ def _process_add_event_tick(
now_seconds,
)
commands.extend(subcommands)
if not handled:
# InputRequiredEvent subclasses are intentionally designed to be handled
# externally by human consumers, not by workflow steps. Don't emit
# UnhandledEvent for these since they're working as intended.
if not isinstance(tick.event, InputRequiredEvent):
event_cls = type(tick.event)
commands.append(
CommandPublishEvent(
UnhandledEvent(
event_type=event_cls.__name__,
qualified_name=f"{event_cls.__module__}.{event_cls.__name__}",
step_name=tick.step_name,
idle=_check_idle_state(state),
)
)
)
return state, commands
@@ -25,6 +25,8 @@ class HandlerQuery:
workflow_name_in: List[str] | None = None
# Matches if the status flag matches
status_in: List[Status] | None = None
# True = only idle handlers, False = only non-idle handlers, None = all
is_idle: bool | None = None
class PersistentHandler(BaseModel):
@@ -0,0 +1,54 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2026 LlamaIndex Inc.
"""Keyed lock utility for per-key mutual exclusion with automatic cleanup."""
import asyncio
from contextlib import asynccontextmanager
from typing import AsyncIterator
class KeyedLock:
"""A collection of locks keyed by string, with automatic cleanup.
Locks are created on-demand and automatically removed when no longer
in use (no waiters or holders). Safe for concurrent asyncio coroutines.
Usage:
locks = KeyedLock()
async with locks("my-key"):
# critical section for "my-key"
pass
"""
def __init__(self) -> None:
self._main_lock = asyncio.Lock()
self._locks: dict[str, asyncio.Lock] = {}
self._refs: dict[str, int] = {}
@asynccontextmanager
async def __call__(self, key: str) -> AsyncIterator[None]:
"""Acquire a lock for the given key.
The lock is created if it doesn't exist and removed when the last
holder/waiter releases it.
"""
# Get or create lock and register interest.
# No await between these lines = atomic in asyncio.
async with self._main_lock:
if key not in self._locks:
self._locks[key] = asyncio.Lock()
self._refs[key] = 0
self._refs[key] += 1
try:
async with self._locks[key]:
yield
finally:
# Deregister and cleanup if last.
# No await between these lines = atomic in asyncio.
async with self._main_lock:
self._refs[key] -= 1
if self._refs[key] == 0:
del self._locks[key]
del self._refs[key]
@@ -27,6 +27,11 @@ def _matches_query(handler: PersistentHandler, query: HandlerQuery) -> bool:
if handler.status not in query.status_in:
return False
if query.is_idle is not None:
handler_is_idle = handler.idle_since is not None
if query.is_idle != handler_is_idle:
return False
return True
@@ -7,7 +7,7 @@ import json
import logging
from contextlib import asynccontextmanager
from dataclasses import dataclass
from datetime import datetime, timezone
from datetime import datetime, timedelta, timezone
from importlib.metadata import version
from pathlib import Path
from typing import Any, AsyncGenerator, Awaitable, Callable, cast
@@ -32,6 +32,7 @@ from workflows.events import (
StepState,
StepStateChanged,
StopEvent,
UnhandledEvent,
WorkflowIdleEvent,
)
from workflows.handler import WorkflowHandler
@@ -58,8 +59,8 @@ from workflows.server.abstract_workflow_store import (
PersistentHandler,
Status,
)
from workflows.server.keyed_lock import KeyedLock
from workflows.server.memory_workflow_store import MemoryWorkflowStore
from workflows.types import RunResultT
# Protocol models are used on the client side; server responds with plain dicts
from workflows.utils import _nanoid as nanoid
@@ -75,17 +76,20 @@ class WorkflowServer:
workflow_store: AbstractWorkflowStore | None = None,
# retry/backoff seconds for persisting the handler state in the store after failures. Configurable mainly for testing.
persistence_backoff: list[float] = [0.5, 3],
# Release idle workflows from memory after this timeout (None = disabled)
idle_release_timeout: timedelta | None = timedelta(seconds=10),
):
# Workflows that this server will run
self._workflows: dict[str, Workflow] = {}
# Additional event schemas supported. Used for serdes lookups. Should be used sparingly if ever.
self._additional_events: dict[str, list[type[Event]] | None] = {}
self._contexts: dict[str, Context] = {}
self._handlers: dict[str, _WorkflowHandler] = {}
self._results: dict[str, RunResultT] = {}
self._workflow_store = (
workflow_store if workflow_store is not None else MemoryWorkflowStore()
)
self._assets_path = Path(__file__).parent / "static"
self._persistence_backoff = list(persistence_backoff)
self._idle_release_timeout = idle_release_timeout
self._middleware = middleware or [
Middleware(
@@ -98,6 +102,9 @@ class WorkflowServer:
)
]
# Lock to prevent concurrent reloads of the same handler
self._reload_lock = KeyedLock()
self._routes = [
Route(
"/workflows",
@@ -192,10 +199,16 @@ class WorkflowServer:
self._additional_events[name] = additional_events
async def start(self) -> "WorkflowServer":
"""Resumes previously running workflows, if they were not complete at last shutdown"""
"""Resumes previously running workflows, if they were not complete at last shutdown.
Idle workflows are not resumed - they remain released and will be
loaded on-demand when events arrive for them.
"""
handlers = await self._workflow_store.query(
HandlerQuery(
status_in=["running"], workflow_name_in=list(self._workflows.keys())
status_in=["running"],
workflow_name_in=list(self._workflows.keys()),
is_idle=False,
)
)
for persistent in handlers:
@@ -252,7 +265,6 @@ class WorkflowServer:
*[self._close_handler(handler) for handler in list(self._handlers.values())]
)
self._handlers.clear()
self._results.clear()
async def serve(
self,
@@ -361,7 +373,7 @@ class WorkflowServer:
"""
---
summary: Health check
description: Returns the server health status.
description: Returns the server health status and workflow counts.
responses:
200:
description: Successful health check
@@ -373,9 +385,28 @@ class WorkflowServer:
status:
type: string
example: healthy
required: [status]
loaded_workflows:
type: integer
description: Number of workflow handlers currently loaded in memory
active_workflows:
type: integer
description: Number of workflow handlers that are active (not idle)
idle_workflows:
type: integer
description: Number of workflow handlers that are idle
required: [status, loaded_workflows, active_workflows, idle_workflows]
"""
return JSONResponse(HealthResponse(status="healthy").model_dump())
loaded = len(self._handlers)
idle = sum(1 for h in self._handlers.values() if h.idle_since is not None)
active = loaded - idle
return JSONResponse(
HealthResponse(
status="healthy",
loaded_workflows=loaded,
active_workflows=active,
idle_workflows=idle,
).model_dump()
)
async def _list_workflows(self, request: Request) -> JSONResponse:
"""
@@ -433,16 +464,15 @@ class WorkflowServer:
if name not in self._workflows:
raise HTTPException(status_code=404, detail=f"Workflow '{name}' not found")
events = self._workflows[name].events
additional_events = self._additional_events.get(name, [])
if additional_events:
events.extend(additional_events)
events = self._workflows[name].events + (
self._additional_events.get(name, []) or []
)
event_objs = []
for event in events:
event_objs.append(event.model_json_schema())
return JSONResponse(WorkflowEventsListResponse(events=event_objs).model_dump())
return JSONResponse(
WorkflowEventsListResponse(
events=[event.model_json_schema() for event in events]
).model_dump()
)
async def _run_workflow(self, request: Request) -> JSONResponse:
"""
@@ -929,14 +959,16 @@ class WorkflowServer:
handler = self._handlers.get(handler_id)
if handler is None:
persisted = await self._workflow_store.query(
HandlerQuery(handler_id_in=[handler_id])
)
if persisted:
status = persisted[0].status
if status in {"completed", "failed", "cancelled"}:
raise HTTPException(detail="Handler is completed", status_code=204)
raise HTTPException(detail="Handler not found", status_code=404)
# Try to reload from persistence (for released idle workflows)
handler, persisted = await self._try_reload_handler(handler_id)
if handler is None:
if persisted:
status = persisted.status
if status in {"completed", "failed", "cancelled"}:
raise HTTPException(
detail="Handler is completed", status_code=204
)
raise HTTPException(detail="Handler not found", status_code=404)
if handler.queue.empty() and handler.task is not None and handler.task.done():
# https://html.spec.whatwg.org/multipage/server-sent-events.html
# Clients will reconnect if the connection is closed; a client can
@@ -1127,16 +1159,27 @@ class WorkflowServer:
if wrapper is not None and is_status_completed(wrapper.status):
raise HTTPException(detail="Workflow already completed", status_code=409)
if wrapper is None:
handler_data = await self._load_handler(handler_id)
if is_status_completed(handler_data.status):
raise HTTPException(
detail="Workflow already completed", status_code=409
)
else:
# this branch is for cases where handler status is running but somehow not in memory
# Ideally, this should never happen. We probably need to revisit when we add pause/expire functionality.
logger.warning(f"Handler {handler_id} is running but not in memory.")
raise HTTPException(detail="Handler expired", status_code=409)
# Try to reload from persistence (for released idle workflows)
wrapper, persisted = await self._try_reload_handler(handler_id)
if wrapper is None:
# Check if it exists but is completed
if persisted and is_status_completed(persisted.status):
raise HTTPException(
detail="Workflow already completed", status_code=409
)
elif persisted is None:
raise HTTPException(detail="Handler not found", status_code=404)
else:
# Shouldn't really happen
raise HTTPException(
detail=f"Failed to resume incomplete handler with status {persisted.status}",
status_code=500,
)
# Immediately mark active to cancel the idle timer before it can fire.
# This prevents a race where the timer releases the handler before we
# finish processing the event.
wrapper.mark_active()
handler = wrapper.run_handler
@@ -1329,6 +1372,7 @@ class WorkflowServer:
handler_id: str,
start_event: StartEvent | None = None,
context: Context | None = None,
idle_since: datetime | None = None,
) -> _WorkflowHandler:
"""Start a workflow and return a wrapper for the handler."""
with instrument_tags({"handler_id": handler_id}):
@@ -1337,12 +1381,16 @@ class WorkflowServer:
start_event=start_event,
)
wrapper = await self._run_workflow_handler(
handler_id, workflow.name, handler
handler_id, workflow.name, handler, idle_since=idle_since
)
return wrapper
async def _run_workflow_handler(
self, handler_id: str, workflow_name: str, handler: WorkflowHandler
self,
handler_id: str,
workflow_name: str,
handler: WorkflowHandler,
idle_since: datetime | None = None,
) -> _WorkflowHandler:
"""
Creates a wrapper for the handler and starts streaming events.
@@ -1362,7 +1410,10 @@ class WorkflowServer:
completed_at=None,
_workflow_store=self._workflow_store,
_persistence_backoff=self._persistence_backoff,
_idle_release_timeout=self._idle_release_timeout,
_on_idle_release=self._release_handler,
)
wrapper.idle_since = idle_since
# Initial checkpoint before registration; fail fast if persistence is unavailable
await wrapper.checkpoint()
# Now register and start streaming
@@ -1370,7 +1421,6 @@ class WorkflowServer:
async def on_finish() -> None:
self._handlers.pop(handler_id, None)
self._results.pop(handler_id, None)
wrapper.start_streaming(on_finish=on_finish)
@@ -1379,21 +1429,107 @@ class WorkflowServer:
async def _close_handler(self, handler: _WorkflowHandler) -> None:
"""Close and cleanup a handler."""
# Cancel the run_handler if not done
if not handler.run_handler.done():
try:
handler.run_handler.cancel()
except Exception:
pass
try:
await handler.run_handler.cancel_run()
except Exception:
pass
if handler.task is not None:
await handler.task
await handler.cancel_handlers_and_tasks()
self._handlers.pop(handler.handler_id, None)
self._results.pop(handler.handler_id, None)
async def _release_handler(self, wrapper: _WorkflowHandler) -> None:
"""Release an idle handler from memory, keeping it in persistence."""
handler_id = wrapper.handler_id
# Use lock to coordinate with _try_reload_handler and prevent race
# where release checkpoint overwrites newer state from reload
async with self._reload_lock(handler_id):
# Check if this handler was already replaced by a reload.
# If a different handler instance is now in _handlers, skip
# checkpointing to avoid overwriting newer state.
current = self._handlers.get(handler_id)
if current is not None and current is not wrapper:
logger.debug(
f"Skipping release checkpoint for {handler_id}: "
"handler was already reloaded"
)
# Mark to skip any further checkpoints (including auto-checkpoint
# in _stream_events when we cancel the task)
# Still need to cancel the old runtime
wrapper._cancel_idle_release_timer(skip_checkpoint=True)
await wrapper.cancel_handlers_and_tasks()
return
# Remove from memory FIRST to prevent race conditions where cancelling
# the handler triggers _stream_events to cancel the timer, which would
# interrupt this method before the pop happens.
self._handlers.pop(handler_id, None)
# Cancel the idle release timer to prevent re-entry
wrapper._cancel_idle_release_timer()
try:
# Final checkpoint to ensure latest state is persisted
await wrapper.checkpoint()
finally:
# Always stop the workflow runtime to prevent memory leaks,
# even if checkpoint fails
await wrapper.cancel_handlers_and_tasks()
logger.info(f"Released idle workflow {handler_id} from memory")
async def _try_reload_handler(
self, handler_id: str
) -> tuple[_WorkflowHandler | None, PersistentHandler | None]:
"""Attempt to reload a released handler from persistence.
Uses per-handler locking to prevent concurrent reloads from creating
duplicate workflow instances.
Returns the persistent handler data if it was fetched from the store, so callers can avoid fetching it again if needed.
"""
async with self._reload_lock(handler_id):
# Check if handler was already reloaded by another concurrent request
if handler_id in self._handlers:
return self._handlers[handler_id], None
found = await self._workflow_store.query(
HandlerQuery(handler_id_in=[handler_id])
)
if not found:
return None, None
handler_data = found[0]
if handler_data.status != "running":
return None, handler_data # Can't reload completed/failed/cancelled
workflow = self._workflows.get(handler_data.workflow_name)
if workflow is None:
logger.warning(
f"Cannot reload {handler_id}: workflow {handler_data.workflow_name} not registered"
)
return None, handler_data
try:
context = Context.from_dict(workflow=workflow, data=handler_data.ctx)
wrapper = await self._start_workflow(
workflow=_NamedWorkflow(
name=handler_data.workflow_name, workflow=workflow
),
handler_id=handler_id,
context=context,
idle_since=handler_data.idle_since,
)
# If workflow was idle when released, start the release timer
# (it will be cancelled if an event wakes the workflow)
if wrapper.idle_since is not None:
wrapper._start_idle_release_timer()
logger.info(f"Reloaded workflow {handler_id} from persistence")
return wrapper, handler_data
except Exception as e:
logger.error(f"Failed to reload handler {handler_id}: {e}")
raise HTTPException(
detail=f"Failed to reload handler: {e}", status_code=500
)
def _event_registry(self, workflow_name: str) -> dict[str, type[Event]]:
items = {e.__name__: e for e in self._workflows[workflow_name].events}
@@ -1429,6 +1565,12 @@ class _WorkflowHandler:
_on_finish: Callable[[], Awaitable[None]] | None = None
idle_since: datetime | None = None
# Idle release support
_idle_release_timeout: timedelta | None = None
_on_idle_release: Callable[["_WorkflowHandler"], Awaitable[None]] | None = None
_idle_release_timer: asyncio.Task[None] | None = None
_skip_checkpoint: bool = False # Set to prevent checkpointing stale handlers
def _as_persistent(self) -> PersistentHandler:
"""Persist the current handler state immediately to the workflow store."""
self.updated_at = datetime.now(timezone.utc)
@@ -1455,6 +1597,9 @@ class _WorkflowHandler:
async def checkpoint(self) -> None:
"""Persist with retry/backoff; cancel handler when retries exhausted."""
if self._skip_checkpoint:
logger.debug(f"Skipping checkpoint for handler {self.handler_id}")
return
backoffs = list(self._persistence_backoff)
try:
persistent = self._as_persistent()
@@ -1561,6 +1706,53 @@ class _WorkflowHandler:
except Exception:
return None
def _start_idle_release_timer(self) -> None:
"""Start a timer to release this handler after the idle timeout."""
if self._idle_release_timeout is None or self._on_idle_release is None:
return
# Cancel any existing timer first
self._cancel_idle_release_timer()
timeout_seconds = self._idle_release_timeout.total_seconds()
async def release_after_timeout() -> None:
try:
await asyncio.sleep(timeout_seconds)
# Only release if still idle and no active stream consumers
if self.idle_since is not None:
if self.consumer_mutex.locked():
# Mutex is locked - reschedule to try again later
self._start_idle_release_timer()
return
if self._on_idle_release is not None:
await self._on_idle_release(self)
except asyncio.CancelledError:
pass # Timer was cancelled, nothing to do
self._idle_release_timer = asyncio.create_task(release_after_timeout())
def _cancel_idle_release_timer(self, skip_checkpoint: bool = False) -> None:
"""Cancel any pending idle release timer."""
if skip_checkpoint:
self._skip_checkpoint = True
if self._idle_release_timer is not None:
self._idle_release_timer.cancel()
self._idle_release_timer = None
def mark_idle(self, idle_since: datetime | None = None) -> None:
self.idle_since = idle_since or datetime.now(timezone.utc)
self._start_idle_release_timer()
def mark_active(self) -> None:
"""Mark this handler as active (not idle).
Call this when an event is being sent to prevent premature release.
"""
if self.idle_since is not None:
self.idle_since = None
self._cancel_idle_release_timer()
def start_streaming(self, on_finish: Callable[[], Awaitable[None]]) -> None:
"""Start streaming events from the handler and managing state."""
self.task = asyncio.create_task(self._stream_events(on_finish=on_finish))
@@ -1571,14 +1763,16 @@ 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
# Track idle state transitions and manage release timer
if isinstance(event, WorkflowIdleEvent):
self.idle_since = datetime.now(timezone.utc)
self.mark_idle()
elif isinstance(event, UnhandledEvent):
self.mark_idle()
elif (
isinstance(event, StepStateChanged)
and event.step_state == StepState.RUNNING
):
self.idle_since = None # Resumed from idle
self.mark_active()
if ( # Watch for a specific internal event that signals the step is complete
isinstance(event, StepStateChanged)
@@ -1595,6 +1789,10 @@ class _WorkflowHandler:
await self.checkpoint()
self.queue.put_nowait(event)
# Workflow is completing - cancel any pending release timer
self._cancel_idle_release_timer()
# done when stream events are complete
try:
await self.run_handler
@@ -1659,6 +1857,24 @@ class _WorkflowHandler:
await self._on_finish()
self.consumer_mutex.release()
async def cancel_handlers_and_tasks(self) -> None:
"""Cancel the handler and release it from the store."""
if not self.run_handler.done():
try:
self.run_handler.cancel()
except Exception:
pass
try:
await self.run_handler.cancel_run()
except Exception:
pass
if self.task and not self.task.done():
self.task.cancel()
try:
await self.task
except asyncio.CancelledError:
pass
class NoLockAvailable(Exception):
"""Raised when no lock is available to acquire after a timeout"""
@@ -0,0 +1,4 @@
PRAGMA user_version=3;
-- Add idle_since column for tracking when a workflow became idle
ALTER TABLE handlers ADD COLUMN idle_since TEXT;
@@ -32,7 +32,7 @@ class SqliteWorkflowStore(AbstractWorkflowStore):
clauses, params = filter_spec
sql = """SELECT handler_id, workflow_name, status, run_id, error, result,
started_at, updated_at, completed_at, ctx FROM handlers"""
started_at, updated_at, completed_at, idle_since, ctx FROM handlers"""
if clauses:
sql = f"{sql} WHERE {' AND '.join(clauses)}"
conn = sqlite3.connect(self.db_path)
@@ -52,8 +52,8 @@ class SqliteWorkflowStore(AbstractWorkflowStore):
cursor.execute(
"""
INSERT INTO handlers (handler_id, workflow_name, status, run_id, error, result,
started_at, updated_at, completed_at, ctx)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
started_at, updated_at, completed_at, idle_since, ctx)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(handler_id) DO UPDATE SET
workflow_name = excluded.workflow_name,
status = excluded.status,
@@ -63,6 +63,7 @@ class SqliteWorkflowStore(AbstractWorkflowStore):
started_at = excluded.started_at,
updated_at = excluded.updated_at,
completed_at = excluded.completed_at,
idle_since = excluded.idle_since,
ctx = excluded.ctx
""",
(
@@ -77,6 +78,7 @@ class SqliteWorkflowStore(AbstractWorkflowStore):
handler.started_at.isoformat() if handler.started_at else None,
handler.updated_at.isoformat() if handler.updated_at else None,
handler.completed_at.isoformat() if handler.completed_at else None,
handler.idle_since.isoformat() if handler.idle_since else None,
json.dumps(handler.ctx),
),
)
@@ -130,6 +132,12 @@ class SqliteWorkflowStore(AbstractWorkflowStore):
return None
add_in_clause("status", query.status_in)
if query.is_idle is not None:
if query.is_idle:
clauses.append("idle_since IS NOT NULL")
else:
clauses.append("idle_since IS NULL")
if not clauses:
return clauses, params
@@ -147,5 +155,6 @@ def _row_to_persistent_handler(row: tuple) -> PersistentHandler:
started_at=datetime.fromisoformat(row[6]) if row[6] else None,
updated_at=datetime.fromisoformat(row[7]) if row[7] else None,
completed_at=datetime.fromisoformat(row[8]) if row[8] else None,
ctx=json.loads(row[9]),
idle_since=datetime.fromisoformat(row[9]) if row[9] else None,
ctx=json.loads(row[10]),
)
@@ -104,5 +104,5 @@ async def test_plugin_with_time_machine() -> AsyncGenerator[
tuple[MockRuntimePlugin, time_machine.Coordinates], None
]:
"""Plugin with time-machine at epoch 1000.0, tick=True."""
with time_machine.travel(1000.0, tick=True) as traveller:
with time_machine.travel("2026-01-07T12:27:00.000-08:00", tick=True) as traveller:
yield MockRuntimePlugin(run_id="test", traveller=traveller), traveller
@@ -10,7 +10,7 @@ testing them in isolation without running the full async control loop.
from __future__ import annotations
from typing import Any
from typing import Any, cast
import pytest
from workflows.decorators import StepConfig
@@ -22,6 +22,7 @@ from workflows.events import (
StepState,
StepStateChanged,
StopEvent,
UnhandledEvent,
WorkflowIdleEvent,
)
from workflows.retry_policy import ConstantDelayRetryPolicy
@@ -56,6 +57,7 @@ from workflows.runtime.types.results import (
AddWaiter,
DeleteCollectedEvent,
DeleteWaiter,
StepFunctionResult,
StepWorkerFailed,
StepWorkerResult,
StepWorkerState,
@@ -131,6 +133,75 @@ def add_worker(state: BrokerState, event: Event, worker_id: int = 0) -> None:
)
def test_add_event_unhandled_emits_internal_event(base_state: BrokerState) -> None:
"""Unhandled events should emit UnhandledEvent with idle status."""
tick = TickAddEvent(event=OtherEvent(data="unused"), step_name=None)
state, commands = _process_add_event_tick(tick, base_state, now_seconds=0.0)
publish_events = [c.event for c in commands if isinstance(c, CommandPublishEvent)]
unhandled = [e for e in publish_events if isinstance(e, UnhandledEvent)]
assert len(unhandled) == 1
assert unhandled[0].event_type == "OtherEvent"
assert unhandled[0].qualified_name.endswith(".OtherEvent")
assert unhandled[0].step_name is None
assert unhandled[0].idle == _check_idle_state(state)
class CustomInputRequired(InputRequiredEvent):
"""Custom InputRequiredEvent subclass for testing."""
prompt: str
def test_add_event_input_required_does_not_emit_unhandled(
base_state: BrokerState,
) -> None:
"""InputRequiredEvent subclasses should NOT emit UnhandledEvent.
InputRequiredEvent events are designed to be handled externally by human
consumers, not by workflow steps. They should not trigger UnhandledEvent.
"""
tick = TickAddEvent(event=CustomInputRequired(prompt="test"), step_name=None)
_, commands = _process_add_event_tick(tick, base_state, now_seconds=0.0)
publish_events = [c.event for c in commands if isinstance(c, CommandPublishEvent)]
unhandled = [e for e in publish_events if isinstance(e, UnhandledEvent)]
assert len(unhandled) == 0
def test_add_event_base_input_required_does_not_emit_unhandled(
base_state: BrokerState,
) -> None:
"""Base InputRequiredEvent should also NOT emit UnhandledEvent."""
tick = TickAddEvent(event=InputRequiredEvent(), step_name=None)
_, commands = _process_add_event_tick(tick, base_state, now_seconds=0.0)
publish_events = [c.event for c in commands if isinstance(c, CommandPublishEvent)]
unhandled = [e for e in publish_events if isinstance(e, UnhandledEvent)]
assert len(unhandled) == 0
def test_add_event_matches_waiter_does_not_emit_unhandled(
base_state: BrokerState,
) -> None:
"""Events that satisfy a waiter should not emit UnhandledEvent."""
base_state.workers["test_step"].collected_waiters.append(
StepWorkerWaiter(
waiter_id="waiter-1",
event=StartEvent(),
waiting_for_event=OtherEvent,
requirements={},
has_requirements=False,
resolved_event=None,
)
)
tick = TickAddEvent(event=OtherEvent(data="hit"), step_name=None)
_, commands = _process_add_event_tick(tick, base_state, now_seconds=0.0)
publish_events = [c.event for c in commands if isinstance(c, CommandPublishEvent)]
assert not any(isinstance(e, UnhandledEvent) for e in publish_events)
@pytest.mark.parametrize(
"result,expected_commands",
[
@@ -253,19 +324,20 @@ def test_waiters(base_state: BrokerState) -> None:
event = MyTestEvent(value=42)
add_worker(base_state, event)
result = AddWaiter(
waiter_id="w1",
waiter_event=InputRequiredEvent(),
requirements={},
timeout=None,
event_type=OtherEvent,
)
# Add waiter
tick: TickStepResult[Any] = TickStepResult(
step_name="test_step",
worker_id=0,
event=event,
result=[
AddWaiter(
waiter_id="w1",
waiter_event=InputRequiredEvent(),
requirements={},
timeout=None,
event_type=OtherEvent,
)
cast(StepFunctionResult[Any], result),
],
)
new_state, _ = _process_step_result_tick(tick, base_state, now_seconds=110.0)
@@ -679,19 +751,19 @@ def test_idle_event_emitted_on_transition_to_idle(base_state: BrokerState) -> No
base_state.workers["test_step"].collected_waiters.append(waiter)
# Process result that completes the worker but leaves waiter active
tick: TickStepResult[Any] = TickStepResult(
result = AddWaiter(
waiter_id="w1",
waiter_event=None,
requirements={},
timeout=None,
event_type=OtherEvent,
)
tick: TickStepResult[Any] = TickStepResult[Any](
step_name="test_step",
worker_id=0,
event=event,
result=[
AddWaiter(
waiter_id="w1",
waiter_event=None,
requirements={},
timeout=None,
event_type=OtherEvent,
)
],
result=[cast(StepFunctionResult[Any], result)],
)
new_state, commands = _process_step_result_tick(tick, base_state, now_seconds=110.0)
@@ -0,0 +1,780 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 LlamaIndex Inc.
"""Tests for idle workflow release and reload functionality."""
import asyncio
from datetime import datetime, timedelta, timezone
from typing import Any, Optional
import pytest
import time_machine
from httpx import ASGITransport, AsyncClient
from workflows import Context, Workflow, step
from workflows.events import HumanResponseEvent, StartEvent, StopEvent
from workflows.server.abstract_workflow_store import (
HandlerQuery,
PersistentHandler,
Status,
)
from workflows.server.memory_workflow_store import MemoryWorkflowStore
from workflows.server.server import WorkflowServer, _WorkflowHandler
from tests.server.util import async_yield, wait_for_passing
class WaitableExternalEvent(HumanResponseEvent):
"""Event sent from external sources."""
response: str
class WaitingWorkflow(Workflow):
"""Workflow that uses ctx.wait_for_event() to properly become idle."""
@step
async def start_and_wait(self, ctx: Context, ev: StartEvent) -> StopEvent:
# Use ctx.wait_for_event() to create a waiter - this makes the workflow idle
external = await ctx.wait_for_event(WaitableExternalEvent)
return StopEvent(result=f"received: {external.response}")
def get_handler_in_memory(server: WorkflowServer, handler_id: str) -> _WorkflowHandler:
wrapper = server._handlers.get(handler_id)
assert wrapper is not None, f"Handler {handler_id} not found in memory"
return wrapper
def assert_handler_in_memory(server: WorkflowServer, handler_id: str) -> None:
get_handler_in_memory(server, handler_id)
def assert_handler_not_in_memory(server: WorkflowServer, handler_id: str) -> None:
assert handler_id not in server._handlers, f"Handler {handler_id} still in memory"
def make_server(
memory_store: MemoryWorkflowStore,
waiting_workflow: Workflow,
idle_release_timeout: Optional[timedelta],
*,
persistence_backoff: Optional[list[float]] = None,
) -> WorkflowServer:
if persistence_backoff is None:
server = WorkflowServer(
workflow_store=memory_store,
idle_release_timeout=idle_release_timeout,
)
else:
server = WorkflowServer(
workflow_store=memory_store,
idle_release_timeout=idle_release_timeout,
persistence_backoff=persistence_backoff,
)
server.add_workflow(
"test", waiting_workflow, additional_events=[WaitableExternalEvent]
)
return server
async def start_waiting_handler(
server: WorkflowServer, handler_id: str
) -> _WorkflowHandler:
handler = server._workflows["test"].run()
await server._run_workflow_handler(handler_id, "test", handler)
await async_yield(20)
wrapper = get_handler_in_memory(server, handler_id)
assert wrapper.idle_since is not None
return wrapper
async def advance_time(traveller: Any, delta: timedelta, iterations: int = 10) -> None:
traveller.shift(delta)
await async_yield(iterations)
async def seed_persistent_handler(
store: MemoryWorkflowStore,
handler_id: str,
*,
idle_since: Optional[datetime],
ctx: dict[str, object],
status: Status = "running",
) -> None:
await store.update(
PersistentHandler(
handler_id=handler_id,
workflow_name="test",
status=status,
idle_since=idle_since,
ctx=ctx,
)
)
@pytest.fixture
def memory_store() -> MemoryWorkflowStore:
return MemoryWorkflowStore()
@pytest.fixture
def waiting_workflow() -> Workflow:
return WaitingWorkflow()
@pytest.mark.asyncio
async def test_is_idle_query_filter_memory_store() -> None:
"""Test that is_idle filter works in MemoryWorkflowStore."""
store = MemoryWorkflowStore()
now = datetime.now(timezone.utc)
# Handler that is idle (has idle_since set)
await seed_persistent_handler(
store,
"idle-1",
idle_since=now - timedelta(minutes=5),
ctx={},
)
# Handler that is not idle (no idle_since)
await seed_persistent_handler(store, "active-1", idle_since=None, ctx={})
# Another idle handler
await seed_persistent_handler(
store,
"idle-2",
idle_since=now - timedelta(seconds=10),
ctx={},
)
# Query for idle handlers
idle_results = await store.query(HandlerQuery(is_idle=True))
assert len(idle_results) == 2
assert {r.handler_id for r in idle_results} == {"idle-1", "idle-2"}
# Query for non-idle handlers
active_results = await store.query(HandlerQuery(is_idle=False))
assert len(active_results) == 1
assert active_results[0].handler_id == "active-1"
# Query without filter returns all
all_results = await store.query(HandlerQuery())
assert len(all_results) == 3
@pytest.mark.asyncio
async def test_workflow_becomes_idle_and_is_released(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Test that a workflow becomes idle after a step and can be released."""
idle_timeout = timedelta(milliseconds=50)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False) as traveller:
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
# Start a workflow
handler_id = "release-test-1"
await start_waiting_handler(server, handler_id)
# Advance time past the idle timeout to trigger the timer
await advance_time(traveller, timedelta(milliseconds=100))
# The handler should be released from memory
assert_handler_not_in_memory(server, handler_id)
# But should still exist in the store with status "running"
persisted = await memory_store.query(
HandlerQuery(handler_id_in=[handler_id])
)
assert len(persisted) == 1
assert persisted[0].status == "running"
assert persisted[0].idle_since is not None
@pytest.mark.asyncio
async def test_released_workflow_is_reloaded_on_event(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Test that a released workflow is reloaded when an event is sent."""
idle_timeout = timedelta(milliseconds=50)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False) as traveller:
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
# Start a workflow
handler_id = "reload-test-1"
await start_waiting_handler(server, handler_id)
# Advance time past the idle timeout to trigger the timer
await advance_time(traveller, timedelta(milliseconds=100))
# Handler should be released
assert_handler_not_in_memory(server, handler_id)
# Now reload by using _try_reload_handler
reloaded, persisted = await server._try_reload_handler(handler_id)
assert reloaded is not None
assert_handler_in_memory(server, handler_id)
assert persisted is not None
assert persisted.status == "running"
# Send the event to complete the workflow
ctx = reloaded.run_handler.ctx
assert ctx is not None
ctx.send_event(WaitableExternalEvent(response="hello"))
result = await reloaded.run_handler
assert result == "received: hello"
@pytest.mark.asyncio
async def test_idle_release_restores_idle_since_on_reload(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Test that idle_since is preserved when reloading via _try_reload_handler."""
server = make_server(
memory_store,
waiting_workflow,
timedelta(minutes=5), # Long timeout
)
# Seed a persisted handler with idle_since
idle_time = datetime.now(timezone.utc) - timedelta(minutes=2)
ctx = Context(waiting_workflow).to_dict()
await seed_persistent_handler(
memory_store,
"idle-restore-1",
idle_since=idle_time,
ctx=ctx,
)
async with server.contextmanager():
# Idle handler should NOT be in memory on startup
assert_handler_not_in_memory(server, "idle-restore-1")
# Reload it (simulating an event arriving)
wrapper, persisted = await server._try_reload_handler("idle-restore-1")
assert wrapper is not None
assert wrapper.idle_since == idle_time
assert persisted is not None
assert persisted.status == "running"
@pytest.mark.asyncio
async def test_reloaded_idle_workflow_is_released_again(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Test that a reloaded workflow that stays idle gets released again."""
# Use very short timeout - this test needs real time for timers
idle_timeout = timedelta(milliseconds=1)
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
# Start a workflow
handler_id = "reload-idle-test-1"
handler = server._workflows["test"].run()
wrapper = await server._run_workflow_handler(handler_id, "test", handler)
async def wrapper_is_idle() -> None:
assert wrapper.idle_since is not None
await wait_for_passing(wrapper_is_idle, interval=0.01, max_duration=1.5)
# Reload the workflow (simulating an event arriving) once the handler
# is released from memory.
async def reload_from_store() -> tuple[_WorkflowHandler, PersistentHandler]:
reloaded, persisted = await server._try_reload_handler(handler_id)
assert reloaded is not None
assert persisted is not None
assert reloaded is not wrapper
return reloaded, persisted
reloaded, persisted = await wait_for_passing(
reload_from_store, interval=0.01, max_duration=1.5
)
assert persisted is not None
assert persisted.status == "running"
# The reloaded handler should have idle_since restored
assert reloaded.idle_since is not None
# Wait for the reloaded handler to be released again by observing
# that a subsequent reload returns a new handler instance.
async def reload_after_release() -> _WorkflowHandler:
reloaded_again, persisted_again = await server._try_reload_handler(
handler_id
)
assert reloaded_again is not None
assert persisted_again is not None
assert reloaded_again is not reloaded
return reloaded_again
reloaded_again = await wait_for_passing(
reload_after_release, interval=0.01, max_duration=1.5
)
await server._close_handler(reloaded_again)
@pytest.mark.asyncio
async def test_idle_handlers_not_resumed_on_server_start(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Test that idle handlers are not loaded into memory on server start."""
# Seed the store with an idle handler
idle_time = datetime.now(timezone.utc) - timedelta(minutes=2)
ctx = Context(waiting_workflow).to_dict()
await seed_persistent_handler(
memory_store,
"idle-on-start-1",
idle_since=idle_time,
ctx=ctx,
)
# Also seed an active (non-idle) handler
await seed_persistent_handler(
memory_store,
"active-on-start-1",
idle_since=None, # Not idle
ctx=ctx,
)
server = make_server(
memory_store,
waiting_workflow,
timedelta(minutes=5),
)
async with server.contextmanager():
# The idle handler should NOT be in memory
assert_handler_not_in_memory(server, "idle-on-start-1")
# The active handler SHOULD be in memory
assert_handler_in_memory(server, "active-on-start-1")
# The idle handler should still exist in the store
persisted = await memory_store.query(
HandlerQuery(handler_id_in=["idle-on-start-1"])
)
assert len(persisted) == 1
assert persisted[0].status == "running"
@pytest.mark.asyncio
async def test_idle_release_disabled_when_timeout_none(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Test that no timer is started when idle_release_timeout is None."""
server = make_server(
memory_store,
waiting_workflow,
None, # Disabled
)
async with server.contextmanager():
# Start a workflow
handler_id = "no-release-test"
wrapper = await start_waiting_handler(server, handler_id)
# Timer should not be set since idle_release_timeout is None
assert wrapper._idle_release_timer is None
@pytest.mark.asyncio
async def test_idle_release_cancels_runtime(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Idle release should stop the workflow runtime.
When a workflow is released from memory, the underlying WorkflowHandler
and its context/broker should be cancelled. Currently, _release_handler
only cancels the server's stream task and removes from _handlers, but
the actual workflow runtime keeps running in the background.
This test captures a reference to the run_handler before release, then
verifies it should be done/cancelled after release.
"""
idle_timeout = timedelta(milliseconds=50)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False) as traveller:
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
# Start a workflow
handler_id = "runtime-leak-test"
wrapper = await start_waiting_handler(server, handler_id)
# Capture reference to the run_handler before release
run_handler = wrapper.run_handler
# Advance time past the idle timeout to trigger the timer
await advance_time(traveller, timedelta(milliseconds=100))
# Handler should be released
assert_handler_not_in_memory(server, handler_id)
assert run_handler.done(), (
"run_handler is still running after release workflow runtime was not stopped, causing memory leak"
)
@pytest.mark.asyncio
async def test_idle_release_waits_for_stream_consumer_then_releases(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Idle release should wait for active consumers and then reschedule.
The idle release timer checks if consumer_mutex is locked and skips release
if so. However, it does NOT reschedule another timer attempt. This means
if a client briefly holds a stream during the timeout window and then
disconnects, the handler will never be released.
This test holds the mutex during the timer, releases it, then verifies
the handler should eventually be released (but won't be due to the issue).
"""
idle_timeout = timedelta(milliseconds=50)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False) as traveller:
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
# Start workflow
handler_id = "mutex-reschedule-test"
wrapper = await start_waiting_handler(server, handler_id)
# Grab the mutex before timer fires, then hold it past the timeout
async with wrapper.consumer_mutex:
# Advance time past the timeout to trigger the timer
await advance_time(traveller, timedelta(milliseconds=100))
# Handler should still be present (timer was blocked by mutex)
assert_handler_in_memory(server, handler_id)
# Now mutex is released - handler should eventually be released
# because the timer reschedules when mutex was locked.
# Advance time to let the rescheduled timer fire.
await advance_time(traveller, timedelta(milliseconds=100))
# Handler should be released because timer was rescheduled
assert_handler_not_in_memory(server, handler_id)
@pytest.mark.asyncio
async def test_mark_active_prevents_release_during_event_post(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Marking active should protect from release during event post.
When an event is posted to wake a workflow, the idle_since is only cleared
when StepStateChanged(RUNNING) is observed in _stream_events. However,
the idle timer can fire in the window between event posting and the
running event being processed, causing the handler to be released while
the workflow is actively processing.
DESIRED IMPLEMENTATION:
We want to support short/immediate idle timeouts (e.g., 0 timeout to release
workflows immediately when idle). However, "revived" workflows (ones that
just received an event but haven't processed it yet) need protection. The
fix should use a minimum grace period for revived workflows - when an event
is sent to an idle workflow, the timer should be rescheduled with at least
a minimum timeout (e.g., 1 second) regardless of the configured timeout.
This prevents immediate release while still allowing fast release of truly
idle workflows.
The fix implemented: The server's _post_event endpoint calls mark_active()
before sending an event, which clears idle_since and cancels the timer.
This test verifies that mark_active() properly protects the handler.
"""
idle_timeout = timedelta(milliseconds=50)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False) as traveller:
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
handler_id = "race-test"
wrapper = await start_waiting_handler(server, handler_id)
# Mark active before sending event (this is what _post_event does)
# This clears idle_since and cancels the timer
wrapper.mark_active()
# Send event to wake the workflow
ctx = wrapper.run_handler.ctx
assert ctx is not None
ctx.send_event(WaitableExternalEvent(response="wake-up"))
# Advance time past the original idle timeout
# The workflow should be protected from release because mark_active was called
await advance_time(traveller, timedelta(milliseconds=100))
# Handler should NOT be released because mark_active()
# cleared idle_since and cancelled the timer before the event was sent
assert_handler_in_memory(server, handler_id)
@pytest.mark.asyncio
async def test_try_reload_is_singleton_under_concurrency(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Concurrent reload requests should only create one handler.
_try_reload_handler has no locking mechanism, so two parallel requests
can both reload and start the same workflow. This creates split-brain
scenarios where multiple workflow instances process events independently.
This test fires multiple reload requests in parallel and checks that
only one succeeds in creating a handler.
"""
server = make_server(
memory_store,
waiting_workflow,
timedelta(minutes=5),
)
# Seed a persisted idle handler
idle_time = datetime.now(timezone.utc) - timedelta(minutes=2)
ctx = Context(waiting_workflow).to_dict()
handler_id = "concurrent-reload-test"
await seed_persistent_handler(
memory_store,
handler_id,
idle_since=idle_time,
ctx=ctx,
)
async with server.contextmanager():
assert_handler_not_in_memory(server, handler_id)
# Track how many workflow instances were created
instances_created: list[object] = []
async def reload_and_track() -> None:
wrapper, _ = await server._try_reload_handler(handler_id)
if wrapper is not None:
# Track the actual run_handler object identity
instances_created.append(id(wrapper.run_handler))
# Fire 5 concurrent reload requests
await asyncio.gather(*[reload_and_track() for _ in range(5)])
# We should have exactly one handler in memory
assert_handler_in_memory(server, handler_id)
# Multiple unique workflow instances may have been created
# (even though only one ends up in _handlers, the others are leaked)
unique_instances = set(instances_created)
assert len(unique_instances) == 1, (
f"{len(unique_instances)} different workflow instances were created "
f"by concurrent reloads. Only one should be created. "
f"The extra instances are leaked and may process events independently."
)
@pytest.mark.asyncio
async def test_release_handler_cancels_runtime_on_checkpoint_failure(
memory_store: MemoryWorkflowStore,
waiting_workflow: Workflow,
) -> None:
"""Release should cancel runtime even if checkpoint fails.
If checkpoint() raises after the handler is already removed from _handlers,
the runtime cancel path never runs, leaking a live workflow with no
in-memory reference.
The current _release_handler structure is:
1. Pop from _handlers (handler removed from memory)
2. Call checkpoint() (if this fails, exception propagates)
3. Cancel runtime (never reached if step 2 fails)
This test triggers a realistic failure by storing a non-serializable object
(a lambda) in the context state. When the idle release timer fires and
tries to checkpoint, serialization fails and the runtime leaks.
The fix should ensure that if ANY part of release fails, the workflow
runtime is still properly cancelled (e.g., use try/finally).
"""
idle_timeout = timedelta(milliseconds=100)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False) as traveller:
server = make_server(
memory_store,
waiting_workflow,
idle_timeout,
persistence_backoff=[], # No retries - fail immediately
)
await server.start()
run_handler = None
try:
handler_id = "checkpoint-fail-test"
wrapper = await start_waiting_handler(server, handler_id)
# Capture reference to verify runtime is stopped
run_handler = wrapper.run_handler
ctx = run_handler.ctx
assert ctx is not None
# Store a non-serializable object in context store
# This will cause checkpoint() to fail when it tries to serialize
await ctx.store.set("bad_data", lambda x: x) # lambdas can't be serialized
# Advance time past the idle timeout to trigger the timer
await advance_time(traveller, timedelta(milliseconds=200), iterations=20)
# Verify handler was removed from _handlers (release started)
assert_handler_not_in_memory(server, handler_id)
# The run_handler should be done/cancelled even if checkpoint failed
# but it's still running because the cancel code was never reached
assert run_handler.done(), (
"run_handler is still running after release attempt failed. "
"_release_handler pops from _handlers before checkpoint(), so if "
"checkpoint() raises (e.g., due to non-serializable state), the "
"cancel path is never reached and the workflow runtime leaks."
)
finally:
# Clean up the leaked runtime manually (since the issue prevents normal cleanup)
if run_handler is not None and not run_handler.done():
run_handler.cancel()
try:
await run_handler.cancel_run()
except Exception:
pass
await server.stop()
@pytest.mark.asyncio
async def test_send_event_clears_idle_state_before_processing(
memory_store: MemoryWorkflowStore, waiting_workflow: Workflow
) -> None:
"""Verify that send_event() clears idle state immediately to prevent race conditions.
mark_active() is called BEFORE processing the event to prevent a race where
the idle release timer fires while we're still handling the request. This means
idle_since is cleared even if send_event() subsequently fails.
"""
idle_timeout = timedelta(milliseconds=100)
with time_machine.travel("2026-01-07T12:00:00Z", tick=False):
server = make_server(memory_store, waiting_workflow, idle_timeout)
async with server.contextmanager():
transport = ASGITransport(app=server.app)
async with AsyncClient(
transport=transport, base_url="http://test"
) as client:
# Start a workflow via HTTP
response = await client.post("/workflows/test/run-nowait", json={})
assert response.status_code == 200
handler_id = response.json()["handler_id"]
# Wait for workflow to become idle (just needs event loop iterations)
await async_yield(20)
wrapper = get_handler_in_memory(server, handler_id)
assert wrapper is not None and wrapper.idle_since is not None
# Post an event with a bad step name via HTTP - this will fail
response = await client.post(
f"/events/{handler_id}",
json={
"event": {
"type": "WaitableExternalEvent",
"value": {"response": "test"},
},
"step": "nonexistent_step", # This step doesn't exist
},
)
# The endpoint returns 400 for bad step
assert response.status_code == 400
# mark_active() is called BEFORE send_event() to prevent race conditions.
# Even though send_event() failed, idle state is cleared.
assert wrapper.idle_since is None, (
"idle_since should be cleared before send_event() is attempted"
)
assert wrapper._idle_release_timer is None, "timer should be cancelled"
@pytest.mark.asyncio
async def test_release_skips_checkpoint_if_handler_was_reloaded(
memory_store: MemoryWorkflowStore,
waiting_workflow: Workflow,
) -> None:
"""Verify that _release_handler skips checkpoint if handler was reloaded.
When release and reload race, _release_handler should detect if a new
handler instance is now in _handlers and skip the checkpoint to avoid
overwriting newer state.
This test verifies:
1. Start handler, wait for idle
2. Reload the handler (simulating it was released earlier)
3. Update the reloaded handler's state
4. Call _release_handler with the OLD wrapper
5. Verify the old release doesn't overwrite the new state
"""
# time_machine with tick=False for fast test execution
with time_machine.travel("2026-01-07T12:27:00.000-08:00", tick=False):
server = WorkflowServer(
workflow_store=memory_store,
idle_release_timeout=None, # Disable auto-release for manual control
)
server.add_workflow(
"test", waiting_workflow, additional_events=[WaitableExternalEvent]
)
async with server.contextmanager():
# Start workflow and wait for it to become idle
handler_id = "race-test-handler"
await start_waiting_handler(server, handler_id)
old_wrapper = get_handler_in_memory(server, handler_id)
original_idle_since = old_wrapper.idle_since
assert original_idle_since is not None
# Checkpoint the old state (simulating what release would have done)
await old_wrapper.checkpoint()
# Simulate the scenario where handler was released and then reloaded:
# Remove old handler from memory
server._handlers.pop(handler_id, None)
# Reload the handler (this gets the persisted state)
reloaded, persisted = await server._try_reload_handler(handler_id)
assert reloaded is not None
assert reloaded is not old_wrapper # Different instance
assert persisted is not None
assert persisted.status == "running"
# The reloaded handler receives an event and becomes active
reloaded.mark_active()
assert reloaded.idle_since is None
# Checkpoint the new state
await reloaded.checkpoint()
# Verify store has the new state (idle_since=None)
stored = await memory_store.query(HandlerQuery(handler_id_in=[handler_id]))
assert stored[0].idle_since is None, "Store should have new state"
# NOW call _release_handler with the OLD wrapper
# This simulates a delayed release that happens after reload
# _release_handler should detect that a different handler is in _handlers
# and skip the checkpoint
await server._release_handler(old_wrapper)
# Verify the store still has the correct (new) state
stored = await memory_store.query(HandlerQuery(handler_id_in=[handler_id]))
assert len(stored) == 1
# The old release should have skipped the checkpoint
# because it detected a different handler instance in _handlers
assert stored[0].idle_since is None, (
f"Release should have skipped checkpoint since handler was reloaded. "
f"Expected idle_since=None (from reload), "
f"but got idle_since={stored[0].idle_since}."
)
@@ -0,0 +1,69 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2026 LlamaIndex Inc.
"""Black box tests for idle release behavior over live HTTP."""
from __future__ import annotations
from datetime import timedelta
import pytest
from workflows import Context, Workflow, step
from workflows.client.client import WorkflowClient
from workflows.events import Event, StartEvent, StopEvent, WorkflowIdleEvent
from workflows.server import WorkflowServer
from workflows.server.memory_workflow_store import MemoryWorkflowStore
from .util import (
live_server, # type: ignore[import]
wait_for_passing, # type: ignore[import]
)
class WaitableExternalEvent(Event):
response: str
class WaitingWorkflow(Workflow):
@step
async def start_and_wait(self, ctx: Context, ev: StartEvent) -> StopEvent:
external = await ctx.wait_for_event(WaitableExternalEvent)
return StopEvent(result=f"received: {external.response}")
@pytest.mark.asyncio
async def test_fast_idle_timeout_does_not_drop_valid_event() -> None:
def make_server() -> WorkflowServer:
server = WorkflowServer(
workflow_store=MemoryWorkflowStore(),
idle_release_timeout=timedelta(milliseconds=1),
)
server.add_workflow(
"waiting",
WaitingWorkflow(),
additional_events=[WaitableExternalEvent],
)
return server
async with live_server(make_server) as (base_url, _server):
client = WorkflowClient(base_url=base_url)
started = await client.run_workflow_nowait("waiting")
handler_id = started.handler_id
async for env in client.get_workflow_events(
handler_id, include_internal_events=True
):
event = env.load_event([WorkflowIdleEvent])
if isinstance(event, WorkflowIdleEvent):
send = await client.send_event(
handler_id, WaitableExternalEvent(response="hello")
)
assert send.status == "sent"
break
async def handler_completed() -> None:
data = await client.get_handler(handler_id)
assert data.status == "completed"
assert data.result is not None
assert data.result.value.get("result") == "received: hello"
await wait_for_passing(handler_completed, max_duration=2.0, interval=0.01)
@@ -0,0 +1,127 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 LlamaIndex Inc.
"""Tests for KeyedLock utility."""
from __future__ import annotations
import asyncio
import pytest
from workflows.server.keyed_lock import KeyedLock
@pytest.fixture
def locks() -> KeyedLock:
return KeyedLock()
async def test_basic_locking(locks: KeyedLock) -> None:
"""Test that basic lock acquisition and release works."""
async with locks("key1"):
pass
# Lock should be reusable after release
async with locks("key1"):
pass
async def test_mutual_exclusion(locks: KeyedLock) -> None:
"""Test that the lock provides mutual exclusion."""
order: list[str] = []
async def task(name: str, delay: float) -> None:
async with locks("shared"):
order.append(f"{name}_enter")
await asyncio.sleep(delay)
order.append(f"{name}_exit")
t1 = asyncio.create_task(task("first", 0.01))
await asyncio.sleep(0.001) # Let first task acquire lock
t2 = asyncio.create_task(task("second", 0.01))
await asyncio.gather(t1, t2)
assert order == ["first_enter", "first_exit", "second_enter", "second_exit"]
async def test_exception_in_critical_section(locks: KeyedLock) -> None:
"""Test that lock is released even if exception occurs."""
with pytest.raises(ValueError, match="test error"):
async with locks("key"):
raise ValueError("test error")
# Lock should still be usable after exception
async with locks("key"):
pass
async def test_cancellation_cleanup(locks: KeyedLock) -> None:
"""Test that lock is released on task cancellation."""
started = asyncio.Event()
async def task() -> None:
async with locks("key"):
started.set()
await asyncio.sleep(10) # Long sleep to be cancelled
t = asyncio.create_task(task())
await started.wait()
t.cancel()
with pytest.raises(asyncio.CancelledError):
await t
# Lock should be usable after cancellation
async with locks("key"):
pass
async def test_waiter_cancelled_while_waiting(locks: KeyedLock) -> None:
"""Test that cancellation while waiting cleans up properly."""
holder_started = asyncio.Event()
waiter_started = asyncio.Event()
async def holder() -> None:
async with locks("key"):
holder_started.set()
await asyncio.sleep(0.1)
async def waiter() -> None:
await holder_started.wait()
waiter_started.set()
async with locks("key"):
pass # Should never get here
t1 = asyncio.create_task(holder())
t2 = asyncio.create_task(waiter())
await waiter_started.wait()
await asyncio.sleep(0.01) # Let waiter register
t2.cancel()
with pytest.raises(asyncio.CancelledError):
await t2
await t1
# Lock should be usable again after cancellation
async with locks("key"):
pass
async def test_parallel_access_mutual_exclusion_with_race_detection(
locks: KeyedLock,
) -> None:
"""Test mutual exclusion using race condition detection."""
shared_value = 0
iterations = 50
async def increment_task() -> None:
nonlocal shared_value
async with locks("key"):
current = shared_value
await asyncio.sleep(0) # Yield to event loop
shared_value = current + 1
tasks = [asyncio.create_task(increment_task()) for _ in range(iterations)]
await asyncio.gather(*tasks)
assert shared_value == iterations
@@ -14,7 +14,6 @@ def test_init() -> None:
server = WorkflowServer()
assert len(server._middleware) == 1
assert server._workflows == {}
assert server._contexts == {}
assert server._handlers == {}
@@ -153,7 +153,11 @@ async def test_run_workflow_with_nonconforming_start_event_type(
async def test_health_check(client: AsyncClient) -> None:
response = await client.get("/health")
assert response.status_code == 200
assert response.json() == {"status": "healthy"}
data = response.json()
assert data["status"] == "healthy"
assert data["loaded_workflows"] == 0
assert data["active_workflows"] == 0
assert data["idle_workflows"] == 0
@pytest.mark.asyncio
@@ -975,6 +979,7 @@ async def test_post_event_context_not_available(
run_handler=SimpleNamespace(done=lambda: False, ctx=None),
workflow_name="test",
status="running",
mark_active=lambda: None,
)
handler_id = "noctx-1"
@@ -5,13 +5,9 @@ from __future__ import annotations
import asyncio
import contextlib
import socket
from contextlib import closing
from typing import AsyncGenerator
import httpx
import pytest
import uvicorn
from workflows import Workflow
from workflows.client.client import WorkflowClient
from workflows.events import StopEvent
@@ -21,16 +17,10 @@ from .conftest import ( # type: ignore[import]
ExternalEvent,
RequestedExternalEvent,
)
from .util import live_server as live_server_ctx # type: ignore[import]
from .util import wait_for_passing # type: ignore[import]
def _get_free_port() -> int:
with closing(socket.socket(socket.AF_INET, socket.SOCK_STREAM)) as s:
s.bind(("127.0.0.1", 0))
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
return int(s.getsockname()[1])
@pytest.fixture
async def live_server(
simple_test_workflow: Workflow,
@@ -38,50 +28,16 @@ async def live_server(
interactive_workflow: Workflow,
error_workflow: Workflow,
) -> AsyncGenerator[tuple[str, WorkflowServer], None]:
port = _get_free_port()
server = WorkflowServer()
server.add_workflow("test", simple_test_workflow)
server.add_workflow("streaming", streaming_workflow)
server.add_workflow("interactive", interactive_workflow)
server.add_workflow("error", error_workflow)
def make_server() -> WorkflowServer:
server = WorkflowServer()
server.add_workflow("test", simple_test_workflow)
server.add_workflow("streaming", streaming_workflow)
server.add_workflow("interactive", interactive_workflow)
server.add_workflow("error", error_workflow)
return server
config = uvicorn.Config(
server.app,
host="127.0.0.1",
port=port,
log_level="error",
loop="asyncio",
)
uv_server = uvicorn.Server(config)
# Start server in background task (lifespan will start workflows)
task = asyncio.create_task(uv_server.serve())
# Wait until server responds on /health or timeout
base_url = f"http://127.0.0.1:{port}"
async with httpx.AsyncClient(base_url=base_url, timeout=1.0) as client:
for _ in range(50): # ~0.5s max wait
try:
resp = await client.get("/health")
if resp.status_code == 200:
break
except Exception:
pass
await asyncio.sleep(0.01)
else:
uv_server.should_exit = True
await task
raise RuntimeError("Live server did not start in time")
try:
async with live_server_ctx(make_server) as (base_url, server):
yield base_url, server
finally:
uv_server.should_exit = True
try:
await task
finally:
# Ensure graceful shutdown of workflow server
await server.stop()
@pytest.mark.asyncio
@@ -1,19 +1,36 @@
import asyncio
import socket
import time
from typing import Awaitable, Callable, TypeVar
from contextlib import asynccontextmanager
from typing import AsyncGenerator, Awaitable, Callable, TypeVar
import httpx
import uvicorn
from workflows.server import WorkflowServer
T = TypeVar("T")
async def async_yield(iterations: int = 10) -> None:
"""Yield to the event loop multiple times to let async tasks run.
Use this when you need to let async tasks make progress without
waiting for real time to pass. Useful with time_machine when time
is frozen (tick=False).
"""
for _ in range(iterations):
await asyncio.sleep(0)
async def wait_for_passing(
func: Callable[[], Awaitable[T]],
max_duration: float = 5.0,
interval: float = 0.05,
) -> T:
start_time = time.time()
start_time = time.monotonic()
last_exception = None
while time.time() - start_time < max_duration:
remaining_duration = max_duration - (time.time() - start_time)
while time.monotonic() - start_time < max_duration:
remaining_duration = max_duration - (time.monotonic() - start_time)
try:
return await asyncio.wait_for(func(), timeout=remaining_duration)
except Exception as e:
@@ -26,3 +43,85 @@ async def wait_for_passing(
raise TimeoutError(
f"Function {func_name} timed out after {max_duration} seconds"
)
@asynccontextmanager
async def live_server(
server_factory: Callable[[], WorkflowServer],
) -> AsyncGenerator[tuple[str, WorkflowServer], None]:
"""Start a live HTTP server for testing with atomic port acquisition.
This context manager handles:
- Atomic port acquisition (no race condition with parallel tests)
- Server startup with health check
- Graceful shutdown
Args:
server_factory: A callable that creates and configures a WorkflowServer.
This allows tests to customize workflows, idle_release_timeout, etc.
Yields:
A tuple of (base_url, server) for making requests and inspecting state.
Example:
def make_server() -> WorkflowServer:
server = WorkflowServer(idle_release_timeout=timedelta(seconds=1))
server.add_workflow("test", MyWorkflow())
return server
async with live_server(make_server) as (base_url, server):
client = WorkflowClient(base_url=base_url)
# ... run tests
"""
# Create socket and bind atomically - prevents race condition in parallel tests
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
try:
sock.bind(("127.0.0.1", 0))
sock.listen(128)
port = sock.getsockname()[1]
server = server_factory()
config = uvicorn.Config(
server.app,
host="127.0.0.1",
port=port,
log_level="error",
loop="asyncio",
)
uv_server = uvicorn.Server(config)
# Start server in background task with our pre-bound socket
task = asyncio.create_task(uv_server.serve(sockets=[sock]))
# Wait until server responds on /health or timeout
base_url = f"http://127.0.0.1:{port}"
async with httpx.AsyncClient(base_url=base_url, timeout=1.0) as client:
for _ in range(50): # ~0.5s max wait
try:
resp = await client.get("/health")
if resp.status_code == 200:
break
except Exception:
pass
await asyncio.sleep(0.01)
else:
uv_server.should_exit = True
await task
raise RuntimeError("Live server did not start in time")
try:
yield base_url, server
finally:
uv_server.should_exit = True
try:
await task
finally:
await server.stop()
finally:
# Socket is managed by uvicorn after serve() starts, but close if we fail early
try:
sock.close()
except Exception:
pass
+1 -1
View File
@@ -6,7 +6,7 @@ build-backend = "uv_build"
dev = [
"pytest>=8.4.2",
"pytest-cov>=7.0.0",
"ty>=0.0.1a26"
"ty>=0.0.1,<0.0.9"
]
[project]
+4 -1
View File
@@ -6,7 +6,7 @@ build-backend = "uv_build"
dev = [
"pytest-xdist>=3.8.0",
"ruff>=0.14.5",
"ty>=0.0.1a26"
"ty>=0.0.1,<0.0.9"
]
[project]
@@ -73,6 +73,9 @@ select = [
"*.ipynb" = ["T201", "ANN001", "ANN002", "ANN003"]
"__main__.py" = ["T201"]
[tool.ty.src]
exclude = ["**/uv.lock", "**/.gitignore"]
[tool.uv.sources]
llama-index-workflows = {workspace = true}
llama-index-utils-workflow = {workspace = true}
Generated
+21 -21
View File
@@ -1581,7 +1581,7 @@ dev = [
dev = [
{ name = "pytest-xdist", specifier = ">=3.8.0" },
{ name = "ruff", specifier = ">=0.14.5" },
{ name = "ty", specifier = ">=0.0.1a26" },
{ name = "ty", specifier = ">=0.0.1,<0.0.9" },
]
[[package]]
@@ -4066,27 +4066,27 @@ wheels = [
[[package]]
name = "ty"
version = "0.0.1a26"
version = "0.0.8"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/39/39/b4b4ecb6ca6d7e937fa56f0b92a8f48d7719af8fe55bdbf667638e9f93e2/ty-0.0.1a26.tar.gz", hash = "sha256:65143f8efeb2da1644821b710bf6b702a31ddcf60a639d5a576db08bded91db4", size = 4432154, upload-time = "2025-11-10T18:02:30.142Z" }
sdist = { url = "https://files.pythonhosted.org/packages/72/9d/59e955cc39206a0d58df5374808785c45ec2a8a2a230eb1638fbb4fe5c5d/ty-0.0.8.tar.gz", hash = "sha256:352ac93d6e0050763be57ad1e02087f454a842887e618ec14ac2103feac48676", size = 4828477, upload-time = "2025-12-29T13:50:07.193Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/cc/6a/661833ecacc4d994f7e30a7f1307bfd3a4a91392a6b03fb6a018723e75b8/ty-0.0.1a26-py3-none-linux_armv6l.whl", hash = "sha256:09208dca99bb548e9200136d4d42618476bfe1f4d2066511f2c8e2e4dfeced5e", size = 9173869, upload-time = "2025-11-10T18:01:46.012Z" },
{ url = "https://files.pythonhosted.org/packages/66/a8/32ea50f064342de391a7267f84349287e2f1c2eb0ad4811d6110916179d6/ty-0.0.1a26-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:91d12b66c91a1b82e698a2aa73fe043a1a9da83ff0dfd60b970500bee0963b91", size = 8973420, upload-time = "2025-11-10T18:01:49.32Z" },
{ url = "https://files.pythonhosted.org/packages/d1/f6/6659d55940cd5158a6740ae46a65be84a7ee9167738033a9b1259c36eef5/ty-0.0.1a26-py3-none-macosx_11_0_arm64.whl", hash = "sha256:c5bc6dfcea5477c81ad01d6a29ebc9bfcbdb21c34664f79c9e1b84be7aa8f289", size = 8528888, upload-time = "2025-11-10T18:01:51.511Z" },
{ url = "https://files.pythonhosted.org/packages/79/c9/4cbe7295013cc412b4f100b509aaa21982c08c59764a2efa537ead049345/ty-0.0.1a26-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:40e5d15635e9918924138e8d3fb1cbf80822dfb8dc36ea8f3e72df598c0c4bea", size = 8801867, upload-time = "2025-11-10T18:01:53.888Z" },
{ url = "https://files.pythonhosted.org/packages/ed/b3/25099b219a6444c4b29f175784a275510c1cd85a23a926d687ab56915027/ty-0.0.1a26-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:86dc147ed0790c7c8fd3f0d6c16c3c5135b01e99c440e89c6ca1e0e592bb6682", size = 8975519, upload-time = "2025-11-10T18:01:56.231Z" },
{ url = "https://files.pythonhosted.org/packages/73/3e/3ad570f4f592cb1d11982dd2c426c90d2aa9f3d38bf77a7e2ce8aa614302/ty-0.0.1a26-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:fbe0e07c9d5e624edfc79a468f2ef191f9435581546a5bb6b92713ddc86ad4a6", size = 9331932, upload-time = "2025-11-10T18:01:58.476Z" },
{ url = "https://files.pythonhosted.org/packages/04/fa/62c72eead0302787f9cc0d613fc671107afeecdaf76ebb04db8f91bb9f7e/ty-0.0.1a26-py3-none-manylinux_2_17_ppc64.manylinux2014_ppc64.whl", hash = "sha256:0dcebbfe9f24b43d98a078f4a41321ae7b08bea40f5c27d81394b3f54e9f7fb5", size = 9921353, upload-time = "2025-11-10T18:02:00.749Z" },
{ url = "https://files.pythonhosted.org/packages/6c/1f/3b329c4b60d878704e09eb9d05467f911f188e699961c044b75932893e0a/ty-0.0.1a26-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:0901b75afc7738224ffc98bbc8ea03a20f167a2a83a4b23a6550115e8b3ddbc6", size = 9700800, upload-time = "2025-11-10T18:02:03.544Z" },
{ url = "https://files.pythonhosted.org/packages/92/24/13fcba20dd86a7c3f83c814279aa3eb6a29c5f1b38a3b3a4a0fd22159189/ty-0.0.1a26-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:4788f34d384c132977958d76fef7f274f8d181b22e33933c4d16cff2bb5ca3b9", size = 9728289, upload-time = "2025-11-10T18:02:06.386Z" },
{ url = "https://files.pythonhosted.org/packages/40/7a/798894ff0b948425570b969be35e672693beeb6b852815b7340bc8de1575/ty-0.0.1a26-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:b98851c11c560ce63cd972ed9728aa079d9cf40483f2cdcf3626a55849bfe107", size = 9279735, upload-time = "2025-11-10T18:02:09.425Z" },
{ url = "https://files.pythonhosted.org/packages/1a/54/71261cc1b8dc7d3c4ad92a83b4d1681f5cb7ea5965ebcbc53311ae8c6424/ty-0.0.1a26-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:c20b4625a20059adecd86fe2c4df87cd6115fea28caee45d3bdcf8fb83d29510", size = 8767428, upload-time = "2025-11-10T18:02:11.956Z" },
{ url = "https://files.pythonhosted.org/packages/8e/07/b248b73a640badba2b301e6845699b7dd241f40a321b9b1bce684d440f70/ty-0.0.1a26-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:d9909e96276f8d16382d285db92ae902174cae842aa953003ec0c06642db2f8a", size = 9009170, upload-time = "2025-11-10T18:02:14.878Z" },
{ url = "https://files.pythonhosted.org/packages/f8/35/ec8353f2bb7fd2f41bca6070b29ecb58e2de9af043e649678b8c132d5439/ty-0.0.1a26-py3-none-musllinux_1_2_i686.whl", hash = "sha256:a76d649ceefe9baa9bbae97d217bee076fd8eeb2a961f66f1dff73cc70af4ac8", size = 9119215, upload-time = "2025-11-10T18:02:18.329Z" },
{ url = "https://files.pythonhosted.org/packages/70/48/db49fe1b7e66edf90dc285869043f99c12aacf7a99c36ee760e297bac6d5/ty-0.0.1a26-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:a0ee0f6366bcf70fae114e714d45335cacc8daa936037441e02998a9110b7a29", size = 9398655, upload-time = "2025-11-10T18:02:21.031Z" },
{ url = "https://files.pythonhosted.org/packages/10/f8/d869492bdbb21ae8cf4c99b02f20812bbbf49aa187cfeb387dfaa03036a8/ty-0.0.1a26-py3-none-win32.whl", hash = "sha256:86689b90024810cac7750bf0c6e1652e4b4175a9de7b82b8b1583202aeb47287", size = 8645669, upload-time = "2025-11-10T18:02:23.23Z" },
{ url = "https://files.pythonhosted.org/packages/b4/18/8a907575d2b335afee7556cb92233ebb5efcefe17752fc9dcab21cffb23b/ty-0.0.1a26-py3-none-win_amd64.whl", hash = "sha256:829e6e6dbd7d9d370f97b2398b4804552554bdcc2d298114fed5e2ea06cbc05c", size = 9442975, upload-time = "2025-11-10T18:02:25.68Z" },
{ url = "https://files.pythonhosted.org/packages/e9/22/af92dcfdd84b78dd97ac6b7154d6a763781f04a400140444885c297cc213/ty-0.0.1a26-py3-none-win_arm64.whl", hash = "sha256:b8f431c784d4cf5b4195a3521b2eca9c15902f239b91154cb920da33f943c62b", size = 8958958, upload-time = "2025-11-10T18:02:28.071Z" },
{ url = "https://files.pythonhosted.org/packages/69/2b/dd61f7e50a69c72f72c625d026e9ab64a0db62b2dd32e7426b520e2429c6/ty-0.0.8-py3-none-linux_armv6l.whl", hash = "sha256:a289d033c5576fa3b4a582b37d63395edf971cdbf70d2d2e6b8c95638d1a4fcd", size = 9853417, upload-time = "2025-12-29T13:50:08.979Z" },
{ url = "https://files.pythonhosted.org/packages/90/72/3f1d3c64a049a388e199de4493689a51fc6aa5ff9884c03dea52b4966657/ty-0.0.8-py3-none-macosx_10_12_x86_64.whl", hash = "sha256:788ea97dc8153a94e476c4d57b2551a9458f79c187c4aba48fcb81f05372924a", size = 9657890, upload-time = "2025-12-29T13:50:27.867Z" },
{ url = "https://files.pythonhosted.org/packages/71/d1/08ac676bd536de3c2baba0deb60e67b3196683a2fabebfd35659d794b5e9/ty-0.0.8-py3-none-macosx_11_0_arm64.whl", hash = "sha256:1b5f1f3d3e230f35a29e520be7c3d90194a5229f755b721e9092879c00842d31", size = 9180129, upload-time = "2025-12-29T13:50:22.842Z" },
{ url = "https://files.pythonhosted.org/packages/af/93/610000e2cfeea1875900f73a375ba917624b0a008d4b8a6c18c894c8dbbc/ty-0.0.8-py3-none-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:6da9ed377fbbcec0a3b60b2ca5fd30496e15068f47cef2344ba87923e78ba996", size = 9683517, upload-time = "2025-12-29T13:50:18.658Z" },
{ url = "https://files.pythonhosted.org/packages/05/04/bef50ba7d8580b0140be597de5cc0ba9a63abe50d3f65560235f23658762/ty-0.0.8-py3-none-manylinux_2_17_armv7l.manylinux2014_armv7l.whl", hash = "sha256:7d0a2bdce5e701d19eb8d46d9da0fe31340f079cecb7c438f5ac6897c73fc5ba", size = 9676279, upload-time = "2025-12-29T13:50:25.207Z" },
{ url = "https://files.pythonhosted.org/packages/aa/b9/2aff1ef1f41b25898bc963173ae67fc8f04ca666ac9439a9c4e78d5cc0ff/ty-0.0.8-py3-none-manylinux_2_17_i686.manylinux2014_i686.whl", hash = "sha256:ef9078799d26d3cc65366e02392e2b78f64f72911b599e80a8497d2ec3117ddb", size = 10073015, upload-time = "2025-12-29T13:50:35.422Z" },
{ url = "https://files.pythonhosted.org/packages/df/0e/9feb6794b6ff0a157c3e6a8eb6365cbfa3adb9c0f7976e2abdc48615dd72/ty-0.0.8-py3-none-manylinux_2_17_ppc64.manylinux2014_ppc64.whl", hash = "sha256:54814ac39b4ab67cf111fc0a236818155cf49828976152378347a7678d30ee89", size = 10961649, upload-time = "2025-12-29T13:49:58.717Z" },
{ url = "https://files.pythonhosted.org/packages/f4/3b/faf7328b14f00408f4f65c9d01efe52e11b9bcc4a79e06187b370457b004/ty-0.0.8-py3-none-manylinux_2_17_ppc64le.manylinux2014_ppc64le.whl", hash = "sha256:c4baf0a80398e8b6c68fa36ff85045a50ede1906cd4edb41fb4fab46d471f1d4", size = 10676190, upload-time = "2025-12-29T13:50:01.11Z" },
{ url = "https://files.pythonhosted.org/packages/64/a5/cfeca780de7eeab7852c911c06a84615a174d23e9ae08aae42a645771094/ty-0.0.8-py3-none-manylinux_2_17_s390x.manylinux2014_s390x.whl", hash = "sha256:ac8e23c3faefc579686799ef1649af8d158653169ad5c3a7df56b152781eeb67", size = 10438641, upload-time = "2025-12-29T13:50:29.664Z" },
{ url = "https://files.pythonhosted.org/packages/0e/8d/8667c7e0ac9f13c461ded487c8d7350f440cd39ba866d0160a8e1b1efd6c/ty-0.0.8-py3-none-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:b558a647a073d0c25540aaa10f8947de826cb8757d034dd61ecf50ab8dbd77bf", size = 10214082, upload-time = "2025-12-29T13:50:31.531Z" },
{ url = "https://files.pythonhosted.org/packages/f8/11/e563229870e2c1d089e7e715c6c3b7605a34436dddf6f58e9205823020c2/ty-0.0.8-py3-none-musllinux_1_2_aarch64.whl", hash = "sha256:8c0104327bf480508bd81f320e22074477df159d9eff85207df39e9c62ad5e96", size = 9664364, upload-time = "2025-12-29T13:50:05.443Z" },
{ url = "https://files.pythonhosted.org/packages/b1/ad/05b79b778bf5237bcd7ee08763b226130aa8da872cbb151c8cfa2e886203/ty-0.0.8-py3-none-musllinux_1_2_armv7l.whl", hash = "sha256:496f1cb87261dd1a036a5609da80ee13de2e6ee4718a661bfa2afb91352fe528", size = 9679440, upload-time = "2025-12-29T13:50:11.289Z" },
{ url = "https://files.pythonhosted.org/packages/12/b5/23ba887769c4a7b8abfd1b6395947dc3dcc87533fbf86379d3a57f87ae8f/ty-0.0.8-py3-none-musllinux_1_2_i686.whl", hash = "sha256:2c488031f92a075ae39d13ac6295fdce2141164ec38c5d47aa8dc24ee3afa37e", size = 9808201, upload-time = "2025-12-29T13:50:21.003Z" },
{ url = "https://files.pythonhosted.org/packages/f8/90/5a82ac0a0707db55376922aed80cd5fca6b2e6d6e9bcd8c286e6b43b4084/ty-0.0.8-py3-none-musllinux_1_2_x86_64.whl", hash = "sha256:90d6f08c5982fa3e802b8918a32e326153519077b827f91c66eea4913a86756a", size = 10313262, upload-time = "2025-12-29T13:50:03.306Z" },
{ url = "https://files.pythonhosted.org/packages/14/f7/ff97f37f0a75db9495ddbc47738ec4339837867c4bfa145bdcfbd0d1eb2f/ty-0.0.8-py3-none-win32.whl", hash = "sha256:d7f460ad6fc9325e9cc8ea898949bbd88141b4609d1088d7ede02ce2ef06e776", size = 9254675, upload-time = "2025-12-29T13:50:33.35Z" },
{ url = "https://files.pythonhosted.org/packages/af/51/eba5d83015e04630002209e3590c310a0ff1d26e1815af204a322617a42e/ty-0.0.8-py3-none-win_amd64.whl", hash = "sha256:1641fb8dedc3d2da43279d21c3c7c1f80d84eae5c264a1e8daa544458e433c19", size = 10131382, upload-time = "2025-12-29T13:50:13.719Z" },
{ url = "https://files.pythonhosted.org/packages/38/1c/0d8454ff0f0f258737ecfe84f6e508729191d29663b404832f98fa5626b7/ty-0.0.8-py3-none-win_arm64.whl", hash = "sha256:ec74f022f315bede478ecae1277a01ab618e6500c1d68450d7883f5cd6ed554a", size = 9636374, upload-time = "2025-12-29T13:50:16.344Z" },
]
[[package]]
@@ -4284,7 +4284,7 @@ requires-dist = [
dev = [
{ name = "pytest", specifier = ">=8.4.2" },
{ name = "pytest-cov", specifier = ">=7.0.0" },
{ name = "ty", specifier = ">=0.0.1a26" },
{ name = "ty", specifier = ">=0.0.1,<0.0.9" },
]
[[package]]