Temporal is a distributed, scalable, durable, and highly available orchestration engine used to execute asynchronous, long-running business logic in a scalable and resilient way.
"Temporal Python SDK" is the framework for authoring workflows and activities using the Python programming language.
Also see:
- Application Development Guide - Once you've tried our
- Python Code Samples
- API Documentation - Complete Temporal Python SDK Package reference.
Type Safe
This library uses the latest typing and MyPy support with generics to ensure all calls can be typed. For example,
starting a workflow with an int parameter when it accepts a str parameter would cause MyPy to fail.
Different Activity Types
The activity worker has been developed to work with async def, threaded, and multiprocess activities. Threaded activities are the initial recommendation, and further guidance can be found in the docs.
Custom asyncio Event Loop
The workflow implementation basically turns async def functions into workflows backed by a distributed, fault-tolerant
event loop. This means task management, sleep, cancellation, etc have all been developed to seamlessly integrate with
asyncio concepts.
See the blog post introducing the Python SDK for an informal introduction to the features and their implementation.
Contents
- Installation - Implementing a Workflow - Running a Workflow - Next Steps - Client - Data Conversion - Pydantic Support - Custom Type Data Conversion - External Storage - Driver Selection - Built-in Drivers - Custom Drivers - Workers - Workflows - Definition - Running - Invoking Activities - Invoking Child Workflows - Timers - Conditions - Asyncio and Determinism - Asyncio Cancellation - Workflow Utilities - Exceptions - Signal and update handlers - External Workflows - Testing - Automatic Time Skipping - Manual Time Skipping - Mocking Activities - Workflow Sandbox - How the Sandbox Works - Avoiding the Sandbox - Customizing the Sandbox - Passthrough Modules - Invalid Module Members - Debugging Workflows withbreakpoint() / pdb
- Known Sandbox Issues
- Global Import/Builtins
- Sandbox is not Secure
- Sandbox Performance
- Extending Restricted Classes
- Certain Standard Library Calls on Restricted Objects
- is_subclass of ABC-based Restricted Classes
- Activities
- Definition
- Types of Activities
- Synchronous Activities
- Synchronous Multithreaded Activities
- Synchronous Multiprocess/Other Activities
- Asynchronous Activities
- Activity Context
- Heartbeating and Cancellation
- Worker Shutdown
- Testing
- Interceptors
- Nexus
- Plugins
- Usage
- Plugin Implementations
- Advanced Plugin Implementations
- Client Plugins
- Worker Plugins
- Workflow Replay
- Observability
- Metrics
- OpenTelemetry Tracing
- OpenTelemetry Metrics
- Protobuf 3.x vs 4.x
- Known Compatibility Issues
- gevent Patching
- Building
- Prepare
- Build
- Use
- FIPS Compliance (Experimental)
- Local SDK development environment
- Testing
- Proto Generation and Testing
- Style
Quick Start
We will guide you through the Temporal basics to create a "hello, world!" script on your machine. It is not intended as one of the ways to use Temporal, but in reality it is very simplified and decidedly not "the only way" to use Temporal. For more information, check out the docs references in "Next Steps" below the quick start.
Installation
Install the temporalio package from PyPI.
These steps can be followed to use with a virtual environment and pip:
- Create a virtual environment
- Update
pip-python -m pip install -U pip
pip may not pick the right wheel
- Install Temporal SDK -
python -m pip install temporalio
NOTE: This README is for the current branch and not necessarily what's released on PyPI.
Implementing a Workflow
Create the following in activities.py:
from temporalio import activity
@activity.defn
def say_hello(name: str) -> str:
return f"Hello, {name}!"
Create the following in workflows.py:
from datetime import timedelta
from temporalio import workflow
Import our activity, passing it through the sandbox
with workflow.unsafe.imports_passed_through():
from activities import say_hello
@workflow.defn
class SayHello:
@workflow.run
async def run(self, name: str) -> str:
return await workflow.execute_activity(
say_hello, name, schedule_to_close_timeout=timedelta(seconds=5)
)
Create the following in run_worker.py:
import asyncio
import concurrent.futures
from temporalio.client import Client
from temporalio.worker import Worker
Import the activity and workflow from our other files
from activities import say_hello
from workflows import SayHello
async def main():
# Create client connected to server at the given address
client = await Client.connect("localhost:7233")
# Run the worker
with concurrent.futures.ThreadPoolExecutor(max_workers=100) as activity_executor:
worker = Worker(
client,
task_queue="my-task-queue",
workflows=[SayHello],
activities=[say_hello],
activity_executor=activity_executor,
)
await worker.run()
if __name__ == "__main__":
asyncio.run(main())
Assuming you have a Temporal server running on localhost, this will run the worker:
python run_worker.py
Running a Workflow
Create the following script at run_workflow.py:
import asyncio
from temporalio.client import Client
Import the workflow from the previous code
from workflows import SayHello
async def main():
# Create client connected to server at the given address
client = await Client.connect("localhost:7233")
# Execute a workflow
result = await client.execute_workflow(SayHello.run, "my name", id="my-workflow-id", task_queue="my-task-queue")
print(f"Result: {result}")
if __name__ == "__main__":
asyncio.run(main())
Assuming you have run_worker.py running from before, this will run the workflow:
python run_workflow.py
The output will be:
Result: Hello, my-name!
Next Steps
Temporal can be implemented in your code in many different ways, to suit your application's needs. The links below will give you much more information about how Temporal works with Python:
- Code Samples - If you want to start with some code, we have provided
- Application Development Guide Our Python specific
- API Documentation - Full Temporal Python SDK package documentation.
Usage
From here, you will find reference documentation about specific pieces of the Temporal Python SDK that were built around Temporal concepts. This section is not intended as a how-to guide -- For more how-to oriented information, check out the links in the Next Steps section above.
Client
A client can be created and used to start a workflow like so:
from temporalio.client import Client
async def main():
# Create client connected to server at the given address and namespace
client = await Client.connect("localhost:7233", namespace="my-namespace")
# Start a workflow
handle = await client.start_workflow(MyWorkflow.run, "some arg", id="my-workflow-id", task_queue="my-task-queue")
# Wait for result
result = await handle.result()
print(f"Result: {result}")
Some things to note about the above code:
- A
Clientdoes not have an explicit "close" - To enable TLS, the
tlsargument toconnectcan be set toTrueor aTLSConfigobject - A single positional argument can be passed to
start_workflow. If there are multiple arguments, only the
start_workflow can be used (i.e. the one accepting a string workflow name) and it must be in
the args keyword argument.
- The
handlerepresents the workflow that was started and can be used for more than just getting the result - Since we are just getting the handle and waiting on the result, we could have called
client.execute_workflowwhich
- Clients can have many more options not shown here (e.g. data converters and interceptors)
- A string can be used instead of the method reference to call a workflow by name (e.g. if defined in another language)
- Clients do not work across forks
client above, this is how to have a client in another namespace:
config = client.config()
config["namespace"] = "my-other-namespace"
other_ns_client = Client(**config)
Data Conversion
Data converters are used to convert raw Temporal payloads to/from actual Python types. A custom data converter of type
temporalio.converter.DataConverter can be set via the data_converter parameter of the Client constructor. Data
converters are a combination of payload converters, external storage, payload codecs, and failure converters. Payload
converters convert Python values to/from serialized bytes. External payload storage optionally stores and retrieves payloads
to/from external storage services using drivers. Payload codecs convert bytes to bytes (e.g. for compression or encryption).
Failure converters convert exceptions to/from serialized failures.
The default data converter supports converting multiple types including:
Nonebytesgoogle.protobuf.message.Message- As JSON when encoding, but has ability to decode binary proto from other languages- Anything that can be converted to JSON including:
json.dump supports natively
* dataclasses
* Iterables including ones JSON dump may not support by default, e.g. set
* IntEnum, StrEnum based enumerates, including enums that mix in int or str
* UUID
* datetime.datetime
To use pydantic model instances, see Pydantic Support.
datetime.date and datetime.time can only be used with the Pydantic data converter.
Although workflows, updates, signals, and queries can all be defined with multiple input parameters, users are strongly
encouraged to use a single dataclass or Pydantic model parameter, so that fields with defaults can be easily added
without breaking compatibility. Similar advice applies to return values.
Classes with generics may not have the generics properly resolved. The current implementation does not have generic type resolution. Users should use concrete types.
##### Pydantic Support
To use Pydantic model instances, install Pydantic and set the Pydantic data converter when creating client instances:
from temporalio.contrib.pydantic import pydantic_data_converter
client = Client(data_converter=pydantic_data_converter, ...)
This data converter supports conversion of all types supported by Pydantic to and from JSON.
In addition to Pydantic models, these include all json.dump-able types, various non-json.dump-able standard library
types such as dataclasses, types from the datetime module, sets, UUID, etc, and custom types composed of any of these.
Pydantic v1 is not supported by this data converter. If you are not yet able to upgrade from Pydantic v1, see https://github.com/temporalio/samples-python/tree/main/pydantic_converter/v1 for limited v1 support.
##### Custom Type Data Conversion
For converting from JSON, the workflow/activity type hint is taken into account to convert to the proper type. Care has
been taken to support all common typings including Optional, Union, all forms of iterables and mappings, NewType,
etc in addition to the regular JSON values mentioned before.
Data converters contain a reference to a payload converter class that is used to convert to/from payloads/values. This
is a class and not an instance because it is instantiated on every workflow run inside the sandbox. The payload
converter is usually a CompositePayloadConverter which contains a multiple EncodingPayloadConverters it uses to try
to serialize/deserialize payloads. Upon serialization, each EncodingPayloadConverter is tried until one succeeds. The
EncodingPayloadConverter provides an "encoding" string serialized onto the payload so that, upon deserialization, the
specific EncodingPayloadConverter for the given "encoding" is used.
The default data converter uses the DefaultPayloadConverter which is simply a CompositePayloadConverter with a known
set of default EncodingPayloadConverters. To implement a custom encoding for a custom type, a new
EncodingPayloadConverter can be created for the new type. For example, to support IPv4Address types:
class IPv4AddressEncodingPayloadConverter(EncodingPayloadConverter):
@property
def encoding(self) -> str:
return "text/ipv4-address"
def to_payload(self, value: Any) -> Optional[Payload]:
if isinstance(value, ipaddress.IPv4Address):
return Payload(
metadata={"encoding": self.encoding.encode()},
data=str(value).encode(),
)
else:
return None
def from_payload(self, payload: Payload, type_hint: Optional[Type] = None) -> Any:
assert not type_hint or type_hint is ipaddress.IPv4Address
return ipaddress.IPv4Address(payload.data.decode())
class IPv4AddressPayloadConverter(CompositePayloadConverter):
def __init__(self) -> None:
# Just add ours as first before the defaults
super().__init__(
IPv4AddressEncodingPayloadConverter(),
*DefaultPayloadConverter.default_encoding_payload_converters,
)
my_data_converter = dataclasses.replace(
DataConverter.default,
payload_converter_class=IPv4AddressPayloadConverter,
)
Imports are left off for brevity.
This is good for many custom types. However, sometimes you want to override the behavior of the just the existing JSON encoding payload converter to support a new type. It is already the last encoding data converter in the list, so it's the fall-through behavior for any otherwise unknown type. Customizing the existing JSON converter has the benefit of making the type work in lists, unions, etc.
The JSONPlainPayloadConverter uses the Python json library with an
advanced JSON encoder by default and a custom value conversion method to turn json.loaded values to their type hints.
The conversion can be customized for serialization with a custom json.JSONEncoder and deserialization with a custom
JSONTypeConverter. For example, to support IPv4Address types in existing JSON conversion:
class IPv4AddressJSONEncoder(AdvancedJSONEncoder):
def default(self, o: Any) -> Any:
if isinstance(o, ipaddress.IPv4Address):
return str(o)
return super().default(o)
class IPv4AddressJSONTypeConverter(JSONTypeConverter):
def to_typed_value(
self, hint: Type, value: Any
) -> Union[Optional[Any], JSONTypeConverterUnhandled]:
if issubclass(hint, ipaddress.IPv4Address):
return ipaddress.IPv4Address(value)
return JSONTypeConverter.Unhandled
class IPv4AddressPayloadConverter(CompositePayloadConverter):
def __init__(self) -> None:
# Replace default JSON plain with our own that has our encoder and type
# converter
json_converter = JSONPlainPayloadConverter(
encoder=IPv4AddressJSONEncoder,
custom_type_converters=[IPv4AddressJSONTypeConverter()],
)
super().__init__(
*[
c if not isinstance(c, JSONPlainPayloadConverter) else json_converter
for c in DefaultPayloadConverter.default_encoding_payload_converters
]
)
my_data_converter = dataclasses.replace(
DataConverter.default,
payload_converter_class=IPv4AddressPayloadConverter,
)
Now IPv4Address can be used in type hints including collections, optionals, etc.
##### External Storage
⚠️ External storage support is currently at an experimental release stage. ⚠️
External storage allows large payloads to be offloaded to an external storage service (such as Amazon S3) rather than stored inline in workflow history. This is useful when workflows or activities work with data that would otherwise exceed Temporal's payload size limits.
External storage is configured via the external_storage parameter on DataConverter. It should be configured on the Client both for clients of your workflow as well as on the worker -- anywhere large payloads may be uploaded or downloaded.
A StorageDriver handles uploading and downloading payloads. Temporal provides built-in drivers for common storage solutions, or you may implement a custom driver. Here's an example using the built-in S3StorageDriver with the SDK's aioboto3 client:
import aioboto3
import dataclasses
from temporalio.client import Client, ClientConfig
from temporalio.contrib.aws.s3driver import S3StorageDriver
from temporalio.contrib.aws.s3driver.aioboto3 import new_aioboto3_client
from temporalio.converter import DataConverter
from temporalio.converter import ExternalStorage
client_config = ClientConfig.load_client_connect_config()
session = aioboto3.Session()
async with session.client("s3") as s3_client:
driver = S3StorageDriver(
client=new_aioboto3_client(s3_client),
bucket="my-bucket",
)
client = await Client.connect(
**client_config,
data_converter=dataclasses.replace(
DataConverter.default,
external_storage=ExternalStorage(drivers=[driver]),
),
)
See the S3 driver README for further details.
Some things to note about external storage:
- Only payloads that meet or exceed
ExternalStorage.payload_size_threshold(default 256 KiB) are offloaded. Smaller payloads are stored inline as normal. - External storage applies transparently to all payloads, whether they are workflow inputs/outputs, activity inputs/outputs, signal inputs, query outputs, update inputs/outputs, or failure details.
DataConverter's payload_codec (if configured) is applied to the payload before* it is handed to the storage driver, so the driver always stores encoded bytes. The reference payload written to workflow history is not encoded by the DataConverter codec.
- Setting
ExternalStorage.payload_size_thresholdto0causes every payload to be considered for external storage regardless of size.
When multiple storage backends are needed, list all drivers in ExternalStorage.drivers and provide a driver_selector to control which driver stores new payloads. Any driver in the list not chosen for storing is still available for retrieval, which is useful when migrating between storage backends.
from temporalio.converter import ExternalStorage
options = ExternalStorage(
drivers=[hot_driver, cold_driver],
driver_selector=lambda context, payload: (
hot_driver if payload.ByteSize() < 5 1024 1024 else cold_driver
),
)
For more complex selection logic, use a plain callable that reads from the StorageDriverStoreContext:
import temporalio.converter
from temporalio.api.common.v1 import Payload
def feature_flag_is_on(workflow_id: str | None) -> bool:
"""Check whether external storage is enabled for this workflow via a feature flag service."""
return workflow_id is not None and len(workflow_id) % 2 == 0
def feature_flag_selector(
context: temporalio.converter.StorageDriverStoreContext, _payload: Payload
) -> temporalio.converter.StorageDriver | None:
workflow_id = (
context.target.id
if isinstance(context.target, temporalio.converter.StorageDriverWorkflowInfo)
else None
)
return my_driver if feature_flag_is_on(workflow_id) else None
options = ExternalStorage(
drivers=[my_driver],
driver_selector=feature_flag_selector,
)
Some things to note about driver selection:
- A
driver_selectoris required when more than one driver is registered. With a single driver,driver_selectormay be omitted and that driver is used for all store operations. - Returning
Nonefrom a selector leaves the payload stored inline in workflow history rather than offloading it. - The driver instance returned by the selector must be one of the instances registered in
ExternalStorage.drivers. If it is not, an error is raised.
- S3 Storage Driver: ⚠️ Experimental ⚠️ Amazon S3 driver. Ships with an aioboto3 client, or bring your own by subclassing
S3StorageDriverClient.
Implement temporalio.converter.StorageDriver to integrate with an external storage system:
from collections.abc import Sequence
from temporalio.converter import StorageDriver, StorageDriverClaim, StorageDriverRetrieveContext, StorageDriverStoreContext
from temporalio.api.common.v1 import Payload
class MyDriver(StorageDriver):
def __init__(self, driver_name: str | None = None):
self._driver_name = driver_name or "my-org:driver:my-driver"
def name(self) -> str:
return self._driver_name
async def store(
self, context: StorageDriverStoreContext, payloads: Sequence[Payload]
) -> list[StorageDriverClaim]:
claims = []
for payload in payloads:
key = await my_storage.put(payload.SerializeToString())
claims.append(StorageDriverClaim(claim_data={"key": key}))
return claims
async def retrieve(
self, context: StorageDriverRetrieveContext, claims: Sequence[StorageDriverClaim]
) -> list[Payload]:
payloads = []
for claim in claims:
data = await my_storage.get(claim.claim_data["key"])
p = Payload()
p.ParseFromString(data)
payloads.append(p)
return payloads
Some things to note about implementing a custom driver:
StorageDriver.name()must return a string that is unique among all drivers inExternalStorage.drivers. This name is embedded in the reference payload stored in workflow history and used to look up the correct driver during retrieval — changing it after payloads have been stored will break retrieval.StorageDriver.type()is automatically implemented to return the name of the class. This can be overridden in subclasses but must remain consistent across all instances of the subclass.- Use
StorageDriverStoreContext.targetinsidestore()when you need workflow or activity identity (namespace, workflow ID, activity ID, etc.) to choose where or how to store payloads.
Workers
Workers host workflows and/or activities. Here's how to run a worker:
import asyncio
import logging
from temporalio.client import Client
from temporalio.worker import Worker
Import your own workflows and activities
from my_workflow_package import MyWorkflow, my_activity
async def run_worker(stop_event: asyncio.Event):
# Create client connected to server at the given address
client = await Client.connect("localhost:7233", namespace="my-namespace")
# Run the worker until the event is set
worker = Worker(client, task_queue="my-task-queue", workflows=[MyWorkflow], activities=[my_activity])
async with worker:
await stop_event.wait()
Some things to note about the above code:
- This creates/uses the same client that is used for starting workflows
- While this example accepts a stop event and uses
async with,run()andshutdown()may be used instead - Workers can have many more options not shown here (e.g. data converters and interceptors)
Workflows
Definition
Workflows are defined as classes decorated with @workflow.defn. The method invoked for the workflow is decorated with
@workflow.run. Methods for signals, queries, and updates are decorated with @workflow.signal, @workflow.query
and @workflow.update respectively. Here's an example of a workflow:
import asyncio
from datetime import timedelta
from temporalio import workflow
Pass the activities through the sandbox
with workflow.unsafe.imports_passed_through():
from .my_activities import GreetingInfo, create_greeting_activity
@workflow.defn
class GreetingWorkflow:
def __init__(self) -> None:
self._current_greeting = "<unset>"
self._greeting_info = GreetingInfo()
self._greeting_info_update = asyncio.Event()
self._complete = asyncio.Event()
@workflow.run
async def run(self, name: str) -> str:
self._greeting_info.name = name
while True:
# Store greeting
self._current_greeting = await workflow.execute_activity(
create_greeting_activity,
self._greeting_info,
start_to_close_timeout=timedelta(seconds=5),
)
workflow.logger.debug("Greeting set to %s", self._current_greeting)
# Wait for salutation update or complete signal (this can be
# cancelled)
await asyncio.wait(
[
asyncio.create_task(self._greeting_info_update.wait()),
asyncio.create_task(self._complete.wait()),
],
return_when=asyncio.FIRST_COMPLETED,
)
if self._complete.is_set():
return self._current_greeting
self._greeting_info_update.clear()
@workflow.signal
async def update_salutation(self, salutation: str) -> None:
self._greeting_info.salutation = salutation
self._greeting_info_update.set()
@workflow.signal
async def complete_with_greeting(self) -> None:
self._complete.set()
@workflow.query
def current_greeting(self) -> str:
return self._current_greeting
@workflow.update
def set_and_get_greeting(self, greeting: str) -> str:
old = self._current_greeting
self._current_greeting = greeting
return old
This assumes there's an activity in my_activities.py like:
from dataclasses import dataclass
from temporalio import workflow
@dataclass
class GreetingInfo:
salutation: str = "Hello"
name: str = "<unknown>"
@activity.defn
def create_greeting_activity(info: GreetingInfo) -> str:
return f"{info.salutation}, {info.name}!"
Some things to note about the above workflow code:
- Workflows run in a sandbox by default.
temporalio imports should usually be "passed through" the sandbox. See the
Workflow Sandbox section for more details.
- This workflow continually updates the queryable current greeting when signalled and can complete with the greeting on
- Workflows are always classes and must have a single
@workflow.runwhich is anasync deffunction - Workflow code must be deterministic. This means no
setiteration, threading, no randomness, no external calls to
asyncio event loop and be
deterministic. Also see the Asyncio and Determinism section later.
@activity.defnis explained in a later section. For normal simple string concatenation, this would just be done in
workflow.execute_activity(create_greeting_activity, ...is actually a typed signature, and MyPy will fail if the
self._greeting_info parameter is not a GreetingInfo
Here are the decorators that can be applied:
@workflow.defn- Defines a workflow class
name param to customize the workflow name, otherwise it defaults to the unqualified class name
* Can have dynamic=True which means all otherwise unhandled workflows fall through to this. If present, cannot have
name argument, and run method must accept a single parameter of Sequence[temporalio.common.RawValue] type. The
payload of the raw value can be converted via workflow.payload_converter().from_payload.
@workflow.run- Defines the primary workflow run method
@workflow.defn, not a base class (but can _also_ be defined on the same
method of a base class)
* Exactly one method name must have this decorator, no more or less
* Must be defined on an async def method
* The method's arguments are the workflow's arguments
* The first parameter must be self, followed by positional arguments. Best practice is to only take a single
argument that is an object/dataclass of fields that can be added to as needed.
@workflow.init- Specifies that the__init__method accepts the workflow's arguments.
__init__ method, the parameters of which must then be identical to those of
the @workflow.run method.
* The purpose of this decorator is to allow operations involving workflow arguments to be performed in the __init__
method, before any signal or update handler has a chance to execute.
@workflow.signal- Defines a method as a signal
async or non-async method at any point in the class hierarchy, but if the decorated method
is overridden, then the override must also be decorated.
* The method's arguments are the signal's arguments.
* Return value is ignored.
* May mutate workflow state, and make calls to other workflow APIs like starting activities, etc.
* Can have a name param to customize the signal name, otherwise it defaults to the unqualified method name.
* Can have dynamic=True which means all otherwise unhandled signals fall through to this. If present, cannot have
name argument, and method parameters must be self, a string signal name, and a
Sequence[temporalio.common.RawValue].
* Non-dynamic method can only have positional arguments. Best practice is to only take a single argument that is an
object/dataclass of fields that can be added to as needed.
* See Signal and update handlers below
@workflow.update- Defines a method as an update
async or non-async method at any point in the class hierarchy, but if the decorated method
is overridden, then the override must also be decorated.
* May accept input and return a value
* The method's arguments are the update's arguments.
* May be async or non-async
* May mutate workflow state, and make calls to other workflow APIs like starting activities, etc.
* Also accepts the name and dynamic parameters like signal, with the same semantics.
* Update handlers may optionally define a validator method by decorating it with @update_handler_method.validator.
To reject an update before any events are written to history, throw an exception in a validator. Validators cannot
be async, cannot mutate workflow state, and return nothing.
* See Signal and update handlers below
@workflow.query- Defines a method as a query
async
* Temporal queries should never mutate anything in the workflow or call any calls that would mutate the workflow
* Also accepts the name and dynamic parameters like signal and update, with the same semantics.
Running
To start a locally-defined workflow from a client, you can simply reference its method like so:
from temporalio.client import Client
from my_workflow_package import GreetingWorkflow
async def create_greeting(client: Client) -> str:
# Start the workflow
handle = await client.start_workflow(GreetingWorkflow.run, "my name", id="my-workflow-id", task_queue="my-task-queue")
# Change the salutation
await handle.signal(GreetingWorkflow.update_salutation, "Aloha")
# Tell it to complete
await handle.signal(GreetingWorkflow.complete_with_greeting)
# Wait and return result
return await handle.result()
Some things to note about the above code:
- This uses the
GreetingWorkflowfrom the previous section - The result of calling this function is
"Aloha, my name!" idandtask_queueare required for running a workflowclient.start_workflowis typed, so MyPy would fail if"my name"were something besides a stringhandle.signalis typed, so MyPy would fail if"Aloha"were something besides a string or if we provided a
complete_with_greeting
handle.resultis typed to the workflow itself, so MyPy would fail if we said thiscreate_greetingreturned
Invoking Activities
- Activities are started with non-async
workflow.start_activity()which accepts either an activity function reference
- A single argument to the activity is positional. Multiple arguments are not supported in the type-safe form of
args keyword argument.
- Activity options are set as keyword arguments after the activity arguments. At least one of
start_to_close_timeout
schedule_to_close_timeout must be provided.
- The result is an activity handle which is an
asyncio.Taskand supports basic task features - An async
workflow.execute_activity()helper is provided which takes the same arguments as
workflow.start_activity() and awaits on the result. This should be used in most cases unless advanced task
capabilities are needed.
- Local activities work very similarly except the functions are
workflow.start_local_activity()and
workflow.execute_local_activity()
- Activities can be methods of a class. Invokers should use
workflow.start_activity_method(),
workflow.execute_activity_method(), workflow.start_local_activity_method(), and
workflow.execute_local_activity_method() instead.
- Activities can callable classes (i.e. that define
__call__). Invokers should useworkflow.start_activity_class(),
workflow.execute_activity_class(), workflow.start_local_activity_class(), and
workflow.execute_local_activity_class() instead.
Invoking Child Workflows
- Child workflows are started with async
workflow.start_child_workflow()which accepts either a workflow run method
- A single argument to the child workflow is positional. Multiple arguments are not supported in the type-safe form of
args keyword argument.
- Child workflow options are set as keyword arguments after the arguments. At least
idmust be provided. - The
awaitof the start does not complete until the start has been accepted by the server - The result is a child workflow handle which is an
asyncio.Taskand supports basic task features. The handle also has
- An async
workflow.execute_child_workflow()helper is provided which takes the same arguments as
workflow.start_child_workflow() and awaits on the result. This should be used in most cases unless advanced task
capabilities are needed.
Timers
- A timer is represented by normal
asyncio.sleep()or aworkflow.sleep()call - Timers are also implicitly started on any
asynciocalls with timeouts (e.g.asyncio.wait_for) - Timers are Temporal server timers, not local ones, so sub-second resolution rarely has value
- Calls that use a specific point in time, e.g.
call_atortimeout_at, should be based on the current loop time
workflow.time()) and not an actual point in time. This is because fixed times are translated to relative ones
by subtracting the current loop time which may not be the actual current time.
Conditions
workflow.wait_conditionis an async function that doesn't return until a provided callback returns true- A
timeoutcan optionally be provided which will throw aasyncio.TimeoutErrorif reached (internally backed by
asyncio.wait_for which uses a timer)
Asyncio and Determinism
Workflows must be deterministic. Workflows are backed by a custom
asyncio event loop. This means many of the common asyncio calls work
as normal. Some asyncio features are disabled such as:
- Thread related calls such as
to_thread(),run_coroutine_threadsafe(),loop.run_in_executor(), etc - Calls that alter the event loop such as
loop.close(),loop.stop(),loop.run_forever(),
loop.set_task_factory(), etc
- Calls that use anything external such as networking, subprocesses, disk IO, etc
asyncio utilities that internally use set() which can make them non-deterministic from one
worker to the next. Therefore the following asyncio functions have workflow-module alternatives that are
deterministic:
asyncio.as_completed()- useworkflow.as_completed()asyncio.wait()- useworkflow.wait()
Asyncio Cancellation
Cancellation is done using asyncio task cancellation.
This means that tasks are requested to be cancelled but can catch the
asyncio.CancelledError, thus
allowing them to perform some cleanup before allowing the cancellation to proceed (i.e. re-raising the error), or to
deny the cancellation entirely. It also means that
asyncio.shield() can be used to
protect tasks against cancellation.
The following tasks, when cancelled, perform a Temporal cancellation:
- Activities - when the task executing an activity is cancelled, a cancellation request may be sent to the activity
- Child workflows - when the task starting or executing a child workflow is cancelled, a cancellation request may be
- Nexus operations - when the task starting or executing a Nexus operation is cancelled, a cancellation request may be
- Timers - when the task executing a timer is cancelled (whether started via sleep or timeout), the timer is cancelled
Task.cancel is called on the main workflow task. Therefore,
asyncio.CancelledError can be caught in order to handle the cancel gracefully.
Workflows follow asyncio cancellation rules exactly which can cause confusion among Python developers. Cancelling a
task doesn't always cancel the thing it created. For example, given
task = asyncio.create_task(workflow.start_child_workflow(..., calling task.cancel does not cancel the child
workflow, it only cancels the starting of it, which has no effect if it has already started. However, cancelling the
result of handle = await workflow.start_child_workflow(... or
task = asyncio.create_task(workflow.execute_child_workflow(... _does_ cancel the child workflow.
Also, due to Temporal rules, a cancellation request is a state not an event. Therefore, repeated cancellation requests are not delivered, only the first. If the workflow chooses swallow a cancellation, it cannot be requested again.
Workflow Utilities
While running in a workflow, in addition to features documented elsewhere, the following items are available from the
temporalio.workflow package:
continue_as_new()- Async function to stop the workflow immediately and continue as newinfo()- Returns information about the current workflowlogger- A logger for use in a workflow (properly skips logging on replay)now()- Returns the "current time" from the workflow's perspective
Exceptions
- Workflows/updates can raise exceptions to fail the workflow or the "workflow task" (i.e. suspend the workflow
- Exceptions that are instances of
temporalio.exceptions.FailureErrorwill fail the workflow with that exception
temporalio.exceptions.ApplicationError. This can
be marked non-retryable or include details as needed.
* Other exceptions that come from activity execution, child execution, cancellation, etc are already instances of
FailureError and will fail the workflow when uncaught.
- Update handlers are special: an instance of
temporalio.exceptions.FailureErrorraised in an update handler will fail
- All other exceptions fail the "workflow task" which means the workflow will continually retry until the workflow is
ApplicationError as mentioned above.
This default can be changed by providing a list of exception types to workflow_failure_exception_types when creating a
Worker or failure_exception_types on the @workflow.defn decorator. If a workflow-thrown exception is an instance
of any type in either list, it will fail the workflow (or update) instead of the workflow task. This means a value of
[Exception] will cause every exception to fail the workflow instead of the workflow task. Also, as a special case, if
temporalio.workflow.NondeterminismError (or any superclass of it) is set, non-deterministic exceptions will fail the
workflow. WARNING: These settings are experimental.
Signal and update handlers
Signal and update handlers are defined using decorated methods as shown in the example above. Client code
sends signals and updates using workflow_handle.signal, workflow_handle.execute_update, or
workflow_handle.start_update. When the workflow receives one of these requests, it starts an asyncio.Task executing
the corresponding handler method with the argument(s) from the request.
The handler methods may be async def and can do all the async operations described above (e.g. invoking activities and
child workflows, and waiting on timers and conditions). Notice that this means that handler tasks will be executing
concurrently with respect to each other and the main workflow task. Use
asyncio.Lock and
asyncio.Semaphore if necessary.
Your main workflow task may finish as a result of successful completion, cancellation, continue-as-new, or failure. You
should ensure that all in-progress signal and update handler tasks have finished before this happens; if you do not, you
will see a warning (the warning can be disabled via the workflow.signal/workflow.update decorators). One way to
ensure that handler tasks have finished is to wait on the workflow.all_handlers_finished condition:
await workflow.wait_condition(workflow.all_handlers_finished)
External Workflows
workflow.get_external_workflow_handle()inside a workflow returns a handle to interact with another workflowworkflow.get_external_workflow_handle_for()can be used instead for a type safe handleawait handle.signal()can be called on the handle to signal the external workflowawait handle.cancel()can be called on the handle to send a cancel to the external workflow
Testing
Workflow testing can be done in an integration-test fashion against a real server, however it is hard to simulate timeouts and other long time-based code. Using the time-skipping workflow test environment can help there.
The time-skipping temporalio.testing.WorkflowEnvironment can be created via the static async start_time_skipping().
This internally downloads the Temporal time-skipping test server to a temporary directory if it doesn't already exist,
then starts the test server which has special APIs for skipping time.
NOTE: The time-skipping test environment does not work on ARM. The SDK will try to download the x64 binary on macOS for use with the Intel emulator, but for Linux or Windows ARM there is no proper time-skipping test server at this time.
##### Automatic Time Skipping
Anytime a workflow result is waited on, the time-skipping server automatically advances to the next event it can. To
manually advance time before waiting on the result of a workflow, the WorkflowEnvironment.sleep method can be used.
Here's a simple example of a workflow that sleeps for 24 hours:
import asyncio
from temporalio import workflow
@workflow.defn
class WaitADayWorkflow:
@workflow.run
async def run(self) -> str:
await asyncio.sleep(24 60 60)
return "all done"
An integration test of this workflow would be way too slow. However the time-skipping server automatically skips to the next event when we wait on the result. Here's a test for that wo
... (README truncated for length)