Temporal과 함께하는 지속 실행
Temporal과 함께하는 지속 실행 (Durable Execution with Temporal)
Temporal은 Pydantic AI가 네이티브로 지원하는 인기 있는 지속 실행 플랫폼이에요. Temporal의 재생(replay) 메커니즘과 TemporalDurability 캐퍼빌리티로 에이전트를 지속으로 만드는 방법을 정리할게요.
출처: 문서
본문
지속 실행 (Durable Execution)
Temporal의 지속 실행 구현에서, 모델이나 API와 상호작용하다 크래시하거나 예외를 만난 프로그램은 성공적으로 완료할 때까지 재시도합니다.
Temporal은 주로 재생 메커니즘에 의존해 실패에서 회복해요. 프로그램이 진행함에 따라 Temporal이 핵심 입력과 결정을 저장해서, 재시작된 프로그램이 정확히 중단한 지점에서 이어받게 합니다.
이것이 동작하게 하는 핵심은 애플리케이션의 반복 가능한(결정적) 부분과 반복 불가능한(비결정적) 부분을 분리하는 것입니다:
- 결정적 조각은 워크플로 (workflows)라고 하며, 같은 입력으로 다시 실행하면 같은 방식으로 실행됩니다.
- 비결정적 조각은 액티비티 (activities)라고 하며, 임의 코드를 실행해 I/O와 다른 모든 연산을 수행할 수 있어요.
워크플로 코드는 오래 실행될 수 있고, 중단되면 정확히 중단한 지점에서 재개됩니다. 결정적으로, 워크플로 코드는 일반적으로 네트워크, 디스크 등 어떤 종류의 I/O도 포함할 수 없어요. 액티비티 코드는 I/O나 외부 상호작용에 제한이 없지만, 액티비티가 도중에 실패하면 처음부터 다시 시작됩니다.
참고
celery에 익숙하다면, Temporal 액티비티를 celery 태스크와 비슷하되 태스크가 완료되어 결과를 얻을 때까지 기다렸다가 워크플로의 다음 스텝으로 진행하는 것으로 생각하면 도움이 될 거예요. 다만 Temporal 워크플로와 액티비티는 celery 태스크보다 훨씬 더 많은 유연성과 기능을 제공합니다.
자세한 내용은 Temporal 문서를 보세요.
Pydantic AI 에이전트의 경우, Temporal과의 통합은 모델 요청, I/O가 필요할 수 있는 도구 호출, MCP 서버 통신이 모두 그 I/O 요구 때문에 Temporal 액티비티로 오프로드되어야 하고, 그것들을 조정하는 로직(즉 에이전트 런)은 워크플로에 살아야 함을 의미해요. 스케줄된 작업이나 웹 요청을 처리하는 코드는 그다음 워크플로를 실행할 수 있고, 워크플로는 필요에 따라 액티비티를 실행합니다.
아래 다이어그램은 Temporal의 에이전트 애플리케이션 전체 아키텍처를 보여줍니다. Temporal Server는 프로그램 실행을 추적하고 관련 상태가 안정적으로 보존되도록(즉 내부 데이터베이스에 저장되고, 가능하면 클라우드 지역에 걸쳐 복제) 책임을 집니다. Temporal Server는 데이터를 암호화된 형태로 관리하므로, 모든 데이터 처리는 워크플로와 액티비티를 실행하는 Worker에서 일어나요.
+---------------------+
| Temporal Server | (Stores workflow state,
+---------------------+ schedules activities,
^ persists progress)
|
Save state, | Schedule Tasks,
progress, | load state on resume
timeouts |
|
+------------------------------------------------------+
| Worker |
| +----------------------------------------------+ |
| | Workflow Code | |
| | (Agent Run Loop) | |
| +----------------------------------------------+ |
| | | | |
| v v v |
| +-----------+ +------------+ +-------------+ |
| | Activity | | Activity | | Activity | |
| | (Tool) | | (MCP Tool) | | (Model API) | |
| +-----------+ +------------+ +-------------+ |
| | | | |
+------------------------------------------------------+
| | |
v v v
[External APIs, services, databases, etc.]
자세한 내용은 Temporal 문서를 보세요.
지속 에이전트 (Durable Agent)
런을 지속으로 만들려면 워커에서 실행되는 Temporal 워크플로 안에서 agent.run()을 호출하세요. 밖에서는 에이전트가 일반적인 비-지속 에이전트로 실행됩니다.
TemporalDurability 캐퍼빌리티를 붙여 어떤 Agent든 지속 실행을 추가할 수 있어요. 에이전트는 어디서나 일반 Agent로 유지됩니다 — 워크플로 밖에서는 투명하게 동작하고, 워크플로 안에서는 캐퍼빌리티가 모델 요청, 도구 호출, MCP 서버 통신을 Temporal 액티비티로 라우팅합니다.
런은 워크플로 안에서만 지속됩니다
TemporalDurability를 붙이는 것만으로는 런이 지속되지 않아요 — 지속성은 여러분의 애플리케이션이 시작한 Temporal 워크플로(아래 예제처럼 Temporal 클라이언트를 통해) 안에서 런을 실행함에서 옵니다. 일반 코드에서 agent.run()이나 스트리밍 메서드를 호출하는 것 — 웹 엔드포인트의 UI 어댑터를 통한 간접 호출을 포함 — 은 평범한 비-지속 에이전트 런입니다. API 뒤에서 지속 에이전트를 제공하려면, 엔드포인트가 워크플로를 시작(또는 신호)하고 그 결과·이벤트를 클라이언트에 연결해 주세요. 엔드포인트 자체에서 에이전트를 실행하지 마세요.
에이전트에 지속 실행을 붙이고, 지속 실행 로직으로 Temporal 워크플로를 만들고, Temporal 서버에 연결하고, 비-지속 코드에서 워크플로를 실행하는 간단하지만 완전한 예제입니다. temporal 옵션 그룹과 함께 Pydantic AI를 설치하기만 하면 됩니다:
pip install pydantic-ai[temporal]
uv add pydantic-ai[temporal]
slim 패키지를 쓰고 있다면 temporal 옵션 그룹으로 설치할 수 있어요:
pip install pydantic-ai-slim[temporal]
uv add pydantic-ai-slim[temporal]
로컬에서 실행 중인 Temporal 서버도 필요해요:
brew install temporal
temporal server start-dev
import uuid
from temporalio import workflow
from temporalio.client import Client
from temporalio.worker import Worker
from pydantic_ai import Agent
from pydantic_ai.durable_exec.temporal import (
PydanticAIPlugin,
PydanticAIWorkflow,
TemporalDurability,
)
agent = Agent(
'openai:gpt-5.6-sol',
instructions="You're an expert in geography.",
name='geography', # (1)
capabilities=[TemporalDurability()], # (2)
)
@workflow.defn
class GeographyWorkflow(PydanticAIWorkflow): # (3)
__pydantic_ai_agents__ = [agent] # (4)
@workflow.run
async def run(self, prompt: str) -> str:
result = await agent.run(prompt) # (5)
return result.output
async def main():
client = await Client.connect( # (6)
'localhost:7233', # (7)
plugins=[PydanticAIPlugin()], # (8)
)
async with Worker( # (9)
client,
task_queue='geography',
workflows=[GeographyWorkflow],
):
output = await client.execute_workflow( # (10)
GeographyWorkflow.run,
args=['What is the capital of Mexico?'],
id=f'geography-{uuid.uuid4()}',
task_queue='geography',
)
print(output)
#> Mexico City (Ciudad de México, CDMX)
- (1) 에이전트의
name은 그 액티비티를 고유하게 식별하는 데 쓰입니다. - (2)
capabilities=[...]로 지속성을 붙입니다. 이 캐퍼빌리티는 에이전트에 바인딩될 때 에이전트의 이름·모델·툴셋을 발견하고 각각 액티비티를 등록합니다. 워크플로 밖에서는 캐퍼빌리티가 투명합니다 — 에이전트가 일반Agent처럼 동작해요. - (3) 워크플로는 I/O가 필요한 연산에 비결정적 액티비티를 쓸 수 있는 결정적 코드 조각을 나타냅니다.
PydanticAIWorkflow서브클래스화는 선택이지만__pydantic_ai_agents__클래스 변수에 적절한 타입을 제공합니다. - (4) 이 워크플로가 쓰는 에이전트를 나열합니다.
PydanticAIPlugin이 각 에이전트의TemporalDurability캐퍼빌리티가 기여한 액티비티를 워커에 자동 등록합니다. 워커 초기화를 수정하는 것이 워크플로 클래스보다 쉽다면,AgentPlugin으로 에이전트의 액티비티를 워커에 직접 등록할 수도 있어요. - (5)
agent.run()이 평소처럼 동작합니다. 워크플로 안에서 모델 요청, 도구 호출, MCP 서버 통신은 Temporal 액티비티로 라우팅됩니다. - (6) 워크플로와 액티비티 실행을 추적하는 Temporal 서버에 연결합니다.
- (7) Temporal 서버가 로컬에서 실행 중이라고 가정합니다.
- (8)
PydanticAIPlugin은 Temporal에 직렬화·역직렬화에 Pydantic을 쓰라고 알리고,__pydantic_ai_agents__에 나열된 에이전트의 액티비티를 자동 등록합니다. 액티비티 재시도 정책은UserError,PydanticUserError,UnexpectedModelBehavior,FallbackExceptionGroup을 비-재시도로 취급하고, Temporal 자신의 초과 페이로드 실패 타입인PayloadsTooLarge와(temporalio1.31 전에는)PayloadSizeError도 (참고: Large Payloads) 비-재시도로 취급합니다. 한편 워커는UserError,PydanticUserError,AgentRunError,UnsupportedEventLoopError를workflow_failure_exception_types로 등록합니다. - (9) 지정된 태스크 큐를 듣고 워크플로·액티비티를 실행할 워커를 시작합니다. 실제 애플리케이션에서는 별도 서비스에서 실행할 수 있어요.
- (10) 지정된 태스크 큐를 듣는 워커에서 워크플로를 실행하라고 서버에 요청합니다.
(이 예제를 실행하려면 asyncio를 임포트하고 asyncio.run(main())을 추가하세요; 다른 변경은 필요 없어요.)
같은 에이전트가 워크플로 안팎에서 동작하므로, TemporalDurability는 각각 Temporal 특정 래퍼 변형 없이도 다른 모든 캐퍼빌리티(계측, SetToolMetadata, ProcessEventStream 등)와 조립됩니다.
워크플로 코드는 async API를 써야 해요. Agent.run_sync()는 이벤트 루프를 스스로 구동하는데, Temporal의 워크플로 이벤트 루프는 그것을 허용하지 않으므로 워크플로 안에서 호출하면 대신 await agent.run()을 하라고 알리는 UserError가 납니다. 워크플로 밖에서는 TemporalDurability가 있는 에이전트가 일반 에이전트처럼 동작하므로 run_sync()이 평소처럼 동작해요.
실제 애플리케이션에서는 에이전트, 워크플로, 워커를 보통 워크플로 실행을 요청하는 코드와 별도로 정의합니다. Temporal 워크플로는 파일 최상위에 정의되어야 하고, 에이전트는 워크플로 안과 워커 시작 시(액티비티 등록용)에 모두 필요하므로, 에이전트도 파일 최상위에 정의되어야 해요.
파이썬 애플리케이션에서 Temporal을 쓰는 법은 Python SDK guide를 보세요.
래퍼 에이전트 경로 (Wrapper-agent path, 더 이상 사용되지 않음)
더 이상 사용되지 않음 (Deprecated)
TemporalAgent는 Temporal 통합의 원래 래퍼 에이전트 경로이고 v3에서 제거됩니다. 새 코드는 위의TemporalDurability캐퍼빌리티를 사용하세요.
TemporalDurability는 TemporalAgent가 기록한 액티비티 이름과 페이로드 모양을 받아들입니다 (에이전트의 name, 툴셋 id, models= 레지스트리 키가 같게 유지되고, TemporalAgent-시대 워크플로가 아직 진행 중인 동안 event_stream_handler=가 TemporalDurability에 남아 있어야 합니다). 그래서 TemporalAgent 아래 시작된 워크플로는 전환 후에도 올바르게 재생됩니다 — 먼저 워크플로를 드레인하거나 버전 관리할 필요가 없어요. 이벤트 스트림 핸들러로 기록된 히스토리는 이벤트당 액티비티를 등록했으므로, 그 워크플로들이 끝나기 전에 ProcessEventStream으로 마이그레이션하면 재생이 깨져요.
어떤 에이전트든 TemporalAgent로 감싸서 Temporal 워크플로 안에서 쓸 수 있는 지속 에이전트 변형을 얻을 수 있습니다. 감싸는 시점에 에이전트의 모델과 툴셋이 고정되고, 각각에 액티비티가 동적으로 만들어지며, 원래 모델·툴셋은 워크플로 안에서 직접 동작을 수행하는 대신 워커에 해당 액티비티를 실행하라고 요청하도록 감싸집니다. 원래 에이전트는 Temporal 워크플로 밖에서 평소처럼 쓸 수 있지만, 감싼 후 모델·툴셋을 바꿔도 지속 에이전트에는 반영되지 않아요.
from pydantic_ai import Agent
from pydantic_ai.durable_exec.temporal import TemporalAgent
agent = Agent('openai:gpt-5.6-sol', name='geography')
temporal_agent = TemporalAgent(agent)
# Use `temporal_agent` from inside a workflow, list it in `__pydantic_ai_agents__`,
# and connect with `PydanticAIPlugin()` exactly like the capability example above.
Temporal 통합 고려사항 (Temporal Integration Considerations)
Temporal을 지속 실행에 쓸 때 에이전트와 툴셋에 특화된 몇 가지 고려사항이 있어요. 에이전트와 툴셋이 Temporal의 워크플로·액티비티 모델과 올바르게 동작하도록 이해하는 것이 중요합니다.
에이전트 이름과 툴셋 ID (Agent Names and Toolset IDs)
Temporal이 액티비티가 실패하거나 중단된 뒤 재시작할 때 무엇을 실행할지 알게 하려면, 그 사이 코드가 바뀌더라도 각 액티비티에 안정적이고 고유한 이름이 필요해요.
TemporalDurability가 에이전트의 모델 요청과 툴셋(특히 자체 도구 나열·호출을 구현하는 것, 즉 FunctionToolset과 MCPToolset)에 액티비티를 동적으로 만들 때, 그 이름은 에이전트의 name과 툴셋의 id에서 유도됩니다. 이 필드들은 보통 선택이지만 Temporal을 쓸 때는 설정이 필요해요. 지속 에이전트가 프로덕션에 배포된 후에는 바꾸면 안 됩니다. 활성 워크플로를 깨뜨릴 수 있으니까요.
이 버전으로 업그레이드하면 args_validator가 있는 도구의 액티비티 순서가 바뀌므로, 그런 도구를 호출하는 진행 중인 워크플로는 Temporal worker versioning이나 patch가 필요해요. args_validator가 있는 도구를 호출하지 않는 워크플로는 영향이 없습니다.
DynamicToolset와 DynamicCapability이 기여한 툴셋은 지원됩니다. 그 팩토리는 도구가 나열·호출될 때 액티비티 안에서 재해석되므로, 런 의존성에 대해 결정적이어야 해요. 다른 감싸진 툴셋처럼, 모든 DynamicToolset은 명시적 id를 요구합니다: 직접 만들 때 id=를 전달하거나, DynamicCapability에 안정적인 캐퍼빌리티 id를 설정하세요. Temporal에서는 per_run_step=False가 존중되지 않는데, 툴셋은 항상 액티비티에서 즉석으로 만들어져야 하기 때문이에요: 액티비티는 워크플로와 다른 워커 프로세스에서 실행될 수 있어, 워크플로가 해석한 툴셋을 재사용할 수 없습니다. DBOS나 Prefect처럼 지속 단위가 워크플로 자신의 프로세스에서 실행되는 엔진은 그것을 재사용해서, 팩토리가 반환한 MCPToolset이 런 동안 하나의 세션을 유지하게 합니다.
같은 이유로, 에이전트에 붙은 MCPToolset은 그 도구를 나열·호출하는 모든 액티비티 안에서 연결됩니다. 액티비티는 워크플로가 가진 세션을 쓸 수 없고, 그 워커가 세션을 가진 프로세스가 아닐 수 있기 때문이에요. 각 액티비티의 세션은 반환될 때 닫히므로, 서버의 cache_tools는 결코 그것을 채운 액티비티만 커버합니다. 워크플로는 그 get_tools 액티비티가 기록한 도구 정의를(cache_tools가 켜져 있는 한) 유지하므로, 발견은 모델 요청당이 아니라 런당 액티비티 하나를 듭니다.
툴셋을 기여하는 캐퍼빌리티 — tools=가 있는 Capability 또는 로컬 실행 MCP 서버 — 은 캐퍼빌리티 자신의 id에서 툴셋의 id를 유도하므로 Capability(id='...', tools=[...]) 또는 MCP(id='...', url='...')을 설정하세요. (MCP는 id가 주어지지 않으면 서버 URL의 호스트·경로에서 유도된 id로 폴백합니다.) toolsets=로 캐퍼빌리티에 전달된 툴셋은 자신의 id를 유지하며, 그것은 툴셋 자체에 설정되어야 해요.
다른 모든 에이전트와 툴셋은 지원됩니다.
도구 인수 검증 (Tool Argument Validation)
도구의 args_validator는 validate_args 액티비티에서 실행되므로 I/O를 수행할 수 있고 동적 도구 검증자는 워크플로 경계를 견뎌냅니다. 없는 도구는 추가 액티비티를 예약하지 않아요. 검증은 승인과 지연 전에 실행되므로 거부된 인수는 승인자에게 닿지 않습니다. 검증자는 ApprovalRequired나 CallDeferred를 발생시킬 수도 있는데, 승인된 호출을 재개하면 tool_call_approved가 설정된 채 검증이 다시 실행됩니다.
각 액티비티는 타입 있는 인수를 독립적으로 유도하므로 스키마 검증자는 멱등이어야 해요. Pydantic 검증 컨텍스트는 Agent에 설정된 validation_context에서 재구성되고, callable 컨텍스트 빌더는 액티비티 자신의 런 컨텍스트로 다시 실행됩니다.
에이전트 런 컨텍스트와 의존성 (Agent Run Context and Dependencies)
워크플로와 액티비티가 별도 프로세스에서 실행되므로, 그 사이에 전달되는 모든 값은 직렬화 가능해야 해요. 이 페이로드들이 워크플로 실행 이벤트 히스토리에 저장되므로, Temporal은 기본적으로 크기를 2MB로 제한합니다 — 그 예산이 실제로 얼마나 살 수 있는지는 Large Payloads를 보세요.
이 제한을 감안해, 액티비티 안에서 실행되는 도구 함수와 이벤트 스트림 핸들러는 에이전트의 RunContext의 제한된 버전을 받으며, Agent.run()에 제공하는 의존성 객체가 Pydantic으로 직렬화 가능하도록 하는 것은 여러분의 책임입니다.
deps만 액티비티로 건너가는 값이 아니에요. model_settings, RunContext의 metadata와 tool_call_metadata, 도구 메타데이터도 건너가며, 모두 Pydantic으로 직렬화 가능해야 합니다. 직렬화할 수 없는 값은 직렬화하지 못한 타입을 이름으로 말하는 UserError를 냅니다. 그 외에는 지원되는 일부 값도 포함됩니다: [model_settings['timeout']][pydantic_ai.settings.ModelSettings.timeout]을 httpx.Timeout 대신 초 단위 숫자로 전달하세요.
Pydantic AI가 경계를 가로질러 타입 없는 사전으로 운반하는 값들 — 도구 메타데이터, RunContext의 metadata와 tool_call_metadata, ApprovalRequired와 CallDeferred의 metadata, 메시지와 그 파트에 실린 metadata, provider_details, vendor_metadata — 은 원래 파이썬 객체가 아니라 JSON 모양으로 도착합니다: set이나 tuple은 list가 되고, dataclass·Pydantic 모델은 dict가 되고, UTF-8 디코딩 가능한 bytes는 str이 되고, 문자열이 아닌 사전 키는 문자열이 됩니다. 임의 바이너리는 통신에 아예 오지 않아요 — 먼저 base64로 인코딩하세요. 통신에는 이 중 무엇도 복원할 타입 정보가 없으므로, 이 페이로드들을 JSON 네이티브로 유지하거나 수신 측에서 기대하는 타입으로 재검증하세요(TypeAdapter 사용). 선언된 타입이 있는 모든 필드 — 각 메시지의 나머지, ToolDefinition, ModelRequestParameters, 그리고 런 컨텍스트의 usage, usage_limits, loaded_capability_ids, discovered_tool_names — 은 충실히 왕복합니다.
영속 페이로드 스키마
Temporal은 현재 배포된 워커에서 사용 가능한 모델·타입 주석으로 영속된 워크플로·액티비티 페이로드를 역직렬화하므로, 이 모델들을 배포 간 내구성 계약으로 취급하세요. 기본값이 있는 선택 필드를 추가하는 것은 호환되지만, 필수 필드를 추가하거나 다른 비호환 변경을 하면 워크플로·액티비티 본문이 실행되기 전에 페이로드 디코딩이 실패할 수 있어요. 이것은 특히 애플리케이션이 소유한 워크플로 입력·의존성 모델과 관련이 큽니다: Pydantic AI는 Temporal 워크플로 히스토리를 소유하거나 마이그레이션하지 않으므로, 장기 실행 워크플로가 있는 애플리케이션은 그것들을 바꿀 때 버전 관리·마이그레이션 전략을 채택해야 해요.
구체적으로, 기본적으로 사용 가능한 필드는 deps, run_id, conversation_id, metadata, retries, tool_call_id, tool_name, tool_call_approved, tool_call_metadata, retry, max_retries, run_step, partial_output, usage, usage_limits, trace_include_content, instrumentation_version, loaded_capability_ids, discovered_tool_names, capability_active뿐입니다. agent와 root_capability는 워커의 에이전트 인스턴스에서 다시 붙고, tool_manager는 None입니다 (그래서 available_tool_names은 액티비티가 디스패치될 때 해석된 스냅샷을 반환하고, serialize_run_context가 그것을 담지 않는 서브클래스에서만 discovered_tool_names로 폴백합니다; active_capability_ids는 같은 방식으로 운반되는데, 이것이 is_tool_available이 캐퍼빌리티가 소유한 도구에 답하게 합니다). pending_messages는 ctx.enqueue()가 액티비티 안에서 오류를 내게 하는 가드를 지니는데, 액티비티의 기록된 결과가 코드를 재실행하지 않고 재생되고 큐에 넣은 메시지가 버려질 것이기 때문입니다.
다른 어떤 필드 — model, prompt, messages, model_settings, tracer, validation_context, capabilities — 에 접근하려 하면 그 필드의 기본값을 반환하는 대신 UserError가 납니다. 그래서 경계를 건너지 못한 필드가 실제 런 상태로 오인될 수 없어요. 멀티모달 prompt는 큰 BinaryContent를 담을 수 있고, 그것을 운반하면 그 콘텐츠가 모든 액티비티 페이로드에 들어가 messages와 같은 Temporal 페이로드 크기 문제를 만듭니다. 이 속성 중 하나 이상이 액티비티 안에서 사용 가능해야 한다면, 커스텀 serialize_run_context와 deserialize_run_context 클래스 메서드로 TemporalRunContext 서브클래스를 만들고 그것을 TemporalDurability의 run_context_type 인수로 전달하세요. 서브클래스는 자신의 프롬프트가 텍스트 전용임을 안다면 prompt 운반에 옵트인할 수 있어요.
액티비티의 RunContext는 직렬화된 페이로드에서 재구성되므로 그 필드는 복사본입니다. 액티비티 안에서 변경해도 런에 영향이 없어요. 특히 usage는 액티비티가 스케줄될 당시 런 사용량의 스냅샷입니다. 도구가 usage=ctx.usage로 다른 에이전트에 위임하면, 그 위임자의 토큰과 요청은 액티비티에 남습니다. 부모 런의 result.usage에서 빠지고 그 사용량 제한에 절대 부과되지 않아요. 위임자 사용량을 반영하려면 직접 운반하세요: 도구에서 위임자의 result.usage를 반환하거나, deps가 닿을 수 있는 외부 저장소에 기록하세요. Temporal은 이 변경을 무조건 잃어요. DBOS와 Prefect는 실시간 RunContext를 인프로세스 지속 단위에 전달하므로, 스텝·태스크 본문이 실제 실행되는 동안 변경은 누적됩니다. 하지만 기록된 결과가 재생(DBOS 워크플로 회복)되거나 재사용(Prefect 태스크 캐시 히트)될 때 본문이 실행되지 않으면 거기도 잃어서, 같은 코드가 런마다 다르게 회계됩니다. 인프로세스 엔진의 동작에 의존하지 마세요. 세 엔진 모두에서 동작하는 반환 채널은 pydantic-ai#6886에서 논의 중입니다.
도구의 prepare 함수는 이 제한의 영향을 받지 않아요. FunctionToolset의 도구(에이전트 자체에 정의된 것 포함)의 경우 워크플로 코드에서 완전한 RunContext로 실행되며, 워크플로 밖처럼 런 스텝당 한 번입니다. 그것이 반환하는 도구 정의는 도구 호출 액티비티에 보내져 그대로 쓰이므로, 모델이 본 도구가 실행되는 도구입니다 — timeout까지요. DynamicToolset의 도구는 예외입니다. 툴셋이 액티비티 안에서 재해석되므로 그 prepare 함수도 거기서 실행되어 제한된 RunContext를 봅니다.
XSearch나 ImageGeneration의 native= 팩토리는 경계 양쪽에서 두 번 해석됩니다: 네이티브 도구를 구성하려고 워크플로 코드에서 한 번, 그리고 폴백 서브에이전트의 도구 호출 액티비티 안에서 제한된 RunContext를 보며 한 번. 거기서는 ctx.messages가 아니라 ctx.deps를 읽으세요.
런타임 캐퍼빌리티 (Capabilities at Runtime)
캐퍼빌리티는 에이전트가 생성될 때 붙여서, TemporalDurability.for_agent()가 워커 시작 전에 그 액티비티를 등록할 수 있게 하세요. 워크플로 안에서 agent.run(capabilities=[...])을 전달하면 UserError가 납니다. 그렇게 늦게 추가된 캐퍼빌리티는 기여하는 툴셋이나 자기 @durable_operation 메서드에 등록된 액티비티가 없기 때문입니다.
런만 관찰하는 캐퍼빌리티는 런별로 붙여도 안전해요. 그 훅은 런 상태를 읽지만 도구·툴셋·지속 연산을 기여하지 않으니까요. Instrumentation이 내장 예제이며 제한에서 면제됩니다. 현재 제한은 더 보수적인데, 서드파티 캐퍼빌리티가 아직 런만 관찰한다고 선언할 수 없기 때문이에요. 캐퍼빌리티가 오버라이드하는 훅에서 이를 유도하는 것은 #5477에서 추적 중입니다. 워크플로 안에서 런별 캐퍼빌리티가 필요하면 그곳에 사용 사례를 공유하세요. 워크플로 밖에서는 지속 캐퍼빌리티가 투명하므로 런별 캐퍼빌리티는 거기서 문제없어요.
큰 페이로드 (Large Payloads)
Temporal은 기록하는 모든 페이로드를 워크플로 실행 이벤트 히스토리에 저장합니다 — 액티비티가 반환하는 것과 스케줄된 인수 둘 다 — 그리고 각각을 기본적으로 2MB로 제한합니다. 두 가지가 자주 그 한도에 걸립니다:
- 바이너리 콘텐츠: 도구가 반환하거나 모델이 응답에 올리는 이미지 같은 것. 바이너리 데이터는 페이로드로 들어가는 길에 base64 인코딩되므로, 원시 바이트의 사용 가능한 예산은 한도의 약 4분의 3 — 대략 1.5MB — 이에요. 크기가 2MB 아래로 읽히는 이미지도 거부될 수 있습니다.
- 큰 의존성 객체: Pydantic AI가 모델·도구·MCP·이벤트 핸들러 액티비티를 스케줄할 때마다 히스토리에 복사됩니다 — 그래서 넓은 도구 팬아웃은 히스토리를 영구히 부풀릴 수도, 페이로드당 한도를 범할 수도 있어요.
도구 반환이나 모델 응답이 너무 클 때
BinaryImage를 반환하는 도구 — ImageGeneration 캐퍼빌리티 뒤의 직접 이미지 생성기든 서브에이전트 폴백이든, 또는 여러분의 local= 도구든 — 는 그 바이트를 액티비티의 결과 페이로드로 보내고, Temporal이 거부하면 Pydantic AI가 도구 이름을 말하는 UserError를 발생시킵니다. 모델이 응답에 올리는 이미지에도 같은 것이 적용되는데, 네이티브 이미지 생성 도구처럼 모델 요청 자체가 액티비티이기 때문입니다. 거기서는 UserError가 모델 이름을 말해요. 스트리밍된 세그먼트는 이미지가 없이도 그 버퍼링된 이벤트만으로 넘칠 수 있습니다.
제어하는 도구의 경우 가장 간단한 수정은 바이트를 히스토리에서 완전히 빼는 것입니다: 이미지를 객체 스토리지에 쓰고 참조를 반환하세요 — ImageUrl 또는 나중에 애플리케이션이 해석하는 키. 작은 페이로드는 워크플로 재생도 빠르게 유지하므로, 이는 한도 아래 유지에 그치지 않고 이점이 있습니다. 바이트가 공급자에게서 돌아오거나, 페이로드가 미디어가 아닌 큰 deps 객체라면 외부 스토리지를 쓰세요.
페이로드를 자신의 스토리지로 오프로드
Temporal은 이를 외부 스토리지로 일반적으로 해결하는데, 크기 임계값 위의 어떤 페이로드든 자신의 저장소로 투명하게 오프로드하고 히스토리에는 참조를 넣습니다. 양방향 모든 페이로드에 적용되므로 UserError가 커버할 수 없는 경우도 다룹니다. DataConverter의 ExternalStorage로 구성하세요:
import dataclasses
from temporalio.client import Client
from temporalio.converter import DataConverter, ExternalStorage
from pydantic_ai.durable_exec.temporal import PydanticAIPlugin
from my_app.storage import my_storage_driver # your `StorageDriver` implementation
async def connect() -> Client:
data_converter = dataclasses.replace(
DataConverter.default,
external_storage=ExternalStorage(drivers=[my_storage_driver]),
)
return await Client.connect(
'localhost:7233',
data_converter=data_converter,
plugins=[PydanticAIPlugin()],
)
PydanticAIPlugin은 Temporal의 기본 페이로드 컨버터 만을 Pydantic 인식 버전으로 교체하므로, 여러분의 external_storage — 커스텀 payload_codec나 failure_converter_class처럼 — 는 보존됩니다.
임계값(기본 256KiB) 아래의 페이로드는 건드리지 않으므로, 의존성이 작은 런에서는 비용이 없어요. 사용 가능한 스토리지 드라이버와 직접 작성법은 Temporal 문서를 보세요. 외부 스토리지는 공개 프리뷰라 API가 아직 바뀔 수 있음에 유의하세요.
저장된 페이로드는 워크플로 히스토리의 일부입니다
빠지거나·만료됐거나·접근 불가능한 페이로드는 워커가 워크플로를 디코딩·재생하지 못하게 합니다. 저장된 페이로드를 최소한 그 히스토리가 재생될 수 있는 기간만큼 사용 가능하게 유지하세요.
오프로드가 페이로드별로 일어나므로, N개 액티비티의 팬아웃은 여전히 N개의 (작은) 참조를 발행합니다. 외부 스토리지는 크기 상한과 히스토리 부풀림을 제거합니다. 단일 워크플로-태스크 완료 내에서 동일한 의존성을 중복 제거하지는 않아요.
공개 프리뷰 API를 채택할 수 없다면, Temporal은 PayloadCodec 위에 만든 손으로 쓴 claim check codec도 문서화하며, PydanticAIPlugin은 같은 방식으로 그것을 보존합니다.
마지막 수단으로, 바이트가 정말 워크플로 히스토리에 있어야 한다면 Temporal 서버에서 limit.blobSize.error 동적 구성을 올리세요. Temporal 자신의 gRPC 메시지 한도는 그 위에 여전히 적용되므로, 이것은 한도를 제거가 아니라 높이는 것임을 유의하세요 — Temporal의 blob size limit 안내를 보세요.
액티비티가 반환하는 것만 친절한 오류를 얻습니다
액티비티 안으로 들어가는 값 — 모델 요청이 운반하는 메시지 히스토리나 event_stream_handler에 넘겨진 스트리밍 이벤트 — 은 워크플로 코드에서 인코딩되며, 거기서 Temporal은 액티비티가 아니라 워크플로 태스크 를 실패시킵니다. 그 실패는 Pydantic AI에 절대 도달하지 않고, 워크플로-태스크 실패이지 액티비티 실패가 아니므로 재시도 정책도 그것을 묶지 않아요. 런은 발생시키는 대신 무기한 재시도합니다. 그래서 맞는 도구 반환도 다음 모델 요청의 히스토리의 일부가 되면 너무 클 수 있어요. 외부 스토리지가 여기서 안정적인 해결책인데, 양방향 모든 페이로드에 적용되니까요.
UserError는 또한 초과 페이로드를 식별하는 데 Temporal의 기본 실패 컨버터에 의존하므로, 자신의 failure_converter_class를 공급하거나 한도를 워커에 보고하지 않는 서버에 대해 실행하면 발화하지 않아요.
이미지 출력 타입은 아예 지원되지 않습니다
output_type이 BinaryImage를 포함하는 에이전트는 크기로 실패하기 전에 — 요청이 만들어지기도 전에 — UserError를 냅니다. 그런 런은 생성된 이미지를 매번 경계를 가로질러 운반해야 하기 때문입니다.
스트리밍 (Streaming)
Agent.run_stream(), Agent.run_stream_events(), Agent.iter()는 Temporal 워크플로 안에서 동작하지만, 그 이벤트는 실시간으로 전달되지 않고 버퍼링됩니다. 모델 스트림은 지속 액티비티 안에서 실행되고, 그 이벤트는 액티비티가 완료된 후 워크플로에 재생됩니다.
I/O 사이드 이펙트가 있는 핸들러는 event_stream_handler=를 TemporalDurability에 전달하세요. 모델 이벤트는 각 모델 요청 액티비티 안에서 실시간 전달되고, 각 도구 이벤트는 자체 이벤트 핸들러 액티비티에서 전달됩니다. 다른 Temporal 액티비티처럼, 액티비티가 재시도되면 핸들러가 두 번 이상 실행될 수 있으니 사이드 이펙트를 멱등으로 유지하세요.
대안으로 ProcessEventStream을 등록하세요. 그 핸들러는 워크플로 코드에서 실행되고, 워크플로 재생 때 다시 실행되므로 결정적이어야 해요. 도구와 최종 출력 이벤트는 실시간으로 도착하고, 실제 캡처된 모델 이벤트는 각 모델 요청이 완료된 후 재생됩니다. 예제는 streaming docs를 보세요.
지속성 event_stream_handler=와 별도로 등록된 ProcessEventStream은 두 개의 별개 핸들러이며 각각 한 번 발동합니다. 지속 핸들러는 지속 액티비티 안에서 실시간 이벤트를 받고, ProcessEventStream은 워크플로 코드의 버퍼링된 재생을 봅니다.
Agent.run(event_stream_handler=...)에 전달된 런별 핸들러도 재생된 모델 이벤트에 대해 워크플로 측에서 실행됩니다.
스트리밍 모델 요청 액티비티, 워크플로, 워크플로 실행 호출이 모두 별도 프로세스에서 일어나므로, 그 사이에 데이터를 전달하려면 주의가 필요합니다:
- 이벤트 스트림 핸들러에 워크플로 호출 지점이나 워크플로에서 데이터를 얻으려면 의존성 객체를 쓸 수 있어요.
- 이벤트 스트림 핸들러에서 워크플로, 워크플로 호출 지점, 프런트엔드로 데이터를 얻으려면 Workflow Stream에 이벤트를 발행하거나(권장, 추가 인프라 없음), 이벤트 스트림 핸들러가 쓰고 이벤트 소비자가 읽을 수 있는 외부 시스템(메시지 큐 같은)을 써요. 의존성 객체로 필요한 곳 모두에서 같은 연결 문자열이나 다른 고유 ID를 사용할 수 있게 하세요.
Workflow Streams로 프런트엔드에 이벤트 스트리밍
메시지 큐를 세우는 대신, Temporal의 내장 Workflow Streams를 전송 수단으로 쓸 수 있어요. 워크플로 자체가, 그 밖의 소비자가 구독하는 지속적이고 오프셋 주소 지정된 채널이 됩니다.
TemporalDurability에 event_stream_topic을 설정하고 워크플로에 AgentEventStream을 호스팅하세요. 토픽을 설정하면 그 자체로 스트리밍이 켜지고, event_stream_handler와 직교합니다. 핸들러도 전달하면 핸들러도 여전히 모든 이벤트를 봅니다.
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
from pydantic_ai import Agent
from pydantic_ai.durable_exec.temporal import AgentEventStream, TemporalDurability
durability = TemporalDurability(event_stream_topic='agent-events')
agent = Agent('openai:gpt-5.6-sol', name='assistant', capabilities=[durability])
@workflow.defn
class AssistantWorkflow:
@workflow.init
def __init__(self, prompt: str) -> None:
self.events = AgentEventStream()
@workflow.run
async def run(self, prompt: str) -> str:
async with self.events:
result = await agent.run(prompt)
return result.output
워크플로 핸들만 가진 소비자는 stream_agent_events()로 런이 진행되는 것을 실시간 관찰합니다:
from uuid import uuid4
from temporalio.client import Client
async def relay_events(client: Client, prompt: str) -> str:
handle = await client.start_workflow(
AssistantWorkflow.run, prompt, id=f'assistant-{uuid4()}', task_queue='my-task-queue'
)
async for event in durability.stream_agent_events(client, handle, output_type=str):
... # forward `event` to the frontend over SSE
return await handle.result()
각 독립 워크플로 체인에 고유한 워크플로 ID를 할당하세요. continue-as-new가 그 ID를 체인 전체에 자동으로 유지하므로, 무관한 워크플로에 재사용하지 마세요.
이것은 사실상 워크플로 경계를 가로지르는 지속적인 run_stream_events()입니다. 런의 결과를 실은 AgentRunResultEvent로 끝나는 같은 AgentStreamEvent들이 순서대로요 — 단, 이벤트들이 다른 프로세스로 건너와 여기 도착했다는 점만 빼요. 그래서 async for는 런이 끝나면 스스로 끝납니다. 언제 멈출지 아는 별도 신호가 필요 없어요.
그것을 가능하게 하는 것은 async with 블록을 떠나는 것입니다. Workflow Stream은 워크플로 자체가 서빙하므로, 워크플로가 반환하면 그 스트림은 더 이상 읽을 수 없고 구독자가 아직 폴링하지 않은 것은 사라져요. 블록은 구독자가 종료 이벤트를 인정할 때까지 — stream_agent_events() 안에서 자동으로 — 워크플로를 열어 둡니다. 아무도 구경하지 않는 것이 오류는 아니에요. 대기는 drain_timeout(기본 30초)으로 제한되고, 어느 쪽이든 워크플로의 결과가 권위를 유지합니다.
UI 프로토콜 구동
스트림이 AgentRunResultEvent로 끝나므로, 그것은 정확히 UIAdapter가 소비하는 것입니다. HTTP 핸들러가 워크플로를 시작하고 토픽에서 직접 UI 이벤트 스트림 프로토콜을 서빙할 수 있어요 — 인프로세스 런처럼 런 결과를 받는 on_complete를 포함해서요.
어댑터의 두 가지 일이 경계를 가로질러 나뉩니다: HTTP 핸들러가 요청을 프로토콜 스트림으로 바꾸고, 워크플로가 같은 요청 본문에서 런 인수를 재구성합니다.
@workflow.defn
class ChatWorkflow:
@workflow.init
def __init__(self, body: bytes) -> None:
self.events = AgentEventStream()
@workflow.run
async def run(self, body: bytes) -> str:
adapter = VercelAIAdapter(agent=agent, run_input=VercelAIAdapter.build_run_input(body))
async with self.events:
result = await agent.run(
message_history=adapter.messages,
deferred_tool_results=adapter.deferred_tool_results,
conversation_id=adapter.conversation_id,
)
return result.output
핸들러는 프런트엔드가 돌아올 수 있는 ID로 그 워크플로를 시작하고 토픽을 스트리밍합니다:
from starlette.requests import Request
from starlette.responses import Response
from pydantic_ai.ui.vercel_ai import VercelAIAdapter
async def post_chat(request: Request) -> Response:
body = await request.body()
run_input = VercelAIAdapter.build_run_input(body)
handle = await client.start_workflow(
ChatWorkflow.run,
body,
id=f'chat-{run_input.id}-{uuid4()}',
task_queue='my-task-queue',
)
adapter = VercelAIAdapter(agent=agent, run_input=run_input)
events = durability.stream_agent_events(client, handle, output_type=str)
return adapter.streaming_response(adapter.transform_stream(events))
워크플로가 큐입니다. 시작하면 런이 Temporal 워커에 넘어가고, HTTP 요청은 그 워커가 만드는 것의 구독자일 뿐입니다. 요청은 런을 건드리지 않고 사라질 수 있고, Temporal이 런의 액티비티를 스스로 재시도하며, 런은 그것을 시작한 프로세스를 견뎌냅니다.
그것은 또한 프런트엔드가 돌아올 수 있다는 뜻입니다. 워크플로 ID를 대화 옆에 — 대화 자체에 이미 강제하는 같은 소유권 아래 — 저장하고, 두 번째 엔드포인트가 재연결 클라이언트를 그것이 시작하지 않은 런에 다시 붙입니다:
async def get_chat(request: Request) -> Response:
# Resolve the run through the conversation the caller owns. A workflow ID taken
# straight off the path would let anyone replay anyone else's conversation, since
# a Temporal handle carries no authorization of its own.
chat = await load_chat(request.path_params['chat_id'], user=request.state.user)
handle = client.get_workflow_handle(chat.workflow_id)
adapter = VercelAIAdapter(agent=agent, run_input=chat.run_input)
events = durability.stream_agent_events(client, handle, output_type=str)
return adapter.streaming_response(adapter.transform_stream(events))
기본 from_offset=0에서 구독하면 런이 지금까지 발행한 모든 이벤트를 재생한 뒤 실시간으로 계속하므로, 새로 고친 탭은 문장 중간에서 이어받지 않고 전체 메시지를 재구성합니다. 이미 받은 것을 들고 있는 클라이언트는 from_offset=offset + 1을 전달해 나머지를 연속으로 가져올 수 있어요. 참고: Reconnecting.
창은 경계가 있습니다. 종료 이벤트를 아무도 인정하지 않는 런은 AgentEventStream의 drain_timeout(기본 30초) 후에 끝나고, 그 스트림도 함께 사라집니다. 그보다 나중에 재연결하는 클라이언트는 토픽이 아니라 handle.result()에서 런의 결과를 얻습니다. 더 긴 창이 필요하면 drain_timeout을 높이세요.
재연결 (Reconnecting)
Workflow Streams는 오프셋 주소 지정되므로, 끊긴 소비자는 중단한 지점에서 재개할 수 있어요 — 이것은 일반 인프로세스 스트리밍이 제공할 수 없는 것이에요. 진행하면서 offset을 체크포인트하고 from_offset=offset + 1로 재연결하세요:
from temporalio.service import RPCError
last_offset = -1
while True:
events = durability.stream_agent_events(client, handle, from_offset=last_offset + 1)
try:
async for event in events:
print(event) # forward to the frontend over SSE
last_offset = events.offset
except RPCError:
continue # the connection dropped: resubscribe at the next offset
break # the iterator ended on its own, so the run is over
return await handle.result() # the output, or the workflow's failure
async for의 깨끗한 종료는 런이 끝났다는 뜻입니다 — 결과를 만들었든 아니든: events.result는 성공 시 런의 AgentRunResult, 워크플로가 다른 방식으로 끝났으면 None이고, 그 경우 handle.result()가 실패·취소를 전파합니다. 중단된 반복만 재구독할 가치가 있어요.
오프셋은 토픽별이 아니라 전체 스트림에 걸쳐 실행되므로, 토픽 필터 구독은 워크플로가 다른 토픽에 발행한 곳마다 공백을 봅니다 — 그래서 세지 않고 스트림에서 읽어야 합니다.
무엇을 발행할지 고르기
모델 스트림은 토큰당 PartDeltaEvent를 발행하고, 발행된 모든 이벤트는 런의 수명 동안 워크플로 상태에 남습니다. 비용과 워크플로 크기를 모두 줄이려면 베어 이름 대신 WorkflowStreamTopic을 전달하고 필터링하세요:
from pydantic_ai.durable_exec.temporal import TemporalDurability, WorkflowStreamTopic
from pydantic_ai.messages import PartDeltaEvent
TemporalDurability(
event_stream_topic=WorkflowStreamTopic(
'agent-events',
events=lambda event: not isinstance(event, PartDeltaEvent), # skip per-token deltas
)
)
종료 AgentRunResultEvent는 항상 발행됩니다. 구독을 끝내는 것이니까요.
내부적으로 실시간 모델 스트림은 workflow_stream_event_handler()가 발행하며, 그것은 직접 구성하거나 감쌀 수 있는 일반 EventStreamHandler를 반환합니다. event_stream_topic을 선호하세요. 그것은 또한 워크플로 코드에서 런의 워크플로 측 이벤트와 구독을 끝내는 종료 이벤트도 발행하는데, 어느 것도 핸들러 혼자서는 할 수 없어요.
주의 사항
- Workflow Streams는 왕복당 약 100ms의 지연을 더합니다(토픽의
batch_interval로 조정 가능). UI 구동에 적합하며, 실시간 음성 같은 초저지연 사용 사례에는 맞지 않아요. - 각
stream_agent_events()구독은 Temporal 서비스에 롱폴을 보유하므로, 프런트엔드 연결당 하나를 열면 Temporal 사용량이 사용자 수에 비례해 커지고 Temporal Cloud의 연결 한도에 부딪힐 수 있어요. 워크플로 런당 소비자 하나를 실행하고 그 이벤트를 프런트엔드 연결들에 직접 팬아웃하세요(예: pub/sub 또는 SSE 브로드캐스트). 구독자당 구독은 개발에서 동작하고 한도에서 무너집니다. - 모델 이벤트는 모델 요청 액티비티 안에서 발행되므로, 그 액티비티가 재시도되면 이벤트가 새 오프셋에서 다시 발행됩니다. 소비자는 중복을 견뎌야 해요. 워크플로 코드에서 발행된 이벤트(도구 이벤트와 종료 이벤트)는 영향이 없습니다. 재생이 로그에 추가하는 것이 아니라 로그를 재구성하니까요.
- 스트림은 지속적이므로, 다른 Pydantic AI 버전이 돌아가는 프로세스가 이벤트를 만들고 소비할 수 있어요. 이벤트 모양은 메이저 버전 내에서 안정적입니다. 생산자와 소비자를 같은 메이저 버전에 두세요.
- 하나의 이터레이터가 하나의 에이전트 런을 덮습니다. 에이전트를 반복 실행하는 워크플로는 런마다 종료 이벤트를 발행하므로, 다음 것을 원하는 소비자는
from_offset=offset + 1로 재연결합니다. - 종료 이벤트는 모든 캐퍼빌리티가 결과를 변형한 후, 하지만 어떤 것도 감싸지 않는 런 라이프사이클의 자체 마무리 전에 발행됩니다. 그 마지막 스텝에서 취소된 런은 결과를 발행한 뒤 실패하므로, 워크플로의 반환 값이 권위 있는 결과로 남습니다.
- 초기 구독은 요청된 워크플로 실행에 고정되지만, continue-as-new 후 Temporal SDK는 고정되지 않은 워크플로 ID로 체인을 따릅니다. 구독자가 여전히 체인을 따라갈 수 있는 동안 그 워크플로 ID를 독립 실행에 재사용하지 마세요. 새 실행에 붙을 수 있으니까요.
- 실시간 이벤트는 워크플로 밖의 소비자만 닿습니다. 워크플로 코드 안의
run_stream_events()는 여전히 버퍼링합니다.
도구나 이벤트 스트림 핸들러에서 ctx.emit()으로 이벤트를 발행하는 것은 현재 지원되지 않습니다. 그것들이 런의 이벤트 스트림에 닿을 수 없는 액티비티 안에서 실행되기 때문이며, 그렇게 하면 UserError가 나요. 이것은 애플리케이션 도구가 발행하는 커스텀 이벤트와 캐퍼빌리티 자신의 도구가 발행하는 캐퍼빌리티 이벤트를 모두 다룹니다. 워크플로에서 실행되는 캐퍼빌리티 훅에서 이벤트를 발행하세요. 액티비티에서의 발행 지원은 pydantic-ai#7971에서 추적합니다.
@on_event로 등록된 캐퍼빌리티 리스너는 지속 단위가 아니라 워크플로 코드에서 실행되므로 매 워크플로 재생마다 다시 실행되고 결정적이어야 해요. I/O는 자신의 액티비티에서 실행되는 지속성 event_stream_handler=에 두세요.
모델 스트림이 액티비티 안에서 소비되므로, 워크플로 측에서 취소하는 것(AgentStream.cancel() 같은)은 지속 경계를 넘어 사용할 수 없어요. 진행 중인 모델 요청을 멈추려면 Temporal 워크플로를 취소하세요. 취소는 (하트비트를 통해) 액티비티에 전달되고, 액티비티는 완료 전에 서버 측 작업을 취소합니다.
전체 런 취소(Cancelling a Run 참고)는 같은 분할을 따르며 Temporal 특유의 결과가 있습니다:
- 워크플로 코드에서
AgentRun.cancel()을 호출하면RunCancelled가 평범한 애플리케이션 결과로 발생합니다: 그것을 잡는 워크플로는 Cancelled 로 끝나는 대신 정상 완료되고, 런은 재생 결정적으로 유지됩니다. 잡지 않은RunCancelled는 타입 있는 애플리케이션 오류로 워크플로를 실패시키고, 런 상태는 실패 경계를 건너지 않습니다 —all_messages()가 필요하면 워크플로 안에서 잡으세요. RunContext.cancel()은 런과 같은 프로세스에 있어야 하므로, 액티비티 안에서 실행되는 도구에서 호출하면 매달리는 대신 명확한UserError가 납니다.CancellationToken도 같은 프로세스 상태이며 Temporal 지속 런에 전달할 수 없습니다. 대신 Temporal 워크플로를 취소하세요.- Temporal 워크플로 자체를 취소하는 것은 외부 취소로 남습니다:
CancelledError가 계속 전파되고 워크플로는 여전히 Cancelled 로 끝납니다.
Agent.run_stream_sync()는 워크플로 코드용이 아니에요. 실행 중인 이벤트 루프가 필요 없고 run_stream()을 감쌉니다. TemporalDurability 아래에서는 위 버퍼링 async 스트리밍 API나 이벤트 스트림 핸들러가 있는 Agent.run()을 쓰세요. 워크플로 밖에서는 TemporalDurability가 있는 에이전트가 일반 에이전트처럼 동작하므로 run_stream_sync()이 평소처럼 동작합니다. (래퍼 TemporalAgent는 워크플로 안에서 run_stream을 금지합니다 — 거기서는 run + 이벤트 스트림 핸들러를 쓰세요.)
일시 정지 턴과 백그라운드 모드 (Suspended Turns and Background Mode)
일부 공급자는 모델 턴을 도중에 일시 정지하거나(Anthropic pause_turn) 준비될 때까지 폴링되는 서버 측 작업으로 실행할 수 있어요(OpenAI background mode). Pydantic AI는 그런 일시 정지된 턴을 완료될 때까지 투명하게 계속합니다. 각 세그먼트는 별도의 모델 요청 액티비티에서 실행되고, 워크플로는 세그먼트 사이에 일시 정지된 ModelResponse와 그 백그라운드 작업 ID를 체크포인트합니다. 최종 응답은 병합되고 사용량은 한 번 기록됩니다. 일시 정지된 응답으로 끝나는 message_history는 그 응답이 첫 액티비티에 전달된 채 재개됩니다.
운영상 몇 가지 의미가 있습니다:
- 타임아웃과 하트비트:
start_to_close_timeout과heartbeat_timeout을 공급자 왕복 한 번으로 잡으세요. 모델 요청 액티비티에는 기본heartbeat_timeout30초가 주어집니다. 다른 액티비티들에서 하트비팅이 어떻게 동작하는지는 Activity Configuration을 보세요. - 재시도와 대기: 실패한 세그먼트는 독립적으로 재시도됩니다. 백그라운드 폴링 사이의 지연은 지속 Temporal 타이머를 쓰고 액티비티 벽시계 시간을 소모하지 않아요.
- 취소: 오류가 일시 정지된 작업을 버리면, 그 공급자 정리는 전용 취소 액티비티에서 실행됩니다.
- 페이로드 크기: 스트리밍 —
event_stream_handler,ProcessEventStream캐퍼빌리티, 또는 런별event_stream_handler— 을 쓸 때마다 각 세그먼트의 버퍼링된 이벤트가 워크플로로 다시 보내져 Temporal의 페이로드 크기 한도(기본 2MB) 안에 맞아야 합니다. 넘치는 세그먼트는 다른 과대 모델 응답과 같은UserError를 냅니다.
참고
자체
serialize_run_context로 커스텀TemporalRunContext서브클래스를 쓴다면usage와usage_limits필드를 계속 포함하세요. 액티비티 안에서 실행되는 도구·캐퍼빌리티가RunContext에서 그것들을 읽어, 예를 들어 런의 남은 사용량 예산에 적응합니다.
런타임 모델 선택 (Model Selection at Runtime)
Agent.run(model=...)은 보통 모델 문자열('openai:gpt-5.6-sol' 같은)과 모델 인스턴스를 모두 지원해요. Temporal 아래에서 모델 인스턴스는 재생 메커니즘을 위해 직렬화할 수 없고, model_id 문자열에서 재구성하면 다른 모델을 만들 것입니다 — 워커 환경이 암시하는 어떤 공급자의 같은 모델 이름이므로, 요청이 다른 자격 증명으로 다른 엔드포인트에 갈 거예요. 그래서 미리 등록되지 않은 인스턴스는 UserError로 거부됩니다. 워크플로 안에서 특정 인스턴스를 쓰는 방법은 두 가지예요:
- 사전 등록:
TemporalDurability에modelsdict를 전달한 뒤 이름으로 참조하거나, 등록된 인스턴스를agent.run(model=...)에 직접 전달; - 모델 이름 문자열을 전달하고
ResolveModelId캐퍼빌리티로 워커에서 인스턴스를 만드세요 — 모델이 런의deps에 의존할 때(예: 사용자별 자격 증명) 올바른 선택입니다.
모델 이름 문자열은 등록이 절대 필요 없어요. 생성 시 설정된 에이전트 자신의 모델은 항상 기본값으로 사용 가능하며, 에이전트는 생성될 때 모델이 설정되어 있어야 합니다.
모델 문자열은 기대한 대로 동작합니다. 모델 문자열이 만들어지는 방식을 커스터마이즈 — 커스텀 공급자, 설정에서 주입된 API 키, 런의 deps에 실린 사용자별 자격 증명 — 하려면 TemporalDurability 앞에 ResolveModelId 캐퍼빌리티를 추가하세요. 그것이 런 설정 시와 액티비티 안에서 모델이 재구성될 때 모두 모든 문자열에 첫 기회를 얻습니다. 거기서 해석기가 런의 실제 deps로 다시 실행됩니다. 해석기가 워커에서 다시 실행되므로, 주어진 (model_id, deps)에 대해 결정적이어야 하며 외부 I/O를 수행하면 안 됩니다 — 자격 증명을 deps에 싣거나 시작 시 로드된 설정을 클로저로 감싸세요.
여러 모델을 사전 등록하고 사용하는 예제는:
import os
from temporalio import workflow
from pydantic_ai import Agent
from pydantic_ai.capabilities import ResolveModelId
from pydantic_ai.durable_exec.temporal import TemporalDurability
from pydantic_ai.models import Model, ModelResolutionContext, infer_model
from pydantic_ai.models.anthropic import AnthropicModel
from pydantic_ai.models.google import GoogleModel
from pydantic_ai.models.openai import OpenAIResponsesModel
from pydantic_ai.providers.openai import OpenAIProvider
# Create models from different providers
default_model = OpenAIResponsesModel('gpt-5.6-sol')
fast_model = AnthropicModel('claude-haiku-4-5')
reasoning_model = GoogleModel('gemini-3-pro-preview')
# Optional: customize how model-name strings are built.
def resolve_model(ctx: ModelResolutionContext[None], model_id: str) -> Model | None:
if model_id.startswith('openai:'):
provider = OpenAIProvider(api_key=os.environ['OPENAI_API_KEY'])
return infer_model(model_id, provider_factory=lambda _: provider)
return None # everything else takes the default `infer_model` path
agent = Agent(
default_model,
name='multi_model_agent',
capabilities=[
ResolveModelId(resolve_model), # Optional
TemporalDurability(
models={
'fast': fast_model,
'reasoning': reasoning_model,
},
),
],
)
@workflow.defn
class MultiModelWorkflow:
@workflow.run
async def run(self, prompt: str, use_reasoning: bool, use_fast: bool) -> str:
if use_reasoning:
# Select by registered name
result = await agent.run(prompt, model='reasoning')
elif use_fast:
# Or pass the registered instance directly
result = await agent.run(prompt, model=fast_model)
else:
# Or pass a model string (resolved by `ResolveModelId` if it matches)
result = await agent.run(prompt, model='openai:gpt-5.6-luna')
return result.output
런타임 툴셋 (Toolsets at Runtime)
지속 감싸기가 필요한 실행 중인 모든 툴셋을 에이전트 생성자에 전달해서 워크플로 실행 전에 액티비티가 워커에 등록되게 하세요. DynamicToolset 포함: 명시적 id를 주고 Agent(toolsets=[...])로 전달합니다. @agent.toolset 데코레이터는 엔진의 지속 단위가 만들어진 후 등록하므로, TemporalDurability 아래 워크플로 안에서 쓰면 UserError가 납니다. 더 이상 사용되지 않는 TemporalAgent는 이 검사를 실행하지 않아요. 워크플로 안에서 감싸는 시점에 동결된 툴셋 목록을 실행하므로, 그렇게 늦게 등록된 툴셋은 조용히 런에서 빠집니다.
추가 툴셋은 agent.run(toolsets=...)으로 런별 전달할 수 있지만, 지속 감싸기가 필요 없는 툴셋만 지원됩니다: 에이전트 런 밖에서 실행되는 도구를 가진 ExternalToolset 같은 비-실행 툴셋과, 모든 도구가 metadata={'temporal': False}로 액티비티 감싸기를 옵트아웃한 FunctionToolset이요. 다른 실행 툴셋(FunctionToolset, MCPToolset)과 동적 툴셋은 에이전트를 생성할 때 설정되어 워크플로 실행 전에 워커에 등록되게 해야 하며, 런타임에 전달하면 UserError가 납니다.
워크플로 안에서 agent.override(toolsets=...)로 바꿔 넣은 툴셋도 같은 규칙에 묶입니다. 그것들도 에이전트의 액티비티가 등록된 후 도착하기 때문이에요. 런타임에 추가된 툴셋은 에이전트를 생성할 때 만든 것의 id를 재사용할 수도 없는데, id가 도구 호출이 어떤 등록된 툴셋의 액티비티로 디스패치되는지 식별하기 때문입니다.
액티비티 설정 (Activity Configuration)
temporalio.workflow.ActivityConfig 객체를 TemporalDurability 생성자에 전달해 타임아웃·재시도 정책 같은 Temporal 액티비티 설정을 커스터마이즈할 수 있어요:
activity_config: 모든 액티비티에 쓸 기본 Temporal 액티비티 설정. 제공하지 않으면start_to_close_timeout60초가 사용됩니다.model_activity_config: 모델 요청 액티비티에 쓸 Temporal 액티비티 설정. 기본 액티비티 설정과 병합됩니다. 모델 요청 액티비티는 지속 실행 재시도 계층을 싣습니다: Temporal 기본RetryPolicy.maximum_attempts0은 모델 요청의 무제한 재실행을 뜻합니다 — SDK 클라이언트·트랜스포트 재시도와 어떻게 쌓이는지는 Retry multiplication을 보세요.event_stream_handler_activity_config: 이벤트 스트림 핸들러 액티비티에 쓸 Temporal 액티비티 설정. 기본 액티비티 설정과 병합됩니다.toolset_activity_config: ID로 식별되는 특정 툴셋의 get-tools·call-tool 액티비티에 쓸 Temporal 액티비티 설정. 기본 액티비티 설정과 병합됩니다.
ActivityConfig는 TypedDict이므로, 철자가 틀리거나 잘못 배치된 키는 런타임에 파이썬이 잡지 못하고, 그 설정이 워크플로 안에서 Temporal에 넘겨질 때만 실패하며, 그 실패는 무기한 재시도됩니다. 그래서 TemporalDurability가 생성될 때 설정 키를 검사해, Temporal이 모르는 키는 UserError를 미리 발생시킵니다.
툴별 액티비티 설정은 도구 자체에 살아있습니다 — 아래 Per-tool activity config를 보세요.
Pydantic AI가 등록하는 모든 액티비티는 실행되는 동안 백그라운드에서 하트비트를 보내서, 오래지만 건강한 액티비티가 크래시된 워커로 오인되지 않고 Temporal 워크플로 취소가 전달될 수 있게 합니다 — 진행 중인 모델 요청이 어떻게 멈추는지는 Streaming을 보세요. 하트비트 타임아웃을 기본으로 받는 것은 모델 요청 액티비티뿐이며 30초입니다. 다른 액티비티에 설정하는 것은 여러분의 몫이에요.
이벤트 루프를 차단할 수 있는 액티비티에 heartbeat_timeout을 설정하지 마세요
하트비트는 백그라운드 태스크에서 발신되므로, 양보 없이 이벤트 루프를 차지하는 액티비티(CPU 바운드 도구 함수 같은)는 박동을 멈추고 타임아웃이 경과하면 그 시도가 실패합니다. 그런 액티비티는 start_to_close_timeout만으로 서빙하는 편이 낫습니다.
툴별 액티비티 설정 (Per-tool activity config)
툴별 액티비티 설정은 도구의 metadata 필드에 살아있는데 — TemporalDurability는 'temporal' 키를 찾습니다. 도구 정의에 직접 메타데이터를 설정하거나, SetToolMetadata 캐퍼빌리티로 도구 선택에 걸쳐 적용할 수 있어요. 전체 선택기 어휘는 capabilities documentation을 보세요.
from datetime import timedelta
from temporalio.workflow import ActivityConfig
from pydantic_ai import Agent
from pydantic_ai.capabilities import SetToolMetadata
from pydantic_ai.durable_exec.temporal import TemporalDurability
from pydantic_ai.toolsets import FunctionToolset
toolset = FunctionToolset(id='research')
@toolset.tool(metadata={'temporal': ActivityConfig(start_to_close_timeout=timedelta(minutes=5))}) # (1)
async def fetch_paper(arxiv_id: str) -> str:
...
@toolset.tool(metadata={'temporal': False}) # (2)
async def now() -> str:
...
agent = Agent(
'openai:gpt-5.6-sol',
name='research',
toolsets=[toolset],
capabilities=[
SetToolMetadata( # (3)
tools=['fetch_paper', 'fetch_dataset'],
temporal=ActivityConfig(start_to_close_timeout=timedelta(minutes=5)),
),
TemporalDurability(),
],
)
- (1) 인라인: 도구 정의와 함께 액티비티 설정을 선언. 툴별 설정은 툴셋·기본 설정 위에 병합됩니다.
- (2) 액티비티 감싸기를 완전히 건너뛰려면
'temporal': False를 설정 (async도구에만 유효 — sync 도구는 스레드가 결정적이지 않으므로 항상 액티비티가 필요합니다). - (3) 선택기 기반:
SetToolMetadata가 도구 선택('all', 이름 목록, dict, 또는 callable)에 걸쳐 같은 메타데이터를 적용.
옵트아웃은 함수·동적 도구에만 적용됩니다. MCP 도구는 I/O를 수행하고 항상 그 Temporal 액티비티에서 실행되므로, MCP 도구의 metadata={'temporal': False}는 UserError를 냅니다.
서드파티 도구 구성
SetToolMetadata는 액티비티 설정이 도구 정의에 있지 않아야 할 때 — 예를 들어 서드파티 패키지에 정의된 도구나, 같은 타임아웃 프로파일을 공유하지만 다른 파일에 있는 도구 그룹 — 권장 경로입니다.
액티비티 재시도 (Activity Retries)
Temporal이 수행하는 요청 실패 자동 재시도 위에, Pydantic AI와 다양한 공급자 API 클라이언트도 자체 요청 재시도 로직이 있어요. 이들을 동시에 켜면 요청이 예상보다 자주 재시도될 수 있고, Retry-After 처리가 부적절해질 수 있어요.
Temporal을 쓸 때는 트랜스포트 재시도를 쓰지 말고 공급자 API 클라이언트의 자체 재시도 로직을 끄는 것을 권장합니다. 예를 들어 커스텀 OpenAIProvider API 클라이언트에 max_retries=0을 설정하세요.
Temporal의 재시도 정책은 액티비티 설정으로 커스터마이즈할 수 있어요.
Logfire로 관측성 (Observability with Logfire)
Temporal은 각 워크플로·액티비티 실행에 텔레메트리 이벤트와 메트릭을 생성하고, Pydantic AI는 각 에이전트 런·모델 요청·도구 호출에 이벤트를 생성합니다. 이들을 Pydantic Logfire로 보내 애플리케이션에서 일어나는 일의 완전한 그림을 얻을 수 있어요.
Logfire를 Temporal과 함께 쓰려면 Temporal의 Client.connect()에 LogfirePlugin 객체를 전달해야 해요:
from temporalio.client import Client
from pydantic_ai.durable_exec.temporal import LogfirePlugin, PydanticAIPlugin
async def main():
client = await Client.connect(
'localhost:7233',
plugins=[PydanticAIPlugin(), LogfirePlugin()],
)
기본적으로 LogfirePlugin은 Temporal(메트릭 포함)과 Pydantic AI를 계측하고 모든 데이터를 Logfire로 보냅니다. Temporal 메트릭은 60초마다 내보내집니다. LogfirePlugin 생성자에 metric_periodicity로 datetime.timedelta를 전달해 간격을 바꿀 수 있어요.
애플리케이션이 이미 logfire.configure()를 호출했다면, 플러그인은 그것을 대체하지 않고 유지하므로 여러분의 스크러빙 옵션·익스포터·샘플링·콘솔 설정이 그대로 남습니다. Logfire 설정·계측을 커스터마이즈하려면 LogfirePlugin 생성자에 setup_logfire 함수를 전달하고 커스텀 Logfire 인스턴스(즉 logfire.configure()의 결과)를 반환할 수 있어요.
Temporal 메트릭을 Logfire로 보내지 않으려면 LogfirePlugin 생성자에 metrics=False를 전달하세요. 이것은 또한 다른 Temporal 텔레메트리 옵션을 구성해야 할 때 Client.connect()에 자신의 Runtime을 공급하게 해 주며, 플러그인은 여전히 트레이싱을 구성합니다.
알려진 문제 (Known Issues)
Pandas
액티비티 안에서 logfire.info를 쓰고 프로젝트 의존성에 pandas 패키지가 있으면, 다음 오류를 만날 수 있는데, 이는 임포트 경쟁 조건의 결과로 보입니다:
AttributeError: partially initialized module 'pandas' has no attribute '_pandas_parser_CAPI' (most likely due to a circular import)
이 문제를 고치려면 temporalio.workflow.unsafe.imports_passed_through() 컨텍스트 매니저를 사용해 패키지를 선제적으로 임포트하고 워크플로 샌드박스에서 다시 로드되지 않게 하세요:
from temporalio import workflow
with workflow.unsafe.imports_passed_through():
import pandas