From da04dc9768e5579b0c52ff15dc815fbd33857b03 Mon Sep 17 00:00:00 2001 From: Justin Anderson <44687433+jmaeagle99@users.noreply.github.com> Date: Tue, 11 Aug 2026 17:17:33 -0700 Subject: [PATCH 1/2] chore: update sdk-core submodule to latest (#1738) * chore: update sdk-core submodule to latest * fix: update sdk-core to pending change sdk-rust * chore: update sdk-core submodule to latest --- temporalio/bridge/Cargo.lock | 214 +----------------- temporalio/bridge/Cargo.toml | 6 +- temporalio/bridge/proto/common/__init__.py | 2 + temporalio/bridge/proto/common/common_pb2.py | 20 +- temporalio/bridge/proto/common/common_pb2.pyi | 53 +++++ .../workflow_completion_pb2.py | 12 +- .../workflow_completion_pb2.pyi | 36 ++- temporalio/bridge/sdk-core | 2 +- temporalio/bridge/src/client.rs | 80 +++---- temporalio/bridge/src/envconfig.rs | 22 +- temporalio/bridge/src/worker.rs | 29 +-- 11 files changed, 195 insertions(+), 281 deletions(-) diff --git a/temporalio/bridge/Cargo.lock b/temporalio/bridge/Cargo.lock index 330916b2b..10632afb2 100644 --- a/temporalio/bridge/Cargo.lock +++ b/temporalio/bridge/Cargo.lock @@ -58,29 +58,6 @@ version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" -[[package]] -name = "aws-lc-rs" -version = "1.17.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "00bdb5da18dac48ca2cc7cd4a98e533e8635a58e2361d13a1a4ee3888e0d72f1" -dependencies = [ - "aws-lc-sys", - "zeroize", -] - -[[package]] -name = "aws-lc-sys" -version = "0.43.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43103168cc76fe62678a375e722fc9cb3a0146159ac5828bc4f0dfd755c2224c" -dependencies = [ - "cc", - "cmake", - "dunce", - "fs_extra", - "pkg-config", -] - [[package]] name = "axum" version = "0.8.9" @@ -209,12 +186,6 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" -[[package]] -name = "cfg_aliases" -version = "0.2.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f079e83a288787bcd14a6aea84cee5c87a67c5a3e660c30f557a3d24761b3527" - [[package]] name = "chacha20" version = "0.10.1" @@ -236,15 +207,6 @@ dependencies = [ "serde", ] -[[package]] -name = "cmake" -version = "0.1.58" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c0f78a02292a74a88ac736019ab962ece0bc380e3f977bf72e376c5d78ff0678" -dependencies = [ - "cc", -] - [[package]] name = "combine" version = "4.6.7" @@ -408,12 +370,6 @@ version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1435fa1053d8b2fbbe9be7e97eca7f33d37b28409959813daefc1446a14247f1" -[[package]] -name = "dunce" -version = "1.0.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "92773504d58c093f6de2459af4af33faa518c13451eb8f2b5698ed3d36e7c813" - [[package]] name = "dyn-clone" version = "1.0.20" @@ -560,12 +516,6 @@ dependencies = [ "futures-core", ] -[[package]] -name = "fs_extra" -version = "1.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" - [[package]] name = "futures" version = "0.3.33" @@ -682,10 +632,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" dependencies = [ "cfg-if", - "js-sys", "libc", "wasi", - "wasm-bindgen", ] [[package]] @@ -707,11 +655,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" dependencies = [ "cfg-if", - "js-sys", "libc", "r-efi 6.0.0", "rand_core 0.10.1", - "wasm-bindgen", ] [[package]] @@ -1159,12 +1105,6 @@ dependencies = [ "hashbrown 0.17.1", ] -[[package]] -name = "lru-slab" -version = "0.1.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" - [[package]] name = "matchers" version = "0.2.0" @@ -1326,7 +1266,7 @@ dependencies = [ "bytes", "http", "opentelemetry", - "reqwest 0.13.4", + "reqwest", ] [[package]] @@ -1341,7 +1281,7 @@ dependencies = [ "opentelemetry-proto", "opentelemetry_sdk", "prost", - "reqwest 0.13.4", + "reqwest", "thiserror", "tokio", "tonic", @@ -1782,63 +1722,6 @@ dependencies = [ "serde", ] -[[package]] -name = "quinn" -version = "0.11.11" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" -dependencies = [ - "bytes", - "cfg_aliases", - "pin-project-lite", - "quinn-proto", - "quinn-udp", - "rustc-hash", - "rustls", - "socket2", - "thiserror", - "tokio", - "tracing", - "web-time", -] - -[[package]] -name = "quinn-proto" -version = "0.11.16" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2f4bfc015262b9df63c8845072ce59068853ff5872180c2ce2f13038b970e560" -dependencies = [ - "aws-lc-rs", - "bytes", - "getrandom 0.4.3", - "lru-slab", - "rand 0.10.2", - "rand_pcg", - "ring", - "rustc-hash", - "rustls", - "rustls-pki-types", - "slab", - "thiserror", - "tinyvec", - "tracing", - "web-time", -] - -[[package]] -name = "quinn-udp" -version = "0.5.15" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "35a133f956daabe89a61a685c2649f13d82d5aa4bd5d12d1277e1072a21c0694" -dependencies = [ - "cfg_aliases", - "libc", - "once_cell", - "socket2", - "tracing", - "windows-sys 0.61.2", -] - [[package]] name = "quote" version = "1.0.47" @@ -1906,15 +1789,6 @@ version = "0.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" -[[package]] -name = "rand_pcg" -version = "0.10.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" -dependencies = [ - "rand_core 0.10.1", -] - [[package]] name = "redox_syscall" version = "0.5.18" @@ -1964,38 +1838,6 @@ version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" -[[package]] -name = "reqwest" -version = "0.12.28" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" -dependencies = [ - "base64", - "bytes", - "futures-core", - "http", - "http-body", - "http-body-util", - "hyper", - "hyper-util", - "js-sys", - "log", - "percent-encoding", - "pin-project-lite", - "serde", - "serde_json", - "serde_urlencoded", - "sync_wrapper", - "tokio", - "tower", - "tower-http", - "tower-service", - "url", - "wasm-bindgen", - "wasm-bindgen-futures", - "web-sys", -] - [[package]] name = "reqwest" version = "0.13.4" @@ -2017,7 +1859,6 @@ dependencies = [ "log", "percent-encoding", "pin-project-lite", - "quinn", "rustls", "rustls-pki-types", "rustls-platform-verifier", @@ -2063,12 +1904,6 @@ dependencies = [ "portable-atomic-util", ] -[[package]] -name = "rustc-hash" -version = "2.1.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6b1e7f9a428571be2dc5bc0505c13fb6bf936822b894ec87abf8a08a4e51742d" - [[package]] name = "rustc_version" version = "0.4.1" @@ -2097,7 +1932,6 @@ version = "0.23.42" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3c54fcab019b409d04215d3a17cb438fd7fbf192ee61461f20f4fe18704bc138" dependencies = [ - "aws-lc-rs", "log", "once_cell", "ring", @@ -2125,7 +1959,6 @@ version = "1.15.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2f4925028c7eb5d1fcdaf196971378ed9d2c1c4efc7dc5d011256f76c99c0a96" dependencies = [ - "web-time", "zeroize", ] @@ -2162,7 +1995,6 @@ version = "0.103.13" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" dependencies = [ - "aws-lc-rs", "ring", "rustls-pki-types", "untrusted", @@ -2509,7 +2341,7 @@ dependencies = [ [[package]] name = "temporalio-client" -version = "0.5.0" +version = "0.6.0" dependencies = [ "anyhow", "async-trait", @@ -2540,7 +2372,7 @@ dependencies = [ [[package]] name = "temporalio-common" -version = "0.5.0" +version = "0.6.0" dependencies = [ "anyhow", "async-trait", @@ -2561,8 +2393,9 @@ dependencies = [ "prometheus", "prost", "prost-types", - "reqwest 0.12.28", + "reqwest", "ringbuf", + "rustls", "serde", "serde_json", "temporalio-common-wasm", @@ -2580,7 +2413,7 @@ dependencies = [ [[package]] name = "temporalio-common-wasm" -version = "0.5.0" +version = "0.6.0" dependencies = [ "anyhow", "async-trait", @@ -2605,7 +2438,7 @@ dependencies = [ [[package]] name = "temporalio-macros" -version = "0.5.0" +version = "0.6.0" dependencies = [ "proc-macro2", "quote", @@ -2614,7 +2447,7 @@ dependencies = [ [[package]] name = "temporalio-protos" -version = "0.5.0" +version = "0.6.0" dependencies = [ "anyhow", "base64", @@ -2635,7 +2468,7 @@ dependencies = [ [[package]] name = "temporalio-sdk-core" -version = "0.5.0" +version = "0.6.0" dependencies = [ "anyhow", "async-trait", @@ -2660,7 +2493,7 @@ dependencies = [ "prost", "prost-wkt-types", "rand 0.10.2", - "reqwest 0.13.4", + "reqwest", "serde", "serde_json", "siphasher", @@ -2726,21 +2559,6 @@ dependencies = [ "zerovec", ] -[[package]] -name = "tinyvec" -version = "1.12.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bb4ebadaa0af04fab11ae01eb5f9fdb5f9c5b875506e210e71c07873528baa7f" -dependencies = [ - "tinyvec_macros", -] - -[[package]] -name = "tinyvec_macros" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" - [[package]] name = "tokio" version = "1.53.1" @@ -3246,16 +3064,6 @@ dependencies = [ "wasm-bindgen", ] -[[package]] -name = "web-time" -version = "1.1.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" -dependencies = [ - "js-sys", - "wasm-bindgen", -] - [[package]] name = "webpki-root-certs" version = "1.0.9" diff --git a/temporalio/bridge/Cargo.toml b/temporalio/bridge/Cargo.toml index 5984d9546..4cf6a02fe 100644 --- a/temporalio/bridge/Cargo.toml +++ b/temporalio/bridge/Cargo.toml @@ -28,11 +28,11 @@ pyo3 = { version = "0.29", features = [ ] } pyo3-async-runtimes = { version = "0.29", features = ["tokio-runtime"] } pythonize = "0.29" -temporalio-client = { version = "0.5", path = "./sdk-core/crates/client" } -temporalio-common = { version = "0.5", path = "./sdk-core/crates/common", features = [ +temporalio-client = { version = "0.6", path = "./sdk-core/crates/client" } +temporalio-common = { version = "0.6", path = "./sdk-core/crates/common", features = [ "envconfig", "otel" ]} -temporalio-sdk-core = { version = "0.5", path = "./sdk-core/crates/sdk-core", features = [ +temporalio-sdk-core = { version = "0.6", path = "./sdk-core/crates/sdk-core", features = [ "ephemeral-server", ] } tokio = "1.26" diff --git a/temporalio/bridge/proto/common/__init__.py b/temporalio/bridge/proto/common/__init__.py index 5622fffb8..a8506090d 100644 --- a/temporalio/bridge/proto/common/__init__.py +++ b/temporalio/bridge/proto/common/__init__.py @@ -1,10 +1,12 @@ from .common_pb2 import ( + ExternalStorageMetrics, NamespacedWorkflowExecution, VersioningIntent, WorkerDeploymentVersion, ) __all__ = [ + "ExternalStorageMetrics", "NamespacedWorkflowExecution", "VersioningIntent", "WorkerDeploymentVersion", diff --git a/temporalio/bridge/proto/common/common_pb2.py b/temporalio/bridge/proto/common/common_pb2.py index c56456fce..481cf216d 100644 --- a/temporalio/bridge/proto/common/common_pb2.py +++ b/temporalio/bridge/proto/common/common_pb2.py @@ -18,7 +18,7 @@ from google.protobuf import duration_pb2 as google_dot_protobuf_dot_duration__pb2 DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile( - b'\n%temporal/sdk/core/common/common.proto\x12\x0e\x63oresdk.common\x1a\x1egoogle/protobuf/duration.proto"U\n\x1bNamespacedWorkflowExecution\x12\x11\n\tnamespace\x18\x01 \x01(\t\x12\x13\n\x0bworkflow_id\x18\x02 \x01(\t\x12\x0e\n\x06run_id\x18\x03 \x01(\t"D\n\x17WorkerDeploymentVersion\x12\x17\n\x0f\x64\x65ployment_name\x18\x01 \x01(\t\x12\x10\n\x08\x62uild_id\x18\x02 \x01(\t*@\n\x10VersioningIntent\x12\x0f\n\x0bUNSPECIFIED\x10\x00\x12\x0e\n\nCOMPATIBLE\x10\x01\x12\x0b\n\x07\x44\x45\x46\x41ULT\x10\x02\x42,\xea\x02)Temporalio::Internal::Bridge::Api::Commonb\x06proto3' + b'\n%temporal/sdk/core/common/common.proto\x12\x0e\x63oresdk.common\x1a\x1egoogle/protobuf/duration.proto"U\n\x1bNamespacedWorkflowExecution\x12\x11\n\tnamespace\x18\x01 \x01(\t\x12\x13\n\x0bworkflow_id\x18\x02 \x01(\t\x12\x0e\n\x06run_id\x18\x03 \x01(\t"D\n\x17WorkerDeploymentVersion\x12\x17\n\x0f\x64\x65ployment_name\x18\x01 \x01(\t\x12\x10\n\x08\x62uild_id\x18\x02 \x01(\t"\x92\x01\n\x16\x45xternalStorageMetrics\x12\x15\n\rpayload_count\x18\x01 \x01(\x04\x12\x18\n\x10total_size_bytes\x18\x02 \x01(\x04\x12\x31\n\x0etotal_duration\x18\x03 \x01(\x0b\x32\x19.google.protobuf.Duration\x12\x14\n\x0c\x64river_names\x18\x04 \x03(\t*@\n\x10VersioningIntent\x12\x0f\n\x0bUNSPECIFIED\x10\x00\x12\x0e\n\nCOMPATIBLE\x10\x01\x12\x0b\n\x07\x44\x45\x46\x41ULT\x10\x02\x42,\xea\x02)Temporalio::Internal::Bridge::Api::Commonb\x06proto3' ) _VERSIONINGINTENT = DESCRIPTOR.enum_types_by_name["VersioningIntent"] @@ -32,6 +32,7 @@ "NamespacedWorkflowExecution" ] _WORKERDEPLOYMENTVERSION = DESCRIPTOR.message_types_by_name["WorkerDeploymentVersion"] +_EXTERNALSTORAGEMETRICS = DESCRIPTOR.message_types_by_name["ExternalStorageMetrics"] NamespacedWorkflowExecution = _reflection.GeneratedProtocolMessageType( "NamespacedWorkflowExecution", (_message.Message,), @@ -54,15 +55,28 @@ ) _sym_db.RegisterMessage(WorkerDeploymentVersion) +ExternalStorageMetrics = _reflection.GeneratedProtocolMessageType( + "ExternalStorageMetrics", + (_message.Message,), + { + "DESCRIPTOR": _EXTERNALSTORAGEMETRICS, + "__module__": "temporal.sdk.core.common.common_pb2", + # @@protoc_insertion_point(class_scope:coresdk.common.ExternalStorageMetrics) + }, +) +_sym_db.RegisterMessage(ExternalStorageMetrics) + if _descriptor._USE_C_DESCRIPTORS == False: DESCRIPTOR._options = None DESCRIPTOR._serialized_options = ( b"\352\002)Temporalio::Internal::Bridge::Api::Common" ) - _VERSIONINGINTENT._serialized_start = 246 - _VERSIONINGINTENT._serialized_end = 310 + _VERSIONINGINTENT._serialized_start = 395 + _VERSIONINGINTENT._serialized_end = 459 _NAMESPACEDWORKFLOWEXECUTION._serialized_start = 89 _NAMESPACEDWORKFLOWEXECUTION._serialized_end = 174 _WORKERDEPLOYMENTVERSION._serialized_start = 176 _WORKERDEPLOYMENTVERSION._serialized_end = 244 + _EXTERNALSTORAGEMETRICS._serialized_start = 247 + _EXTERNALSTORAGEMETRICS._serialized_end = 393 # @@protoc_insertion_point(module_scope) diff --git a/temporalio/bridge/proto/common/common_pb2.pyi b/temporalio/bridge/proto/common/common_pb2.pyi index 739a129e1..8862fa036 100644 --- a/temporalio/bridge/proto/common/common_pb2.pyi +++ b/temporalio/bridge/proto/common/common_pb2.pyi @@ -4,10 +4,13 @@ isort:skip_file """ import builtins +import collections.abc import sys import typing import google.protobuf.descriptor +import google.protobuf.duration_pb2 +import google.protobuf.internal.containers import google.protobuf.internal.enum_type_wrapper import google.protobuf.message @@ -121,3 +124,53 @@ class WorkerDeploymentVersion(google.protobuf.message.Message): ) -> None: ... global___WorkerDeploymentVersion = WorkerDeploymentVersion + +class ExternalStorageMetrics(google.protobuf.message.Message): + """Metrics for a set of external payload storage operations (all uploads and downloads) + performed while processing a task, so core can emit unified logging and metrics. + """ + + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + PAYLOAD_COUNT_FIELD_NUMBER: builtins.int + TOTAL_SIZE_BYTES_FIELD_NUMBER: builtins.int + TOTAL_DURATION_FIELD_NUMBER: builtins.int + DRIVER_NAMES_FIELD_NUMBER: builtins.int + payload_count: builtins.int + """Number of payloads stored or retrieved externally.""" + total_size_bytes: builtins.int + """Total size in bytes of the externally stored or retrieved payloads.""" + @property + def total_duration(self) -> google.protobuf.duration_pb2.Duration: + """Wall-clock time spent on the external storage operations.""" + @property + def driver_names( + self, + ) -> google.protobuf.internal.containers.RepeatedScalarFieldContainer[builtins.str]: + """Names of the drivers that participated in the operations.""" + def __init__( + self, + *, + payload_count: builtins.int = ..., + total_size_bytes: builtins.int = ..., + total_duration: google.protobuf.duration_pb2.Duration | None = ..., + driver_names: collections.abc.Iterable[builtins.str] | None = ..., + ) -> None: ... + def HasField( + self, field_name: typing_extensions.Literal["total_duration", b"total_duration"] + ) -> builtins.bool: ... + def ClearField( + self, + field_name: typing_extensions.Literal[ + "driver_names", + b"driver_names", + "payload_count", + b"payload_count", + "total_duration", + b"total_duration", + "total_size_bytes", + b"total_size_bytes", + ], + ) -> None: ... + +global___ExternalStorageMetrics = ExternalStorageMetrics diff --git a/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.py b/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.py index ce26b220d..057b301e4 100644 --- a/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.py +++ b/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.py @@ -31,7 +31,7 @@ ) DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile( - b'\n?temporal/sdk/core/workflow_completion/workflow_completion.proto\x12\x1b\x63oresdk.workflow_completion\x1a%temporal/api/failure/v1/message.proto\x1a(temporal/api/enums/v1/failed_cause.proto\x1a$temporal/api/enums/v1/workflow.proto\x1a%temporal/sdk/core/common/common.proto\x1a;temporal/sdk/core/workflow_commands/workflow_commands.proto"\xac\x01\n\x1cWorkflowActivationCompletion\x12\x0e\n\x06run_id\x18\x01 \x01(\t\x12:\n\nsuccessful\x18\x02 \x01(\x0b\x32$.coresdk.workflow_completion.SuccessH\x00\x12\x36\n\x06\x66\x61iled\x18\x03 \x01(\x0b\x32$.coresdk.workflow_completion.FailureH\x00\x42\x08\n\x06status"\xac\x01\n\x07Success\x12<\n\x08\x63ommands\x18\x01 \x03(\x0b\x32*.coresdk.workflow_commands.WorkflowCommand\x12\x1b\n\x13used_internal_flags\x18\x06 \x03(\r\x12\x46\n\x13versioning_behavior\x18\x07 \x01(\x0e\x32).temporal.api.enums.v1.VersioningBehavior"\x81\x01\n\x07\x46\x61ilure\x12\x31\n\x07\x66\x61ilure\x18\x01 \x01(\x0b\x32 .temporal.api.failure.v1.Failure\x12\x43\n\x0b\x66orce_cause\x18\x02 \x01(\x0e\x32..temporal.api.enums.v1.WorkflowTaskFailedCauseB8\xea\x02\x35Temporalio::Internal::Bridge::Api::WorkflowCompletionb\x06proto3' + b'\n?temporal/sdk/core/workflow_completion/workflow_completion.proto\x12\x1b\x63oresdk.workflow_completion\x1a%temporal/api/failure/v1/message.proto\x1a(temporal/api/enums/v1/failed_cause.proto\x1a$temporal/api/enums/v1/workflow.proto\x1a%temporal/sdk/core/common/common.proto\x1a;temporal/sdk/core/workflow_commands/workflow_commands.proto"\xbe\x02\n\x1cWorkflowActivationCompletion\x12\x0e\n\x06run_id\x18\x01 \x01(\t\x12:\n\nsuccessful\x18\x02 \x01(\x0b\x32$.coresdk.workflow_completion.SuccessH\x00\x12\x36\n\x06\x66\x61iled\x18\x03 \x01(\x0b\x32$.coresdk.workflow_completion.FailureH\x00\x12H\n\x18payload_download_metrics\x18\x04 \x01(\x0b\x32&.coresdk.common.ExternalStorageMetrics\x12\x46\n\x16payload_upload_metrics\x18\x05 \x01(\x0b\x32&.coresdk.common.ExternalStorageMetricsB\x08\n\x06status"\xac\x01\n\x07Success\x12<\n\x08\x63ommands\x18\x01 \x03(\x0b\x32*.coresdk.workflow_commands.WorkflowCommand\x12\x1b\n\x13used_internal_flags\x18\x06 \x03(\r\x12\x46\n\x13versioning_behavior\x18\x07 \x01(\x0e\x32).temporal.api.enums.v1.VersioningBehavior"\x81\x01\n\x07\x46\x61ilure\x12\x31\n\x07\x66\x61ilure\x18\x01 \x01(\x0b\x32 .temporal.api.failure.v1.Failure\x12\x43\n\x0b\x66orce_cause\x18\x02 \x01(\x0e\x32..temporal.api.enums.v1.WorkflowTaskFailedCauseB8\xea\x02\x35Temporalio::Internal::Bridge::Api::WorkflowCompletionb\x06proto3' ) @@ -79,9 +79,9 @@ b"\352\0025Temporalio::Internal::Bridge::Api::WorkflowCompletion" ) _WORKFLOWACTIVATIONCOMPLETION._serialized_start = 316 - _WORKFLOWACTIVATIONCOMPLETION._serialized_end = 488 - _SUCCESS._serialized_start = 491 - _SUCCESS._serialized_end = 663 - _FAILURE._serialized_start = 666 - _FAILURE._serialized_end = 795 + _WORKFLOWACTIVATIONCOMPLETION._serialized_end = 634 + _SUCCESS._serialized_start = 637 + _SUCCESS._serialized_end = 809 + _FAILURE._serialized_start = 812 + _FAILURE._serialized_end = 941 # @@protoc_insertion_point(module_scope) diff --git a/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.pyi b/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.pyi index 5b438f360..8e12736aa 100644 --- a/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.pyi +++ b/temporalio/bridge/proto/workflow_completion/workflow_completion_pb2.pyi @@ -14,6 +14,7 @@ import google.protobuf.message import temporalio.api.enums.v1.failed_cause_pb2 import temporalio.api.enums.v1.workflow_pb2 import temporalio.api.failure.v1.message_pb2 +import temporalio.bridge.proto.common.common_pb2 import temporalio.bridge.proto.workflow_commands.workflow_commands_pb2 if sys.version_info >= (3, 8): @@ -31,23 +32,52 @@ class WorkflowActivationCompletion(google.protobuf.message.Message): RUN_ID_FIELD_NUMBER: builtins.int SUCCESSFUL_FIELD_NUMBER: builtins.int FAILED_FIELD_NUMBER: builtins.int + PAYLOAD_DOWNLOAD_METRICS_FIELD_NUMBER: builtins.int + PAYLOAD_UPLOAD_METRICS_FIELD_NUMBER: builtins.int run_id: builtins.str """The run id from the workflow activation you are completing""" @property def successful(self) -> global___Success: ... @property def failed(self) -> global___Failure: ... + @property + def payload_download_metrics( + self, + ) -> temporalio.bridge.proto.common.common_pb2.ExternalStorageMetrics: + """Metrics for external payload storage downloads (retrievals) performed while processing + this activation. Only set when external storage retrieved payloads. + """ + @property + def payload_upload_metrics( + self, + ) -> temporalio.bridge.proto.common.common_pb2.ExternalStorageMetrics: + """Metrics for external payload storage uploads (stores) performed while processing this + activation. Only set when external storage stored payloads. + """ def __init__( self, *, run_id: builtins.str = ..., successful: global___Success | None = ..., failed: global___Failure | None = ..., + payload_download_metrics: temporalio.bridge.proto.common.common_pb2.ExternalStorageMetrics + | None = ..., + payload_upload_metrics: temporalio.bridge.proto.common.common_pb2.ExternalStorageMetrics + | None = ..., ) -> None: ... def HasField( self, field_name: typing_extensions.Literal[ - "failed", b"failed", "status", b"status", "successful", b"successful" + "failed", + b"failed", + "payload_download_metrics", + b"payload_download_metrics", + "payload_upload_metrics", + b"payload_upload_metrics", + "status", + b"status", + "successful", + b"successful", ], ) -> builtins.bool: ... def ClearField( @@ -55,6 +85,10 @@ class WorkflowActivationCompletion(google.protobuf.message.Message): field_name: typing_extensions.Literal[ "failed", b"failed", + "payload_download_metrics", + b"payload_download_metrics", + "payload_upload_metrics", + b"payload_upload_metrics", "run_id", b"run_id", "status", diff --git a/temporalio/bridge/sdk-core b/temporalio/bridge/sdk-core index 21fcc2952..78ac17e4b 160000 --- a/temporalio/bridge/sdk-core +++ b/temporalio/bridge/sdk-core @@ -1 +1 @@ -Subproject commit 21fcc2952257478489bd6e74cd0a65eb0b2c63be +Subproject commit 78ac17e4bdaffdb54eb1c925a07d1fd40ce7b7f6 diff --git a/temporalio/bridge/src/client.rs b/temporalio/bridge/src/client.rs index 37bfc796d..aec5c0ef8 100644 --- a/temporalio/bridge/src/client.rs +++ b/temporalio/bridge/src/client.rs @@ -280,10 +280,12 @@ impl ClientConfig { .maybe_http_connect_proxy(self.http_connect_proxy_config.map(Into::into)) .dns_load_balancing(dns_load_balancing) .grpc_compression(grpc_compression_from_str(&self.grpc_compression)?) - .payload_limits(temporalio_client::PayloadLimitsOptions { - payloads_warn_size: self.payloads_warn_size, - memo_warn_size: self.memo_warn_size, - }) + .payload_limits( + temporalio_client::PayloadLimitsOptions::builder() + .payloads_warn_size(self.payloads_warn_size) + .memo_warn_size(self.memo_warn_size) + .build(), + ) .headers(ascii_headers) .binary_headers(binary_headers) .maybe_api_key(self.api_key) @@ -325,25 +327,26 @@ impl TryFrom for temporalio_client::TlsOptions { Some(fixed_server_name_verifier(&name, &ca_cert)?) } }; - Ok(temporalio_client::TlsOptions { - server_root_ca_cert, - domain: conf.domain, - client_tls_options: match (conf.client_cert, conf.client_private_key) { - (None, None) => None, - (Some(client_cert), Some(client_private_key)) => { - Some(temporalio_client::ClientTlsOptions { - client_cert, - client_private_key, - }) - } - _ => { - return Err(PyValueError::new_err( - "Must have both client cert and private key or neither", - )) - } - }, - server_cert_verifier, - }) + let client_tls_options = match (conf.client_cert, conf.client_private_key) { + (None, None) => None, + (Some(client_cert), Some(client_private_key)) => Some( + temporalio_client::ClientTlsOptions::builder() + .client_cert(client_cert) + .client_private_key(client_private_key) + .build(), + ), + _ => { + return Err(PyValueError::new_err( + "Must have both client cert and private key or neither", + )) + } + }; + Ok(temporalio_client::TlsOptions::builder() + .maybe_server_root_ca_cert(server_root_ca_cert) + .maybe_domain(conf.domain) + .maybe_client_tls_options(client_tls_options) + .maybe_server_cert_verifier(server_cert_verifier) + .build()) } } @@ -430,32 +433,31 @@ impl ServerCertVerifier for FixedServerNameVerifier { impl From for RetryOptions { fn from(conf: ClientRetryConfig) -> Self { - RetryOptions { - initial_interval: Duration::from_millis(conf.initial_interval_millis), - randomization_factor: conf.randomization_factor, - multiplier: conf.multiplier, - max_interval: Duration::from_millis(conf.max_interval_millis), - max_elapsed_time: conf.max_elapsed_time_millis.map(Duration::from_millis), - max_retries: conf.max_retries, - } + RetryOptions::builder() + .initial_interval(Duration::from_millis(conf.initial_interval_millis)) + .randomization_factor(conf.randomization_factor) + .multiplier(conf.multiplier) + .max_interval(Duration::from_millis(conf.max_interval_millis)) + .max_elapsed_time(conf.max_elapsed_time_millis.map(Duration::from_millis)) + .max_retries(conf.max_retries) + .build() } } impl From for CoreClientKeepAliveConfig { fn from(conf: ClientKeepAliveConfig) -> Self { - CoreClientKeepAliveConfig { - interval: Duration::from_millis(conf.interval_millis), - timeout: Duration::from_millis(conf.timeout_millis), - } + CoreClientKeepAliveConfig::builder() + .interval(Duration::from_millis(conf.interval_millis)) + .timeout(Duration::from_millis(conf.timeout_millis)) + .build() } } impl From for HttpConnectProxyOptions { fn from(conf: ClientHttpConnectProxyConfig) -> Self { - HttpConnectProxyOptions { - target_addr: conf.target_host, - basic_auth: conf.basic_auth, - } + HttpConnectProxyOptions::new(conf.target_host) + .maybe_basic_auth(conf.basic_auth) + .build() } } diff --git a/temporalio/bridge/src/envconfig.rs b/temporalio/bridge/src/envconfig.rs index 40ef9b638..7ce2f4664 100644 --- a/temporalio/bridge/src/envconfig.rs +++ b/temporalio/bridge/src/envconfig.rs @@ -91,10 +91,10 @@ fn load_client_config_inner( config_file_strict: bool, env_vars: Option>, ) -> PyResult> { - let options = LoadClientConfigOptions { - config_source, - config_file_strict, - }; + let options = LoadClientConfigOptions::builder() + .maybe_config_source(config_source) + .config_file_strict(config_file_strict) + .build(); let core_config = core_load_client_config(options, env_vars.as_ref()) .map_err(|e| ConfigError::new_err(format!("{e}")))?; @@ -110,13 +110,13 @@ fn load_client_connect_config_inner( config_file_strict: bool, env_vars: Option>, ) -> PyResult> { - let options = LoadClientConfigProfileOptions { - config_source, - config_file_profile: profile, - config_file_strict, - disable_file, - disable_env, - }; + let options = LoadClientConfigProfileOptions::builder() + .maybe_config_source(config_source) + .maybe_config_file_profile(profile) + .config_file_strict(config_file_strict) + .disable_file(disable_file) + .disable_env(disable_env) + .build(); let profile = core_load_client_config_profile(options, env_vars.as_ref()) .map_err(|e| ConfigError::new_err(format!("{e}")))?; diff --git a/temporalio/bridge/src/worker.rs b/temporalio/bridge/src/worker.rs index d0b007dba..321dc6560 100644 --- a/temporalio/bridge/src/worker.rs +++ b/temporalio/bridge/src/worker.rs @@ -899,24 +899,25 @@ fn convert_versioning_strategy( }, WorkerVersioningStrategy::DeploymentBased(options) => { temporalio_sdk_core::WorkerVersioningStrategy::WorkerDeploymentBased( - temporalio_common::worker::WorkerDeploymentOptions { - version: temporalio_common::worker::WorkerDeploymentVersion { + temporalio_common::worker::WorkerDeploymentOptions::new( + temporalio_common::worker::WorkerDeploymentVersion { deployment_name: options.version.deployment_name, build_id: options.version.build_id, }, - use_worker_versioning: options.use_worker_versioning, - default_versioning_behavior: if options.use_worker_versioning { - Some( - temporalio_common::protos::temporal::api::enums::v1::VersioningBehavior::try_from( - options.default_versioning_behavior, - ) - .unwrap_or_default() - .into(), + ) + .use_worker_versioning(options.use_worker_versioning) + .maybe_default_versioning_behavior(if options.use_worker_versioning { + Some( + temporalio_common::protos::temporal::api::enums::v1::VersioningBehavior::try_from( + options.default_versioning_behavior, ) - } else { - None - }, - }, + .unwrap_or_default() + .into(), + ) + } else { + None + }) + .build(), ) } WorkerVersioningStrategy::LegacyBuildIdBased(lb) => { From f35f4a1ef6c293cf1032eb05a811447b3f606ada Mon Sep 17 00:00:00 2001 From: Ribhav Jain Date: Wed, 12 Aug 2026 04:27:15 +0400 Subject: [PATCH 2/2] Add workflow.uuid7() (#1733) * Add workflow.uuid7() * Suppress unreachable-code warning on pre-3.14 type checks --------- Co-authored-by: Tim Conley --- CHANGELOG.md | 8 +++ .../worker/workflow_sandbox/_restrictions.py | 4 +- temporalio/workflow/__init__.py | 2 + temporalio/workflow/_context.py | 28 ++++++++ tests/worker/test_workflow.py | 67 +++++++++++++++++++ tests/worker/workflow_sandbox/test_runner.py | 4 ++ 6 files changed, 112 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d146f045b..614fd8e01 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,6 +20,14 @@ to include examples, links to docs, or any other relevant information. ### Added +- `temporalio.workflow.uuid7()` generates a determinism-safe, time-sortable + UUIDv7 (RFC 9562) from workflow time and the workflow's deterministic random + generator, complementing the existing `workflow.uuid4()` + ([#1450](https://github.com/temporalio/sdk-python/issues/1450)). The + workflow sandbox now also restricts the non-deterministic `uuid.uuid7()` + added to the standard library in Python 3.14, matching the existing + `uuid.uuid1()`/`uuid.uuid4()` restrictions. + ### Changed - `temporalio.contrib.pydantic` converters now reuse Pydantic type adapters diff --git a/temporalio/worker/workflow_sandbox/_restrictions.py b/temporalio/worker/workflow_sandbox/_restrictions.py index 34f9efee2..23774fcd2 100644 --- a/temporalio/worker/workflow_sandbox/_restrictions.py +++ b/temporalio/worker/workflow_sandbox/_restrictions.py @@ -769,7 +769,9 @@ def _public_callables(parent: Any, *, exclude: set[str] = set()) -> set[str]: "urllib": SandboxMatcher( children={"request": SandboxMatcher.all_uses}, ), - "uuid": SandboxMatcher(use={"uuid1", "uuid4"}, only_runtime=True), + # uuid7 only exists in the stdlib on Python 3.14+; matching a + # nonexistent attribute is harmless on older versions + "uuid": SandboxMatcher(use={"uuid1", "uuid4", "uuid7"}, only_runtime=True), "webbrowser": SandboxMatcher.all_uses, "xmlrpc": SandboxMatcher.all_uses, "zipfile": SandboxMatcher( diff --git a/temporalio/workflow/__init__.py b/temporalio/workflow/__init__.py index 3d5a65c77..fa2681139 100644 --- a/temporalio/workflow/__init__.py +++ b/temporalio/workflow/__init__.py @@ -94,6 +94,7 @@ upsert_memo, upsert_search_attributes, uuid4, + uuid7, wait_condition, ) from ._definition import ( @@ -225,6 +226,7 @@ "upsert_memo", "upsert_search_attributes", "uuid4", + "uuid7", "wait_condition", "DynamicWorkflowConfig", "defn", diff --git a/temporalio/workflow/_context.py b/temporalio/workflow/_context.py index d2ec0d633..b33f83150 100644 --- a/temporalio/workflow/_context.py +++ b/temporalio/workflow/_context.py @@ -65,6 +65,7 @@ "upsert_memo", "upsert_search_attributes", "uuid4", + "uuid7", "wait_condition", ] @@ -901,6 +902,33 @@ def uuid4() -> uuid.UUID: return uuid.UUID(bytes=random().getrandbits(16 * 8).to_bytes(16, "big"), version=4) +def uuid7() -> uuid.UUID: + """Get a new, determinism-safe v7 UUID based on :py:func:`time_ns` and + :py:func:`random`. + + Per RFC 9562, the UUID's leading 48 bits are the current workflow time as + milliseconds since the epoch, so UUIDs from successive workflow tasks sort + by creation time. The remaining 74 bits are random. UUIDs generated within + the same workflow task share the same workflow time and are not guaranteed + to be monotonically ordered with respect to one another. + + Note, this UUID is not cryptographically safe and should not be used for + security purposes. + + Returns: + A deterministically-seeded v7 UUID. + """ + # uuid.UUID's version parameter only accepts 1-5 before Python 3.14, so + # the version and variant bits are set manually. + unix_ts_ms = (time_ns() // 1_000_000) & 0xFFFF_FFFF_FFFF + rand = random() + rand_a = rand.getrandbits(12) + rand_b = rand.getrandbits(62) + return uuid.UUID( + int=(unix_ts_ms << 80) | (0x7 << 76) | (rand_a << 64) | (0b10 << 62) | rand_b + ) + + async def sleep(duration: float | timedelta, *, summary: str | None = None) -> None: """Sleep for the given duration. diff --git a/tests/worker/test_workflow.py b/tests/worker/test_workflow.py index 283357c2d..39e648e60 100644 --- a/tests/worker/test_workflow.py +++ b/tests/worker/test_workflow.py @@ -3934,6 +3934,73 @@ async def test_workflow_uuid(client: Client): assert handle2_query_result == await handle2.query(UUIDWorkflow.result) +@workflow.defn +class UUID7Workflow: + def __init__(self) -> None: + self._result = "" + self._time_ms = -1 + + @workflow.run + async def run(self) -> None: + self._time_ms = workflow.time_ns() // 1_000_000 + self._result = str(workflow.uuid7()) + + @workflow.query + def result(self) -> str: + return self._result + + @workflow.query + def time_ms(self) -> int: + return self._time_ms + + +async def test_workflow_uuid7(client: Client): + task_queue = str(uuid.uuid4()) + async with new_worker( + client, UUID7Workflow, task_queue=task_queue, max_cached_workflows=0 + ): + # Get two handle UUID results. Need to disable workflow cache since we + # restart the worker and don't want to pay the sticky queue penalty. + handle1 = await client.start_workflow( + UUID7Workflow.run, id=f"workflow-{uuid.uuid4()}", task_queue=task_queue + ) + await handle1.result() + handle1_query_result = await handle1.query(UUID7Workflow.result) + + handle2 = await client.start_workflow( + UUID7Workflow.run, + id=f"workflow-{uuid.uuid4()}", + task_queue=task_queue, + ) + await handle2.result() + handle2_query_result = await handle2.query(UUID7Workflow.result) + + # Confirm they aren't equal to each other but they are equal to retries + # of the same query + assert handle1_query_result != handle2_query_result + assert handle1_query_result == await handle1.query(UUID7Workflow.result) + assert handle2_query_result == await handle2.query(UUID7Workflow.result) + + # Confirm RFC 9562 shape: version 7, RFC variant, and the leading 48 + # bits are the workflow time in milliseconds at generation + for handle, query_result in ( + (handle1, handle1_query_result), + (handle2, handle2_query_result), + ): + result_uuid = uuid.UUID(query_result) + assert result_uuid.version == 7 + assert result_uuid.variant == uuid.RFC_4122 + workflow_time_ms = await handle.query(UUID7Workflow.time_ms) + assert int(result_uuid) >> 80 == workflow_time_ms + + # Now confirm those results are the same even on a new worker + async with new_worker( + client, UUID7Workflow, task_queue=task_queue, max_cached_workflows=0 + ): + assert handle1_query_result == await handle1.query(UUID7Workflow.result) + assert handle2_query_result == await handle2.query(UUID7Workflow.result) + + @activity.defn(name="custom-name") class CallableClassActivity: def __init__(self, orig_field1: str) -> None: diff --git a/tests/worker/workflow_sandbox/test_runner.py b/tests/worker/workflow_sandbox/test_runner.py index 288da0861..024ab7273 100644 --- a/tests/worker/workflow_sandbox/test_runner.py +++ b/tests/worker/workflow_sandbox/test_runner.py @@ -197,6 +197,10 @@ async def test_workflow_sandbox_restrictions(client: Client): if sys.version_info < (3, 14): invalid_code_to_check.append("import os.path\nos.path.abspath('foo')") # type: ignore[reportUnreachable] + # uuid7 was only added to the stdlib in 3.14 + if sys.version_info >= (3, 14): + invalid_code_to_check.append("import uuid\nuuid.uuid7()") # type: ignore[reportUnreachable] + for code in invalid_code_to_check: with pytest.raises(WorkflowFailureError) as err: await client.execute_workflow(