diff --git a/deepagents_plugin/README.md b/deepagents_plugin/README.md index 7f5d435b..da532e18 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) | 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 @@ -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 their [setup instructions](../external_storage): ```bash uv run deepagents_plugin/langsmith_tracing/main.py +uv run deepagents_plugin/external_storage/main.py ``` ## Key Features Demonstrated @@ -109,6 +114,8 @@ uv run deepagents_plugin/langsmith_tracing/main.py - **Streaming** — forward model chunks to external subscribers while keeping the durable result unchanged. - **Observability** — compose with `LangSmithPlugin` for tracing. +- **Large payload storage** — configure native `ExternalStorage` on the client's + data converter; the standalone plugin composes it with its LangChain converter. ## Related 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_storage/client.py b/deepagents_plugin/external_storage/client.py new file mode 100644 index 00000000..b6b66487 --- /dev/null +++ b/deepagents_plugin/external_storage/client.py @@ -0,0 +1,49 @@ +"""Configure native S3 External Storage for the Deep Agents plugin.""" + +from collections.abc import AsyncIterator, Callable +from contextlib import asynccontextmanager +from dataclasses import replace + +import aioboto3 +from temporalio.client import Client +from temporalio.contrib.aws.s3driver import S3StorageDriver +from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client +from temporalio.converter import DataConverter, ExternalStorage +from temporalio.deepagents import DeepAgentsPlugin +from temporalio.envconfig import ClientConfig + +S3_ENDPOINT = "http://localhost:5000" +S3_BUCKET = "temporal-payloads" + + +@asynccontextmanager +async def connect_client( + plugin_factory: Callable[[DataConverter], DeepAgentsPlugin], + *, + target_host: str | None = None, + s3_endpoint: str = S3_ENDPOINT, +) -> AsyncIterator[Client]: + """Keep the S3 client alive while Temporal reads and writes payloads.""" + config = ClientConfig.load_client_connect_config() + config.setdefault("target_host", "localhost:7233") + if target_host is not None: + config["target_host"] = target_host + + session = aioboto3.Session() + async with session.client( + "s3", + endpoint_url=s3_endpoint, + aws_access_key_id="test", + aws_secret_access_key="test", + region_name="us-east-1", + ) as s3_client: + driver = S3StorageDriver( + client=new_aioboto3_client(s3_client), + bucket=S3_BUCKET, + ) + # Small on purpose: this demo also externalizes serialized conversation + # context, while the large tool result shows the same path at scale. + storage = ExternalStorage(drivers=[driver], payload_size_threshold=128) + data_converter = replace(DataConverter.default, external_storage=storage) + plugin = plugin_factory(data_converter) + yield await Client.connect(**config, plugins=[plugin]) diff --git a/deepagents_plugin/external_storage/main.py b/deepagents_plugin/external_storage/main.py new file mode 100644 index 00000000..2394d40f --- /dev/null +++ b/deepagents_plugin/external_storage/main.py @@ -0,0 +1,81 @@ +"""Run one scripted conversation with native Temporal S3 External Storage.""" + +import asyncio +import hashlib +import uuid + +from langchain_core.messages import AIMessage +from temporalio.converter import DataConverter +from temporalio.deepagents import DeepAgentsPlugin +from temporalio.deepagents.testing import mock_model_provider +from temporalio.worker import Worker + +from deepagents_plugin.external_storage.client import connect_client +from deepagents_plugin.external_storage.workflow import ( + PAYLOAD_BYTES, + TASK_QUEUE, + ExternalStorageAgent, +) + + +def create_plugin(data_converter: DataConverter) -> DeepAgentsPlugin: + """Script the model so the demo needs no LLM provider or API key.""" + return DeepAgentsPlugin( + data_converter=data_converter, + model_provider=mock_model_provider( + [ + AIMessage( + content="", + tool_calls=[ + { + "name": "read_file", + "args": {"file_path": "large-document.txt"}, + "id": "call-read", + } + ], + ), + AIMessage(content="I received the complete document."), + AIMessage(content="The project is called Cedar."), + AIMessage(content="Cedar uses Python."), + AIMessage(content="Cedar is due Friday."), + AIMessage(content="Cedar is a Python project due Friday."), + ] + ), + ) + + +async def main() -> None: + questions = [ + "Read large-document.txt and acknowledge receiving it.", + "I am working on a project called Cedar.", + "The project uses Python.", + "It is due Friday.", + "Remind me of the project name, language, and deadline.", + ] + async with connect_client(create_plugin) as client: + async with Worker( + client, + task_queue=TASK_QUEUE, + workflows=[ExternalStorageAgent], + max_cached_workflows=0, + ): + workflow_id = f"deepagents-external-storage-{uuid.uuid4()}" + result = await client.execute_workflow( + ExternalStorageAgent.run, + questions, + id=workflow_id, + task_queue=TASK_QUEUE, + ) + + assert result.tool_result_bytes == PAYLOAD_BYTES + assert result.tool_result_sha256 == hashlib.sha256(b"x" * PAYLOAD_BYTES).hexdigest() + print(f"Workflow: {workflow_id}") + print( + f"Tool result: {result.tool_result_bytes:,} bytes (6 MiB), integrity verified" + ) + print(f"Conversation turns: {result.turns}") + print(f"Final answer: {result.last_answer}") + + +if __name__ == "__main__": + asyncio.run(main()) 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_storage_test.py b/tests/deepagents_plugin/external_storage_test.py new file mode 100644 index 00000000..65f250c4 --- /dev/null +++ b/tests/deepagents_plugin/external_storage_test.py @@ -0,0 +1,125 @@ +"""Verify large tool and conversation payloads use native S3 External Storage.""" + +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.converter import DataConverter +from temporalio.worker import Replayer, Worker + +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, + ExternalStorageAgent, + ExternalStorageResult, +) +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 + 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: + 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, + questions, + id=queue, + task_queue=queue, + execution_timeout=timedelta(seconds=60), + ) + result = await handle.result() + history = await handle.fetch_history() + + 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] == 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_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() + + # A fresh client can decode the externally stored payloads for replay. + async with connect_client( + create_plugin, target_host=target_host, s3_endpoint=s3_endpoint + ) as replay_client: + await Replayer( + workflows=[ExternalStorageAgent], + data_converter=replay_client.data_converter, + plugins=[create_plugin(replay_client.data_converter)], + ).replay_workflow(history)