mirror of
https://github.com/run-llama/workflows-py.git
synced 2026-08-24 20:01:34 -04:00
drop and restore idle workflows in workflow server (#266)
This commit is contained in:
@@ -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.
|
||||
@@ -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"""
|
||||
|
||||
+4
@@ -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;
|
||||
+13
-4
@@ -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
|
||||
|
||||
@@ -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
@@ -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}
|
||||
|
||||
@@ -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]]
|
||||
|
||||
Reference in New Issue
Block a user