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
214 changes: 11 additions & 203 deletions temporalio/bridge/Cargo.lock

Large diffs are not rendered by default.

6 changes: 3 additions & 3 deletions temporalio/bridge/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
2 changes: 2 additions & 0 deletions temporalio/bridge/proto/common/__init__.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
from .common_pb2 import (
ExternalStorageMetrics,
NamespacedWorkflowExecution,
VersioningIntent,
WorkerDeploymentVersion,
)

__all__ = [
"ExternalStorageMetrics",
"NamespacedWorkflowExecution",
"VersioningIntent",
"WorkerDeploymentVersion",
Expand Down
20 changes: 17 additions & 3 deletions temporalio/bridge/proto/common/common_pb2.py

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

53 changes: 53 additions & 0 deletions temporalio/bridge/proto/common/common_pb2.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand All @@ -31,30 +32,63 @@ 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(
self,
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",
Expand Down
2 changes: 1 addition & 1 deletion temporalio/bridge/sdk-core
Submodule sdk-core updated 121 files
80 changes: 41 additions & 39 deletions temporalio/bridge/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -325,25 +327,26 @@ impl TryFrom<ClientTlsConfig> 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())
}
}

Expand Down Expand Up @@ -430,32 +433,31 @@ impl ServerCertVerifier for FixedServerNameVerifier {

impl From<ClientRetryConfig> 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<ClientKeepAliveConfig> 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<ClientHttpConnectProxyConfig> 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()
}
}

Expand Down
22 changes: 11 additions & 11 deletions temporalio/bridge/src/envconfig.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,10 +91,10 @@ fn load_client_config_inner(
config_file_strict: bool,
env_vars: Option<HashMap<String, String>>,
) -> PyResult<Py<PyAny>> {
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}")))?;

Expand All @@ -110,13 +110,13 @@ fn load_client_connect_config_inner(
config_file_strict: bool,
env_vars: Option<HashMap<String, String>>,
) -> PyResult<Py<PyAny>> {
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}")))?;
Expand Down
Loading
Loading