Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 11 additions & 4 deletions deepagents_plugin/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ one side.
| [subagents](subagents) | Durability propagates across the agent tree — sub-agent model calls become activities with no per-sub-agent wiring. |
| [streaming](streaming) | Stream model chunks to external subscribers via `streaming_topic` + `WorkflowStream`, keeping the durable result identical. |
| [langsmith_tracing](langsmith_tracing) | Compose `DeepAgentsPlugin` with `LangSmithPlugin` for durable execution + LLM tracing. |
| [external_storage](external_storage) | Use the Temporal data converter's native S3 External Storage for a 6 MiB tool result and multi-turn conversation payloads; scripted model, no API keys. |

## Prerequisites

Expand All @@ -45,8 +46,10 @@ one side.
> Temporal 1.34.0 or later and Python >= 3.11 — on older interpreters the group resolves to nothing
> and the samples are skipped.

2. Configure a model provider. The samples use
`anthropic:claude-sonnet-4-5`, which needs an Anthropic API key:
2. Configure a model provider. Most samples use
`anthropic:claude-sonnet-4-5`, which needs an Anthropic API key. The
[external_storage](external_storage) sample uses a scripted model and needs
no provider credentials:

```bash
export ANTHROPIC_API_KEY=...
Expand Down Expand Up @@ -85,11 +88,13 @@ uv run deepagents_plugin/hello_world/run_worker.py
uv run deepagents_plugin/hello_world/run_workflow.py
```

The `langsmith_tracing` sample instead bundles the worker and starter into a
single driver:
The `langsmith_tracing` and `external_storage` samples instead bundle the worker
and starter into a single driver. The storage sample also needs a mock S3 service;
see their [setup instructions](../external_storage):

```bash
uv run deepagents_plugin/langsmith_tracing/main.py
uv run deepagents_plugin/external_storage/main.py
```

## Key Features Demonstrated
Expand All @@ -109,6 +114,8 @@ uv run deepagents_plugin/langsmith_tracing/main.py
- **Streaming** — forward model chunks to external subscribers while keeping the
durable result unchanged.
- **Observability** — compose with `LangSmithPlugin` for tracing.
- **Large payload storage** — configure native `ExternalStorage` on the client's
data converter; the standalone plugin composes it with its LangChain converter.

## Related

Expand Down
75 changes: 75 additions & 0 deletions deepagents_plugin/external_storage/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
# External Storage

This sample combines two cases for Temporal's native `ExternalStorage` data
converter:

* A Deep Agent tool returns a **6 MiB** document. The tool result and the next
model activity input are stored in S3, with small references in workflow
history.
* The workflow carries a multi-turn conversation. Each model call receives the
earlier turns as context. The same converter externalizes those activity
payloads when they exceed the configured threshold.

The sample does not use a DeepAgents backend or manage conversation objects
itself. Temporal's converter externalizes large serialized payloads; it is not
a key-value conversation store. Small conversation state is reconstructed by
workflow replay, while large activity payloads are stored by the configured
S3 driver. The demo threshold is deliberately low (128 bytes) so conversation
payloads also exercise the converter. In a real application, choose a threshold
appropriate for your payloads.

The S3 driver uses the local mock service from the
[general External Storage sample](../../external_storage). No AWS credentials,
Docker, or LLM API key is needed. The model replies are scripted; the real Deep
Agents loop, tool activity, model activities, and S3 transfers still run.

## Run it

Use Python **3.11 or later**:

```bash
uv sync --python 3.13 --group deepagents --group external-storage
```

Start Temporal in one terminal:

```bash
temporal server start-dev
```

Start the mock S3 service in another terminal. It listens on port 5000 and
creates the `temporal-payloads` bucket:

```bash
uv run external_storage/s3.py
```

Run the demo:

```bash
uv run deepagents_plugin/external_storage/main.py
```

The output includes the workflow ID, verified document size and SHA-256, number
of conversation turns, and final answer. Inspect the workflow history with:

```bash
temporal workflow show --workflow-id <printed-workflow-id>
```

Activity inputs and results above the threshold appear as native storage
references. To view externally stored payloads in the UI, run the
[storage-aware codec server](../../external_storage#5-optional-run-the-codec-server)
and set the UI's Remote Codec Endpoint to `http://localhost:8081`.

The client configures `ExternalStorage` on the SDK's default data converter and
passes that converter to `DeepAgentsPlugin`. The worker inherits the converter;
application code does not manually upload or download Temporal payloads. For
real S3, use an existing bucket and normal AWS credentials. Configure compatible
storage drivers on every client, worker, and replayer that reads the history,
and retain objects for as long as execution, reset, or replay may need them.

External Storage reduces bytes in Temporal history and transport. It does not
reduce the number of history events, the model context size, or the total number
of externally stored bytes. Carrying the full transcript into each model call
still increases total conversation payload storage over time.
49 changes: 49 additions & 0 deletions deepagents_plugin/external_storage/client.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
"""Configure native S3 External Storage for the Deep Agents plugin."""

from collections.abc import AsyncIterator, Callable
from contextlib import asynccontextmanager
from dataclasses import replace

import aioboto3
from temporalio.client import Client
from temporalio.contrib.aws.s3driver import S3StorageDriver
from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client
from temporalio.converter import DataConverter, ExternalStorage
from temporalio.deepagents import DeepAgentsPlugin
from temporalio.envconfig import ClientConfig

S3_ENDPOINT = "http://localhost:5000"
S3_BUCKET = "temporal-payloads"


@asynccontextmanager
async def connect_client(
plugin_factory: Callable[[DataConverter], DeepAgentsPlugin],
*,
target_host: str | None = None,
s3_endpoint: str = S3_ENDPOINT,
) -> AsyncIterator[Client]:
"""Keep the S3 client alive while Temporal reads and writes payloads."""
config = ClientConfig.load_client_connect_config()
config.setdefault("target_host", "localhost:7233")
if target_host is not None:
config["target_host"] = target_host

session = aioboto3.Session()
async with session.client(
"s3",
endpoint_url=s3_endpoint,
aws_access_key_id="test",
aws_secret_access_key="test",
region_name="us-east-1",
) as s3_client:
driver = S3StorageDriver(
client=new_aioboto3_client(s3_client),
bucket=S3_BUCKET,
)
# Small on purpose: this demo also externalizes serialized conversation
# context, while the large tool result shows the same path at scale.
storage = ExternalStorage(drivers=[driver], payload_size_threshold=128)
data_converter = replace(DataConverter.default, external_storage=storage)
plugin = plugin_factory(data_converter)
yield await Client.connect(**config, plugins=[plugin])
81 changes: 81 additions & 0 deletions deepagents_plugin/external_storage/main.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
"""Run one scripted conversation with native Temporal S3 External Storage."""

import asyncio
import hashlib
import uuid

from langchain_core.messages import AIMessage
from temporalio.converter import DataConverter
from temporalio.deepagents import DeepAgentsPlugin
from temporalio.deepagents.testing import mock_model_provider
from temporalio.worker import Worker

from deepagents_plugin.external_storage.client import connect_client
from deepagents_plugin.external_storage.workflow import (
PAYLOAD_BYTES,
TASK_QUEUE,
ExternalStorageAgent,
)


def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin:
"""Script the model so the demo needs no LLM provider or API key."""
return DeepAgentsPlugin(
data_converter=data_converter,
model_provider=mock_model_provider(
[
AIMessage(
content="",
tool_calls=[
{
"name": "read_file",
"args": {"file_path": "large-document.txt"},
"id": "call-read",
}
],
),
AIMessage(content="I received the complete document."),
AIMessage(content="The project is called Cedar."),
AIMessage(content="Cedar uses Python."),
AIMessage(content="Cedar is due Friday."),
AIMessage(content="Cedar is a Python project due Friday."),
]
),
)


async def main() -> None:
questions = [
"Read large-document.txt and acknowledge receiving it.",
"I am working on a project called Cedar.",
"The project uses Python.",
"It is due Friday.",
"Remind me of the project name, language, and deadline.",
]
async with connect_client(create_plugin) as client:
async with Worker(
client,
task_queue=TASK_QUEUE,
workflows=[ExternalStorageAgent],
max_cached_workflows=0,
):
workflow_id = f"deepagents-external-storage-{uuid.uuid4()}"
result = await client.execute_workflow(
ExternalStorageAgent.run,
questions,
id=workflow_id,
task_queue=TASK_QUEUE,
)

assert result.tool_result_bytes == PAYLOAD_BYTES
assert result.tool_result_sha256 == hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest()
print(f"Workflow: {workflow_id}")
print(
f"Tool result: {result.tool_result_bytes:,} bytes (6 MiB), integrity verified"
)
print(f"Conversation turns: {result.turns}")
print(f"Final answer: {result.last_answer}")


if __name__ == "__main__":
asyncio.run(main())
95 changes: 95 additions & 0 deletions deepagents_plugin/external_storage/workflow.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
"""Exercise native External Storage with a large tool result and a conversation."""

import hashlib
from dataclasses import dataclass
from datetime import timedelta

from langchain_core.messages import ToolMessage
from langchain_core.tools import tool
from temporalio import workflow
from temporalio.common import RetryPolicy
from temporalio.deepagents import create_temporal_deep_agent, tool_as_activity

PAYLOAD_BYTES = 6 * 1024 * 1024
MODEL = "fake:external-storage"
TASK_QUEUE = "deepagents-external-storage"


@dataclass
class ConversationTurn:
user: str
assistant: str


@dataclass
class ExternalStorageResult:
tool_result_bytes: int
tool_result_sha256: str
turns: int
last_answer: str


@tool
async def read_file(file_path: str) -> str:
"""Read a mock document containing exactly 6 MiB of text."""
del file_path
return "x" * PAYLOAD_BYTES


@workflow.defn
class ExternalStorageAgent:
@workflow.run
async def run(self, questions: list[str]) -> ExternalStorageResult:
reader = tool_as_activity(
read_file,
start_to_close_timeout=timedelta(seconds=30),
activity_options={"retry_policy": RetryPolicy(maximum_attempts=1)},
)
turns: list[ConversationTurn] = []
tool_result: str | None = None
last_answer = ""

for index, question in enumerate(questions):
history = "\n".join(
f"User: {turn.user}\nAssistant: {turn.assistant}" for turn in turns
)
system_prompt = (
"Answer the latest message using the conversation history. "
"For the first request, read the requested document and acknowledge "
"receiving it.\n\n"
f"Conversation history:\n{history or '(no earlier turns)'}"
)
agent = create_temporal_deep_agent(
model=MODEL,
tools=[reader] if index == 0 else [],
system_prompt=system_prompt,
activity_options={"start_to_close_timeout": timedelta(seconds=30)},
)
result = await agent.ainvoke(
{"messages": [{"role": "user", "content": question}]}
)
if index == 0:
message = next(
message
for message in result["messages"]
if isinstance(message, ToolMessage) and message.name == "read_file"
)
if not isinstance(message.content, str):
raise TypeError("Expected the document tool result to be text")
tool_result = message.content

answer = result["messages"][-1].content
if not isinstance(answer, str):
raise TypeError("Expected a text response from the model")
turns.append(ConversationTurn(user=question, assistant=answer))
last_answer = answer

if tool_result is None:
raise ValueError("At least one question is required")
tool_bytes = tool_result.encode()
return ExternalStorageResult(
tool_result_bytes=len(tool_bytes),
tool_result_sha256=hashlib.sha256(tool_bytes).hexdigest(),
turns=len(turns),
last_answer=last_answer,
)
Loading
Loading