Skip to content

Commit 8b0ea82

Browse files
committed
Add Deep Agents external storage sample
1 parent ae5df2f commit 8b0ea82

6 files changed

Lines changed: 467 additions & 4 deletions

File tree

‎deepagents_plugin/README.md‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ one side.
2626
| [subagents](subagents) | Durability propagates across the agent tree — sub-agent model calls become activities with no per-sub-agent wiring. |
2727
| [streaming](streaming) | Stream model chunks to external subscribers via `streaming_topic` + `WorkflowStream`, keeping the durable result identical. |
2828
| [langsmith_tracing](langsmith_tracing) | Compose `DeepAgentsPlugin` with `LangSmithPlugin` for durable execution + LLM tracing. |
29+
| [external_storage](external_storage) | Carry a 6 MiB tool result and the next model input through native S3 External Storage; scripted model, no API keys. |
2930

3031
## Prerequisites
3132

@@ -45,8 +46,10 @@ one side.
4546
> Python >= 3.11 — on older interpreters the group resolves to nothing
4647
> and the samples are skipped.
4748
48-
2. Configure a model provider. The samples use
49-
`anthropic:claude-sonnet-4-5`, which needs an Anthropic API key:
49+
2. Configure a model provider. Most samples use
50+
`anthropic:claude-sonnet-4-5`, which needs an Anthropic API key. The
51+
[external_storage](external_storage) sample uses a scripted model and needs
52+
no provider credentials:
5053

5154
```bash
5255
export ANTHROPIC_API_KEY=...
@@ -85,11 +88,13 @@ uv run deepagents_plugin/hello_world/run_worker.py
8588
uv run deepagents_plugin/hello_world/run_workflow.py
8689
```
8790

88-
The `langsmith_tracing` sample instead bundles the worker and starter into a
89-
single driver:
91+
The `langsmith_tracing` and `external_storage` samples instead bundle the worker
92+
and starter into a single driver. The storage sample also needs a mock S3 service;
93+
see its [setup instructions](external_storage):
9094

9195
```bash
9296
uv run deepagents_plugin/langsmith_tracing/main.py
97+
uv run deepagents_plugin/external_storage/main.py
9398
```
9499

95100
## Key Features Demonstrated
@@ -109,6 +114,8 @@ uv run deepagents_plugin/langsmith_tracing/main.py
109114
- **Streaming** — forward model chunks to external subscribers while keeping the
110115
durable result unchanged.
111116
- **Observability** — compose with `LangSmithPlugin` for tracing.
117+
- **Large payload storage** — compose native `ExternalStorage` with the plugin's
118+
data converter, keeping full tool results in S3 and small references in history.
112119

113120
## Related
114121

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
1+
# External Storage
2+
3+
Run a Deep Agent whose `read_file` tool returns **6 MiB** of text, exceeding
4+
Temporal's default 2 MiB payload limit. The SDK's native `ExternalStorage`
5+
offloads the tool output and the next model activity's input to S3, recording
6+
small references in workflow history and retrieving the full bytes on decode.
7+
The workflow returns only a byte count, SHA-256, and a short final answer.
8+
9+
Native External Storage is in **Public Preview**. This sample uses the SDK's
10+
`S3StorageDriver` with the local mock S3 service from the
11+
[general External Storage sample](../../external_storage). Neither AWS credentials,
12+
Docker, nor an LLM API key is needed. The model replies are scripted; the real
13+
Deep Agents loop, tool activity, model activities, and S3 transfers still run.
14+
15+
## Running the sample
16+
17+
Use Python **3.11 or later**. Run these commands from the repository root:
18+
19+
```bash
20+
uv sync --python 3.13 --group deepagents --group external-storage
21+
```
22+
23+
Start Temporal with its default payload limits in one terminal:
24+
25+
```bash
26+
temporal server start-dev
27+
```
28+
29+
In a second terminal, start the existing mock S3 service. It listens on port
30+
5000 and creates the `temporal-payloads` bucket:
31+
32+
```bash
33+
uv run external_storage/s3.py
34+
```
35+
36+
In a third terminal, run the demo. It starts a worker, executes one workflow,
37+
checks the result's integrity, and shuts down the worker:
38+
39+
```bash
40+
uv run deepagents_plugin/external_storage/main.py
41+
```
42+
43+
The output includes:
44+
45+
```text
46+
Tool result: 6,291,456 bytes (6 MiB), integrity verified
47+
Agent: Received the complete document.
48+
```
49+
50+
The script also prints the workflow ID and complete SHA-256. Inspect the history
51+
in the Temporal UI or with:
52+
53+
```bash
54+
temporal workflow show --workflow-id <printed-workflow-id>
55+
```
56+
57+
The completed `deepagents.invoke_tool` activity and the following
58+
`deepagents.invoke_model` input contain native storage references. The first
59+
model input and final workflow result remain small and inline.
60+
61+
To view externally stored payloads in the UI, reuse the
62+
[storage-aware codec server](../../external_storage#5-optional-run-the-codec-server):
63+
run `uv run external_storage/codec_server.py` and set the UI's Remote Codec
64+
Endpoint to `http://localhost:8081`. Its gzip decoder passes uncompressed payloads
65+
through, so it can also read this sample's references.
66+
67+
## How the configuration works
68+
69+
`client.py` creates `ExternalStorage(drivers=[driver],
70+
payload_size_threshold=256 * 1024)`. It uses the SDK's `SimplePlugin` converter
71+
hook to apply this configuration **after** `DeepAgentsPlugin` installs its
72+
LangChain-aware payload converter:
73+
74+
```python
75+
def configure(converter: DataConverter | None) -> DataConverter:
76+
return replace(converter or DataConverter.default, external_storage=storage)
77+
78+
client = await Client.connect(
79+
"localhost:7233",
80+
plugins=[
81+
create_plugin(),
82+
SimplePlugin("ExternalStorage", data_converter=configure),
83+
],
84+
)
85+
```
86+
87+
This order also supports the older plugin version in this repository's lockfile,
88+
which otherwise replaces the client's converter. The worker inherits the final
89+
converter. The S3 client stays open while the worker runs and the starter decodes
90+
results. No application code uploads or downloads payloads manually.
91+
92+
This demo deliberately uses **no compression codec**: repeated `x` characters
93+
would compress well below the storage threshold. It overrides the built-in
94+
`read_file` tool with an activity-backed mock bulk reader. Deep Agents excludes
95+
`read_file` from tool-result eviction, so the complete result crosses both
96+
activity boundaries. Applications using real models should page documents or
97+
retrieve excerpts to manage context separately. Storage reduces transport and
98+
history bytes; it does not reduce model tokens or decoded workflow memory.
99+
100+
With storage removed, this tool result fails with `PayloadsTooLarge` at the
101+
default limits. Raising the blob limit alone still leaves the gRPC message
102+
limit described in [Hanyu Liu's post](https://hanyuliu.me/posts/temporal-deep-agent-payload-limits/).
103+
Native storage avoids transmitting the large bytes through either limit.
104+
105+
For real S3, use an existing bucket and normal AWS credentials instead of this
106+
sample's local endpoint and mock credentials. Configure compatible storage
107+
drivers on every client, worker, and replayer that reads the history. Retain
108+
objects for as long as execution, reset, or replay may need them; Temporal does
109+
not delete stored payloads when an activity finishes.
110+
111+
## Tests
112+
113+
The integration test starts an isolated mock S3 service and verifies the full
114+
result, both native references, activity routing, and replay with a new S3
115+
client after worker shutdown.
116+
No manually running S3 service or model credentials are required:
117+
118+
```bash
119+
uv sync --python 3.13 --group deepagents --group external-storage --group dev
120+
uv run pytest tests/deepagents_plugin/external_storage_test.py
121+
```
Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,64 @@
1+
"""Compose native S3 storage with the Deep Agents data converter."""
2+
3+
from collections.abc import AsyncIterator
4+
from contextlib import asynccontextmanager
5+
from dataclasses import replace
6+
7+
import aioboto3
8+
from temporalio.client import Client
9+
from temporalio.contrib.aws.s3driver import S3StorageDriver
10+
from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client
11+
from temporalio.contrib.deepagents import DeepAgentsPlugin
12+
from temporalio.converter import DataConverter, ExternalStorage
13+
from temporalio.envconfig import ClientConfig
14+
from temporalio.plugin import SimplePlugin
15+
16+
S3_ENDPOINT = "http://localhost:5000"
17+
S3_BUCKET = "temporal-payloads"
18+
19+
20+
# @@@SNIPSTART python-deepagents-external-storage-converter
21+
def storage_plugin(storage: ExternalStorage) -> SimplePlugin:
22+
"""Add storage after the agent plugin has installed its payload converter."""
23+
24+
def configure(converter: DataConverter | None) -> DataConverter:
25+
return replace(converter or DataConverter.default, external_storage=storage)
26+
27+
# Applying this after DeepAgentsPlugin also supports older plugin versions
28+
# that replace the converter rather than composing with the client's one.
29+
return SimplePlugin("ExternalStorage", data_converter=configure)
30+
31+
32+
@asynccontextmanager
33+
async def connect_client(
34+
plugin: DeepAgentsPlugin,
35+
*,
36+
target_host: str | None = None,
37+
s3_endpoint: str = S3_ENDPOINT,
38+
) -> AsyncIterator[Client]:
39+
"""Keep the S3 client alive throughout workflow execution and payload reads."""
40+
config = ClientConfig.load_client_connect_config()
41+
config.setdefault("target_host", "localhost:7233")
42+
if target_host is not None:
43+
config["target_host"] = target_host
44+
45+
session = aioboto3.Session()
46+
async with session.client(
47+
"s3",
48+
endpoint_url=s3_endpoint,
49+
aws_access_key_id="test",
50+
aws_secret_access_key="test",
51+
region_name="us-east-1",
52+
) as s3_client:
53+
driver = S3StorageDriver(
54+
client=new_aioboto3_client(s3_client),
55+
bucket=S3_BUCKET,
56+
)
57+
storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024)
58+
yield await Client.connect(
59+
**config,
60+
plugins=[plugin, storage_plugin(storage)],
61+
)
62+
63+
64+
# @@@SNIPEND
Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
1+
"""Run one Deep Agent with a 6 MiB tool result and native S3 storage.
2+
3+
Both the worker and starter run in this process. The two model replies are
4+
scripted, so the demo needs no LLM credentials or provider context budget.
5+
The model and tool still run as real Temporal activities.
6+
"""
7+
8+
import asyncio
9+
import hashlib
10+
import uuid
11+
12+
from langchain_core.messages import AIMessage
13+
from temporalio.contrib.deepagents import DeepAgentsPlugin
14+
from temporalio.contrib.deepagents.testing import mock_model_provider
15+
from temporalio.worker import Worker
16+
17+
from deepagents_plugin.external_storage.client import connect_client
18+
from deepagents_plugin.external_storage.workflow import (
19+
PAYLOAD_BYTES,
20+
TASK_QUEUE,
21+
ExternalStorageAgent,
22+
)
23+
24+
25+
def create_plugin() -> DeepAgentsPlugin:
26+
"""Script a tool request and acknowledgment for this single-workflow demo."""
27+
return DeepAgentsPlugin(
28+
model_provider=mock_model_provider(
29+
[
30+
AIMessage(
31+
content="",
32+
tool_calls=[
33+
{
34+
"name": "read_file",
35+
"args": {"file_path": "large-document.txt"},
36+
"id": "call-read",
37+
}
38+
],
39+
),
40+
AIMessage(content="Received the complete document."),
41+
]
42+
)
43+
)
44+
45+
46+
async def main() -> None:
47+
async with connect_client(create_plugin()) as client:
48+
async with Worker(
49+
client,
50+
task_queue=TASK_QUEUE,
51+
workflows=[ExternalStorageAgent],
52+
max_cached_workflows=0,
53+
):
54+
workflow_id = f"deepagents-external-storage-{uuid.uuid4()}"
55+
result = await client.execute_workflow(
56+
ExternalStorageAgent.run,
57+
"Read large-document.txt and acknowledge receiving it.",
58+
id=workflow_id,
59+
task_queue=TASK_QUEUE,
60+
)
61+
62+
assert result.result_bytes == PAYLOAD_BYTES
63+
assert result.sha256 == hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest()
64+
print(f"Workflow: {workflow_id}")
65+
print(f"Tool result: {result.result_bytes:,} bytes (6 MiB), integrity verified")
66+
print(f"SHA-256: {result.sha256}")
67+
print(f"Agent: {result.answer}")
68+
69+
70+
if __name__ == "__main__":
71+
asyncio.run(main())
Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,75 @@
1+
"""A Deep Agent whose tool returns 6 MiB, carried by native External Storage."""
2+
3+
import hashlib
4+
from dataclasses import dataclass
5+
from datetime import timedelta
6+
7+
from langchain_core.messages import ToolMessage
8+
from langchain_core.tools import tool
9+
from temporalio import workflow
10+
from temporalio.common import RetryPolicy
11+
from temporalio.contrib.deepagents import (
12+
create_temporal_deep_agent,
13+
tool_as_activity,
14+
)
15+
16+
PAYLOAD_BYTES = 6 * 1024 * 1024
17+
MODEL = "fake:external-storage"
18+
TASK_QUEUE = "deepagents-external-storage"
19+
20+
21+
@dataclass
22+
class AgentResult:
23+
"""Small integrity evidence; the large document stays out of the final result."""
24+
25+
result_bytes: int
26+
sha256: str
27+
answer: str
28+
29+
30+
@tool
31+
async def read_file(file_path: str) -> str:
32+
"""Read a mock document containing exactly 6 MiB of text."""
33+
# Override the built-in read_file with an activity-backed bulk reader.
34+
# Deep Agents excludes read_file from tool-result eviction, so its complete
35+
# content reaches the next model activity. Real file I/O would go here.
36+
del file_path
37+
return "x" * PAYLOAD_BYTES
38+
39+
40+
# @@@SNIPSTART python-deepagents-external-storage-workflow
41+
@workflow.defn
42+
class ExternalStorageAgent:
43+
@workflow.run
44+
async def run(self, question: str) -> AgentResult:
45+
reader = tool_as_activity(
46+
read_file,
47+
start_to_close_timeout=timedelta(seconds=30),
48+
activity_options={"retry_policy": RetryPolicy(maximum_attempts=1)},
49+
)
50+
agent = create_temporal_deep_agent(
51+
model=MODEL,
52+
tools=[reader],
53+
system_prompt="Read the document, then acknowledge receiving the result.",
54+
activity_options={"start_to_close_timeout": timedelta(seconds=30)},
55+
)
56+
result = await agent.ainvoke(
57+
{"messages": [{"role": "user", "content": question}]}
58+
)
59+
tool_result = next(
60+
message
61+
for message in result["messages"]
62+
if isinstance(message, ToolMessage) and message.name == "read_file"
63+
)
64+
assert isinstance(tool_result.content, str)
65+
data = tool_result.content.encode()
66+
answer = result["messages"][-1].content
67+
assert isinstance(answer, str)
68+
return AgentResult(
69+
result_bytes=len(data),
70+
sha256=hashlib.sha256(data).hexdigest(),
71+
answer=answer,
72+
)
73+
74+
75+
# @@@SNIPEND

0 commit comments

Comments
 (0)