From 155a16b45af06acf229c8d78877d71ffd4777244 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Thu, 8 Oct 2026 21:56:55 -0700 Subject: [PATCH 1/7] Add Deep Agents external storage sample --- deepagents_plugin/README.md | 15 ++- deepagents_plugin/external_storage/README.md | 121 +++++++++++++++++ deepagents_plugin/external_storage/client.py | 64 +++++++++ deepagents_plugin/external_storage/main.py | 71 ++++++++++ .../external_storage/workflow.py | 75 +++++++++++ .../external_storage_test.py | 125 ++++++++++++++++++ 6 files changed, 467 insertions(+), 4 deletions(-) create mode 100644 deepagents_plugin/external_storage/README.md create mode 100644 deepagents_plugin/external_storage/client.py create mode 100644 deepagents_plugin/external_storage/main.py create mode 100644 deepagents_plugin/external_storage/workflow.py create mode 100644 tests/deepagents_plugin/external_storage_test.py diff --git a/deepagents_plugin/README.md b/deepagents_plugin/README.md index 7f5d435b..d5fb9029 100644 --- a/deepagents_plugin/README.md +++ b/deepagents_plugin/README.md @@ -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) | Carry a 6 MiB tool result and the next model input through native S3 External Storage; scripted model, no API keys. | ## Prerequisites @@ -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=... @@ -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 its [setup instructions](external_storage): ```bash uv run deepagents_plugin/langsmith_tracing/main.py +uv run deepagents_plugin/external_storage/main.py ``` ## Key Features Demonstrated @@ -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** — compose native `ExternalStorage` with the plugin's + data converter, keeping full tool results in S3 and small references in history. ## Related diff --git a/deepagents_plugin/external_storage/README.md b/deepagents_plugin/external_storage/README.md new file mode 100644 index 00000000..3fe8af72 --- /dev/null +++ b/deepagents_plugin/external_storage/README.md @@ -0,0 +1,121 @@ +# External Storage + +Run a Deep Agent whose `read_file` tool returns **6 MiB** of text, exceeding +Temporal's default 2 MiB payload limit. The SDK's native `ExternalStorage` +offloads the tool output and the next model activity's input to S3, recording +small references in workflow history and retrieving the full bytes on decode. +The workflow returns only a byte count, SHA-256, and a short final answer. + +Native External Storage is in **Public Preview**. This sample uses the SDK's +`S3StorageDriver` with the local mock S3 service from the +[general External Storage sample](../../external_storage). Neither AWS credentials, +Docker, nor an LLM API key is needed. The model replies are scripted; the real +Deep Agents loop, tool activity, model activities, and S3 transfers still run. + +## Running the sample + +Use Python **3.11 or later**. Run these commands from the repository root: + +```bash +uv sync --python 3.13 --group deepagents --group external-storage +``` + +Start Temporal with its default payload limits in one terminal: + +```bash +temporal server start-dev +``` + +In a second terminal, start the existing mock S3 service. It listens on port +5000 and creates the `temporal-payloads` bucket: + +```bash +uv run external_storage/s3.py +``` + +In a third terminal, run the demo. It starts a worker, executes one workflow, +checks the result's integrity, and shuts down the worker: + +```bash +uv run deepagents_plugin/external_storage/main.py +``` + +The output includes: + +```text +Tool result: 6,291,456 bytes (6 MiB), integrity verified +Agent: Received the complete document. +``` + +The script also prints the workflow ID and complete SHA-256. Inspect the history +in the Temporal UI or with: + +```bash +temporal workflow show --workflow-id +``` + +The completed `deepagents.invoke_tool` activity and the following +`deepagents.invoke_model` input contain native storage references. The first +model input and final workflow result remain small and inline. + +To view externally stored payloads in the UI, reuse the +[storage-aware codec server](../../external_storage#5-optional-run-the-codec-server): +run `uv run external_storage/codec_server.py` and set the UI's Remote Codec +Endpoint to `http://localhost:8081`. Its gzip decoder passes uncompressed payloads +through, so it can also read this sample's references. + +## How the configuration works + +`client.py` creates `ExternalStorage(drivers=[driver], +payload_size_threshold=256 * 1024)`. It uses the SDK's `SimplePlugin` converter +hook to apply this configuration **after** `DeepAgentsPlugin` installs its +LangChain-aware payload converter: + +```python +def configure(converter: DataConverter | None) -> DataConverter: + return replace(converter or DataConverter.default, external_storage=storage) + +client = await Client.connect( + "localhost:7233", + plugins=[ + create_plugin(), + SimplePlugin("ExternalStorage", data_converter=configure), + ], +) +``` + +This order also supports the older plugin version in this repository's lockfile, +which otherwise replaces the client's converter. The worker inherits the final +converter. The S3 client stays open while the worker runs and the starter decodes +results. No application code uploads or downloads payloads manually. + +This demo deliberately uses **no compression codec**: repeated `x` characters +would compress well below the storage threshold. It overrides the built-in +`read_file` tool with an activity-backed mock bulk reader. Deep Agents excludes +`read_file` from tool-result eviction, so the complete result crosses both +activity boundaries. Applications using real models should page documents or +retrieve excerpts to manage context separately. Storage reduces transport and +history bytes; it does not reduce model tokens or decoded workflow memory. + +With storage removed, this tool result fails with `PayloadsTooLarge` at the +default limits. Raising the blob limit alone still leaves the gRPC message +limit described in [Hanyu Liu's post](https://hanyuliu.me/posts/temporal-deep-agent-payload-limits/). +Native storage avoids transmitting the large bytes through either limit. + +For real S3, use an existing bucket and normal AWS credentials instead of this +sample's local endpoint and mock credentials. Configure compatible storage +drivers on every client, worker, and replayer that reads the history. Retain +objects for as long as execution, reset, or replay may need them; Temporal does +not delete stored payloads when an activity finishes. + +## Tests + +The integration test starts an isolated mock S3 service and verifies the full +result, both native references, activity routing, and replay with a new S3 +client after worker shutdown. +No manually running S3 service or model credentials are required: + +```bash +uv sync --python 3.13 --group deepagents --group external-storage --group dev +uv run pytest tests/deepagents_plugin/external_storage_test.py +``` diff --git a/deepagents_plugin/external_storage/client.py b/deepagents_plugin/external_storage/client.py new file mode 100644 index 00000000..92fdca48 --- /dev/null +++ b/deepagents_plugin/external_storage/client.py @@ -0,0 +1,64 @@ +"""Compose native S3 storage with the Deep Agents data converter.""" + +from collections.abc import AsyncIterator +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.contrib.deepagents import DeepAgentsPlugin +from temporalio.converter import DataConverter, ExternalStorage +from temporalio.envconfig import ClientConfig +from temporalio.plugin import SimplePlugin + +S3_ENDPOINT = "http://localhost:5000" +S3_BUCKET = "temporal-payloads" + + +# @@@SNIPSTART python-deepagents-external-storage-converter +def storage_plugin(storage: ExternalStorage) -> SimplePlugin: + """Add storage after the agent plugin has installed its payload converter.""" + + def configure(converter: DataConverter | None) -> DataConverter: + return replace(converter or DataConverter.default, external_storage=storage) + + # Applying this after DeepAgentsPlugin also supports older plugin versions + # that replace the converter rather than composing with the client's one. + return SimplePlugin("ExternalStorage", data_converter=configure) + + +@asynccontextmanager +async def connect_client( + plugin: DeepAgentsPlugin, + *, + target_host: str | None = None, + s3_endpoint: str = S3_ENDPOINT, +) -> AsyncIterator[Client]: + """Keep the S3 client alive throughout workflow execution and payload reads.""" + 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, + ) + storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) + yield await Client.connect( + **config, + plugins=[plugin, storage_plugin(storage)], + ) + + +# @@@SNIPEND diff --git a/deepagents_plugin/external_storage/main.py b/deepagents_plugin/external_storage/main.py new file mode 100644 index 00000000..57a7be05 --- /dev/null +++ b/deepagents_plugin/external_storage/main.py @@ -0,0 +1,71 @@ +"""Run one Deep Agent with a 6 MiB tool result and native S3 storage. + +Both the worker and starter run in this process. The two model replies are +scripted, so the demo needs no LLM credentials or provider context budget. +The model and tool still run as real Temporal activities. +""" + +import asyncio +import hashlib +import uuid + +from langchain_core.messages import AIMessage +from temporalio.contrib.deepagents import DeepAgentsPlugin +from temporalio.contrib.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() -> DeepAgentsPlugin: + """Script a tool request and acknowledgment for this single-workflow demo.""" + return DeepAgentsPlugin( + model_provider=mock_model_provider( + [ + AIMessage( + content="", + tool_calls=[ + { + "name": "read_file", + "args": {"file_path": "large-document.txt"}, + "id": "call-read", + } + ], + ), + AIMessage(content="Received the complete document."), + ] + ) + ) + + +async def main() -> None: + 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, + "Read large-document.txt and acknowledge receiving it.", + id=workflow_id, + task_queue=TASK_QUEUE, + ) + + assert result.result_bytes == PAYLOAD_BYTES + assert result.sha256 == hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest() + print(f"Workflow: {workflow_id}") + print(f"Tool result: {result.result_bytes:,} bytes (6 MiB), integrity verified") + print(f"SHA-256: {result.sha256}") + print(f"Agent: {result.answer}") + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/deepagents_plugin/external_storage/workflow.py b/deepagents_plugin/external_storage/workflow.py new file mode 100644 index 00000000..2c339b81 --- /dev/null +++ b/deepagents_plugin/external_storage/workflow.py @@ -0,0 +1,75 @@ +"""A Deep Agent whose tool returns 6 MiB, carried by native External Storage.""" + +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.contrib.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 AgentResult: + """Small integrity evidence; the large document stays out of the final result.""" + + result_bytes: int + sha256: str + answer: str + + +@tool +async def read_file(file_path: str) -> str: + """Read a mock document containing exactly 6 MiB of text.""" + # Override the built-in read_file with an activity-backed bulk reader. + # Deep Agents excludes read_file from tool-result eviction, so its complete + # content reaches the next model activity. Real file I/O would go here. + del file_path + return "x" * PAYLOAD_BYTES + + +# @@@SNIPSTART python-deepagents-external-storage-workflow +@workflow.defn +class ExternalStorageAgent: + @workflow.run + async def run(self, question: str) -> AgentResult: + reader = tool_as_activity( + read_file, + start_to_close_timeout=timedelta(seconds=30), + activity_options={"retry_policy": RetryPolicy(maximum_attempts=1)}, + ) + agent = create_temporal_deep_agent( + model=MODEL, + tools=[reader], + system_prompt="Read the document, then acknowledge receiving the result.", + activity_options={"start_to_close_timeout": timedelta(seconds=30)}, + ) + result = await agent.ainvoke( + {"messages": [{"role": "user", "content": question}]} + ) + tool_result = next( + message + for message in result["messages"] + if isinstance(message, ToolMessage) and message.name == "read_file" + ) + assert isinstance(tool_result.content, str) + data = tool_result.content.encode() + answer = result["messages"][-1].content + assert isinstance(answer, str) + return AgentResult( + result_bytes=len(data), + sha256=hashlib.sha256(data).hexdigest(), + answer=answer, + ) + + +# @@@SNIPEND diff --git a/tests/deepagents_plugin/external_storage_test.py b/tests/deepagents_plugin/external_storage_test.py new file mode 100644 index 00000000..c9aa3870 --- /dev/null +++ b/tests/deepagents_plugin/external_storage_test.py @@ -0,0 +1,125 @@ +"""Verify the sample against real Temporal activities and a local S3 service.""" + +import hashlib +import uuid +from collections.abc import Iterator +from datetime import timedelta + +import boto3 +import pytest +from moto.server import ThreadedMotoServer +from temporalio.api.enums.v1 import EventType +from temporalio.api.sdk.v1 import ExternalStorageReference +from temporalio.client import Client +from temporalio.contrib.deepagents import DeepAgentsPlugin +from temporalio.converter import DataConverter +from temporalio.worker import Replayer, Worker + +from deepagents_plugin.external_storage.client import ( + S3_BUCKET, + connect_client, + storage_plugin, +) +from deepagents_plugin.external_storage.main import create_plugin +from deepagents_plugin.external_storage.workflow import ( + PAYLOAD_BYTES, + AgentResult, + ExternalStorageAgent, +) +from tests.deepagents_plugin.helpers import ( + INVOKE_MODEL, + INVOKE_TOOL, + count_scheduled_activities, +) + + +@pytest.fixture +def s3_endpoint() -> Iterator[str]: + server = ThreadedMotoServer(ip_address="127.0.0.1", port=0, verbose=False) + server.start() + try: + host, port = server.get_host_and_port() + endpoint = f"http://{host}:{port}" + s3 = boto3.client( + "s3", + endpoint_url=endpoint, + aws_access_key_id="test", + aws_secret_access_key="test", + region_name="us-east-1", + ) + s3.create_bucket(Bucket=S3_BUCKET) + yield endpoint + finally: + server.stop() + + +async def test_external_storage(client: Client, s3_endpoint: str) -> None: + queue = f"deepagents-storage-{uuid.uuid4()}" + target_host = client.service_client.config.target_host + async with connect_client( + create_plugin(), target_host=target_host, s3_endpoint=s3_endpoint + ) as stored_client: + assert stored_client.data_converter.external_storage is not None + async with Worker( + stored_client, + task_queue=queue, + workflows=[ExternalStorageAgent], + max_cached_workflows=0, + ): + handle = await stored_client.start_workflow( + ExternalStorageAgent.run, + "Read large-document.txt and acknowledge receiving it.", + id=queue, + task_queue=queue, + execution_timeout=timedelta(seconds=60), + ) + result = await handle.result() + history = await handle.fetch_history() + + assert result == AgentResult( + result_bytes=PAYLOAD_BYTES, + sha256=hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest(), + answer="Received the complete document.", + ) + counts = await count_scheduled_activities(handle) + assert counts[INVOKE_MODEL] == 2 + assert counts[INVOKE_TOOL] == 1 + + scheduled = { + event.event_id: event.activity_task_scheduled_event_attributes + for event in history.events + if event.event_type == EventType.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED + } + tool_result = next( + event.activity_task_completed_event_attributes.result.payloads[0] + for event in history.events + if event.event_type == EventType.EVENT_TYPE_ACTIVITY_TASK_COMPLETED + and scheduled[ + event.activity_task_completed_event_attributes.scheduled_event_id + ].activity_type.name + == INVOKE_TOOL + ) + model_input = list(scheduled.values())[-1].input.payloads[0] + for payload in (tool_result, model_input): + reference = DataConverter.default.payload_converter.from_payload(payload) + assert isinstance(reference, ExternalStorageReference) + assert reference.driver_name == "aws.s3driver" + assert payload.ByteSize() < 1024 + assert "x" * 1024 not in history.to_json() + + # The original worker and S3 client are closed. Retrieve persisted claims + # through a fresh client and driver, without executing the activities again. + async with connect_client( + DeepAgentsPlugin(), target_host=target_host, s3_endpoint=s3_endpoint + ) as replay_client: + storage = replay_client.data_converter.external_storage + assert storage is not None + await Replayer( + workflows=[ExternalStorageAgent], + data_converter=replay_client.data_converter, + plugins=[ + DeepAgentsPlugin(), + # Preserve storage after the older plugin's converter override. + storage_plugin(storage), + ], + ).replay_workflow(history) From d9ef09fcdeb04d3dd7dd8ac1c5ca50c37bd38898 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Thu, 8 Oct 2026 22:07:19 -0700 Subject: [PATCH 2/7] Use standalone Deep Agents plugin for external storage --- deepagents_plugin/README.md | 4 +-- deepagents_plugin/external_storage/README.md | 36 ++++++------------- deepagents_plugin/external_storage/client.py | 20 +++-------- deepagents_plugin/external_storage/main.py | 4 +-- .../external_storage/workflow.py | 2 +- .../external_storage_test.py | 11 ++---- 6 files changed, 22 insertions(+), 55 deletions(-) diff --git a/deepagents_plugin/README.md b/deepagents_plugin/README.md index d5fb9029..688eaee4 100644 --- a/deepagents_plugin/README.md +++ b/deepagents_plugin/README.md @@ -114,8 +114,8 @@ uv run deepagents_plugin/external_storage/main.py - **Streaming** — forward model chunks to external subscribers while keeping the durable result unchanged. - **Observability** — compose with `LangSmithPlugin` for tracing. -- **Large payload storage** — compose native `ExternalStorage` with the plugin's - data converter, keeping full tool results in S3 and small references in history. +- **Large payload storage** — configure native `ExternalStorage` on the client's + data converter; the standalone plugin composes it with its LangChain converter. ## Related diff --git a/deepagents_plugin/external_storage/README.md b/deepagents_plugin/external_storage/README.md index 3fe8af72..a4cca761 100644 --- a/deepagents_plugin/external_storage/README.md +++ b/deepagents_plugin/external_storage/README.md @@ -66,28 +66,13 @@ through, so it can also read this sample's references. ## How the configuration works -`client.py` creates `ExternalStorage(drivers=[driver], -payload_size_threshold=256 * 1024)`. It uses the SDK's `SimplePlugin` converter -hook to apply this configuration **after** `DeepAgentsPlugin` installs its -LangChain-aware payload converter: - -```python -def configure(converter: DataConverter | None) -> DataConverter: - return replace(converter or DataConverter.default, external_storage=storage) - -client = await Client.connect( - "localhost:7233", - plugins=[ - create_plugin(), - SimplePlugin("ExternalStorage", data_converter=configure), - ], -) -``` - -This order also supports the older plugin version in this repository's lockfile, -which otherwise replaces the client's converter. The worker inherits the final -converter. The S3 client stays open while the worker runs and the starter decodes -results. No application code uploads or downloads payloads manually. +`client.py` opens an S3 client, creates an SDK `S3StorageDriver`, and passes a +`DataConverter` with `ExternalStorage` directly to `Client.connect`. The standalone +`temporalio-deepagents` plugin composes its LangChain payload converter with that +converter, preserving the storage configuration. The worker uses the client's +converter, and the sample keeps the S3 client open while the worker runs and the +starter decodes results. No application code uploads or downloads payloads +manually. This demo deliberately uses **no compression codec**: repeated `x` characters would compress well below the storage threshold. It overrides the built-in @@ -105,15 +90,14 @@ Native storage avoids transmitting the large bytes through either limit. For real S3, use an existing bucket and normal AWS credentials instead of this sample's local endpoint and mock credentials. Configure compatible storage drivers on every client, worker, and replayer that reads the history. Retain -objects for as long as execution, reset, or replay may need them; Temporal does -not delete stored payloads when an activity finishes. +objects for as long as execution, reset, or replay may need them; Temporal does not delete stored payloads when an activity finishes. ## Tests The integration test starts an isolated mock S3 service and verifies the full result, both native references, activity routing, and replay with a new S3 -client after worker shutdown. -No manually running S3 service or model credentials are required: +client after worker shutdown. No manually running S3 service or model credentials +are required: ```bash uv sync --python 3.13 --group deepagents --group external-storage --group dev diff --git a/deepagents_plugin/external_storage/client.py b/deepagents_plugin/external_storage/client.py index 92fdca48..4cbb8ab8 100644 --- a/deepagents_plugin/external_storage/client.py +++ b/deepagents_plugin/external_storage/client.py @@ -1,4 +1,4 @@ -"""Compose native S3 storage with the Deep Agents data converter.""" +"""Configure native S3 storage for the Deep Agents client.""" from collections.abc import AsyncIterator from contextlib import asynccontextmanager @@ -8,27 +8,15 @@ from temporalio.client import Client from temporalio.contrib.aws.s3driver import S3StorageDriver from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client -from temporalio.contrib.deepagents import DeepAgentsPlugin from temporalio.converter import DataConverter, ExternalStorage from temporalio.envconfig import ClientConfig -from temporalio.plugin import SimplePlugin +from temporalio.deepagents import DeepAgentsPlugin S3_ENDPOINT = "http://localhost:5000" S3_BUCKET = "temporal-payloads" # @@@SNIPSTART python-deepagents-external-storage-converter -def storage_plugin(storage: ExternalStorage) -> SimplePlugin: - """Add storage after the agent plugin has installed its payload converter.""" - - def configure(converter: DataConverter | None) -> DataConverter: - return replace(converter or DataConverter.default, external_storage=storage) - - # Applying this after DeepAgentsPlugin also supports older plugin versions - # that replace the converter rather than composing with the client's one. - return SimplePlugin("ExternalStorage", data_converter=configure) - - @asynccontextmanager async def connect_client( plugin: DeepAgentsPlugin, @@ -55,9 +43,11 @@ async def connect_client( bucket=S3_BUCKET, ) storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) + converter = replace(DataConverter.default, external_storage=storage) yield await Client.connect( **config, - plugins=[plugin, storage_plugin(storage)], + data_converter=converter, + plugins=[plugin], ) diff --git a/deepagents_plugin/external_storage/main.py b/deepagents_plugin/external_storage/main.py index 57a7be05..e0bc3f22 100644 --- a/deepagents_plugin/external_storage/main.py +++ b/deepagents_plugin/external_storage/main.py @@ -10,8 +10,8 @@ import uuid from langchain_core.messages import AIMessage -from temporalio.contrib.deepagents import DeepAgentsPlugin -from temporalio.contrib.deepagents.testing import mock_model_provider +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 diff --git a/deepagents_plugin/external_storage/workflow.py b/deepagents_plugin/external_storage/workflow.py index 2c339b81..49842b3b 100644 --- a/deepagents_plugin/external_storage/workflow.py +++ b/deepagents_plugin/external_storage/workflow.py @@ -8,7 +8,7 @@ from langchain_core.tools import tool from temporalio import workflow from temporalio.common import RetryPolicy -from temporalio.contrib.deepagents import ( +from temporalio.deepagents import ( create_temporal_deep_agent, tool_as_activity, ) diff --git a/tests/deepagents_plugin/external_storage_test.py b/tests/deepagents_plugin/external_storage_test.py index c9aa3870..4d28dbdf 100644 --- a/tests/deepagents_plugin/external_storage_test.py +++ b/tests/deepagents_plugin/external_storage_test.py @@ -11,14 +11,13 @@ from temporalio.api.enums.v1 import EventType from temporalio.api.sdk.v1 import ExternalStorageReference from temporalio.client import Client -from temporalio.contrib.deepagents import DeepAgentsPlugin +from temporalio.deepagents import DeepAgentsPlugin from temporalio.converter import DataConverter from temporalio.worker import Replayer, Worker from deepagents_plugin.external_storage.client import ( S3_BUCKET, connect_client, - storage_plugin, ) from deepagents_plugin.external_storage.main import create_plugin from deepagents_plugin.external_storage.workflow import ( @@ -112,14 +111,8 @@ async def test_external_storage(client: Client, s3_endpoint: str) -> None: async with connect_client( DeepAgentsPlugin(), target_host=target_host, s3_endpoint=s3_endpoint ) as replay_client: - storage = replay_client.data_converter.external_storage - assert storage is not None await Replayer( workflows=[ExternalStorageAgent], data_converter=replay_client.data_converter, - plugins=[ - DeepAgentsPlugin(), - # Preserve storage after the older plugin's converter override. - storage_plugin(storage), - ], + plugins=[DeepAgentsPlugin()], ).replay_workflow(history) From f8b152c13e697049400d1f73952330c346072416 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Thu, 8 Oct 2026 22:10:46 -0700 Subject: [PATCH 3/7] Fix import ordering for sample CI --- deepagents_plugin/external_storage/client.py | 2 +- tests/deepagents_plugin/external_storage_test.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/deepagents_plugin/external_storage/client.py b/deepagents_plugin/external_storage/client.py index 4cbb8ab8..feb54d89 100644 --- a/deepagents_plugin/external_storage/client.py +++ b/deepagents_plugin/external_storage/client.py @@ -9,8 +9,8 @@ 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.envconfig import ClientConfig from temporalio.deepagents import DeepAgentsPlugin +from temporalio.envconfig import ClientConfig S3_ENDPOINT = "http://localhost:5000" S3_BUCKET = "temporal-payloads" diff --git a/tests/deepagents_plugin/external_storage_test.py b/tests/deepagents_plugin/external_storage_test.py index 4d28dbdf..a56d84e4 100644 --- a/tests/deepagents_plugin/external_storage_test.py +++ b/tests/deepagents_plugin/external_storage_test.py @@ -11,8 +11,8 @@ from temporalio.api.enums.v1 import EventType from temporalio.api.sdk.v1 import ExternalStorageReference from temporalio.client import Client -from temporalio.deepagents import DeepAgentsPlugin from temporalio.converter import DataConverter +from temporalio.deepagents import DeepAgentsPlugin from temporalio.worker import Replayer, Worker from deepagents_plugin.external_storage.client import ( From 581b1c372d90ee4a6eaf12387ea8a47129d4aeb2 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Fri, 9 Oct 2026 10:06:06 -0700 Subject: [PATCH 4/7] Rename Deep Agents external storage samples --- deepagents_plugin/README.md | 13 +- .../external_conversation_storage/README.md | 59 +++++++ .../external_conversation_storage/client.py | 52 ++++++ .../external_conversation_storage/main.py | 70 ++++++++ .../external_conversation_storage/workflow.py | 159 ++++++++++++++++++ .../README.md | 35 ++-- .../client.py | 14 +- .../main.py | 4 +- .../workflow.py | 4 - ...st.py => external_payload_storage_test.py} | 13 +- 10 files changed, 386 insertions(+), 37 deletions(-) create mode 100644 deepagents_plugin/external_conversation_storage/README.md create mode 100644 deepagents_plugin/external_conversation_storage/client.py create mode 100644 deepagents_plugin/external_conversation_storage/main.py create mode 100644 deepagents_plugin/external_conversation_storage/workflow.py rename deepagents_plugin/{external_storage => external_payload_storage}/README.md (81%) rename deepagents_plugin/{external_storage => external_payload_storage}/client.py (85%) rename deepagents_plugin/{external_storage => external_payload_storage}/main.py (94%) rename deepagents_plugin/{external_storage => external_payload_storage}/workflow.py (96%) rename tests/deepagents_plugin/{external_storage_test.py => external_payload_storage_test.py} (91%) diff --git a/deepagents_plugin/README.md b/deepagents_plugin/README.md index 688eaee4..e4ffd3aa 100644 --- a/deepagents_plugin/README.md +++ b/deepagents_plugin/README.md @@ -26,7 +26,8 @@ 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) | Carry a 6 MiB tool result and the next model input through native S3 External Storage; scripted model, no API keys. | +| [external_payload_storage](external_payload_storage) | Carry a 6 MiB tool result and the next model input through native S3 External Storage; scripted model, no API keys. | +| [external_conversation_storage](external_conversation_storage) | Store each exchange in S3 and include the full conversation history in each model call; native External Storage also protects large Temporal payloads. | ## Prerequisites @@ -48,7 +49,7 @@ one side. 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 + [external_payload_storage](external_payload_storage) sample uses a scripted model and needs no provider credentials: ```bash @@ -88,13 +89,15 @@ uv run deepagents_plugin/hello_world/run_worker.py uv run deepagents_plugin/hello_world/run_workflow.py ``` -The `langsmith_tracing` and `external_storage` samples instead bundle the worker +The `langsmith_tracing`, `external_payload_storage`, and +`external_conversation_storage` samples instead bundle the worker and starter into a single driver. The storage sample also needs a mock S3 service; -see its [setup instructions](external_storage): +see their [setup instructions](../external_storage): ```bash uv run deepagents_plugin/langsmith_tracing/main.py -uv run deepagents_plugin/external_storage/main.py +uv run deepagents_plugin/external_payload_storage/main.py +uv run deepagents_plugin/external_conversation_storage/main.py ``` ## Key Features Demonstrated diff --git a/deepagents_plugin/external_conversation_storage/README.md b/deepagents_plugin/external_conversation_storage/README.md new file mode 100644 index 00000000..6ee10c20 --- /dev/null +++ b/deepagents_plugin/external_conversation_storage/README.md @@ -0,0 +1,59 @@ +# External Storage and Conversation History + +This companion demonstrates the same native S3 payload configuration pattern as +[external_payload_storage](../external_payload_storage), but has its own client +setup and can be run independently. It runs a four-turn conversation. It writes +each user/assistant exchange to a separate S3 object, then loads every previous +exchange for each model call. The full transcript is available to the agent +throughout the conversation. + +The S3 transcript is application-managed conversation memory. Temporal's native +`ExternalStorage` is configured too, so large workflow and activity payloads can +be stored as S3 references in Temporal history. These solve different problems: +the transcript design makes the full history available to the model, while +native External Storage bounds payload bytes in Temporal history and transport. +The scripted run uses small prompts, so it does not need to cross the native +storage threshold. The full-history approach increases model input size on every +turn; total tokens across a long conversation can grow roughly quadratically. + +External storage preserves and retrieves the transcript; it does not reduce +model tokens. For long conversations, a summary or selective retrieval strategy +can keep the model context smaller while still drawing on older turns. + +## Run it + +Use Python 3.11 or later. Install dependencies: + +```bash +uv sync --python 3.13 --group deepagents --group external-storage +``` + +Start the Temporal dev server in one terminal: + +```bash +temporal server start-dev +``` + +Start the repository's mock S3 service in another terminal: + +```bash +uv run external_storage/s3.py +``` + +Then run the sample: + +```bash +uv run deepagents_plugin/external_conversation_storage/main.py +``` + +It uses a scripted model and needs no provider API key. Each turn's full exchange +is stored under a workflow-specific S3 prefix. The final turn demonstrates that +the complete history carries earlier project details into the answer. + +The sample uses the same `temporal-payloads` bucket and local mock S3 service as +the [large-payload sample](../external_payload_storage), but shares no code with +it. +For production, use durable storage with access controls and retention +appropriate for conversation data. +Keep the storage driver configured on clients, workers, and replayers that need +to decode native Temporal payload references. diff --git a/deepagents_plugin/external_conversation_storage/client.py b/deepagents_plugin/external_conversation_storage/client.py new file mode 100644 index 00000000..44201aab --- /dev/null +++ b/deepagents_plugin/external_conversation_storage/client.py @@ -0,0 +1,52 @@ +"""Connect a Deep Agents client with Temporal native S3 External Storage.""" + +from collections.abc import AsyncIterator +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.deepagents import DeepAgentsPlugin +from temporalio.converter import DataConverter, ExternalStorage +from temporalio.envconfig import ClientConfig + +S3_ENDPOINT = "http://localhost:5000" +S3_BUCKET = "temporal-payloads" + + +@asynccontextmanager +async def connect_client( + plugin: 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, + ) + storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) + plugin.data_converter = replace( + plugin.data_converter or DataConverter.default, + external_storage=storage, + ) + yield await Client.connect( + **config, + plugins=[plugin], + ) diff --git a/deepagents_plugin/external_conversation_storage/main.py b/deepagents_plugin/external_conversation_storage/main.py new file mode 100644 index 00000000..448d7b16 --- /dev/null +++ b/deepagents_plugin/external_conversation_storage/main.py @@ -0,0 +1,70 @@ +"""Run a multi-turn agent that supplies its full transcript to the model.""" + +import asyncio +import uuid + +from langchain_core.messages import AIMessage +from temporalio.deepagents import DeepAgentsPlugin +from temporalio.deepagents.testing import mock_model_provider +from temporalio.worker import Worker + +from deepagents_plugin.external_conversation_storage.client import connect_client +from deepagents_plugin.external_conversation_storage.workflow import ( + TASK_QUEUE, + StoredConversationAgent, + load_conversation_history, + save_turn, +) + + +def create_plugin() -> DeepAgentsPlugin: + """Script replies so the sample runs without an LLM provider or API key.""" + return DeepAgentsPlugin( + model_provider=mock_model_provider( + [ + AIMessage(content="The project name is Cedar."), + AIMessage(content="The project is called Cedar and uses Python."), + AIMessage(content="Cedar uses Python and is due Friday."), + AIMessage(content="Cedar is a Python project due Friday."), + ] + ) + ) + + +async def main() -> None: + user_turns = [ + "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=[StoredConversationAgent], + activities=[load_conversation_history, save_turn], + max_cached_workflows=0, + ): + workflow_id = f"deepagents-conversation-{uuid.uuid4()}" + handle = await client.start_workflow( + StoredConversationAgent.run, + id=workflow_id, + task_queue=TASK_QUEUE, + ) + for question in user_turns: + await handle.signal(StoredConversationAgent.submit_turn, question) + await handle.signal(StoredConversationAgent.finish) + result = await handle.result() + + print(f"Workflow: {workflow_id}") + print(f"Turns processed: {result.turns}") + print(f"First-turn prompt context: {result.first_context_chars:,} characters") + print(f"Last-turn prompt context: {result.last_context_chars:,} characters") + print(f"Final answer: {result.last_answer}") + print(f"Full transcript stored externally: {result.transcript_turns} turns") + print("Each model call included the complete conversation history.") + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/deepagents_plugin/external_conversation_storage/workflow.py b/deepagents_plugin/external_conversation_storage/workflow.py new file mode 100644 index 00000000..79ccc81f --- /dev/null +++ b/deepagents_plugin/external_conversation_storage/workflow.py @@ -0,0 +1,159 @@ +"""A multi-turn agent that loads its full S3 transcript on each turn.""" + +from dataclasses import dataclass +from datetime import timedelta + +from temporalio import activity, workflow +from temporalio.common import RetryPolicy +from temporalio.deepagents import create_temporal_deep_agent + +MODEL = "fake:external-conversation-storage" +TASK_QUEUE = "deepagents-external-conversation-storage" +S3_BUCKET = "temporal-payloads" + + +@dataclass +class Turn: + user: str + assistant: str + + +@dataclass +class ConversationResult: + last_answer: str + turns: int + first_context_chars: int + last_context_chars: int + transcript_turns: int + + +def _conversation_prefix(workflow_id: str) -> str: + """Keep user-provided workflow IDs out of S3 key syntax.""" + import hashlib + + return f"conversation-memory/{hashlib.sha256(workflow_id.encode()).hexdigest()}/" + + +@activity.defn +async def load_conversation_history(workflow_id: str, turn_index: int) -> list[Turn]: + """Read every earlier exchange from the external transcript.""" + import aioboto3 + + session = aioboto3.Session() + async with session.client( + "s3", + endpoint_url="http://localhost:5000", + aws_access_key_id="test", + aws_secret_access_key="test", + region_name="us-east-1", + ) as client: + turns = [] + prefix = _conversation_prefix(workflow_id) + for index in range(turn_index): + response = await client.get_object( + Bucket=S3_BUCKET, Key=f"{prefix}turn-{index:08d}.json" + ) + import json + + turns.append(Turn(**json.loads(await response["Body"].read()))) + return turns + + +@activity.defn +async def save_turn(workflow_id: str, turn_index: int, turn: Turn) -> None: + """Append one exchange as its own immutable S3 object.""" + import aioboto3 + import json + from dataclasses import asdict + + session = aioboto3.Session() + async with session.client( + "s3", + endpoint_url="http://localhost:5000", + aws_access_key_id="test", + aws_secret_access_key="test", + region_name="us-east-1", + ) as client: + await client.put_object( + Bucket=S3_BUCKET, + Key=f"{_conversation_prefix(workflow_id)}turn-{turn_index:08d}.json", + Body=json.dumps(asdict(turn)).encode(), + ContentType="application/json", + ) + + +@workflow.defn +class StoredConversationAgent: + def __init__(self) -> None: + self.pending_turns: list[str] = [] + self.finished = False + + @workflow.signal + def submit_turn(self, question: str) -> None: + self.pending_turns.append(question) + + @workflow.signal + def finish(self) -> None: + self.finished = True + + @workflow.run + async def run(self) -> ConversationResult: + workflow_id = workflow.info().workflow_id + turn_count = 0 + first_context_chars = 0 + last_context_chars = 0 + last_answer = "" + + while True: + await workflow.wait_condition(lambda: self.pending_turns or self.finished) + if not self.pending_turns: + break + question = self.pending_turns.pop(0) + turn_index = turn_count + history = await workflow.execute_activity( + load_conversation_history, + args=[workflow_id, turn_index], + start_to_close_timeout=timedelta(seconds=30), + retry_policy=RetryPolicy(maximum_attempts=1), + ) + memory = "\n".join( + f"User: {turn.user}\nAssistant: {turn.assistant}" + for turn in history + ) + system_prompt = ( + "Answer the user's latest message. Use the full conversation " + "when relevant. This is the complete conversation so far.\n\n" + "Conversation history:\n" + f"{memory or '(no earlier turns)'}" + ) + # A fresh agent invocation receives the complete prior transcript + # and current prompt; prior turns are not carried in graph state. + agent = create_temporal_deep_agent( + model=MODEL, + system_prompt=system_prompt, + activity_options={"start_to_close_timeout": timedelta(seconds=30)}, + ) + context_chars = len(system_prompt) + len(question) + if turn_count == 0: + first_context_chars = context_chars + result = await agent.ainvoke({"messages": [{"role": "user", "content": question}]}) + answer = result["messages"][-1].content + if not isinstance(answer, str): + raise TypeError("Expected a text response from the model") + last_answer = answer + last_context_chars = context_chars + await workflow.execute_activity( + save_turn, + args=[workflow_id, turn_index, Turn(user=question, assistant=answer)], + start_to_close_timeout=timedelta(seconds=30), + retry_policy=RetryPolicy(maximum_attempts=1), + ) + turn_count += 1 + + return ConversationResult( + last_answer=last_answer, + turns=turn_count, + first_context_chars=first_context_chars, + last_context_chars=last_context_chars, + transcript_turns=turn_count, + ) diff --git a/deepagents_plugin/external_storage/README.md b/deepagents_plugin/external_payload_storage/README.md similarity index 81% rename from deepagents_plugin/external_storage/README.md rename to deepagents_plugin/external_payload_storage/README.md index a4cca761..6105768b 100644 --- a/deepagents_plugin/external_storage/README.md +++ b/deepagents_plugin/external_payload_storage/README.md @@ -37,7 +37,7 @@ In a third terminal, run the demo. It starts a worker, executes one workflow, checks the result's integrity, and shuts down the worker: ```bash -uv run deepagents_plugin/external_storage/main.py +uv run deepagents_plugin/external_payload_storage/main.py ``` The output includes: @@ -66,13 +66,23 @@ through, so it can also read this sample's references. ## How the configuration works -`client.py` opens an S3 client, creates an SDK `S3StorageDriver`, and passes a -`DataConverter` with `ExternalStorage` directly to `Client.connect`. The standalone -`temporalio-deepagents` plugin composes its LangChain payload converter with that -converter, preserving the storage configuration. The worker uses the client's -converter, and the sample keeps the S3 client open while the worker runs and the -starter decodes results. No application code uploads or downloads payloads -manually. +`client.py` creates `ExternalStorage(drivers=[driver], +payload_size_threshold=256 * 1024)` and adds it to the LangChain-aware converter +provided by `DeepAgentsPlugin`: + +```python +plugin = create_plugin() +plugin.data_converter = replace(plugin.data_converter, external_storage=storage) + +client = await Client.connect( + "localhost:7233", + plugins=[plugin], +) +``` + +The worker inherits the configured converter. The S3 client stays open while +the worker runs and the starter decodes results. No application code uploads or +downloads Temporal payloads manually. This demo deliberately uses **no compression codec**: repeated `x` characters would compress well below the storage threshold. It overrides the built-in @@ -90,16 +100,17 @@ Native storage avoids transmitting the large bytes through either limit. For real S3, use an existing bucket and normal AWS credentials instead of this sample's local endpoint and mock credentials. Configure compatible storage drivers on every client, worker, and replayer that reads the history. Retain -objects for as long as execution, reset, or replay may need them; Temporal does not delete stored payloads when an activity finishes. +objects for as long as execution, reset, or replay may need them; Temporal does +not delete stored payloads when an activity finishes. ## Tests The integration test starts an isolated mock S3 service and verifies the full result, both native references, activity routing, and replay with a new S3 -client after worker shutdown. No manually running S3 service or model credentials -are required: +client after worker shutdown. +No manually running S3 service or model credentials are required: ```bash uv sync --python 3.13 --group deepagents --group external-storage --group dev -uv run pytest tests/deepagents_plugin/external_storage_test.py +uv run pytest tests/deepagents_plugin/external_payload_storage_test.py ``` diff --git a/deepagents_plugin/external_storage/client.py b/deepagents_plugin/external_payload_storage/client.py similarity index 85% rename from deepagents_plugin/external_storage/client.py rename to deepagents_plugin/external_payload_storage/client.py index feb54d89..df667d0d 100644 --- a/deepagents_plugin/external_storage/client.py +++ b/deepagents_plugin/external_payload_storage/client.py @@ -1,4 +1,4 @@ -"""Configure native S3 storage for the Deep Agents client.""" +"""Compose native S3 storage with the Deep Agents data converter.""" from collections.abc import AsyncIterator from contextlib import asynccontextmanager @@ -8,15 +8,14 @@ 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.converter import DataConverter, ExternalStorage from temporalio.envconfig import ClientConfig S3_ENDPOINT = "http://localhost:5000" S3_BUCKET = "temporal-payloads" -# @@@SNIPSTART python-deepagents-external-storage-converter @asynccontextmanager async def connect_client( plugin: DeepAgentsPlugin, @@ -43,12 +42,11 @@ async def connect_client( bucket=S3_BUCKET, ) storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) - converter = replace(DataConverter.default, external_storage=storage) + plugin.data_converter = replace( + plugin.data_converter or DataConverter.default, + external_storage=storage, + ) yield await Client.connect( **config, - data_converter=converter, plugins=[plugin], ) - - -# @@@SNIPEND diff --git a/deepagents_plugin/external_storage/main.py b/deepagents_plugin/external_payload_storage/main.py similarity index 94% rename from deepagents_plugin/external_storage/main.py rename to deepagents_plugin/external_payload_storage/main.py index e0bc3f22..fb520aca 100644 --- a/deepagents_plugin/external_storage/main.py +++ b/deepagents_plugin/external_payload_storage/main.py @@ -14,8 +14,8 @@ 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 ( +from deepagents_plugin.external_payload_storage.client import connect_client +from deepagents_plugin.external_payload_storage.workflow import ( PAYLOAD_BYTES, TASK_QUEUE, ExternalStorageAgent, diff --git a/deepagents_plugin/external_storage/workflow.py b/deepagents_plugin/external_payload_storage/workflow.py similarity index 96% rename from deepagents_plugin/external_storage/workflow.py rename to deepagents_plugin/external_payload_storage/workflow.py index 49842b3b..1e4a2e35 100644 --- a/deepagents_plugin/external_storage/workflow.py +++ b/deepagents_plugin/external_payload_storage/workflow.py @@ -37,7 +37,6 @@ async def read_file(file_path: str) -> str: return "x" * PAYLOAD_BYTES -# @@@SNIPSTART python-deepagents-external-storage-workflow @workflow.defn class ExternalStorageAgent: @workflow.run @@ -70,6 +69,3 @@ async def run(self, question: str) -> AgentResult: sha256=hashlib.sha256(data).hexdigest(), answer=answer, ) - - -# @@@SNIPEND diff --git a/tests/deepagents_plugin/external_storage_test.py b/tests/deepagents_plugin/external_payload_storage_test.py similarity index 91% rename from tests/deepagents_plugin/external_storage_test.py rename to tests/deepagents_plugin/external_payload_storage_test.py index a56d84e4..f4157a7d 100644 --- a/tests/deepagents_plugin/external_storage_test.py +++ b/tests/deepagents_plugin/external_payload_storage_test.py @@ -11,16 +11,16 @@ from temporalio.api.enums.v1 import EventType from temporalio.api.sdk.v1 import ExternalStorageReference from temporalio.client import Client -from temporalio.converter import DataConverter from temporalio.deepagents import DeepAgentsPlugin +from temporalio.converter import DataConverter from temporalio.worker import Replayer, Worker -from deepagents_plugin.external_storage.client import ( +from deepagents_plugin.external_payload_storage.client import ( S3_BUCKET, connect_client, ) -from deepagents_plugin.external_storage.main import create_plugin -from deepagents_plugin.external_storage.workflow import ( +from deepagents_plugin.external_payload_storage.main import create_plugin +from deepagents_plugin.external_payload_storage.workflow import ( PAYLOAD_BYTES, AgentResult, ExternalStorageAgent, @@ -108,11 +108,12 @@ async def test_external_storage(client: Client, s3_endpoint: str) -> None: # The original worker and S3 client are closed. Retrieve persisted claims # through a fresh client and driver, without executing the activities again. + replay_plugin = DeepAgentsPlugin() async with connect_client( - DeepAgentsPlugin(), target_host=target_host, s3_endpoint=s3_endpoint + replay_plugin, target_host=target_host, s3_endpoint=s3_endpoint ) as replay_client: await Replayer( workflows=[ExternalStorageAgent], data_converter=replay_client.data_converter, - plugins=[DeepAgentsPlugin()], + plugins=[replay_plugin], ).replay_workflow(history) From 5970487261c6eb5dd0424d8b788a346b5305e256 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Fri, 9 Oct 2026 10:08:48 -0700 Subject: [PATCH 5/7] Fix Deep Agents external storage lint --- .../external_conversation_storage/client.py | 2 +- .../external_conversation_storage/workflow.py | 14 +++++++++----- .../external_payload_storage/client.py | 2 +- .../external_payload_storage_test.py | 2 +- 4 files changed, 12 insertions(+), 8 deletions(-) diff --git a/deepagents_plugin/external_conversation_storage/client.py b/deepagents_plugin/external_conversation_storage/client.py index 44201aab..91d9e1d3 100644 --- a/deepagents_plugin/external_conversation_storage/client.py +++ b/deepagents_plugin/external_conversation_storage/client.py @@ -8,8 +8,8 @@ from temporalio.client import Client from temporalio.contrib.aws.s3driver import S3StorageDriver from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client -from temporalio.deepagents import DeepAgentsPlugin from temporalio.converter import DataConverter, ExternalStorage +from temporalio.deepagents import DeepAgentsPlugin from temporalio.envconfig import ClientConfig S3_ENDPOINT = "http://localhost:5000" diff --git a/deepagents_plugin/external_conversation_storage/workflow.py b/deepagents_plugin/external_conversation_storage/workflow.py index 79ccc81f..ea6a9e06 100644 --- a/deepagents_plugin/external_conversation_storage/workflow.py +++ b/deepagents_plugin/external_conversation_storage/workflow.py @@ -62,10 +62,11 @@ async def load_conversation_history(workflow_id: str, turn_index: int) -> list[T @activity.defn async def save_turn(workflow_id: str, turn_index: int, turn: Turn) -> None: """Append one exchange as its own immutable S3 object.""" - import aioboto3 import json from dataclasses import asdict + import aioboto3 + session = aioboto3.Session() async with session.client( "s3", @@ -105,7 +106,9 @@ async def run(self) -> ConversationResult: last_answer = "" while True: - await workflow.wait_condition(lambda: self.pending_turns or self.finished) + await workflow.wait_condition( + lambda: bool(self.pending_turns) or self.finished + ) if not self.pending_turns: break question = self.pending_turns.pop(0) @@ -117,8 +120,7 @@ async def run(self) -> ConversationResult: retry_policy=RetryPolicy(maximum_attempts=1), ) memory = "\n".join( - f"User: {turn.user}\nAssistant: {turn.assistant}" - for turn in history + f"User: {turn.user}\nAssistant: {turn.assistant}" for turn in history ) system_prompt = ( "Answer the user's latest message. Use the full conversation " @@ -136,7 +138,9 @@ async def run(self) -> ConversationResult: context_chars = len(system_prompt) + len(question) if turn_count == 0: first_context_chars = context_chars - result = await agent.ainvoke({"messages": [{"role": "user", "content": question}]}) + result = await agent.ainvoke( + {"messages": [{"role": "user", "content": question}]} + ) answer = result["messages"][-1].content if not isinstance(answer, str): raise TypeError("Expected a text response from the model") diff --git a/deepagents_plugin/external_payload_storage/client.py b/deepagents_plugin/external_payload_storage/client.py index df667d0d..26daaa64 100644 --- a/deepagents_plugin/external_payload_storage/client.py +++ b/deepagents_plugin/external_payload_storage/client.py @@ -8,8 +8,8 @@ from temporalio.client import Client from temporalio.contrib.aws.s3driver import S3StorageDriver from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client -from temporalio.deepagents import DeepAgentsPlugin from temporalio.converter import DataConverter, ExternalStorage +from temporalio.deepagents import DeepAgentsPlugin from temporalio.envconfig import ClientConfig S3_ENDPOINT = "http://localhost:5000" diff --git a/tests/deepagents_plugin/external_payload_storage_test.py b/tests/deepagents_plugin/external_payload_storage_test.py index f4157a7d..6305ec45 100644 --- a/tests/deepagents_plugin/external_payload_storage_test.py +++ b/tests/deepagents_plugin/external_payload_storage_test.py @@ -11,8 +11,8 @@ from temporalio.api.enums.v1 import EventType from temporalio.api.sdk.v1 import ExternalStorageReference from temporalio.client import Client -from temporalio.deepagents import DeepAgentsPlugin from temporalio.converter import DataConverter +from temporalio.deepagents import DeepAgentsPlugin from temporalio.worker import Replayer, Worker from deepagents_plugin.external_payload_storage.client import ( From 64f9a7643b92e0d2aa0d3293310f977397f30652 Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Fri, 9 Oct 2026 10:17:00 -0700 Subject: [PATCH 6/7] Pass S3 converter into Deep Agents plugin --- .../external_conversation_storage/client.py | 10 ++++------ .../external_conversation_storage/main.py | 8 +++++--- deepagents_plugin/external_payload_storage/README.md | 12 ++++++++---- deepagents_plugin/external_payload_storage/client.py | 10 ++++------ deepagents_plugin/external_payload_storage/main.py | 8 +++++--- .../external_payload_storage_test.py | 6 ++++-- 6 files changed, 30 insertions(+), 24 deletions(-) diff --git a/deepagents_plugin/external_conversation_storage/client.py b/deepagents_plugin/external_conversation_storage/client.py index 91d9e1d3..1d6a963e 100644 --- a/deepagents_plugin/external_conversation_storage/client.py +++ b/deepagents_plugin/external_conversation_storage/client.py @@ -1,6 +1,6 @@ """Connect a Deep Agents client with Temporal native S3 External Storage.""" -from collections.abc import AsyncIterator +from collections.abc import AsyncIterator, Callable from contextlib import asynccontextmanager from dataclasses import replace @@ -18,7 +18,7 @@ @asynccontextmanager async def connect_client( - plugin: DeepAgentsPlugin, + plugin_factory: Callable[[DataConverter], DeepAgentsPlugin], *, target_host: str | None = None, s3_endpoint: str = S3_ENDPOINT, @@ -42,10 +42,8 @@ async def connect_client( bucket=S3_BUCKET, ) storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) - plugin.data_converter = replace( - plugin.data_converter or DataConverter.default, - external_storage=storage, - ) + data_converter = replace(DataConverter.default, external_storage=storage) + plugin = plugin_factory(data_converter) yield await Client.connect( **config, plugins=[plugin], diff --git a/deepagents_plugin/external_conversation_storage/main.py b/deepagents_plugin/external_conversation_storage/main.py index 448d7b16..0a7edf62 100644 --- a/deepagents_plugin/external_conversation_storage/main.py +++ b/deepagents_plugin/external_conversation_storage/main.py @@ -4,6 +4,7 @@ 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 @@ -17,9 +18,10 @@ ) -def create_plugin() -> DeepAgentsPlugin: +def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin: """Script replies so the sample runs without an LLM provider or API key.""" return DeepAgentsPlugin( + data_converter=data_converter, model_provider=mock_model_provider( [ AIMessage(content="The project name is Cedar."), @@ -27,7 +29,7 @@ def create_plugin() -> DeepAgentsPlugin: AIMessage(content="Cedar uses Python and is due Friday."), AIMessage(content="Cedar is a Python project due Friday."), ] - ) + ), ) @@ -38,7 +40,7 @@ async def main() -> None: "It is due Friday.", "Remind me of the project name, language, and deadline.", ] - async with connect_client(create_plugin()) as client: + async with connect_client(create_plugin) as client: async with Worker( client, task_queue=TASK_QUEUE, diff --git a/deepagents_plugin/external_payload_storage/README.md b/deepagents_plugin/external_payload_storage/README.md index 6105768b..e4120f5a 100644 --- a/deepagents_plugin/external_payload_storage/README.md +++ b/deepagents_plugin/external_payload_storage/README.md @@ -67,12 +67,12 @@ through, so it can also read this sample's references. ## How the configuration works `client.py` creates `ExternalStorage(drivers=[driver], -payload_size_threshold=256 * 1024)` and adds it to the LangChain-aware converter -provided by `DeepAgentsPlugin`: +payload_size_threshold=256 * 1024)`, adds it to the SDK's default converter with +`dataclasses.replace`, and passes that converter into the Deep Agents plugin: ```python -plugin = create_plugin() -plugin.data_converter = replace(plugin.data_converter, external_storage=storage) +data_converter = replace(DataConverter.default, external_storage=storage) +plugin = DeepAgentsPlugin(data_converter=data_converter) client = await Client.connect( "localhost:7233", @@ -80,6 +80,10 @@ client = await Client.connect( ) ``` +The plugin upgrades the default payload converter for LangChain types while +preserving the converter's external-storage configuration. This uses the +standalone plugin's supported `data_converter` constructor argument. + The worker inherits the configured converter. The S3 client stays open while the worker runs and the starter decodes results. No application code uploads or downloads Temporal payloads manually. diff --git a/deepagents_plugin/external_payload_storage/client.py b/deepagents_plugin/external_payload_storage/client.py index 26daaa64..dbbd290e 100644 --- a/deepagents_plugin/external_payload_storage/client.py +++ b/deepagents_plugin/external_payload_storage/client.py @@ -1,6 +1,6 @@ """Compose native S3 storage with the Deep Agents data converter.""" -from collections.abc import AsyncIterator +from collections.abc import AsyncIterator, Callable from contextlib import asynccontextmanager from dataclasses import replace @@ -18,7 +18,7 @@ @asynccontextmanager async def connect_client( - plugin: DeepAgentsPlugin, + plugin_factory: Callable[[DataConverter], DeepAgentsPlugin], *, target_host: str | None = None, s3_endpoint: str = S3_ENDPOINT, @@ -42,10 +42,8 @@ async def connect_client( bucket=S3_BUCKET, ) storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) - plugin.data_converter = replace( - plugin.data_converter or DataConverter.default, - external_storage=storage, - ) + data_converter = replace(DataConverter.default, external_storage=storage) + plugin = plugin_factory(data_converter) yield await Client.connect( **config, plugins=[plugin], diff --git a/deepagents_plugin/external_payload_storage/main.py b/deepagents_plugin/external_payload_storage/main.py index fb520aca..07d8bd97 100644 --- a/deepagents_plugin/external_payload_storage/main.py +++ b/deepagents_plugin/external_payload_storage/main.py @@ -10,6 +10,7 @@ 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 @@ -22,9 +23,10 @@ ) -def create_plugin() -> DeepAgentsPlugin: +def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin: """Script a tool request and acknowledgment for this single-workflow demo.""" return DeepAgentsPlugin( + data_converter=data_converter, model_provider=mock_model_provider( [ AIMessage( @@ -39,12 +41,12 @@ def create_plugin() -> DeepAgentsPlugin: ), AIMessage(content="Received the complete document."), ] - ) + ), ) async def main() -> None: - async with connect_client(create_plugin()) as client: + async with connect_client(create_plugin) as client: async with Worker( client, task_queue=TASK_QUEUE, diff --git a/tests/deepagents_plugin/external_payload_storage_test.py b/tests/deepagents_plugin/external_payload_storage_test.py index 6305ec45..ecfc095b 100644 --- a/tests/deepagents_plugin/external_payload_storage_test.py +++ b/tests/deepagents_plugin/external_payload_storage_test.py @@ -56,7 +56,7 @@ async def test_external_storage(client: Client, s3_endpoint: str) -> None: queue = f"deepagents-storage-{uuid.uuid4()}" target_host = client.service_client.config.target_host async with connect_client( - create_plugin(), target_host=target_host, s3_endpoint=s3_endpoint + create_plugin, target_host=target_host, s3_endpoint=s3_endpoint ) as stored_client: assert stored_client.data_converter.external_storage is not None async with Worker( @@ -110,7 +110,9 @@ async def test_external_storage(client: Client, s3_endpoint: str) -> None: # through a fresh client and driver, without executing the activities again. replay_plugin = DeepAgentsPlugin() async with connect_client( - replay_plugin, target_host=target_host, s3_endpoint=s3_endpoint + lambda data_converter: DeepAgentsPlugin(data_converter=data_converter), + target_host=target_host, + s3_endpoint=s3_endpoint, ) as replay_client: await Replayer( workflows=[ExternalStorageAgent], From 95c24f556b5247db976eb7f1f05e9fcbe3a9fe8d Mon Sep 17 00:00:00 2001 From: Brian Strauch Date: Fri, 9 Oct 2026 12:21:07 -0700 Subject: [PATCH 7/7] Combine Deep Agents external storage samples --- deepagents_plugin/README.md | 11 +- .../external_conversation_storage/README.md | 59 ------- .../external_conversation_storage/main.py | 72 -------- .../external_conversation_storage/workflow.py | 163 ------------------ .../external_payload_storage/README.md | 120 ------------- .../external_payload_storage/client.py | 50 ------ .../external_payload_storage/workflow.py | 71 -------- deepagents_plugin/external_storage/README.md | 75 ++++++++ .../client.py | 11 +- .../main.py | 42 +++-- .../external_storage/workflow.py | 95 ++++++++++ ...orage_test.py => external_storage_test.py} | 90 +++++----- 12 files changed, 251 insertions(+), 608 deletions(-) delete mode 100644 deepagents_plugin/external_conversation_storage/README.md delete mode 100644 deepagents_plugin/external_conversation_storage/main.py delete mode 100644 deepagents_plugin/external_conversation_storage/workflow.py delete mode 100644 deepagents_plugin/external_payload_storage/README.md delete mode 100644 deepagents_plugin/external_payload_storage/client.py delete mode 100644 deepagents_plugin/external_payload_storage/workflow.py create mode 100644 deepagents_plugin/external_storage/README.md rename deepagents_plugin/{external_conversation_storage => external_storage}/client.py (83%) rename deepagents_plugin/{external_payload_storage => external_storage}/main.py (53%) create mode 100644 deepagents_plugin/external_storage/workflow.py rename tests/deepagents_plugin/{external_payload_storage_test.py => external_storage_test.py} (50%) diff --git a/deepagents_plugin/README.md b/deepagents_plugin/README.md index e4ffd3aa..da532e18 100644 --- a/deepagents_plugin/README.md +++ b/deepagents_plugin/README.md @@ -26,8 +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_payload_storage](external_payload_storage) | Carry a 6 MiB tool result and the next model input through native S3 External Storage; scripted model, no API keys. | -| [external_conversation_storage](external_conversation_storage) | Store each exchange in S3 and include the full conversation history in each model call; native External Storage also protects large Temporal payloads. | +| [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 @@ -49,7 +48,7 @@ one side. 2. Configure a model provider. Most samples use `anthropic:claude-sonnet-4-5`, which needs an Anthropic API key. The - [external_payload_storage](external_payload_storage) sample uses a scripted model and needs + [external_storage](external_storage) sample uses a scripted model and needs no provider credentials: ```bash @@ -89,15 +88,13 @@ uv run deepagents_plugin/hello_world/run_worker.py uv run deepagents_plugin/hello_world/run_workflow.py ``` -The `langsmith_tracing`, `external_payload_storage`, and -`external_conversation_storage` samples instead bundle the worker +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_payload_storage/main.py -uv run deepagents_plugin/external_conversation_storage/main.py +uv run deepagents_plugin/external_storage/main.py ``` ## Key Features Demonstrated diff --git a/deepagents_plugin/external_conversation_storage/README.md b/deepagents_plugin/external_conversation_storage/README.md deleted file mode 100644 index 6ee10c20..00000000 --- a/deepagents_plugin/external_conversation_storage/README.md +++ /dev/null @@ -1,59 +0,0 @@ -# External Storage and Conversation History - -This companion demonstrates the same native S3 payload configuration pattern as -[external_payload_storage](../external_payload_storage), but has its own client -setup and can be run independently. It runs a four-turn conversation. It writes -each user/assistant exchange to a separate S3 object, then loads every previous -exchange for each model call. The full transcript is available to the agent -throughout the conversation. - -The S3 transcript is application-managed conversation memory. Temporal's native -`ExternalStorage` is configured too, so large workflow and activity payloads can -be stored as S3 references in Temporal history. These solve different problems: -the transcript design makes the full history available to the model, while -native External Storage bounds payload bytes in Temporal history and transport. -The scripted run uses small prompts, so it does not need to cross the native -storage threshold. The full-history approach increases model input size on every -turn; total tokens across a long conversation can grow roughly quadratically. - -External storage preserves and retrieves the transcript; it does not reduce -model tokens. For long conversations, a summary or selective retrieval strategy -can keep the model context smaller while still drawing on older turns. - -## Run it - -Use Python 3.11 or later. Install dependencies: - -```bash -uv sync --python 3.13 --group deepagents --group external-storage -``` - -Start the Temporal dev server in one terminal: - -```bash -temporal server start-dev -``` - -Start the repository's mock S3 service in another terminal: - -```bash -uv run external_storage/s3.py -``` - -Then run the sample: - -```bash -uv run deepagents_plugin/external_conversation_storage/main.py -``` - -It uses a scripted model and needs no provider API key. Each turn's full exchange -is stored under a workflow-specific S3 prefix. The final turn demonstrates that -the complete history carries earlier project details into the answer. - -The sample uses the same `temporal-payloads` bucket and local mock S3 service as -the [large-payload sample](../external_payload_storage), but shares no code with -it. -For production, use durable storage with access controls and retention -appropriate for conversation data. -Keep the storage driver configured on clients, workers, and replayers that need -to decode native Temporal payload references. diff --git a/deepagents_plugin/external_conversation_storage/main.py b/deepagents_plugin/external_conversation_storage/main.py deleted file mode 100644 index 0a7edf62..00000000 --- a/deepagents_plugin/external_conversation_storage/main.py +++ /dev/null @@ -1,72 +0,0 @@ -"""Run a multi-turn agent that supplies its full transcript to the model.""" - -import asyncio -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_conversation_storage.client import connect_client -from deepagents_plugin.external_conversation_storage.workflow import ( - TASK_QUEUE, - StoredConversationAgent, - load_conversation_history, - save_turn, -) - - -def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin: - """Script replies so the sample runs without an LLM provider or API key.""" - return DeepAgentsPlugin( - data_converter=data_converter, - model_provider=mock_model_provider( - [ - AIMessage(content="The project name is Cedar."), - AIMessage(content="The project is called Cedar and uses Python."), - AIMessage(content="Cedar uses Python and is due Friday."), - AIMessage(content="Cedar is a Python project due Friday."), - ] - ), - ) - - -async def main() -> None: - user_turns = [ - "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=[StoredConversationAgent], - activities=[load_conversation_history, save_turn], - max_cached_workflows=0, - ): - workflow_id = f"deepagents-conversation-{uuid.uuid4()}" - handle = await client.start_workflow( - StoredConversationAgent.run, - id=workflow_id, - task_queue=TASK_QUEUE, - ) - for question in user_turns: - await handle.signal(StoredConversationAgent.submit_turn, question) - await handle.signal(StoredConversationAgent.finish) - result = await handle.result() - - print(f"Workflow: {workflow_id}") - print(f"Turns processed: {result.turns}") - print(f"First-turn prompt context: {result.first_context_chars:,} characters") - print(f"Last-turn prompt context: {result.last_context_chars:,} characters") - print(f"Final answer: {result.last_answer}") - print(f"Full transcript stored externally: {result.transcript_turns} turns") - print("Each model call included the complete conversation history.") - - -if __name__ == "__main__": - asyncio.run(main()) diff --git a/deepagents_plugin/external_conversation_storage/workflow.py b/deepagents_plugin/external_conversation_storage/workflow.py deleted file mode 100644 index ea6a9e06..00000000 --- a/deepagents_plugin/external_conversation_storage/workflow.py +++ /dev/null @@ -1,163 +0,0 @@ -"""A multi-turn agent that loads its full S3 transcript on each turn.""" - -from dataclasses import dataclass -from datetime import timedelta - -from temporalio import activity, workflow -from temporalio.common import RetryPolicy -from temporalio.deepagents import create_temporal_deep_agent - -MODEL = "fake:external-conversation-storage" -TASK_QUEUE = "deepagents-external-conversation-storage" -S3_BUCKET = "temporal-payloads" - - -@dataclass -class Turn: - user: str - assistant: str - - -@dataclass -class ConversationResult: - last_answer: str - turns: int - first_context_chars: int - last_context_chars: int - transcript_turns: int - - -def _conversation_prefix(workflow_id: str) -> str: - """Keep user-provided workflow IDs out of S3 key syntax.""" - import hashlib - - return f"conversation-memory/{hashlib.sha256(workflow_id.encode()).hexdigest()}/" - - -@activity.defn -async def load_conversation_history(workflow_id: str, turn_index: int) -> list[Turn]: - """Read every earlier exchange from the external transcript.""" - import aioboto3 - - session = aioboto3.Session() - async with session.client( - "s3", - endpoint_url="http://localhost:5000", - aws_access_key_id="test", - aws_secret_access_key="test", - region_name="us-east-1", - ) as client: - turns = [] - prefix = _conversation_prefix(workflow_id) - for index in range(turn_index): - response = await client.get_object( - Bucket=S3_BUCKET, Key=f"{prefix}turn-{index:08d}.json" - ) - import json - - turns.append(Turn(**json.loads(await response["Body"].read()))) - return turns - - -@activity.defn -async def save_turn(workflow_id: str, turn_index: int, turn: Turn) -> None: - """Append one exchange as its own immutable S3 object.""" - import json - from dataclasses import asdict - - import aioboto3 - - session = aioboto3.Session() - async with session.client( - "s3", - endpoint_url="http://localhost:5000", - aws_access_key_id="test", - aws_secret_access_key="test", - region_name="us-east-1", - ) as client: - await client.put_object( - Bucket=S3_BUCKET, - Key=f"{_conversation_prefix(workflow_id)}turn-{turn_index:08d}.json", - Body=json.dumps(asdict(turn)).encode(), - ContentType="application/json", - ) - - -@workflow.defn -class StoredConversationAgent: - def __init__(self) -> None: - self.pending_turns: list[str] = [] - self.finished = False - - @workflow.signal - def submit_turn(self, question: str) -> None: - self.pending_turns.append(question) - - @workflow.signal - def finish(self) -> None: - self.finished = True - - @workflow.run - async def run(self) -> ConversationResult: - workflow_id = workflow.info().workflow_id - turn_count = 0 - first_context_chars = 0 - last_context_chars = 0 - last_answer = "" - - while True: - await workflow.wait_condition( - lambda: bool(self.pending_turns) or self.finished - ) - if not self.pending_turns: - break - question = self.pending_turns.pop(0) - turn_index = turn_count - history = await workflow.execute_activity( - load_conversation_history, - args=[workflow_id, turn_index], - start_to_close_timeout=timedelta(seconds=30), - retry_policy=RetryPolicy(maximum_attempts=1), - ) - memory = "\n".join( - f"User: {turn.user}\nAssistant: {turn.assistant}" for turn in history - ) - system_prompt = ( - "Answer the user's latest message. Use the full conversation " - "when relevant. This is the complete conversation so far.\n\n" - "Conversation history:\n" - f"{memory or '(no earlier turns)'}" - ) - # A fresh agent invocation receives the complete prior transcript - # and current prompt; prior turns are not carried in graph state. - agent = create_temporal_deep_agent( - model=MODEL, - system_prompt=system_prompt, - activity_options={"start_to_close_timeout": timedelta(seconds=30)}, - ) - context_chars = len(system_prompt) + len(question) - if turn_count == 0: - first_context_chars = context_chars - result = await agent.ainvoke( - {"messages": [{"role": "user", "content": question}]} - ) - answer = result["messages"][-1].content - if not isinstance(answer, str): - raise TypeError("Expected a text response from the model") - last_answer = answer - last_context_chars = context_chars - await workflow.execute_activity( - save_turn, - args=[workflow_id, turn_index, Turn(user=question, assistant=answer)], - start_to_close_timeout=timedelta(seconds=30), - retry_policy=RetryPolicy(maximum_attempts=1), - ) - turn_count += 1 - - return ConversationResult( - last_answer=last_answer, - turns=turn_count, - first_context_chars=first_context_chars, - last_context_chars=last_context_chars, - transcript_turns=turn_count, - ) diff --git a/deepagents_plugin/external_payload_storage/README.md b/deepagents_plugin/external_payload_storage/README.md deleted file mode 100644 index e4120f5a..00000000 --- a/deepagents_plugin/external_payload_storage/README.md +++ /dev/null @@ -1,120 +0,0 @@ -# External Storage - -Run a Deep Agent whose `read_file` tool returns **6 MiB** of text, exceeding -Temporal's default 2 MiB payload limit. The SDK's native `ExternalStorage` -offloads the tool output and the next model activity's input to S3, recording -small references in workflow history and retrieving the full bytes on decode. -The workflow returns only a byte count, SHA-256, and a short final answer. - -Native External Storage is in **Public Preview**. This sample uses the SDK's -`S3StorageDriver` with the local mock S3 service from the -[general External Storage sample](../../external_storage). Neither AWS credentials, -Docker, nor an LLM API key is needed. The model replies are scripted; the real -Deep Agents loop, tool activity, model activities, and S3 transfers still run. - -## Running the sample - -Use Python **3.11 or later**. Run these commands from the repository root: - -```bash -uv sync --python 3.13 --group deepagents --group external-storage -``` - -Start Temporal with its default payload limits in one terminal: - -```bash -temporal server start-dev -``` - -In a second terminal, start the existing mock S3 service. It listens on port -5000 and creates the `temporal-payloads` bucket: - -```bash -uv run external_storage/s3.py -``` - -In a third terminal, run the demo. It starts a worker, executes one workflow, -checks the result's integrity, and shuts down the worker: - -```bash -uv run deepagents_plugin/external_payload_storage/main.py -``` - -The output includes: - -```text -Tool result: 6,291,456 bytes (6 MiB), integrity verified -Agent: Received the complete document. -``` - -The script also prints the workflow ID and complete SHA-256. Inspect the history -in the Temporal UI or with: - -```bash -temporal workflow show --workflow-id -``` - -The completed `deepagents.invoke_tool` activity and the following -`deepagents.invoke_model` input contain native storage references. The first -model input and final workflow result remain small and inline. - -To view externally stored payloads in the UI, reuse the -[storage-aware codec server](../../external_storage#5-optional-run-the-codec-server): -run `uv run external_storage/codec_server.py` and set the UI's Remote Codec -Endpoint to `http://localhost:8081`. Its gzip decoder passes uncompressed payloads -through, so it can also read this sample's references. - -## How the configuration works - -`client.py` creates `ExternalStorage(drivers=[driver], -payload_size_threshold=256 * 1024)`, adds it to the SDK's default converter with -`dataclasses.replace`, and passes that converter into the Deep Agents plugin: - -```python -data_converter = replace(DataConverter.default, external_storage=storage) -plugin = DeepAgentsPlugin(data_converter=data_converter) - -client = await Client.connect( - "localhost:7233", - plugins=[plugin], -) -``` - -The plugin upgrades the default payload converter for LangChain types while -preserving the converter's external-storage configuration. This uses the -standalone plugin's supported `data_converter` constructor argument. - -The worker inherits the configured converter. The S3 client stays open while -the worker runs and the starter decodes results. No application code uploads or -downloads Temporal payloads manually. - -This demo deliberately uses **no compression codec**: repeated `x` characters -would compress well below the storage threshold. It overrides the built-in -`read_file` tool with an activity-backed mock bulk reader. Deep Agents excludes -`read_file` from tool-result eviction, so the complete result crosses both -activity boundaries. Applications using real models should page documents or -retrieve excerpts to manage context separately. Storage reduces transport and -history bytes; it does not reduce model tokens or decoded workflow memory. - -With storage removed, this tool result fails with `PayloadsTooLarge` at the -default limits. Raising the blob limit alone still leaves the gRPC message -limit described in [Hanyu Liu's post](https://hanyuliu.me/posts/temporal-deep-agent-payload-limits/). -Native storage avoids transmitting the large bytes through either limit. - -For real S3, use an existing bucket and normal AWS credentials instead of this -sample's local endpoint and mock credentials. Configure compatible storage -drivers on every client, worker, and replayer that reads the history. Retain -objects for as long as execution, reset, or replay may need them; Temporal does -not delete stored payloads when an activity finishes. - -## Tests - -The integration test starts an isolated mock S3 service and verifies the full -result, both native references, activity routing, and replay with a new S3 -client after worker shutdown. -No manually running S3 service or model credentials are required: - -```bash -uv sync --python 3.13 --group deepagents --group external-storage --group dev -uv run pytest tests/deepagents_plugin/external_payload_storage_test.py -``` diff --git a/deepagents_plugin/external_payload_storage/client.py b/deepagents_plugin/external_payload_storage/client.py deleted file mode 100644 index dbbd290e..00000000 --- a/deepagents_plugin/external_payload_storage/client.py +++ /dev/null @@ -1,50 +0,0 @@ -"""Compose native S3 storage with the Deep Agents data converter.""" - -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 throughout workflow execution and payload reads.""" - 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, - ) - storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) - data_converter = replace(DataConverter.default, external_storage=storage) - plugin = plugin_factory(data_converter) - yield await Client.connect( - **config, - plugins=[plugin], - ) diff --git a/deepagents_plugin/external_payload_storage/workflow.py b/deepagents_plugin/external_payload_storage/workflow.py deleted file mode 100644 index 1e4a2e35..00000000 --- a/deepagents_plugin/external_payload_storage/workflow.py +++ /dev/null @@ -1,71 +0,0 @@ -"""A Deep Agent whose tool returns 6 MiB, carried by native External Storage.""" - -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 AgentResult: - """Small integrity evidence; the large document stays out of the final result.""" - - result_bytes: int - sha256: str - answer: str - - -@tool -async def read_file(file_path: str) -> str: - """Read a mock document containing exactly 6 MiB of text.""" - # Override the built-in read_file with an activity-backed bulk reader. - # Deep Agents excludes read_file from tool-result eviction, so its complete - # content reaches the next model activity. Real file I/O would go here. - del file_path - return "x" * PAYLOAD_BYTES - - -@workflow.defn -class ExternalStorageAgent: - @workflow.run - async def run(self, question: str) -> AgentResult: - reader = tool_as_activity( - read_file, - start_to_close_timeout=timedelta(seconds=30), - activity_options={"retry_policy": RetryPolicy(maximum_attempts=1)}, - ) - agent = create_temporal_deep_agent( - model=MODEL, - tools=[reader], - system_prompt="Read the document, then acknowledge receiving the result.", - activity_options={"start_to_close_timeout": timedelta(seconds=30)}, - ) - result = await agent.ainvoke( - {"messages": [{"role": "user", "content": question}]} - ) - tool_result = next( - message - for message in result["messages"] - if isinstance(message, ToolMessage) and message.name == "read_file" - ) - assert isinstance(tool_result.content, str) - data = tool_result.content.encode() - answer = result["messages"][-1].content - assert isinstance(answer, str) - return AgentResult( - result_bytes=len(data), - sha256=hashlib.sha256(data).hexdigest(), - answer=answer, - ) diff --git a/deepagents_plugin/external_storage/README.md b/deepagents_plugin/external_storage/README.md new file mode 100644 index 00000000..1c563a73 --- /dev/null +++ b/deepagents_plugin/external_storage/README.md @@ -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 +``` + +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. diff --git a/deepagents_plugin/external_conversation_storage/client.py b/deepagents_plugin/external_storage/client.py similarity index 83% rename from deepagents_plugin/external_conversation_storage/client.py rename to deepagents_plugin/external_storage/client.py index 1d6a963e..b6b66487 100644 --- a/deepagents_plugin/external_conversation_storage/client.py +++ b/deepagents_plugin/external_storage/client.py @@ -1,4 +1,4 @@ -"""Connect a Deep Agents client with Temporal native S3 External Storage.""" +"""Configure native S3 External Storage for the Deep Agents plugin.""" from collections.abc import AsyncIterator, Callable from contextlib import asynccontextmanager @@ -41,10 +41,9 @@ async def connect_client( client=new_aioboto3_client(s3_client), bucket=S3_BUCKET, ) - storage = ExternalStorage(drivers=[driver], payload_size_threshold=256 * 1024) + # 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], - ) + yield await Client.connect(**config, plugins=[plugin]) diff --git a/deepagents_plugin/external_payload_storage/main.py b/deepagents_plugin/external_storage/main.py similarity index 53% rename from deepagents_plugin/external_payload_storage/main.py rename to deepagents_plugin/external_storage/main.py index 07d8bd97..2394d40f 100644 --- a/deepagents_plugin/external_payload_storage/main.py +++ b/deepagents_plugin/external_storage/main.py @@ -1,9 +1,4 @@ -"""Run one Deep Agent with a 6 MiB tool result and native S3 storage. - -Both the worker and starter run in this process. The two model replies are -scripted, so the demo needs no LLM credentials or provider context budget. -The model and tool still run as real Temporal activities. -""" +"""Run one scripted conversation with native Temporal S3 External Storage.""" import asyncio import hashlib @@ -15,8 +10,8 @@ from temporalio.deepagents.testing import mock_model_provider from temporalio.worker import Worker -from deepagents_plugin.external_payload_storage.client import connect_client -from deepagents_plugin.external_payload_storage.workflow import ( +from deepagents_plugin.external_storage.client import connect_client +from deepagents_plugin.external_storage.workflow import ( PAYLOAD_BYTES, TASK_QUEUE, ExternalStorageAgent, @@ -24,7 +19,7 @@ def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin: - """Script a tool request and acknowledgment for this single-workflow demo.""" + """Script the model so the demo needs no LLM provider or API key.""" return DeepAgentsPlugin( data_converter=data_converter, model_provider=mock_model_provider( @@ -39,13 +34,24 @@ def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin: } ], ), - AIMessage(content="Received the complete document."), + 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, @@ -56,17 +62,19 @@ async def main() -> None: workflow_id = f"deepagents-external-storage-{uuid.uuid4()}" result = await client.execute_workflow( ExternalStorageAgent.run, - "Read large-document.txt and acknowledge receiving it.", + questions, id=workflow_id, task_queue=TASK_QUEUE, ) - assert result.result_bytes == PAYLOAD_BYTES - assert result.sha256 == hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest() - print(f"Workflow: {workflow_id}") - print(f"Tool result: {result.result_bytes:,} bytes (6 MiB), integrity verified") - print(f"SHA-256: {result.sha256}") - print(f"Agent: {result.answer}") + 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__": diff --git a/deepagents_plugin/external_storage/workflow.py b/deepagents_plugin/external_storage/workflow.py new file mode 100644 index 00000000..afe63009 --- /dev/null +++ b/deepagents_plugin/external_storage/workflow.py @@ -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, + ) diff --git a/tests/deepagents_plugin/external_payload_storage_test.py b/tests/deepagents_plugin/external_storage_test.py similarity index 50% rename from tests/deepagents_plugin/external_payload_storage_test.py rename to tests/deepagents_plugin/external_storage_test.py index ecfc095b..65f250c4 100644 --- a/tests/deepagents_plugin/external_payload_storage_test.py +++ b/tests/deepagents_plugin/external_storage_test.py @@ -1,4 +1,4 @@ -"""Verify the sample against real Temporal activities and a local S3 service.""" +"""Verify large tool and conversation payloads use native S3 External Storage.""" import hashlib import uuid @@ -12,18 +12,14 @@ from temporalio.api.sdk.v1 import ExternalStorageReference from temporalio.client import Client from temporalio.converter import DataConverter -from temporalio.deepagents import DeepAgentsPlugin from temporalio.worker import Replayer, Worker -from deepagents_plugin.external_payload_storage.client import ( - S3_BUCKET, - connect_client, -) -from deepagents_plugin.external_payload_storage.main import create_plugin -from deepagents_plugin.external_payload_storage.workflow import ( +from deepagents_plugin.external_storage.client import S3_BUCKET, connect_client +from deepagents_plugin.external_storage.main import create_plugin +from deepagents_plugin.external_storage.workflow import ( PAYLOAD_BYTES, - AgentResult, ExternalStorageAgent, + ExternalStorageResult, ) from tests.deepagents_plugin.helpers import ( INVOKE_MODEL, @@ -55,6 +51,13 @@ def s3_endpoint() -> Iterator[str]: async def test_external_storage(client: Client, s3_endpoint: str) -> None: queue = f"deepagents-storage-{uuid.uuid4()}" target_host = client.service_client.config.target_host + 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, target_host=target_host, s3_endpoint=s3_endpoint ) as stored_client: @@ -67,7 +70,7 @@ async def test_external_storage(client: Client, s3_endpoint: str) -> None: ): handle = await stored_client.start_workflow( ExternalStorageAgent.run, - "Read large-document.txt and acknowledge receiving it.", + questions, id=queue, task_queue=queue, execution_timeout=timedelta(seconds=60), @@ -75,47 +78,48 @@ async def test_external_storage(client: Client, s3_endpoint: str) -> None: result = await handle.result() history = await handle.fetch_history() - assert result == AgentResult( - result_bytes=PAYLOAD_BYTES, - sha256=hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest(), - answer="Received the complete document.", + assert result == ExternalStorageResult( + tool_result_bytes=PAYLOAD_BYTES, + tool_result_sha256=hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest(), + turns=5, + last_answer="Cedar is a Python project due Friday.", ) counts = await count_scheduled_activities(handle) - assert counts[INVOKE_MODEL] == 2 + assert counts[INVOKE_MODEL] == 6 assert counts[INVOKE_TOOL] == 1 - scheduled = { - event.event_id: event.activity_task_scheduled_event_attributes - for event in history.events - if event.event_type == EventType.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED - } - tool_result = next( - event.activity_task_completed_event_attributes.result.payloads[0] - for event in history.events - if event.event_type == EventType.EVENT_TYPE_ACTIVITY_TASK_COMPLETED - and scheduled[ - event.activity_task_completed_event_attributes.scheduled_event_id - ].activity_type.name - == INVOKE_TOOL - ) - model_input = list(scheduled.values())[-1].input.payloads[0] - for payload in (tool_result, model_input): - reference = DataConverter.default.payload_converter.from_payload(payload) - assert isinstance(reference, ExternalStorageReference) - assert reference.driver_name == "aws.s3driver" - assert payload.ByteSize() < 1024 - assert "x" * 1024 not in history.to_json() + scheduled = { + event.event_id: event.activity_task_scheduled_event_attributes + for event in history.events + if event.event_type == EventType.EVENT_TYPE_ACTIVITY_TASK_SCHEDULED + } + tool_result = next( + event.activity_task_completed_event_attributes.result.payloads[0] + for event in history.events + if event.event_type == EventType.EVENT_TYPE_ACTIVITY_TASK_COMPLETED + and scheduled[ + event.activity_task_completed_event_attributes.scheduled_event_id + ].activity_type.name + == INVOKE_TOOL + ) + model_inputs = [ + attrs.input.payloads[0] + for attrs in scheduled.values() + if attrs.activity_type.name == INVOKE_MODEL + ] + for payload in [tool_result, *model_inputs]: + reference = DataConverter.default.payload_converter.from_payload(payload) + assert isinstance(reference, ExternalStorageReference) + assert reference.driver_name == "aws.s3driver" + assert payload.ByteSize() < 1024 + assert "x" * 1024 not in history.to_json() - # The original worker and S3 client are closed. Retrieve persisted claims - # through a fresh client and driver, without executing the activities again. - replay_plugin = DeepAgentsPlugin() + # A fresh client can decode the externally stored payloads for replay. async with connect_client( - lambda data_converter: DeepAgentsPlugin(data_converter=data_converter), - target_host=target_host, - s3_endpoint=s3_endpoint, + create_plugin, target_host=target_host, s3_endpoint=s3_endpoint ) as replay_client: await Replayer( workflows=[ExternalStorageAgent], data_converter=replay_client.data_converter, - plugins=[replay_plugin], + plugins=[create_plugin(replay_client.data_converter)], ).replay_workflow(history)