Skip to content
Pyrula

Agents

Generated from the type stubs and docstrings. Do not edit by hand.

FastAPI dependency that validates a static Bearer token.

Usage::

router = RunRouter(worker=..., config=RunRouterConfig(auth_config=AuthConfig(provider=APIKeyAuth("sk-proj-xxx"))))

Async HTTP client for Pyrula agents.

submit(self, agent_name: str, metadata: Optional[dict[str, Any]] = None, idempotency_key: Optional[str] = None, headers: Optional[dict[str, str]] = None, **params: Any) -> TurnHandle
stream(self, agent_name: str, turn_id: str, last_event_id: Optional[str] = None, headers: Optional[dict[str, str]] = None, idle_timeout_seconds: Optional[float] = None) -> AsyncIterator[Event]
stream_from_path(self, stream_path: str, last_event_id: Optional[str] = None, headers: Optional[dict[str, str]] = None, idle_timeout_seconds: Optional[float] = None) -> AsyncIterator[Event]
replay(self, agent_name: str, turn_id: str, headers: Optional[dict[str, str]] = None) -> AsyncIterator[Event]
stream_with_replay_fallback(self, agent_name: str, turn_id: str, last_event_id: Optional[str] = None, headers: Optional[dict[str, str]] = None, include_overflow_event: bool = False) -> AsyncIterator[Event]

Stream live events, falling back to replay when the live cursor is stale.

This is useful when reconnecting with a saved Last-Event-ID that may have fallen behind the retained live stream window. In that case the live stream emits stream:overflow; this helper can switch to replay() automatically.

run(self, agent_name: str, metadata: Optional[dict[str, Any]] = None, idempotency_key: Optional[str] = None, headers: Optional[dict[str, str]] = None, **params: Any) -> list[Event]

Submit a turn and collect all events until the terminal event.

Follows ctx.continue_as_new chains transparently: a run:continued event means the turn handed off to a successor run in the same thread (one logical turn - see pyrula.workflows.continuation). The successor’s stream is a separate SSE connection (a different turn_id under the hood), so reconnect and keep collecting into the same event list rather than returning the chain’s RUN_CONTINUED hop as if it were the final result.

Bounded the same way resolve_continuation is: a self-loop (continued_to pointing back at the turn we just streamed) raises immediately, and any other chain - cyclic or merely runaway - raises after _CONTINUATION_MAX_HOPS reconnects. Both raise ContinuationChainError so a pathological chain fails loud instead of reconnecting forever.

cancel(self, agent_name: str, turn_id: str, headers: Optional[dict[str, str]] = None) -> None
status(self, agent_name: str, turn_id: str, headers: Optional[dict[str, str]] = None) -> TurnStatusInfo

Agent execution context extending BaseContext with LLM/stream/tool-loop.

Inherits from pyrula.workflows BaseContext: step, sleep, signal, timers, determinism, checkpoint, emit, emit_status, interrupt.

Adds agent-only: call_agent, gather, memory, thread, llm.

run_id and metadata default to ""/{} so the context can be constructed from a bound TurnState alone (testing and low-level usage). __post_init__ back-fills run_id from state when omitted (an empty string is never a valid run id); metadata keeps its empty-dict default.

Fields:

  • run_id: str = ''
  • metadata: dict[str, Any] = field(default_factory=dict)
  • llm: Any = None
  • state: Optional[TurnState] = None
  • cancelled: bool
  • turn_id: str
  • memory: Any
  • thread: Optional[ThreadHandle]
mcp_tools(self, server: Optional[str] = None) -> list[Any]

Discovered MCP tool callables to splat into ctx.llm.stream(tools=[…]).

Live: the worker pool’s tools for one server (or all when server is None). Replay with a dead server: stub callables named from the run’s recorded tool_use names, so tool_map matches the recorded run.

server=None means “all MCP tools”: the live pool’s tools, or on replay the recorded MCP tool_use names, or [] when no MCP is configured (asking for all tools is empty, not an error). A named, unconnected server raises.

emit_status(self, text: str) -> TurnEvent
emit(self, kind: str, payload: dict[str, Any], id: Optional[str] = None) -> TurnEvent
call_agent(self, agent_fn: Any, llm: Optional[Any] = None, id: Optional[str] = None, **kwargs: Any) -> Any
gather_subagents(self, *calls: tuple[Any, dict[str, Any]], id: Optional[str] = None, poll_interval: float = 0.5) -> list[Any]
spawn_workflow(self, name: str, params: Optional[dict[str, Any]] = None) -> str

Submit a @workflow run; returns run_id immediately (fire-and-forget).

The run_id is replay-stable (ctx.uuid) and the submit is idempotent, so a crash-and-replay re-submits the SAME run rather than spawning a duplicate. To await a child’s result durably, use ctx.run_child instead.

wait_for_workflow(self, name: str, run_id: str, poll_interval: float = 0.5, timeout: Optional[float] = None) -> Any

Poll the store until a spawned workflow completes.

Returns the workflow result on success, raises PyrulaError if the workflow ended with an error event, and WaitTimeoutError if timeout seconds elapse first (None waits indefinitely).

Recognises run:complete/run:error event kinds so it works regardless of which runner executes the spawned workflow.


Fields:

  • max_tool_calls: Optional[int] = None
  • max_llm_calls: Optional[int] = DEFAULT_MAX_LLM_CALLS
  • max_agent_depth: Optional[int] = None
  • max_parallel_subagents: Optional[int] = None
  • max_event_count_per_turn: Optional[int] = None
  • max_event_payload_bytes: Optional[int] = None
  • max_turn_wall_seconds: Optional[float] = None
  • context_window: Optional[int] = None
  • reserved_output_tokens: Optional[int] = None
  • compaction_fraction: Optional[float] = None
  • history_compactor: Any = UNSET
  • estimate_tokens: Optional[Any] = None


FastAPI router that exposes pyrula.agents over HTTP with SSE streaming.

Subclasses RunRouter to inherit generic durable-run infrastructure (attach, _classify_status, _setup_auth, ping_store, draining). Overrides lifespan (per-agent recovery), route registration (agents URL patterns), and auth checks (authorize_turn semantics).

Fields:

  • get_stream = get_stream
  • get_replay = get_replay
  • get_status = get_status
  • delete_cancel = delete_cancel
  • post_resume = post_resume
  • post_signal = post_signal
  • post_create_thread = post_create_thread
  • get_thread = get_thread
  • delete_thread = delete_thread
  • get_thread_turns = get_thread_turns
  • get_agents = get_agents
  • get_runs = get_runs
  • connector_state: Optional[str] Pyrula Cloud connector state, or None when cloud is not active.
lifespan(self) -> AsyncGenerator[None, None]

Embedded-mode startup (per-agent recovery) and shutdown (drain + cancel).

apply_redaction_hook(store: Any, hook: Callable[[str, dict[str, Any]], dict[str, Any]]) -> None

Patch store.append_event_sequenced in-place to fire hook for every event.

Patches the canonical allocator (append_event delegates to it, and submit_run writes RUN_INIT through it directly), so every persisted event is redacted. Only intercepts appends through this store instance. Raises ValueError if a hook is already installed to prevent silent hook-chaining.


Configuration for AgentRouter - extends RunRouterConfig with agent-only fields.

Fields:

  • reporter: Any = None
  • cloud: bool = True
  • deps: Any = None
  • deps_factory: Optional[Callable[..., Any]] = None
  • heartbeat_config: Optional[ReaderHeartbeatConfig] = None
  • auth_config: Optional[AuthConfig] = None

Anthropic Messages API adapter (the native Pyrula LLM client for @agent).

To keep your existing anthropic SDK code and get durability by swapping one import, use the drop-in pyrula.agents.compat.anthropic.AsyncAnthropic instead.

Fields:

  • provider = 'anthropic'
  • model = model
  • streaming_mode = streaming_mode
  • max_buffer_bytes = max_buffer_bytes
  • max_retries = max_retries
  • retry_delay = retry_delay
  • max_tokens = max_tokens
  • prompt_caching = prompt_caching
  • cache_ttl = cache_ttl
  • llm_call_timeout = llm_call_timeout
  • rate_limit = rate_limit
stream(self, messages: list[dict[str, Any]], tools: Optional[list[ToolDefinition]] = None, system: Optional[str] = None, tool_choice: Optional[dict[str, Any]] = None) -> AsyncGenerator[ContentBlock, None]
complete(self, messages: list[dict[str, Any]], system: Optional[str] = None) -> CompletionResult

Fields:

  • prefix: str = '/agents'
  • auth: Any = UNSET
  • deployment_mode: Optional[str] = None
  • sse_keepalive_interval: float = 15.0
  • cancel_on_reader_disconnect: bool = False
  • include_health: bool = True
  • include_ready: bool = True
  • health_path: str = '/healthz'
  • ready_path: str = '/readyz'
  • reporter: Optional[Reporter] = None
  • max_queue_depth: Optional[int] = None
  • deps: Any = None
  • cloud: bool = True
  • redaction_hook: Optional[Callable[[str, dict[str, Any]], dict[str, Any]]] = None
  • idempotency_ttl_s: float = 86400.0
  • cors_origins: Optional[list[str]] = None


Fields:

  • STORE_UNAVAILABLE = 'store_unavailable'
  • STREAM_TOO_LARGE = 'stream_too_large'
  • TURN_TIMEOUT = 'turn_timeout'
  • VERSION_SKEW = 'version_skew'
  • CANCELLED = 'cancelled'
  • PAUSED = 'paused'
  • TOOL_ERROR = 'tool_error'
  • BUDGET_EXCEEDED = 'budget_exceeded'
  • AGENT_ERROR = 'agent_error'
  • BUFFER_OVERFLOW = 'buffer_overflow'
  • LLM_ERROR = 'llm_error'
  • LLM_RATE_LIMITED = 'llm_rate_limited'
  • LLM_PROVIDER_UNAVAILABLE = 'llm_provider_unavailable'
  • POISON_QUARANTINE = 'poison_quarantine'
  • MAX_RETRIES = 'max_retries'
  • OUTPUT_VALIDATION = 'output_validation'
  • AUTH_MISSING = 'auth_missing'
  • AUTH_INVALID = 'auth_invalid'
  • AUTH_FORBIDDEN = 'auth_forbidden'
  • NOT_FOUND = 'not_found'
  • INVALID_REQUEST = 'invalid_request'
  • BACKPRESSURE_REJECTED = 'backpressure_rejected'
  • IDEMPOTENCY_CONFLICT = 'idempotency_conflict'
  • INTERNAL = 'internal'
  • INTERRUPT_TIMEOUT = 'interrupt_timeout'
  • CONTEXT_WINDOW_EXCEEDED = 'context_window_exceeded'
  • AUTH_UNAVAILABLE = 'auth_unavailable'
  • WRITE_TOOL_BLOCKED = 'write_tool_blocked'
  • CONTINUE_AS_NEW_NOT_SUPPORTED = _ContractErrorCode.CONTINUE_AS_NEW_NOT_SUPPORTED.value

In-memory store for testing only - workflow MemoryStore + agent threads.

Does not persist across process restarts. Use ValkeyStore for production.

Fields:

  • max_thread_turns = max_thread_turns

Run-neutral lifecycle + event-log view over a pyrula.agents turn-scoped store.

Implements both core ABCs in one object. The engine holds a single reference; mixin methods (cache_get, cache_set, set_thread_completed, get_sleeping_turns, wake_turn, …) pass through via __getattr__.

Idempotent: LifecycleAdapter(LifecycleAdapter(store)) returns the existing adapter (no double-wrap).

When the underlying store is a core store (with run_* naming), the run_* methods delegate directly; turn_* aliases still work by remapping to run_*.

Fields:

  • raw: Store
submit_turn(self, name: str, turn_id: str, params: dict[str, object], **kwargs: Any) -> str
claim_next_turn(self, name: str, worker_id: str, worker_version: Optional[str] = None) -> Optional[str]
claim_run(self, name: str, run_id: str, worker_id: str) -> None
release_turn(self, name: str, turn_id: str, **kwargs: Any) -> None
requeue_turn(self, name: str, turn_id: str, **kwargs: Any) -> None
expire_turn(self, name: str, turn_id: str, **kwargs: Any) -> bool
set_turn_owner(self, name: str, turn_id: str, owner: str) -> None
resume_turn(self, name: str, turn_id: str, action: str, value: Any) -> None
recover_pending_turns(self, name: str, now: Optional[float] = None) -> list[str]
submit_run(self, name: str, run_id: str, params: dict[str, object], metadata: Optional[dict[str, Any]] = None, owner: Optional[str] = None, idempotency_key: Optional[str] = None, idempotency_body_hash: Optional[str] = None, idempotency_ttl: int = 86400, version: str = '', version_behavior: str = 'pinned') -> str
claim_next_run(self, name: str, worker_id: str, worker_version: Optional[str] = None) -> Optional[str]
release_run(self, name: str, run_id: str, worker_id: Optional[str] = None) -> None
requeue_run(self, name: str, run_id: str, worker_id: Optional[str] = None) -> None
recover_pending_runs(self, name: str, now: Optional[float] = None) -> list[str]
get_pending_count(self, name: str) -> int
heartbeat(self, name: str, run_id: str, worker_id: str) -> None
request_cancel(self, name: str, run_id: str) -> str
get_cancel_requested_at(self, name: str, run_id: str) -> Optional[str]
request_pause(self, name: str, run_id: str) -> str
get_pause_requested_at(self, name: str, run_id: str) -> Optional[str]
clear_pause(self, name: str, run_id: str) -> None
resume_paused_run(self, name: str, run_id: str) -> None
expire_run(self, name: str, run_id: str, ttl_seconds: int = 360) -> bool
set_run_owner(self, name: str, run_id: str, owner: str) -> None
get_run_owner(self, name: str, run_id: str) -> Optional[str]
get_run_version(self, name: str, run_id: str) -> Optional[tuple[str, str]]
adopt_run(self, name: str, run_id: str, new_version: str) -> None
submitted_at(self, name: str, run_id: str) -> Optional[float]
apply_retention(self, name: str, run_id: str) -> None
ping(self) -> None
list_runs(self, name: str, limit: int = 50, cursor: Optional[str] = None, owner: Optional[str] = None) -> tuple[list[dict[str, Any]], Optional[str]]
set_interrupted(self, name: str, run_id: str, event_id: str, payload: dict[str, Any]) -> None
resume_run(self, name: str, run_id: str, action: str, value: Any) -> None
get_interrupt_response(self, name: str, run_id: str) -> Optional[dict[str, Any]]
complete_run(self, name: str, run_id: str) -> None
append_event(self, name: str, run_id: str, event: RunEvent) -> None
append_event_sequenced(self, name: str, run_id: str, event: RunEvent) -> int
append_events(self, name: str, run_id: str, events: list[RunEvent]) -> None
append_events_sequenced(self, name: str, run_id: str, events: list[RunEvent], retain: bool = True) -> list[int]
load_run(self, name: str, run_id: str, max_events: Optional[int] = None) -> list[RunEvent]
replay_run(self, name: str, run_id: str, cursor: Optional[Any] = None, count: Optional[int] = None) -> tuple[list[RunEvent], Optional[Any]]
get_last_event(self, name: str, run_id: str) -> Optional[RunEvent]
read_entries_from(self, name: str, run_id: str, cursor: Optional[str], count: int) -> tuple[list[tuple[str, RunEvent]], Optional[str]]

LiteLLM adapter - routes to any supported provider.

Model names follow LiteLLM format: "provider/model".

complete(self, messages: list[dict[str, Any]], system: Optional[str] = None) -> CompletionResult

OpenAI Chat Completions adapter (the native Pyrula LLM client for @agent).

Named OpenAILLM (matching AnthropicLLM). If you instead want to keep your existing openai SDK code and get durability by swapping one import, use the OpenAI-SDK-shaped drop-in pyrula.agents.compat.openai.AsyncOpenAI.

base_url can point to any OpenAI-compatible API (Ollama, Groq, Together AI, OpenRouter, Azure, etc.).

Fields:

  • provider = 'openai'
  • base_url = base_url
complete(self, messages: list[dict[str, Any]], system: Optional[str] = None) -> CompletionResult

Fields:

  • max_retries: int = 3
  • retry_ttl_s: float = 3600.0

Base class for Pyrula exceptions.

Fields:

  • code: Optional[ErrorCode] = None

Simplified entry point for running a Pyrula worker process.

Fields:

  • connector_state: Optional[str]
run(self) -> None
run_async(self) -> None
submit(self, agent_name: str, metadata: Optional[dict[str, Any]] = None, **params: object) -> str

Store operation failed after exhausting the retry budget.

Fields:

  • code = ErrorCode.STORE_UNAVAILABLE

Fields:

  • http: httpx.AsyncClient
  • deps: Any = None
  • state: Optional[TurnState] = None
  • tool_name: Optional[str] = None
  • tool_use_id: Optional[str] = None
  • cache_getter: Optional[Callable[[str], Optional[Any]]] = None
  • cache_setter: Optional[Callable[[str, Any, float], None]] = None
  • idempotency_key: Optional[str] Stable per-tool-call key, turn_id:tool_use_id, deterministic across replay (both components are journaled). Tool effects are at-least-once - a crash between executing a tool and journaling its result re-fires the tool on replay. Pass this key to the external system (e.g. a Stripe idempotency key, a dedup table, a conditional write) so the re-fire is a no-op. None when the context is unbound (no turn / no tool_use_id).
  • cache: CacheProxy
  • cancelled: bool
  • turn_id: str
emit_status(self, text: str) -> TurnEvent

Fields:

  • SLEEPING_INDEX_KEY: str = 'sleeping:index'
  • max_thread_turns = max_thread_turns

Fields:

  • max_concurrent_turns: int = 50
  • orphan_scan_interval: float = 60.0
  • shutdown_grace: float = 30.0
  • cancel_timeout: float = 5.0
  • claim_poll_interval: float = 0.1
  • reporter: Optional[Reporter] = None
  • timer_poll_interval_s: float = 5.0
  • timer_max_batch: int = 100
  • schedule_poll_interval_s: float = 30.0
  • poison: PoisonOptions = field(default_factory=PoisonOptions)
agent(fn: Optional[Callable[..., Any]] = None, name: Optional[str] = None, timeout: float = 600.0, hide_thinking: bool = False, llm: Any = None, limits: Optional[AgentLimits] = None, output_type: Optional[Any] = None, deps_type: Optional[Any] = None, capture: str = CaptureLevel.lifecycle, mcp_servers: Optional[list[Any]] = None, version: Optional[str] = None, version_behavior: VersionBehavior = 'pinned') -> AgentCallable | Callable[[Callable[..., Any]], AgentCallable]
create_app(agents: Optional[list[Callable[..., Any]]] = None, store: Optional[Store] = None, llm: Any, options: Optional[AppOptions] = None) -> FastAPI
data_step(ctx: BaseContext, step_id: str, op: Any, *args: Any, **kwargs: Any) -> Any

Durable data operation with codec-driven replay.

registered_agents() -> list[AgentCallable]
tool(fn: Optional[Callable[..., Any]] = None, retry: int = 0, timeout: float = 30.0, name: Optional[str] = None, description: Optional[str] = None, schema: Optional[dict[str, Any]] = None, cache_ttl: Optional[float] = None, writes: bool = False) -> ToolCallable | Callable[[Callable[..., Any]], ToolCallable]