Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,12 @@ adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## Unreleased

FIXED

- Asynchronous Azure Blob payload uploads and downloads no longer run gzip
compression or decompression on the calling event loop, keeping concurrent
async operations responsive during large transfers.

## v1.10.1

FIXED
Expand Down
12 changes: 12 additions & 0 deletions azure-functions-durable/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## Unreleased

ADDED

- Added `DFApp.configure_large_payloads(payload_store=...)` to externalize large
durable payloads to Azure Blob Storage or a custom payload store.
- Configured clients automatically hydrate stored payloads, including
orchestration history and entity operation inputs and results.

FIXED

- Preserved the application's source directory in `context.function_directory`
for decorated activities and durable-client functions.

## v2.0.0rc1

CHANGED
Expand Down
123 changes: 123 additions & 0 deletions azure-functions-durable/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,129 @@ Key capabilities include durable orchestrations and sub-orchestrations, durable
timers, external events, durable entities, retries, versioning, durable HTTP
calls (`context.call_http(...)`), recurring scheduled tasks, and history export.

## Large payloads

Configure a `durabletask.payload.PayloadStore` once at app startup to store large
serialized payloads outside orchestration history. For Azure Blob Storage, install
the optional dependencies:

```bash
pip install azure-functions-durable "durabletask[azure-blob-payloads]" aiohttp
```

In your Function app, configure the root `DFApp` before any invocations:

```python
import os

import azure.durable_functions as df
from durabletask.extensions.azure_blob_payloads import (
BlobPayloadStore,
BlobPayloadStoreOptions,
)

app = df.DFApp()
app.configure_large_payloads(
payload_store=BlobPayloadStore(BlobPayloadStoreOptions(
connection_string=os.environ["PAYLOAD_STORAGE_CONNECTION_STRING"],
container_name="durable-payloads",
threshold_bytes=256 * 1024,
))
)
```

Set `PAYLOAD_STORAGE_CONNECTION_STRING` in your Function app settings (or in
`local.settings.json` for local development). The store automatically uploads
serialized payloads above the threshold and downloads their contents when the
SDK consumes them. Orchestration and activity inputs and outputs, custom status,
external events, and entity inputs, results, and state use the configured store.
Sub-orchestrations and continue-as-new use it as well. The default maximum stored
payload size is 10 MiB; `max_stored_payload_bytes` can configure this limit.

Configuration applies to both synchronous and asynchronous durable clients and
all registered blueprints, including blueprints imported before configuration.
There is one store per Python worker process. Registering the same store object
again is allowed; registering a different object raises `ValueError`. Configure
every scaled-out worker with access to the same backing storage and retain that
access across deployments. Keep the store open for the process lifetime.

> [!WARNING]
> Keep payload blobs for as long as any retained orchestration history or entity
> state references them, including histories needed for replay. Purging an
> orchestration does not delete its payload blobs; manage retention separately.

This is SDK-managed storage, separate from the Azure Storage backend's automatic
large-message handling. Without configuration, the SDK keeps payloads inline.
Use the configured Python clients to retrieve hydrated payloads. Host management
HTTP endpoints and other consumers that do not use this configuration can expose
reference strings instead. Applications exchanging externalized payloads must
agree on the store and reference encoding; Functions references are JSON strings.

> [!WARNING]
> Storage failures can fail durable invocations, including orchestrations.
> Storage transport retries are separate from durable activity retry policies.
> This SDK does not add an activity retry policy or guarantee that the Functions
> host abandons and redelivers a work item after a storage failure. A transient
> storage error can therefore become a terminal orchestration failure.

Registered orchestration and entity handlers await the store's async methods
before and after execution. Orchestrators remain synchronous generators, and
orchestration replay and entity code run on execution threads with their
invocation logging context preserved. Their payload downloads and uploads do
not occupy those threads, and serialization does not access storage during
replay. Custom stores must implement genuinely nonblocking async methods to
benefit from this behavior.

For orchestration and entity execution, the SDK reuses the Functions runtime's
thread pool when the runtime exposes it; otherwise it uses a process-wide SDK
pool. Both honor `PYTHON_THREADPOOL_THREAD_COUNT`.

Activities retain their synchronous or asynchronous calling convention.
Synchronous activities use synchronous storage inside the host-managed execution
thread and remain directly callable without `await`; async activities await
async storage. Direct calls to decorated activities return ordinary Python values
without accessing payload storage; transport processing applies only to host
binding invocations. Binding converters perform no storage I/O. Synchronous functions
still receive the synchronous durable client, and synchronous client APIs use
synchronous storage. Direct `Orchestrator.handle()` and `Orchestrator.create()`
adapters also remain synchronous.

Both client history APIs hydrate entity operation inputs and results, including
values nested in the host's entity protocol envelopes. During orchestration
replay, nested entity results are hydrated, but historical nested request inputs
are not downloaded because replay only needs their correlation metadata.
Historical scheduled activity inputs are also not downloaded during replay;
explicit history retrieval continues to hydrate those inputs.

> [!WARNING]
> With payload storage configured, whole payload strings recognized by the
> store's `is_known_token()` are reserved references, not literal application
> data. For `BlobPayloadStore`, this includes strings of the form
> `blob:v1:<container>:<blobName>`, with nonempty container and blob names.
> Functions recognizes both raw and JSON-quoted references. A matching string
> is treated as already externalized on output and downloaded on input, even
> below the size threshold. Missing or inaccessible references raise errors;
> they do not fall back to literal strings.
Comment thread
andystaples marked this conversation as resolved.

To pass a reference as application data for later retrieval, wrap it in an
object, for example `{"reference": "blob:v1:container:blob"}`. Reference detection
does not recursively inspect strings inside application JSON objects. The
wrapper preserves the literal reference whether the object stays inline or is
itself externalized. Keep the wrapper whenever passing that value across a
durable payload boundary; passing its string field alone opts back into reference
interpretation. Custom payload stores define their own reserved token syntax.

> [!WARNING]
> Recognized references are trusted transport inputs, not authorization
> boundaries. `BlobPayloadStore` reads from the container named in the token
> using its configured credentials; `container_name` selects the upload
> container and does not restrict downloads. An account-wide connection string
> can therefore allow reads outside that container. Use least-privilege
> credentials scoped to the intended payload storage, and explicitly decide
> whether external callers may supply references. Reject untrusted references
> or validate their allowed storage locations before passing them into durable
> APIs; token recognition alone does not authorize a read.

## Unit testing entities

Use `execute_entity()` to run one entity operation in-process without a
Expand Down
30 changes: 29 additions & 1 deletion azure-functions-durable/azure/durable_functions/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,11 +8,12 @@
import threading

from datetime import datetime, timedelta
from typing import Any, Mapping, Optional, Union, cast
from typing import Any, Mapping, Optional, Union, cast, override
from warnings import deprecated
import azure.functions as func
from urllib.parse import urlparse, quote

from durabletask import history
from durabletask.client import (
AsyncTaskHubGrpcClient,
OrchestrationQuery,
Expand All @@ -26,6 +27,11 @@
AzureFunctionsDefaultClientInterceptorImpl,
)
from .internal.serialization import DEFAULT_FUNCTIONS_DATA_CONVERTER
from .internal.payloads import (
get_transport_payload_store,
hydrate_entity_history,
hydrate_entity_history_async,
)
from .http.http_management_payload import HttpManagementPayload, replace_url_origin
from .internal.compat.durable_orchestration_status import DurableOrchestrationStatus
from .internal.compat.entity_state_response import EntityStateResponse
Expand Down Expand Up @@ -178,6 +184,7 @@ def __init__(self, client_as_string: str):
interceptors=interceptors,
channel_options=channel_options,
data_converter=DEFAULT_FUNCTIONS_DATA_CONVERTER,
payload_store=get_transport_payload_store(),
Comment thread
andystaples marked this conversation as resolved.
emit_trace_spans=False,
logger=_LOGGER)

Expand All @@ -193,6 +200,16 @@ def __init__(self, client_as_string: str):
self._creation_loop = None
self._close_scheduled = False

@override
async def get_orchestration_history(
self, instance_id: str, *, execution_id: str | None = None,
for_work_item_processing: bool = False) -> list[history.HistoryEvent]:
events = await super().get_orchestration_history(
instance_id, execution_id=execution_id,
for_work_item_processing=for_work_item_processing)
await hydrate_entity_history_async(events, self._payload_store, instance_id)
return events

def schedule_close(self) -> None:
"""Schedule the underlying gRPC channel to close after the invocation.

Expand Down Expand Up @@ -664,9 +681,20 @@ def __init__(self, client_as_string: str):
interceptors=interceptors,
channel_options=channel_options,
data_converter=DEFAULT_FUNCTIONS_DATA_CONVERTER,
payload_store=get_transport_payload_store(),
emit_trace_spans=False,
logger=_LOGGER)

@override
def get_orchestration_history(
self, instance_id: str, *, execution_id: str | None = None,
for_work_item_processing: bool = False) -> list[history.HistoryEvent]:
events = super().get_orchestration_history(
instance_id, execution_id=execution_id,
for_work_item_processing=for_work_item_processing)
hydrate_entity_history(events, self._payload_store, instance_id)
return events

@classmethod
def get_cached(cls, client_as_string: str) -> "SyncDurableFunctionsClient":
"""Get the process-wide client for a durable-client binding configuration.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from azure.functions.decorators.function_app import DecoratorApi, FunctionBuilder

from durabletask import task
from durabletask.payload import PayloadStore

from .metadata import OrchestrationTrigger, ActivityTrigger, EntityTrigger, \
DurableClient
Expand All @@ -19,9 +20,9 @@
builtin_http_activity,
builtin_http_poll_orchestrator,
)
from ..internal.compat.activity import wrap_activity
from ..internal.compat.activity import wrap_activity, wrap_activity_payloads
from ..internal.invocation import preserve_function_source, wrap_invocation
from ..worker import DurableFunctionsWorker
from ..orchestrator import Orchestrator


class Blueprint(TriggerApi, BindingApi):
Expand Down Expand Up @@ -166,6 +167,7 @@ def configure_history_export(self, writer: Any) -> None:
def _configure_orchestrator_callable(
self,
wrap: Callable[[Callable[..., Any]], FunctionBuilder],
context_name: str,
input_type: Optional[type] = None
) -> Callable[[task.Orchestrator[Any, Any]], FunctionBuilder]:
"""Obtain decorator to construct an Orchestrator class from a user-defined Function.
Expand Down Expand Up @@ -193,7 +195,13 @@ def decorator(orchestrator_func: task.Orchestrator[Any, Any]) -> FunctionBuilder
# feed it to a v1-style ``context.get_input()``.
orchestrator_func._df_input_type = input_type # type: ignore[attr-defined] # noqa: E501

handle = Orchestrator.create(orchestrator_func)
worker = DurableFunctionsWorker()

async def handle(context: func.OrchestrationContext) -> str:
return await worker.execute_orchestration_request_async(orchestrator_func, context)

handle.orchestrator_function = orchestrator_func # pyright: ignore[reportFunctionMemberAccess]
handle = wrap_invocation(handle, "context", registered_trigger_name=context_name)

# invoke next decorator, with the Orchestrator as input
handle.__name__ = orchestrator_func.__name__
Expand All @@ -204,6 +212,7 @@ def decorator(orchestrator_func: task.Orchestrator[Any, Any]) -> FunctionBuilder
def _configure_entity_callable(
self,
wrap: Callable[[Callable[..., Any]], FunctionBuilder],
context_name: str,
entity_name: Optional[str] = None
) -> Callable[[task.Entity[Any, Any]], FunctionBuilder]:
"""Obtain decorator to construct an Entity class from a user-defined Function.
Expand Down Expand Up @@ -235,18 +244,16 @@ def decorator(entity_func: task.Entity[Any, Any]) -> FunctionBuilder:
# Construct an orchestrator based on the end-user code
worker = DurableFunctionsWorker()

# TODO: Because this handle method is the one actually exposed to the Functions SDK decorator,
# the parameter name will always be "context" here, even if the user specified a different name.
# We need to find a way to allow custom context names (like "ctx").
# The generated handle is what the Azure Functions host registers,
# so its ``context`` parameter must be annotated with
# ``azure.functions.EntityContext`` for the host's entityTrigger
# binding converter to accept it; at runtime the host passes that
# transport context (exposing ``.body``).
def handle(context: func.EntityContext) -> str:
return worker.execute_entity_batch_request(entity_func, context)
async def handle(context: func.EntityContext) -> str:
return await worker.execute_entity_batch_request_async(entity_func, context)

handle.entity_function = entity_func # pyright: ignore[reportFunctionMemberAccess]
handle = wrap_invocation(handle, "context", registered_trigger_name=context_name)

# invoke next decorator, with the Entity as input
handle.__name__ = entity_func.__name__
Expand Down Expand Up @@ -294,14 +301,15 @@ def orchestration_trigger(self, context_name: str,
def wrap(fb: FunctionBuilder) -> FunctionBuilder:

def decorator() -> FunctionBuilder:
registered = fb._function._func # pyright: ignore[reportPrivateUsage]
fb.add_trigger(
trigger=OrchestrationTrigger(name=context_name,
trigger=OrchestrationTrigger(name=getattr(registered, "_df_trigger_name", context_name),
orchestration=orchestration))
return fb

return decorator()

return self._configure_orchestrator_callable(wrap, input_type=input_type)
return self._configure_orchestrator_callable(wrap, context_name, input_type=input_type)

def activity_trigger(self, input_name: str,
activity: Optional[str] = None
Expand All @@ -326,7 +334,13 @@ def decorator(user_fn: Callable[..., Any]) -> FunctionBuilder:
# Adapt a durabletask-native two-argument activity ((ctx, input))
# to the host's single-input convention; one-argument activities
# pass through unchanged.
return wrap(wrap_activity(user_fn, input_name))
function = (user_fn._function._func # pyright: ignore[reportPrivateUsage]
if isinstance(user_fn, FunctionBuilder) else user_fn)
registered = wrap_activity_payloads(wrap_activity(function, input_name), input_name)
Comment thread
andystaples marked this conversation as resolved.
if isinstance(user_fn, FunctionBuilder):
user_fn._function._func = registered # pyright: ignore[reportPrivateUsage]
return wrap(user_fn)
return wrap(registered)

return decorator

Expand All @@ -347,14 +361,15 @@ def entity_trigger(self,
@self._build_function
def wrap(fb: FunctionBuilder) -> FunctionBuilder:
def decorator() -> FunctionBuilder:
registered = fb._function._func # pyright: ignore[reportPrivateUsage]
fb.add_trigger(
trigger=EntityTrigger(name=context_name,
trigger=EntityTrigger(name=getattr(registered, "_df_trigger_name", context_name),
entity_name=entity_name))
return fb

return decorator()

return self._configure_entity_callable(wrap, entity_name)
return self._configure_entity_callable(wrap, context_name, entity_name)

def durable_client_input(self,
client_name: str,
Expand Down Expand Up @@ -439,6 +454,7 @@ def set_client_metadata(client_bound: Callable[..., Any]) -> None:
annotations[client_name] = str
client_bound.__annotations__ = annotations
setattr(client_bound, "client_function", function)
preserve_function_source(client_bound, function)

if is_async_function:
@wraps(function)
Expand Down Expand Up @@ -477,6 +493,19 @@ class DFApp(Blueprint, FunctionRegister):
Exports the decorators required to declare and index DF Function-types.
"""

def configure_large_payloads(self, *, payload_store: PayloadStore) -> None:
"""Enable payload externalization for this app and its blueprints.

Call once at app startup in every worker process, before invocations.
The store is shared by all durable clients, orchestrations, entities,
and activities in the process. All scaled-out workers must have access
to the same backing storage. Registering a different store in the same
process raises ValueError. Re-registering the same object is allowed.
"""
from ..internal.payloads import configure_payload_store

configure_payload_store(payload_store)

def register_functions(self, function_container: DecoratorApi) -> None:
"""Register the functions of a blueprint into this app.

Expand Down
Loading
Loading