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
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 53 files
+10 −0 CHANGELOG.md
+312 −217 crates/client/src/async_activity_handle.rs
+10 −0 crates/client/src/dns.rs
+3 −0 crates/client/src/errors.rs
+1,520 −0 crates/client/src/interceptors.rs
+954 −156 crates/client/src/lib.rs
+59 −6 crates/client/src/options_structs.rs
+239 −0 crates/client/src/rpc_options.rs
+744 −255 crates/client/src/schedules.rs
+35 −5 crates/client/src/test_helpers.rs
+130 −2 crates/client/src/worker.rs
+465 −235 crates/client/src/workflow_handle.rs
+15 −12 crates/common-wasm/src/data_converters.rs
+11 −10 crates/common/build.rs
+156 −66 crates/common/src/payload_visitor.rs
+13 −0 crates/protos/protos/local/temporal/sdk/core/common/common.proto
+6 −0 crates/protos/protos/local/temporal/sdk/core/workflow_completion/workflow_completion.proto
+5 −0 crates/protos/src/protos/mod.rs
+70 −24 crates/sdk-core-c-bridge/src/worker.rs
+1 −1 crates/sdk-core/src/core_tests/mod.rs
+65 −5 crates/sdk-core/src/core_tests/workers.rs
+3 −3 crates/sdk-core/src/core_tests/workflow_tasks.rs
+9 −11 crates/sdk-core/src/pollers/poll_buffer.rs
+30 −2 crates/sdk-core/src/replay/mod.rs
+12 −2 crates/sdk-core/src/test_help/integ_helpers.rs
+1 −0 crates/sdk-core/src/test_help/unit_helpers.rs
+5 −8 crates/sdk-core/src/worker/client/mocks.rs
+67 −4 crates/sdk-core/src/worker/heartbeat.rs
+525 −301 crates/sdk-core/src/worker/mod.rs
+8 −7 crates/sdk-core/src/worker/tuner/resource_based.rs
+4 −0 crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs
+198 −4 crates/sdk-core/src/worker/workflow/managed_run.rs
+16 −0 crates/sdk-core/src/worker/workflow/mod.rs
+6 −3 crates/sdk-core/src/worker/workflow/wft_poller.rs
+7 −2 crates/sdk-core/src/worker/workflow/workflow_stream.rs
+10 −8 crates/sdk-core/tests/heavy_tests.rs
+11 −2 crates/sdk-core/tests/integ_tests/async_activity_client_tests.rs
+166 −7 crates/sdk-core/tests/integ_tests/data_converter_tests.rs
+2 −2 crates/sdk-core/tests/integ_tests/metrics_tests.rs
+4 −4 crates/sdk-core/tests/integ_tests/polling_tests.rs
+62 −37 crates/sdk-core/tests/integ_tests/schedule_tests.rs
+4 −4 crates/sdk-core/tests/integ_tests/worker_heartbeat_tests.rs
+3 −1 crates/sdk-core/tests/integ_tests/worker_tests.rs
+2 −0 crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs
+74 −4 crates/sdk-core/tests/integ_tests/workflow_client_tests.rs
+2 −1 crates/sdk-core/tests/integ_tests/workflow_tests.rs
+8 −5 crates/sdk-core/tests/integ_tests/workflow_tests/activities.rs
+94 −4 crates/sdk-core/tests/integ_tests/workflow_tests/client_interactions.rs
+1 −1 crates/sdk-core/tests/integ_tests/workflow_tests/stickyness.rs
+10 −10 crates/sdk-core/tests/manual_tests.rs
+5 −3 crates/sdk/examples/schedules/starter.rs
+211 −32 crates/sdk/src/lib.rs
+1 −0 crates/sdk/src/workflow_future.rs
105 changes: 24 additions & 81 deletions temporalio/worker/_workflow.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,13 +9,13 @@
import os
import sys
import threading
import time
from collections.abc import Awaitable, Callable, MutableMapping, Sequence
from dataclasses import dataclass
from datetime import timedelta, timezone
from datetime import timezone
from types import TracebackType

import temporalio.api.common.v1
import temporalio.bridge.proto.common
import temporalio.bridge.proto.workflow_activation
import temporalio.bridge.proto.workflow_completion
import temporalio.bridge.runtime
Expand Down Expand Up @@ -64,6 +64,17 @@
_DEFAULT_WORKFLOW_TASK_EXTERNAL_STORAGE_CONCURRENCY: int = 3


def _set_external_storage_metrics(
target: temporalio.bridge.proto.common.ExternalStorageMetrics,
metrics: temporalio.converter._extstore.StorageOperationMetrics,
) -> None:
"""Populate a proto ``ExternalStorageMetrics`` from measured storage metrics."""
target.payload_count = metrics.payload_count
target.total_size_bytes = metrics.total_size
target.total_duration.FromTimedelta(metrics.total_duration)
target.driver_names.extend(sorted(metrics.driver_names))


class _WorkflowWorker: # type:ignore[reportUnusedClass]
def __init__(
self,
Expand Down Expand Up @@ -325,7 +336,6 @@ async def _handle_activation(
completion.successful.SetInParent()
workflow = None
data_converter = self._data_converter
task_start_time = time.monotonic()
download_metrics = temporalio.converter._extstore.StorageOperationMetrics()
try:
if LOG_PROTOS:
Expand Down Expand Up @@ -500,6 +510,17 @@ async def _handle_activation(
completion.failed.Clear()
completion.failed.failure.message = f"Failed encoding completion: {err}"

# Reported on the completion so core can include them in its workflow-task duration
# log; core measures the duration itself.
if download_metrics.payload_count > 0:
_set_external_storage_metrics(
completion.payload_download_metrics, download_metrics
)
if upload_metrics.payload_count > 0:
_set_external_storage_metrics(
completion.payload_upload_metrics, upload_metrics
)

# Send off completion
if LOG_PROTOS:
logger.debug("Sending workflow completion:\n%s", completion)
Expand All @@ -511,84 +532,6 @@ async def _handle_activation(
"Failed completing activation on workflow with run ID %s", act.run_id
)

# Log workflow task duration with external storage metrics
self._log_workflow_task_duration(
act, workflow, task_start_time, download_metrics, upload_metrics
)

def _log_workflow_task_duration(
self,
act: temporalio.bridge.proto.workflow_activation.WorkflowActivation,
workflow: _RunningWorkflow | None,
task_start_time: float,
download_metrics: temporalio.converter._extstore.StorageOperationMetrics,
upload_metrics: temporalio.converter._extstore.StorageOperationMetrics,
) -> None:
task_duration = timedelta(seconds=time.monotonic() - task_start_time)

def _fmt_duration(td: timedelta) -> str:
secs = td.total_seconds()
if secs >= 1:
return f"{secs:.3f}s"
return f"{secs * 1000:.3f}ms"

completed_event_id = act.history_length + 1
_info = workflow.get_info() if workflow is not None else None
attempt = _info.attempt if _info is not None else "unknown"
log_id = f"{act.run_id}:{completed_event_id}:{attempt}"
msg_details, extra = temporalio.workflow._build_log_context(
_info._logger_details() if _info is not None else None,
full_workflow_info=_info,
)
msg_details["event_id"] = completed_event_id
msg_details["workflow_task_duration"] = _fmt_duration(task_duration)
msg_details["workflow_history_size"] = act.history_size_bytes
extra["event_id"] = completed_event_id
extra["workflow_task_duration"] = task_duration
extra["workflow_history_size"] = act.history_size_bytes
if download_metrics.payload_count > 0:
msg_details["payload_download_count"] = download_metrics.payload_count
msg_details["payload_download_size"] = download_metrics.total_size
msg_details["payload_download_duration"] = _fmt_duration(
download_metrics.total_duration
)
msg_details["payload_download_drivers"] = sorted(
download_metrics.driver_names
)
extra["payload_download_count"] = download_metrics.payload_count
extra["payload_download_size"] = download_metrics.total_size
extra["payload_download_duration"] = download_metrics.total_duration
extra["payload_download_drivers"] = sorted(download_metrics.driver_names)
if upload_metrics.payload_count > 0:
msg_details["payload_upload_count"] = upload_metrics.payload_count
msg_details["payload_upload_size"] = upload_metrics.total_size
msg_details["payload_upload_duration"] = _fmt_duration(
upload_metrics.total_duration
)
msg_details["payload_upload_drivers"] = sorted(upload_metrics.driver_names)
extra["payload_upload_count"] = upload_metrics.payload_count
extra["payload_upload_size"] = upload_metrics.total_size
extra["payload_upload_duration"] = upload_metrics.total_duration
extra["payload_upload_drivers"] = sorted(upload_metrics.driver_names)
if task_duration.total_seconds() > 10:
logger.warning(
f"[TMPRL1104] {log_id} Workflow task exceeded 10 seconds (%s)",
msg_details,
extra=extra,
)
elif task_duration.total_seconds() > 5:
logger.info(
f"[TMPRL1104] {log_id} Workflow task exceeded 5 seconds (%s)",
msg_details,
extra=extra,
)
else:
logger.debug(
f"[TMPRL1104] {log_id} Workflow task duration information (%s)",
msg_details,
extra=extra,
)

async def _handle_cache_eviction(
self,
act: temporalio.bridge.proto.workflow_activation.WorkflowActivation,
Expand Down
Loading
Loading