Metadata-Version: 2.5
Name: pydantic-ai-waymark
Version: 0.2.2
Summary: Run Pydantic AI graphs with durable Waymark model and tool actions
Requires-Python: >=3.12
Requires-Dist: mountaineer-di<0.3,>=0.2
Requires-Dist: pydantic-ai-slim<2.22,>=2.21
Requires-Dist: waymark<0.31,>=0.30.0.dev7
Provides-Extra: openai
Requires-Dist: pydantic-ai-slim[openai]<2.22,>=2.21; extra == 'openai'
Provides-Extra: web
Requires-Dist: fastapi<1,>=0.116; extra == 'web'
Requires-Dist: uvicorn<1,>=0.35; extra == 'web'
Description-Content-Type: text/markdown

# pydantic-ai-waymark

AI agents are outgrowing the request-response cycle. They increasingly run for
days, weeks, or months. Some of that is active work, but much of it is waiting:
for a user to approve an action, for another system to respond, or for a
two-week onboarding period to end. An agent should not have to stay in memory
and occupy a worker through those gaps. Over that lifetime, workers will
restart, code will be deployed, and temporary failures will happen.

Saving the current step in a database is easy. The harder part is making the
whole control loop reliable: recording which model and tool actions completed,
applying retries and timeouts at the right boundaries, scheduling durable
timers, and resuming the right run after a restart. Queues, schedulers, and
state tables can solve each piece, but stitching them together becomes an
orchestration system embedded in every agent.

This is what durable execution provides. [Waymark](https://github.com/piercefreeman/waymark)
compiles ordinary Python control flow into a durable state machine, checkpoints
progress at action boundaries, and stores waits instead of holding a live
process. Workers can come and go while the workflow continues from its last
completed step.

Waymark solves that problem for general Python workflows. `pydantic-ai-waymark`
applies it to Pydantic AI agents: keep the model, tools, retries, timeouts, and
low-level control you already have while Waymark handles persistence, wakeups,
and scalable orchestration. It is more infrastructure than a short, one-shot
agent needs; it earns its keep when the agent must outlive the process running
it.

## Install

Add the package to your project from PyPI:

```bash
uv add pydantic-ai-waymark
```

Include the OpenAI provider used by the example with the `openai` extra:

```bash
uv add "pydantic-ai-waymark[openai]"
```

## Define an agent

This library wraps your existing Pydantic AI agents, so define agents and tools
as you normally would. The only change is to wrap the completed `Agent(...)`
initialization in `waymark_agent(...)`:

```python
from pydantic import BaseModel
from pydantic_ai import Agent
from pydantic_ai_waymark import AIRequestBase, waymark_agent


class Reply(BaseModel):
    answer: str
    needs_human: bool


support_agent = waymark_agent(
    Agent(
        "openai:gpt-5.2",
        name="support_agent",
        instructions="Answer concisely.",
        output_type=Reply,
        defer_model_check=True,
    )
)


@support_agent.tool_plain
def lookup_policy(topic: str) -> str:
    """Look up the support policy for a topic."""
    return f"Policy for {topic}: escalate account changes."
```

## Compile the agent into a workflow

Parameterize `PydanticAIWorkflow` with the request type, implement the Waymark
entrypoint, and call `run_agent` from it:

```python
from waymark import workflow
from pydantic_ai_waymark import AIRequestBase, PydanticAIWorkflow


class SupportRequest(AIRequestBase[None]):
    agent = support_agent


@workflow
class SupportWorkflow(PydanticAIWorkflow[SupportRequest]):
    async def run(self, request: SupportRequest) -> Reply:
        return (await self.run_agent(request)).output


reply = await SupportWorkflow().run(
    SupportRequest(prompt="How do I update my account?")
)
```

The request parameter may be a union such as
`PydanticAIWorkflow[SupportRequest | SalesRequest]`. This simply acts as a typehint
for the run_agent function. You can similarly nest these values within a large request blob:

```python
class Agent1Request(APIRequestBase[None]):
    agent = agent_1

class Agent2Request(APIRequestBase[None]):
    agent = agent_2

class MainRequest(BaseModel):
    request_1: Agent1Request
    request_2: Agent2Request

@workflow
class MultiAgentWorkflow(PydanticAIWorkflow[Agent1Request | Agent2Request]):
    async def run(self, request: MainRequest) -> None:
        response_1 = await self.run_agent(request.request_1)
        response_2 = await self.run_agent(request.request_2)

reply = await MultiAgentWorkflow().run(
    MainRequest(
        request_1=SupportRequest(prompt="What's your name?")
        request_2=SupportRequest(prompt="What's your name?")
    )
)
```

`AIRequestBase` also accepts `message_history`, `deps`, `model`,
`conversation_id`, and `run_id`. Its serialized representation includes the stable
agent reference needed by the worker.

## Durable payload codecs

Large tool values such as images should not be copied into every Waymark snapshot.
Register a serializer/deserializer pair to replace them with a small durable reference
before an action result, graph state, or message history is persisted:

```python
from mountaineer_di import Depends
from pydantic_ai_waymark import Payload, SerializedPayload


async def serialize_payload(
    payload: Payload,
    db=Depends(get_db_connection),
) -> SerializedPayload:
    value = await replace_large_values_with_database_refs(db, payload.to_python())
    return payload.serialized(value)


async def deserialize_payload(
    payload: SerializedPayload,
    db=Depends(get_db_connection),
) -> Payload:
    value = await restore_database_refs(db, payload.value)
    return payload.deserialized(value)


support_agent = waymark_agent(
    Agent(..., name="support_agent"),
    serializer=serialize_payload,
    deserializer=deserialize_payload,
)
```

`Payload` is a discriminated union covering graph state, messages, agent output, tool
output, tool action results, deferred tool results, and user prompts. Match on
`payload.kind` when a context needs special handling. `payload.to_python()` produces
the plain-Python tree to transform; `serialized(...)` and `deserialized(...)` preserve
and validate its kind across the round trip. Codecs may declare additional
`mountaineer_di.Depends(...)` parameters. Sync and async codecs are supported, and
generator dependencies remain open for the call and are cleaned up afterward.

The workflow keeps one current checkpoint and stores message history as ordered deltas.
Pending model-action results contain only the messages added by that action, preventing
Waymark's retained action results from accumulating successively larger history copies.
Graph-state payloads therefore exclude message history; `messages` payloads are
reassembled before each agent step. Use stable keys or upserts when storing large values
so repeated serialization reuses the same durable reference.

Waymark actions use the same dependency resolver, so application I/O can use the same
providers without putting database or storage clients in Pydantic AI's per-run `deps`:

```python
from waymark import action


@action
async def save_result(result: dict, db = Depends(get_db_connection)) -> None:
    await db.insert(result)
```

## Extras

We make our best effort to wrap `pydantic-ai`'s features 1:1 - just with the addition of the magic of durable execution. For instance, you can use Pydantic AI's existing retry and timeout settings as usual:

```python
@support_agent.tool_plain(
    retries=3,
    timeout=120,
)
def lookup_policy(topic: str) -> str:
    return f"Policy for {topic}: escalate account changes."
```

Timeouts and `ModelRetry` responses follow Pydantic AI's retry flow across
durable Waymark actions. Other exceptions fail the workflow immediately.

Tools run in parallel by default. Mark a tool as sequential when it must run
alone:

```python
@support_agent.tool_plain
async def read_profile() -> str:
    return "Profile loaded."


@support_agent.tool_plain(sequential=True)
async def update_account() -> str:
    return "Account updated."


@support_agent.tool_plain
async def send_confirmation() -> str:
    return "Confirmation sent."
```

`sequential=True` acts as a barrier. If the model calls `read_profile`,
`lookup_policy`, `update_account`, and `send_confirmation` in that order,
Waymark resolves them as follows:

1. `read_profile` and `lookup_policy` run in parallel.
2. `update_account` runs alone after both finish.
3. `send_confirmation` starts after `update_account` finishes.

Calls after the barrier can run in parallel again until the next sequential
tool call.

To let a tool request a durable wait, raise `DurableSleep`. The tool action
records the request, the workflow performs the timer, and the supplied result
is returned to the model under the original tool-call ID:

```python
from pydantic_ai_waymark import DurableSleep


@support_agent.tool_plain
def wait_for_follow_up(seconds: float = 5) -> str:
    raise DurableSleep(seconds, result="Follow-up wait completed.")
```

## Docker Compose example

The example includes Postgres, Waymark workers, the Waymark dashboard, and a
small FastAPI form. Put the OpenAI key in the repository's `.env` file and run:

```bash
cp .env.example .env
docker compose -f examples/docker-compose.yml up --build
```

Open [http://localhost:8000](http://localhost:8000). The Waymark dashboard is at
[http://localhost:24119](http://localhost:24119).

The support form calls three action tools without sleeping. A separate sleep
form passes a configurable duration through the agent's dependencies and
visibly exercises a `DurableSleep` timer before the model resumes.

Stop the stack and remove its example database with:

```bash
docker compose -f examples/docker-compose.yml down -v
```

## Lifecycle hooks

Sometimes you want to keep users informed about an agent's progress. The easiest
way to do that is to register hooks for state changes, such as when an agent receives
a tool call and when that tool finishes. This mirrors agent harnesses like Codex and
Claude Code, which show the currently running tool in shimmering text beneath the
conversation history.

Use these hooks to save the current state to a database for polling, or push updates
through a websocket service for broadcast, as shown below. `PydanticAIWorkflow`
provides no-op `on_agent_start`, `on_agent_end`, `on_message`, `on_tool_start`, and
`on_tool_end` methods. Override only the hooks you need, and put external I/O in a
Waymark action so the side effect remains durable:

```python
from typing import Any

from waymark import action, workflow


@action
async def publish_event(event: str, payload: dict[str, Any]) -> None:
    await websocket_service.publish(event, payload)


@workflow
class ObservableSupportWorkflow(PydanticAIWorkflow[SupportRequest]):
    async def on_message(
        self,
        agent_request: SupportRequest,
        message: str,
    ) -> None:
        await self.run_action(
            publish_event(
                event="message.received",
                payload={"message": message},
            )
        )

    async def on_tool_start(
        self,
        agent_request: SupportRequest,
        tool_id: str,
        tool_args: object,
    ) -> None:
        await self.run_action(
            publish_event(
                event="tool.started",
                payload={"tool_id": tool_id, "args": tool_args},
            )
        )

    async def on_tool_end(
        self,
        agent_request: SupportRequest,
        tool_id: str,
        payload: object,
    ) -> None:
        await self.run_action(
            publish_event(
                event="tool.ended",
                payload={"tool_id": tool_id, "result": payload},
            )
        )

    async def run(self, request: SupportRequest) -> Reply:
        return (await self.run_agent(request)).output
```

The end-hook payloads are the complete `AgentResult` and tool-result dictionary.
Use `tool_id` as an idempotency key when the receiving service may see retries.

## Workers

The module containing registered agents must be importable by each worker:

```bash
export WAYMARK_DATABASE_URL=postgresql://postgres:postgres@localhost:5432/waymark
export WAYMARK_USER_MODULE=examples.support_agent
export OPENAI_API_KEY=...
uv run waymark-start-workers
```

## Checks

```bash
uv run pytest -q
uv run ruff check .
uv run ty check
```
