Skip to content
Open
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
78 changes: 78 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,84 @@ to include examples, links to docs, or any other relevant information.
through `ExternalOutputStreamClient`. Workflow output is staged outside
History and becomes readable only after its compact Workflow Task marker is
committed.
### Changed

- Standalone Activities are now generally available (GA). (Standalone Activities as Nexus operations
and Standalone Activities operator commands remain experimental. Operator commands are `pause`,
`unpause`, `updateOptions`, `restoreOriginal`.)
- System Nexus Signal-with-Start Workflow operations now use the typed
`WorkflowOutboundInterceptor.start_signal_with_start_workflow` interception point instead of
the generic `WorkflowOutboundInterceptor.start_nexus_operation` method.
- System Nexus Signal-with-Start Workflow operations now invoke
`WorkflowOutboundInterceptor.start_system_nexus_operation` after their typed interception
point. They continue not to invoke `WorkflowOutboundInterceptor.start_nexus_operation`.
- The experimental `GetNexusOperationResultInput` now includes the Nexus endpoint, service, and
operation.

### :boom: Breaking Changes

- Experimental external storage: `ExternalStorage.driver_selector` is now called with a
`StorageDriverSelectContext` instead of a `StorageDriverStoreContext`. Update the annotation;
the new type carries the same `target` field. Since selectors are plain callables, a stale
annotation fails type checking rather than at runtime.
- `client.ActivityExecution` and `client.ActivityExecutionDescription` had some fields removed or renamed
to match RPC API.
- Dataclass parameters for these types were changed to `frozen=True, eq=False, kw_only=True`.
- `scheduled_time` was renamed `schedule_time`.
- `last_failure` was changed from field to method that runs data converter on demand.
- `state_transition_count`, `eager_execution_requested`, `paused` and `long_poll_token` were removed.
- ActivityHandle.describe() long-poll token was removed. The functionality can still be used manually
through raw gRPC API.

### Fixed

- `temporalio.contrib.google_genai` now requires `google-genai` 2.21.0 or later
and supports its file download API, including video inputs and download
destinations.
- `temporalio.contrib.deepagents` no longer dedups repeated identical tool,
model, and backend-op calls: each dispatch runs its own Activity, and the
continue-as-new result cache is retired for new executions (a continued run
resumes from the carried transcript and never re-executes prior dispatches,
so a carried cache entry could only serve stale results). Patch-gated
(`deepagents.retire-result-cache`), so histories recorded before this change
replay unchanged; note that deferring the patch keeps the full legacy dedup
cache — including the stale-result behavior this entry describes — and that
a chain upgraded mid-continue-as-new re-executes rather than reuses a
repeated identical call (the conservative direction).
- `contrib.deepagents`: summarization middleware configured with a model name string now routes its LLM calls through Activities instead of running them in the Workflow.

- **Experimental**: External storage metrics now report the wall-clock time storage was in flight.
Previously each batch's duration was summed, over-reporting the time whenever storage operations
ran concurrently.
- System Nexus Signal-with-Start workflow operations now give custom payload
converters the target workflow's serialization context when encoding their
inner request payloads.
- Cancelling an activity from a signal while the workflow itself is cancelled
no longer causes a nondeterminism error from duplicate activity-cancellation
commands.
- `StrandsPlugin` now disables Botocore retries for its default Bedrock model so
model request retries are handled exclusively by Temporal.
- `temporalio.contrib.openai_agents` now honors the `retry-after-ms` and
`retry-after` headers when OpenAI returns `x-should-retry: true`. Previously
the delay the server asked for was discarded on that path and the activity
retried on its configured interval instead.
- Nexus-context workflow/activity starts no longer set `on_conflict_options` when there are no links
or callbacks to attach.
- The workflow sandbox now passes `pydantic_core` through by default, alongside `pydantic`.

## [1.32.0] - 2026-08-24

### Added

- Added `temporalio.converter.create_payload_validation_error` to create the
non-retryable application error used when a converted payload fails validation.
- Added experimental `temporalio.contrib.opentelemetry.ReplaySafeMeterProvider` and
`ReplaySafeLoggerProvider` (and exported `ReplaySafeTracerProvider`): wrap an
OpenTelemetry provider so metrics and log events recorded from workflow code (e.g. by
Google ADK) are not duplicated on replay. `GoogleAdkPlugin` warns when a global OTel
provider is not replay-safe.
- Added `LoggingConfig.format` to select compact, pretty, or newline-delimited JSON output for
Core logs written to the console.

- Added the `Runtime(disable_environment_info=...)` option to control whether
runtime, hosting, and platform information is included in worker heartbeats.
Expand Down
7 changes: 6 additions & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ lambda-worker-otel = [
"opentelemetry-sdk-extension-aws>=2.0.0,<3",
]
aioboto3 = ["aioboto3>=10.4.0", "types-aioboto3[s3]>=10.4.0"]
google-genai = ["google-genai>=2.10.0,<3.0.0"]
google-genai = ["google-genai>=2.21.0,<3.0.0"]
strands-agents = ["strands-agents>=1.39.0"]

[project.urls]
Expand Down Expand Up @@ -188,8 +188,13 @@ exclude = [
# Ignore generated code
'temporalio/api',
'temporalio/bridge/proto',
'temporalio/nexus/system/workflow_service',
]

[[tool.mypy.overrides]]
module = "temporalio.nexus.system.workflow_service.*"
ignore_errors = true

[tool.pydocstyle]
convention = "google"
# https://github.com/PyCQA/pydocstyle/issues/363#issuecomment-625563088
Expand Down
25 changes: 18 additions & 7 deletions scripts/gen_nexus_system_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,15 +35,28 @@
/ "v1"
/ "request_response.proto"
)
NEX_GEN_VERSION = "0.2.4"


def nex_gen_command() -> list[str]:
if bin_path := os.environ.get("NEX_GEN_BIN"):
return [bin_path]

if shutil.which("nex-gen") is None:
subprocess.check_call(["cargo", "install", "--locked", "nex-gen", "--force"])
return ["nex-gen"]
if shutil.which("nexgen") is None:
subprocess.check_call(
[
"cargo",
"install",
"--locked",
"nexgen",
"--version",
NEX_GEN_VERSION,
"--features",
"advanced",
"--force",
]
)
return ["nexgen"]


def build_descriptor_set(descriptor_path: Path) -> None:
Expand Down Expand Up @@ -114,13 +127,11 @@ def generate_nexus_system_api() -> None:
subprocess.check_call(
[
*command,
"generate",
"--lang",
"python",
"--input",
str(wit_path),
"--input",
str(wit_deps_dir),
"--native-api",
"--system-nexus",
"--support-file",
str(python_support_path),
"--descriptors",
Expand Down
37 changes: 25 additions & 12 deletions scripts/gen_payload_visitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import sys
from importlib.util import module_from_spec, spec_from_file_location
from pathlib import Path
from typing import cast
from typing import cast, get_args, get_origin, get_type_hints

import google.protobuf.message
import nexusrpc
Expand All @@ -19,14 +19,15 @@
from temporalio.bridge.proto.workflow_completion.workflow_completion_pb2 import (
WorkflowActivationCompletion,
)
from temporalio.converter._payload_converter import _get_transfer_type_converter


def discover_system_nexus_roots() -> list[Descriptor]:
module_path = (
base_dir / "temporalio" / "nexus" / "system" / "workflow_service" / "service.py"
)
module_path = base_dir / "temporalio" / "nexus" / "system" / "workflow_service"
spec = spec_from_file_location(
"temporalio_nexus_system_workflow_service", module_path
"temporalio_nexus_system_workflow_service",
module_path / "__init__.py",
submodule_search_locations=[str(module_path)],
)
if spec is None or spec.loader is None:
raise RuntimeError(f"Cannot load generated system service from {module_path}")
Expand All @@ -35,10 +36,14 @@ def discover_system_nexus_roots() -> list[Descriptor]:
spec.loader.exec_module(module)

roots: list[Descriptor] = []
for operation in vars(module.WorkflowService).values():
if not isinstance(operation, nexusrpc.Operation):
for annotation in get_type_hints(module._services.WorkflowService).values():
if get_origin(annotation) is not nexusrpc.Operation:
continue
for proto_type in (operation.input_type, operation.output_type):
for operation_type in get_args(annotation):
converter = _get_transfer_type_converter(operation_type)
proto_type = (
converter.transfer_type if converter is not None else operation_type
)
if isinstance(proto_type, type) and issubclass(
proto_type, google.protobuf.message.Message
):
Expand Down Expand Up @@ -119,12 +124,20 @@ def generate(self, roots: list[Descriptor]) -> str:

The generated code defines async visitor functions for each reachable
protobuf message type starting from WorkflowActivation, including support
for repeated fields and map entries, and a convenience entrypoint
function `visit`.
for repeated fields and map entries. Payload-free roots get no-op methods
so the `visit` entrypoint recognizes them as supported.
"""

for r in roots:
self.walk(r)
for root in roots:
if not self.walk(root):
self.methods.append(
f"""\
async def _visit_{name_for(root)}(
self, fs: VisitorFunctions, o: Any
) -> None:
pass
"""
)

header = """
from __future__ import annotations
Expand Down
93 changes: 86 additions & 7 deletions scripts/nex_gen_support.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,12 @@
import temporalio.api.workflow.v1
import temporalio.common
import temporalio.converter
import temporalio.nexus.system


class SignalWithStartWorkflowModelRequest(typing.Protocol):
namespace: str
id: str


def retry_policy_from_proto(
Expand Down Expand Up @@ -52,6 +58,12 @@ def workflow_type_to_proto(
return common_pb2.WorkflowType(name=workflow_function_name(workflow_type))


def workflow_type_from_proto(
proto: common_pb2.WorkflowType,
) -> str:
return proto.name


def task_queue_from_proto(
proto: taskqueue_pb2.TaskQueue,
) -> str:
Expand All @@ -70,12 +82,33 @@ def workflow_namespace() -> str:
return info().namespace


def signal_with_start_workflow_serialization_context(
request: SignalWithStartWorkflowModelRequest,
) -> temporalio.converter.WorkflowSerializationContext:
return temporalio.converter.WorkflowSerializationContext(
namespace=request.namespace,
workflow_id=request.id,
)


def payloads_to_proto(
values: collections.abc.Sequence[typing.Any],
) -> common_pb2.Payloads:
from temporalio.workflow import payload_converter
return (
temporalio.nexus.system._current_user_payload_converter().to_payloads_wrapper(
values
)
)

return payload_converter().to_payloads_wrapper(values)

def payloads_from_proto(
proto: common_pb2.Payloads,
) -> list[object]:
return list(
temporalio.nexus.system._current_user_payload_converter().from_payloads_wrapper(
proto
)
)


def _clone_payload(payload: common_pb2.Payload) -> common_pb2.Payload:
Expand All @@ -87,20 +120,23 @@ def _clone_payload(payload: common_pb2.Payload) -> common_pb2.Payload:
def _value_to_payload(value: object | common_pb2.Payload) -> common_pb2.Payload:
if isinstance(value, common_pb2.Payload):
return _clone_payload(value)
from temporalio.workflow import payload_converter

payloads = payload_converter().to_payloads_wrapper([value])
payloads = (
temporalio.nexus.system._current_user_payload_converter().to_payloads_wrapper(
[value]
)
)
return _clone_payload(payloads.payloads[0])


def _payload_to_value(payload: common_pb2.Payload) -> object:
wrapper = common_pb2.Payloads()
wrapper.payloads.add().CopyFrom(payload)
from temporalio.workflow import payload_converter

return typing.cast(
object,
payload_converter().from_payloads_wrapper(wrapper)[0],
temporalio.nexus.system._current_user_payload_converter().from_payloads_wrapper(
wrapper
)[0],
)


Expand Down Expand Up @@ -131,6 +167,21 @@ def memo_to_proto(
return message


def header_from_proto(
proto: common_pb2.Header,
) -> collections.abc.Mapping[str, object]:
return {key: _payload_to_value(value) for key, value in proto.fields.items()}


def header_to_proto(
header: collections.abc.Mapping[str, object],
) -> common_pb2.Header:
message = common_pb2.Header()
for key, value in header.items():
message.fields[key].CopyFrom(_value_to_payload(value))
return message


def duration_from_proto(proto: google.protobuf.duration_pb2.Duration) -> timedelta:
return proto.ToTimedelta()

Expand Down Expand Up @@ -177,6 +228,12 @@ def search_attributes_to_proto(
return proto


def search_attributes_from_proto(
proto: common_pb2.SearchAttributes,
) -> temporalio.common.TypedSearchAttributes:
return temporalio.converter.decode_typed_search_attributes(proto)


def priority_from_proto(
proto: common_pb2.Priority,
) -> temporalio.common.Priority:
Expand All @@ -193,3 +250,25 @@ def versioning_override_to_proto(
versioning_override: temporalio.common.VersioningOverride,
) -> temporalio.api.workflow.v1.VersioningOverride:
return versioning_override._to_proto() # pyright: ignore[reportPrivateUsage]


def versioning_override_from_proto(
proto: temporalio.api.workflow.v1.VersioningOverride,
) -> temporalio.common.VersioningOverride:
if proto.HasField("pinned") and proto.pinned.HasField("version"):
version = proto.pinned.version
return temporalio.common.PinnedVersioningOverride(
temporalio.common.WorkerDeploymentVersion(
deployment_name=version.deployment_name,
build_id=version.build_id,
)
)
if proto.pinned_version:
return temporalio.common.PinnedVersioningOverride(
temporalio.common.WorkerDeploymentVersion.from_canonical_string(
proto.pinned_version
)
)
if proto.auto_upgrade:
return temporalio.common.AutoUpgradeVersioningOverride()
raise ValueError("unknown versioning override proto shape")
5 changes: 5 additions & 0 deletions temporalio/bridge/_visitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -621,3 +621,8 @@ async def _visit_temporal_api_workflowservice_v1_SignalWithStartWorkflowExecutio
await self._visit_temporal_api_common_v1_Header(fs, o.header)
if o.HasField("user_metadata"):
await self._visit_temporal_api_sdk_v1_UserMetadata(fs, o.user_metadata)

async def _visit_temporal_api_workflowservice_v1_SignalWithStartWorkflowExecutionResponse(
self, fs: VisitorFunctions, o: Any
) -> None:
pass
Loading
Loading