From ce69d10f0e80ec154264c3a7ed395af1e18aa796 Mon Sep 17 00:00:00 2001 From: jmaeagle99 <44687433+jmaeagle99@users.noreply.github.com> Date: Tue, 28 Jul 2026 09:20:50 -0700 Subject: [PATCH 1/3] feat(core): emit workflow task duration logging --- crates/common/src/payload_visitor.rs | 2 + .../temporal/sdk/core/common/common.proto | 13 ++ .../workflow_completion.proto | 6 + crates/protos/src/protos/mod.rs | 5 + crates/sdk-core/src/test_help/unit_helpers.rs | 1 + .../workflow/machines/workflow_machines.rs | 4 + .../src/worker/workflow/managed_run.rs | 202 +++++++++++++++++- crates/sdk-core/src/worker/workflow/mod.rs | 16 ++ .../src/worker/workflow/workflow_stream.rs | 9 +- .../tests/integ_tests/worker_tests.rs | 2 + .../integ_tests/worker_versioning_tests.rs | 2 + crates/sdk/src/workflow_future.rs | 1 + 12 files changed, 257 insertions(+), 6 deletions(-) diff --git a/crates/common/src/payload_visitor.rs b/crates/common/src/payload_visitor.rs index 8fd647dbc..317d01c50 100644 --- a/crates/common/src/payload_visitor.rs +++ b/crates/common/src/payload_visitor.rs @@ -432,6 +432,7 @@ mod tests { ..Default::default() }, )), + ..Default::default() }; encode_payloads( @@ -613,6 +614,7 @@ mod tests { ..Default::default() }, )), + ..Default::default() }; encode_payloads( diff --git a/crates/protos/protos/local/temporal/sdk/core/common/common.proto b/crates/protos/protos/local/temporal/sdk/core/common/common.proto index 302050cb0..0833e4b93 100644 --- a/crates/protos/protos/local/temporal/sdk/core/common/common.proto +++ b/crates/protos/protos/local/temporal/sdk/core/common/common.proto @@ -34,4 +34,17 @@ enum VersioningIntent { message WorkerDeploymentVersion { string deployment_name = 1; string build_id = 2; +} + +// 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. +message ExternalStorageMetrics { + // Number of payloads stored or retrieved externally. + uint64 payload_count = 1; + // Total size in bytes of the externally stored or retrieved payloads. + uint64 total_size_bytes = 2; + // Wall-clock time spent on the external storage operations. + google.protobuf.Duration total_duration = 3; + // Names of the drivers that participated in the operations. + repeated string driver_names = 4; } \ No newline at end of file diff --git a/crates/protos/protos/local/temporal/sdk/core/workflow_completion/workflow_completion.proto b/crates/protos/protos/local/temporal/sdk/core/workflow_completion/workflow_completion.proto index 78637f640..afe7d6442 100644 --- a/crates/protos/protos/local/temporal/sdk/core/workflow_completion/workflow_completion.proto +++ b/crates/protos/protos/local/temporal/sdk/core/workflow_completion/workflow_completion.proto @@ -17,6 +17,12 @@ message WorkflowActivationCompletion { Success successful = 2; Failure failed = 3; } + // Metrics for external payload storage downloads (retrievals) performed while processing + // this activation. Only set when external storage retrieved payloads. + common.ExternalStorageMetrics payload_download_metrics = 4; + // Metrics for external payload storage uploads (stores) performed while processing this + // activation. Only set when external storage stored payloads. + common.ExternalStorageMetrics payload_upload_metrics = 5; } // Successful workflow activation with a list of commands generated by the workflow execution diff --git a/crates/protos/src/protos/mod.rs b/crates/protos/src/protos/mod.rs index ab89d5be4..c7dbfb303 100644 --- a/crates/protos/src/protos/mod.rs +++ b/crates/protos/src/protos/mod.rs @@ -108,6 +108,7 @@ pub mod coresdk { Self { run_id: run_id.into(), status: Some(workflow_activation_completion::Status::Successful(success)), + ..Default::default() } } @@ -120,6 +121,7 @@ pub mod coresdk { Self { run_id: run_id.into(), status: Some(workflow_activation_completion::Status::Successful(success)), + ..Default::default() } } @@ -129,6 +131,7 @@ pub mod coresdk { Self { run_id: run_id.into(), status: Some(workflow_activation_completion::Status::Successful(success)), + ..Default::default() } } @@ -146,6 +149,7 @@ pub mod coresdk { as i32, }, )), + ..Default::default() } } @@ -269,6 +273,7 @@ pub mod coresdk { WorkflowActivationCompletion { run_id, status: Some(workflow_activation_completion::Status::Successful(success)), + ..Default::default() } } } diff --git a/crates/sdk-core/src/test_help/unit_helpers.rs b/crates/sdk-core/src/test_help/unit_helpers.rs index 00767f9a9..4d3220c13 100644 --- a/crates/sdk-core/src/test_help/unit_helpers.rs +++ b/crates/sdk-core/src/test_help/unit_helpers.rs @@ -135,6 +135,7 @@ pub(crate) async fn poll_and_reply_clears_outstanding_evicts<'a>( WorkflowActivationCompletion { run_id: res.run_id.clone(), status: Some(reply.clone()), + ..Default::default() } }; diff --git a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs index e0f13a23e..ffb790ead 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -406,6 +406,10 @@ impl WorkflowMachines { self.current_started_event_id } + pub(crate) fn history_size_bytes(&self) -> u64 { + self.history_size_bytes + } + pub(crate) fn prepare_for_wft_response(&mut self) -> MachinesWFTResponseContent<'_> { MachinesWFTResponseContent { replaying: self.replaying, diff --git a/crates/sdk-core/src/worker/workflow/managed_run.rs b/crates/sdk-core/src/worker/workflow/managed_run.rs index 545e0235d..a2001e4da 100644 --- a/crates/sdk-core/src/worker/workflow/managed_run.rs +++ b/crates/sdk-core/src/worker/workflow/managed_run.rs @@ -12,8 +12,8 @@ use crate::{ FailedActivationWFTReport, HeartbeatTimeoutMsg, HistoryUpdate, LocalActivityRequestSink, LocalResolution, NextPageReq, OutstandingActivation, OutstandingTask, PermittedWFT, RequestEvictMsg, RunBasics, - ServerCommandsWithWorkflowInfo, WFCommand, WFCommandVariant, WFMachinesError, - WFT_HEARTBEAT_TIMEOUT_FRACTION, WFTReportStatus, WorkflowTaskInfo, + ServerCommandsWithWorkflowInfo, TaskStorageMetrics, WFCommand, WFCommandVariant, + WFMachinesError, WFT_HEARTBEAT_TIMEOUT_FRACTION, WFTReportStatus, WorkflowTaskInfo, history_update::HistoryPaginator, machines::{MachinesWFTResponseContent, WorkflowMachines}, }, @@ -31,6 +31,7 @@ use std::{ use temporalio_common::protos::{ TaskToken, coresdk::{ + common::ExternalStorageMetrics, workflow_activation::{ WorkflowActivation, create_evict_activation, query_to_job, remove_from_cache::EvictionReason, workflow_activation_job, @@ -292,6 +293,7 @@ impl ManagedRun { pub(super) fn mark_wft_complete( &mut self, report_status: WFTReportStatus, + task_storage_metrics: &TaskStorageMetrics, ) -> Option { debug!("Marking WFT completed"); let retme = self.wft.take(); @@ -299,7 +301,17 @@ impl ManagedRun { if let Some(ot) = &retme && let Some(ct) = report_status.completion_time() { - self.metrics.wf_task_latency(ct.sub(ot.start_time)); + let task_duration = ct.sub(ot.start_time); + self.metrics.wf_task_latency(task_duration); + log_workflow_task_duration( + &self.wfm.machines.run_id, + &self.wfm.machines.workflow_type, + self.wfm.machines.last_processed_event + 1, + ot.info.attempt, + self.wfm.machines.history_size_bytes(), + task_duration, + task_storage_metrics, + ); } if let WFTReportStatus::Reported { @@ -1529,13 +1541,195 @@ impl From for RunUpdateErr { } } +#[allow(clippy::too_many_arguments)] +fn log_workflow_task_duration( + run_id: &str, + workflow_type: &str, + event_id: i64, + attempt: u32, + history_size_bytes: u64, + duration: Duration, + storage: &TaskStorageMetrics, +) { + let dl = storage.download.as_ref(); + let ul = storage.upload.as_ref(); + let duration_millis = |d: Duration| d.as_millis() as u64; + let storage_millis = |m: Option<&ExternalStorageMetrics>| -> u64 { + m.and_then(|m| m.total_duration) + .and_then(|d| Duration::try_from(d).ok()) + .map(duration_millis) + .unwrap_or_default() + }; + + macro_rules! emit { + ($lvl:ident, $msg:literal) => { + $lvl!( + workflow_type = %workflow_type, + event_id = event_id, + attempt = attempt, + workflow_task_duration = duration_millis(duration), + workflow_history_size = history_size_bytes, + payload_download_count = dl.map(|m| m.payload_count).unwrap_or_default(), + payload_download_size = dl.map(|m| m.total_size_bytes).unwrap_or_default(), + payload_download_duration = storage_millis(dl), + payload_download_drivers = ?dl.map(|m| sorted(&m.driver_names)).unwrap_or_default(), + payload_upload_count = ul.map(|m| m.payload_count).unwrap_or_default(), + payload_upload_size = ul.map(|m| m.total_size_bytes).unwrap_or_default(), + payload_upload_duration = storage_millis(ul), + payload_upload_drivers = ?ul.map(|m| sorted(&m.driver_names)).unwrap_or_default(), + $msg, + ) + }; + } + + if duration > Duration::from_secs(5) { + emit!( + warn, + "[TMPRL1104] {run_id}:{event_id}:{attempt} Workflow task duration exceeded 5 seconds." + ); + } else { + emit!(trace, "Workflow task duration information."); + } +} + +fn sorted(names: &[String]) -> Vec { + let mut v = names.to_vec(); + v.sort(); + v +} + #[cfg(test)] mod tests { + use super::{TaskStorageMetrics, log_workflow_task_duration}; use crate::worker::workflow::{WFCommand, WFCommandVariant}; - use std::mem::{Discriminant, discriminant}; + use std::{ + fmt::Write, + mem::{Discriminant, discriminant}, + sync::{Arc, Mutex}, + time::Duration, + }; + use temporalio_common::protos::coresdk::common::ExternalStorageMetrics; + use tracing::{Event, Level, Metadata, Subscriber, field::Field, field::Visit, span}; use command_utils::*; + #[derive(Default)] + struct CapturedEvent { + level: Option, + fields: String, + } + #[derive(Default, Clone)] + struct CapturingSub { + events: Arc>>, + } + struct FieldVisitor(String); + impl Visit for FieldVisitor { + fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) { + let _ = write!(self.0, "{}={:?};", field.name(), value); + } + fn record_u64(&mut self, field: &Field, value: u64) { + let _ = write!(self.0, "{}={};", field.name(), value); + } + fn record_i64(&mut self, field: &Field, value: i64) { + let _ = write!(self.0, "{}={};", field.name(), value); + } + fn record_str(&mut self, field: &Field, value: &str) { + let _ = write!(self.0, "{}={};", field.name(), value); + } + } + impl Subscriber for CapturingSub { + fn enabled(&self, _: &Metadata<'_>) -> bool { + true + } + fn new_span(&self, _: &span::Attributes<'_>) -> span::Id { + span::Id::from_u64(1) + } + fn record(&self, _: &span::Id, _: &span::Record<'_>) {} + fn record_follows_from(&self, _: &span::Id, _: &span::Id) {} + fn event(&self, event: &Event<'_>) { + let mut v = FieldVisitor(String::new()); + event.record(&mut v); + self.events.lock().unwrap().push(CapturedEvent { + level: Some(*event.metadata().level()), + fields: v.0, + }); + } + fn enter(&self, _: &span::Id) {} + fn exit(&self, _: &span::Id) {} + } + + fn capture(duration: Duration, storage: &TaskStorageMetrics) -> CapturedEvent { + let sub = CapturingSub::default(); + tracing::subscriber::with_default(sub.clone(), || { + log_workflow_task_duration("run-1", "MyWorkflow", 12, 3, 4096, duration, storage); + }); + sub.events + .lock() + .unwrap() + .drain(..) + .next() + .expect("one event emitted") + } + + #[test] + fn tmprl1104_levels_follow_thresholds() { + let none = TaskStorageMetrics::default(); + let trace_ev = capture(Duration::from_secs(2), &none); + assert_eq!(trace_ev.level, Some(Level::TRACE)); + assert!( + !trace_ev.fields.contains("TMPRL1104"), + "fields: {}", + trace_ev.fields + ); + let warn_ev = capture(Duration::from_secs(7), &none); + assert_eq!(warn_ev.level, Some(Level::WARN)); + assert!( + warn_ev.fields.contains("[TMPRL1104]"), + "fields: {}", + warn_ev.fields + ); + } + + #[test] + fn tmprl1104_fields_present() { + let storage = TaskStorageMetrics { + download: Some(ExternalStorageMetrics { + payload_count: 2, + total_size_bytes: 1024, + total_duration: Some(prost_types::Duration { + seconds: 0, + nanos: 5_000_000, + }), + driver_names: vec!["s3".to_string()], + }), + upload: None, + }; + // Only the warn-level message carries the identifier, so exceed the 5s threshold. + let ev = capture(Duration::from_secs(6), &storage); + assert!(ev.fields.contains("attempt=3"), "fields: {}", ev.fields); + assert!( + ev.fields.contains("[TMPRL1104] run-1:12:3"), + "fields: {}", + ev.fields + ); + assert!( + ev.fields.contains("workflow_history_size=4096"), + "fields: {}", + ev.fields + ); + assert!( + ev.fields.contains("payload_download_count=2"), + "fields: {}", + ev.fields + ); + // No upload occurred; that group must still be present as zero. + assert!( + ev.fields.contains("payload_upload_count=0"), + "fields: {}", + ev.fields + ); + } + #[rstest::rstest] #[case::empty( vec![], diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index b9366a1c2..4b3340fd3 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -64,6 +64,7 @@ use temporalio_common::{ protos::{ TaskToken, coresdk::{ + common::ExternalStorageMetrics, workflow_activation::{ QueryWorkflow, WorkflowActivation, WorkflowActivationJob, remove_from_cache::EvictionReason, workflow_activation_job, @@ -302,6 +303,7 @@ impl Workflows { status: Some( workflow_completion::Success::from_variants(vec![]).into(), ), + ..Default::default() }, true, // We need to say a type, but the type is irrelevant, so imagine some @@ -322,6 +324,7 @@ impl Workflows { WorkflowActivationCompletion { run_id, status: Some(machines_err.as_failure().into()), + ..Default::default() }, true, Option::>::None, @@ -546,6 +549,10 @@ impl Workflows { post_activate_hook: Option, ) -> Result<(), CompleteWfError> { let is_empty_completion = completion.is_empty(); + let task_storage_metrics = TaskStorageMetrics { + download: completion.payload_download_metrics.clone(), + upload: completion.payload_upload_metrics.clone(), + }; let completion = validate_completion(completion, is_autocomplete)?; let run_id = completion.run_id().to_string(); let (tx, rx) = oneshot::channel(); @@ -615,6 +622,7 @@ impl Workflows { wft_report_status, wft_from_complete: maybe_pwft, is_autocomplete, + task_storage_metrics, }); Ok(()) @@ -1131,7 +1139,15 @@ struct PostActivationMsg { wft_report_status: WFTReportStatus, wft_from_complete: Option, is_autocomplete: bool, + task_storage_metrics: TaskStorageMetrics, +} + +#[derive(Debug, Default, Clone)] +struct TaskStorageMetrics { + download: Option, + upload: Option, } + #[derive(Debug, Clone)] struct RequestEvictMsg { run_id: String, diff --git a/crates/sdk-core/src/worker/workflow/workflow_stream.rs b/crates/sdk-core/src/worker/workflow/workflow_stream.rs index 0c86f3fe5..ad9bb0b48 100644 --- a/crates/sdk-core/src/worker/workflow/workflow_stream.rs +++ b/crates/sdk-core/src/worker/workflow/workflow_stream.rs @@ -340,7 +340,11 @@ impl WFStream { let mut res = None; - let maybe_t = self.complete_wft(run_id, report.wft_report_status); + let maybe_t = self.complete_wft( + run_id, + report.wft_report_status, + &report.task_storage_metrics, + ); // Augment the WFT from complete with the permit if both exist let wft_from_complete = wft_from_complete.and_then(|wft| { maybe_t.map(|t| PermittedWFT { @@ -467,6 +471,7 @@ impl WFStream { &mut self, run_id: &str, wft_report_status: WFTReportStatus, + task_storage_metrics: &TaskStorageMetrics, ) -> Option { // If the WFT completion wasn't sent to the server, but we did see the final event, we still // want to clear the workflow task. This can really only happen in replay testing, where we @@ -493,7 +498,7 @@ impl WFStream { return None; } - rh.mark_wft_complete(wft_report_status) + rh.mark_wft_complete(wft_report_status, task_storage_metrics) } else { None } diff --git a/crates/sdk-core/tests/integ_tests/worker_tests.rs b/crates/sdk-core/tests/integ_tests/worker_tests.rs index b26dd9a5f..386e67ef7 100644 --- a/crates/sdk-core/tests/integ_tests/worker_tests.rs +++ b/crates/sdk-core/tests/integ_tests/worker_tests.rs @@ -148,6 +148,7 @@ async fn worker_handles_unknown_workflow_types_gracefully() { .. })), run_id, + .. } if message == "Workflow type unregistered not found" && *run_id == self.run_id ) { self.unregistered_failure_seen.set(true); @@ -158,6 +159,7 @@ async fn worker_handles_unknown_workflow_types_gracefully() { WorkflowActivationCompletion { status: Some(Status::Successful(..)), run_id, + .. } if self.unregistered_failure_seen.get() && *run_id == self.run_id ) { // Shutdown the worker diff --git a/crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs b/crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs index 2f99ce51c..0c2e10426 100644 --- a/crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs +++ b/crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs @@ -69,6 +69,7 @@ async fn sets_deployment_info_on_task_responses(#[values(true, false)] use_defau core.complete_workflow_activation(WorkflowActivationCompletion { run_id: res.run_id.clone(), status: Some(success_complete.into()), + ..Default::default() }) .await .unwrap(); @@ -290,6 +291,7 @@ async fn versioning_off_with_custom_build_id() { ]) .into(), ), + ..Default::default() }) .await .unwrap(); diff --git a/crates/sdk/src/workflow_future.rs b/crates/sdk/src/workflow_future.rs index 21011031a..4be240f78 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -194,6 +194,7 @@ impl WorkflowFuture { versioning_behavior: VersioningBehavior::Unspecified.into(), }, )), + ..Default::default() }) .expect("Completion channel intact"); } From dbbe5a4e8fdac92d412a8c22a145375ae92c2edf Mon Sep 17 00:00:00 2001 From: jmaeagle99 <44687433+jmaeagle99@users.noreply.github.com> Date: Thu, 30 Jul 2026 09:28:40 -0700 Subject: [PATCH 2/3] test: verify storage metrics reach duration log --- .../sdk-core/src/core_tests/workflow_tasks.rs | 98 +++++++++++++++++++ .../sdk-core/src/test_help/integ_helpers.rs | 5 + crates/sdk-core/src/worker/mod.rs | 8 ++ 3 files changed, 111 insertions(+) diff --git a/crates/sdk-core/src/core_tests/workflow_tasks.rs b/crates/sdk-core/src/core_tests/workflow_tasks.rs index 4eebc17df..ef2283afe 100644 --- a/crates/sdk-core/src/core_tests/workflow_tasks.rs +++ b/crates/sdk-core/src/core_tests/workflow_tasks.rs @@ -1097,6 +1097,104 @@ async fn complete_after_eviction() { core.shutdown().await; } +/// The external-storage metrics a lang SDK reports on a workflow activation completion must +/// reach core's workflow-task duration log. This exercises the full extraction/threading path +/// (completion proto -> `TaskStorageMetrics` -> `mark_wft_complete`), which the direct-call unit +/// tests of `log_workflow_task_duration` deliberately bypass. +#[tokio::test] +async fn tmprl1104_completion_storage_metrics_reach_log() { + use std::sync::Mutex; + use temporalio_common::{ + protos::coresdk::common::ExternalStorageMetrics, + telemetry::{ + CoreLog, CoreLogConsumer, Logger, TelemetryOptions, construct_filter_string, + telemetry_init, + }, + }; + use tracing::Level; + + #[derive(Debug)] + struct CaptureConsumer(Mutex>); + impl CoreLogConsumer for CaptureConsumer { + fn on_log(&self, log: CoreLog) { + self.0.lock().unwrap().push(log); + } + } + + let consumer = Arc::new(CaptureConsumer(Mutex::new(vec![]))); + let instance = telemetry_init( + TelemetryOptions::builder() + .logging(Logger::Push { + // The routine duration line is emitted at TRACE (only the >5s line is WARN). + filter: construct_filter_string(Level::TRACE, Level::WARN), + consumer: consumer.clone(), + }) + .build(), + ) + .unwrap(); + let sub = instance.trace_subscriber().expect("subscriber present"); + + let wfid = "fake_wf_id"; + // A workflow that completes in a single WFT, so that lone WFT is the reported one that + // reaches `mark_wft_complete` and carries the metrics we attach below. + let mut t = TestHistoryBuilder::default(); + t.add_by_type(EventType::WorkflowExecutionStarted); + t.add_workflow_task_scheduled_and_started(); + let mock = mock_worker_client(); + let mut mock = single_hist_mock_sg(wfid, t, [1], mock, true); + mock.set_trace_subscriber(sub); + let core = mock_worker(mock); + + // Attach download/upload metrics to the completion, as a lang SDK would. + let activation = core.poll_workflow_activation().await.unwrap(); + let mut completion = WorkflowActivationCompletion::from_cmds( + activation.run_id, + vec![CompleteWorkflowExecution { result: None }.into()], + ); + completion.payload_download_metrics = Some(ExternalStorageMetrics { + payload_count: 2, + total_size_bytes: 1024, + total_duration: Some(prost_types::Duration { + seconds: 0, + nanos: 5_000_000, + }), + driver_names: vec!["s3".to_string()], + }); + completion.payload_upload_metrics = Some(ExternalStorageMetrics { + payload_count: 3, + total_size_bytes: 2048, + total_duration: Some(prost_types::Duration { + seconds: 1, + nanos: 0, + }), + driver_names: vec!["gcs".to_string()], + }); + core.complete_workflow_activation(completion).await.unwrap(); + core.shutdown().await; + + // The lone WFT's duration log must carry the completion's metrics, keeping download and + // upload distinct (a swap or dropped field would fail here). + let logs = consumer.0.lock().unwrap(); + let ev = logs + .iter() + .find(|l| l.message == "Workflow task duration information.") + .expect("workflow task duration log"); + let u64_field = |name: &str| ev.fields.get(name).and_then(|v| v.as_u64()); + let str_field = |name: &str| { + ev.fields + .get(name) + .and_then(|v| v.as_str()) + .unwrap_or_default() + .to_string() + }; + assert_eq!(u64_field("payload_download_count"), Some(2)); + assert_eq!(u64_field("payload_download_size"), Some(1024)); + assert_eq!(u64_field("payload_upload_count"), Some(3)); + assert_eq!(u64_field("payload_upload_size"), Some(2048)); + assert!(str_field("payload_download_drivers").contains("s3")); + assert!(str_field("payload_upload_drivers").contains("gcs")); +} + #[tokio::test] async fn sends_appropriate_sticky_task_queue_responses() { // This test verifies that when completions are sent with sticky queues enabled, that they diff --git a/crates/sdk-core/src/test_help/integ_helpers.rs b/crates/sdk-core/src/test_help/integ_helpers.rs index 5847f83df..66f41a6bb 100644 --- a/crates/sdk-core/src/test_help/integ_helpers.rs +++ b/crates/sdk-core/src/test_help/integ_helpers.rs @@ -275,6 +275,11 @@ impl MocksHolder { self.worker_telemetry = Some(WorkerTelemetry::from_meter(meter)); } + #[cfg(test)] + pub(crate) fn set_trace_subscriber(&mut self, sub: Arc) { + self.worker_telemetry = Some(WorkerTelemetry::from_trace_subscriber(sub)); + } + /// Can be used for tests that need to avoid auto-shutdown due to running out of mock responses pub fn make_wft_stream_interminable(&mut self) { if let Some(old_stream) = self.inputs.wft_stream.take() { diff --git a/crates/sdk-core/src/worker/mod.rs b/crates/sdk-core/src/worker/mod.rs index 739e01d98..ab823a8d7 100644 --- a/crates/sdk-core/src/worker/mod.rs +++ b/crates/sdk-core/src/worker/mod.rs @@ -583,6 +583,14 @@ impl WorkerTelemetry { trace_subscriber: None, } } + + #[cfg(test)] + pub(crate) fn from_trace_subscriber(sub: Arc) -> Self { + Self { + temporal_metric_meter: None, + trace_subscriber: Some(sub), + } + } } impl Worker { From 1a43a3ba04ec924cd8c364427696c465a04ca41e Mon Sep 17 00:00:00 2001 From: jmaeagle99 <44687433+jmaeagle99@users.noreply.github.com> Date: Thu, 30 Jul 2026 16:16:47 -0700 Subject: [PATCH 3/3] fix: only warn over default 5 sec threshold --- .../sdk-core/src/core_tests/workflow_tasks.rs | 98 ------------- .../sdk-core/src/test_help/integ_helpers.rs | 5 - crates/sdk-core/src/worker/mod.rs | 8 -- .../src/worker/workflow/managed_run.rs | 129 +++++++++++------- crates/sdk-core/src/worker/workflow/mod.rs | 38 +++++- 5 files changed, 116 insertions(+), 162 deletions(-) diff --git a/crates/sdk-core/src/core_tests/workflow_tasks.rs b/crates/sdk-core/src/core_tests/workflow_tasks.rs index ef2283afe..4eebc17df 100644 --- a/crates/sdk-core/src/core_tests/workflow_tasks.rs +++ b/crates/sdk-core/src/core_tests/workflow_tasks.rs @@ -1097,104 +1097,6 @@ async fn complete_after_eviction() { core.shutdown().await; } -/// The external-storage metrics a lang SDK reports on a workflow activation completion must -/// reach core's workflow-task duration log. This exercises the full extraction/threading path -/// (completion proto -> `TaskStorageMetrics` -> `mark_wft_complete`), which the direct-call unit -/// tests of `log_workflow_task_duration` deliberately bypass. -#[tokio::test] -async fn tmprl1104_completion_storage_metrics_reach_log() { - use std::sync::Mutex; - use temporalio_common::{ - protos::coresdk::common::ExternalStorageMetrics, - telemetry::{ - CoreLog, CoreLogConsumer, Logger, TelemetryOptions, construct_filter_string, - telemetry_init, - }, - }; - use tracing::Level; - - #[derive(Debug)] - struct CaptureConsumer(Mutex>); - impl CoreLogConsumer for CaptureConsumer { - fn on_log(&self, log: CoreLog) { - self.0.lock().unwrap().push(log); - } - } - - let consumer = Arc::new(CaptureConsumer(Mutex::new(vec![]))); - let instance = telemetry_init( - TelemetryOptions::builder() - .logging(Logger::Push { - // The routine duration line is emitted at TRACE (only the >5s line is WARN). - filter: construct_filter_string(Level::TRACE, Level::WARN), - consumer: consumer.clone(), - }) - .build(), - ) - .unwrap(); - let sub = instance.trace_subscriber().expect("subscriber present"); - - let wfid = "fake_wf_id"; - // A workflow that completes in a single WFT, so that lone WFT is the reported one that - // reaches `mark_wft_complete` and carries the metrics we attach below. - let mut t = TestHistoryBuilder::default(); - t.add_by_type(EventType::WorkflowExecutionStarted); - t.add_workflow_task_scheduled_and_started(); - let mock = mock_worker_client(); - let mut mock = single_hist_mock_sg(wfid, t, [1], mock, true); - mock.set_trace_subscriber(sub); - let core = mock_worker(mock); - - // Attach download/upload metrics to the completion, as a lang SDK would. - let activation = core.poll_workflow_activation().await.unwrap(); - let mut completion = WorkflowActivationCompletion::from_cmds( - activation.run_id, - vec![CompleteWorkflowExecution { result: None }.into()], - ); - completion.payload_download_metrics = Some(ExternalStorageMetrics { - payload_count: 2, - total_size_bytes: 1024, - total_duration: Some(prost_types::Duration { - seconds: 0, - nanos: 5_000_000, - }), - driver_names: vec!["s3".to_string()], - }); - completion.payload_upload_metrics = Some(ExternalStorageMetrics { - payload_count: 3, - total_size_bytes: 2048, - total_duration: Some(prost_types::Duration { - seconds: 1, - nanos: 0, - }), - driver_names: vec!["gcs".to_string()], - }); - core.complete_workflow_activation(completion).await.unwrap(); - core.shutdown().await; - - // The lone WFT's duration log must carry the completion's metrics, keeping download and - // upload distinct (a swap or dropped field would fail here). - let logs = consumer.0.lock().unwrap(); - let ev = logs - .iter() - .find(|l| l.message == "Workflow task duration information.") - .expect("workflow task duration log"); - let u64_field = |name: &str| ev.fields.get(name).and_then(|v| v.as_u64()); - let str_field = |name: &str| { - ev.fields - .get(name) - .and_then(|v| v.as_str()) - .unwrap_or_default() - .to_string() - }; - assert_eq!(u64_field("payload_download_count"), Some(2)); - assert_eq!(u64_field("payload_download_size"), Some(1024)); - assert_eq!(u64_field("payload_upload_count"), Some(3)); - assert_eq!(u64_field("payload_upload_size"), Some(2048)); - assert!(str_field("payload_download_drivers").contains("s3")); - assert!(str_field("payload_upload_drivers").contains("gcs")); -} - #[tokio::test] async fn sends_appropriate_sticky_task_queue_responses() { // This test verifies that when completions are sent with sticky queues enabled, that they diff --git a/crates/sdk-core/src/test_help/integ_helpers.rs b/crates/sdk-core/src/test_help/integ_helpers.rs index 66f41a6bb..5847f83df 100644 --- a/crates/sdk-core/src/test_help/integ_helpers.rs +++ b/crates/sdk-core/src/test_help/integ_helpers.rs @@ -275,11 +275,6 @@ impl MocksHolder { self.worker_telemetry = Some(WorkerTelemetry::from_meter(meter)); } - #[cfg(test)] - pub(crate) fn set_trace_subscriber(&mut self, sub: Arc) { - self.worker_telemetry = Some(WorkerTelemetry::from_trace_subscriber(sub)); - } - /// Can be used for tests that need to avoid auto-shutdown due to running out of mock responses pub fn make_wft_stream_interminable(&mut self) { if let Some(old_stream) = self.inputs.wft_stream.take() { diff --git a/crates/sdk-core/src/worker/mod.rs b/crates/sdk-core/src/worker/mod.rs index ab823a8d7..739e01d98 100644 --- a/crates/sdk-core/src/worker/mod.rs +++ b/crates/sdk-core/src/worker/mod.rs @@ -583,14 +583,6 @@ impl WorkerTelemetry { trace_subscriber: None, } } - - #[cfg(test)] - pub(crate) fn from_trace_subscriber(sub: Arc) -> Self { - Self { - temporal_metric_meter: None, - trace_subscriber: Some(sub), - } - } } impl Worker { diff --git a/crates/sdk-core/src/worker/workflow/managed_run.rs b/crates/sdk-core/src/worker/workflow/managed_run.rs index a2001e4da..923de33be 100644 --- a/crates/sdk-core/src/worker/workflow/managed_run.rs +++ b/crates/sdk-core/src/worker/workflow/managed_run.rs @@ -1551,6 +1551,10 @@ fn log_workflow_task_duration( duration: Duration, storage: &TaskStorageMetrics, ) { + let threshold = wft_duration_warn_threshold(); + if duration <= threshold { + return; + } let dl = storage.download.as_ref(); let ul = storage.upload.as_ref(); let duration_millis = |d: Duration| d.as_millis() as u64; @@ -1560,36 +1564,41 @@ fn log_workflow_task_duration( .map(duration_millis) .unwrap_or_default() }; + warn!( + workflow_type = %workflow_type, + event_id = event_id, + attempt = attempt, + workflow_task_duration = duration_millis(duration), + workflow_history_size = history_size_bytes, + payload_download_count = dl.map(|m| m.payload_count).unwrap_or_default(), + payload_download_size = dl.map(|m| m.total_size_bytes).unwrap_or_default(), + payload_download_duration = storage_millis(dl), + payload_download_drivers = ?dl.map(|m| sorted(&m.driver_names)).unwrap_or_default(), + payload_upload_count = ul.map(|m| m.payload_count).unwrap_or_default(), + payload_upload_size = ul.map(|m| m.total_size_bytes).unwrap_or_default(), + payload_upload_duration = storage_millis(ul), + payload_upload_drivers = ?ul.map(|m| sorted(&m.driver_names)).unwrap_or_default(), + "[TMPRL1104] {run_id}:{event_id}:{attempt} Workflow task duration exceeded {} seconds.", + threshold.as_secs() + ); +} - macro_rules! emit { - ($lvl:ident, $msg:literal) => { - $lvl!( - workflow_type = %workflow_type, - event_id = event_id, - attempt = attempt, - workflow_task_duration = duration_millis(duration), - workflow_history_size = history_size_bytes, - payload_download_count = dl.map(|m| m.payload_count).unwrap_or_default(), - payload_download_size = dl.map(|m| m.total_size_bytes).unwrap_or_default(), - payload_download_duration = storage_millis(dl), - payload_download_drivers = ?dl.map(|m| sorted(&m.driver_names)).unwrap_or_default(), - payload_upload_count = ul.map(|m| m.payload_count).unwrap_or_default(), - payload_upload_size = ul.map(|m| m.total_size_bytes).unwrap_or_default(), - payload_upload_duration = storage_millis(ul), - payload_upload_drivers = ?ul.map(|m| sorted(&m.driver_names)).unwrap_or_default(), - $msg, - ) - }; - } +fn wft_duration_warn_threshold() -> Duration { + static THRESHOLD: std::sync::OnceLock = std::sync::OnceLock::new(); + *THRESHOLD.get_or_init(|| { + parse_wft_duration_warn_threshold( + std::env::var("TEMPORAL_WORKFLOW_TASK_DURATION_WARN_SECONDS").ok(), + ) + }) +} - if duration > Duration::from_secs(5) { - emit!( - warn, - "[TMPRL1104] {run_id}:{event_id}:{attempt} Workflow task duration exceeded 5 seconds." - ); - } else { - emit!(trace, "Workflow task duration information."); - } +// Separated from the env read so the parse + default fallback can be unit-tested without mutating +// the process environment. +fn parse_wft_duration_warn_threshold(value: Option) -> Duration { + value + .and_then(|s| s.parse::().ok()) + .map(Duration::from_secs) + .unwrap_or(Duration::from_secs(5)) } fn sorted(names: &[String]) -> Vec { @@ -1600,7 +1609,9 @@ fn sorted(names: &[String]) -> Vec { #[cfg(test)] mod tests { - use super::{TaskStorageMetrics, log_workflow_task_duration}; + use super::{ + TaskStorageMetrics, log_workflow_task_duration, parse_wft_duration_warn_threshold, + }; use crate::worker::workflow::{WFCommand, WFCommandVariant}; use std::{ fmt::Write, @@ -1658,30 +1669,19 @@ mod tests { fn exit(&self, _: &span::Id) {} } - fn capture(duration: Duration, storage: &TaskStorageMetrics) -> CapturedEvent { + fn capture(duration: Duration, storage: &TaskStorageMetrics) -> Option { let sub = CapturingSub::default(); tracing::subscriber::with_default(sub.clone(), || { log_workflow_task_duration("run-1", "MyWorkflow", 12, 3, 4096, duration, storage); }); - sub.events - .lock() - .unwrap() - .drain(..) - .next() - .expect("one event emitted") + sub.events.lock().unwrap().drain(..).next() } #[test] - fn tmprl1104_levels_follow_thresholds() { + fn tmprl1104_warns_only_over_threshold() { let none = TaskStorageMetrics::default(); - let trace_ev = capture(Duration::from_secs(2), &none); - assert_eq!(trace_ev.level, Some(Level::TRACE)); - assert!( - !trace_ev.fields.contains("TMPRL1104"), - "fields: {}", - trace_ev.fields - ); - let warn_ev = capture(Duration::from_secs(7), &none); + assert!(capture(Duration::from_secs(2), &none).is_none()); + let warn_ev = capture(Duration::from_secs(7), &none).expect("warn emitted"); assert_eq!(warn_ev.level, Some(Level::WARN)); assert!( warn_ev.fields.contains("[TMPRL1104]"), @@ -1690,6 +1690,36 @@ mod tests { ); } + #[test] + fn tmprl1104_threshold_parsing() { + assert_eq!( + parse_wft_duration_warn_threshold(None), + Duration::from_secs(5) + ); + assert_eq!( + parse_wft_duration_warn_threshold(Some("10".to_string())), + Duration::from_secs(10) + ); + assert_eq!( + parse_wft_duration_warn_threshold(Some("0".to_string())), + Duration::from_secs(0) + ); + // Unparseable / empty / negative values fall back to the default (parsed as u64, so a + // negative can never yield a threshold). + assert_eq!( + parse_wft_duration_warn_threshold(Some("nope".to_string())), + Duration::from_secs(5) + ); + assert_eq!( + parse_wft_duration_warn_threshold(Some(String::new())), + Duration::from_secs(5) + ); + assert_eq!( + parse_wft_duration_warn_threshold(Some("-5".to_string())), + Duration::from_secs(5) + ); + } + #[test] fn tmprl1104_fields_present() { let storage = TaskStorageMetrics { @@ -1704,14 +1734,19 @@ mod tests { }), upload: None, }; - // Only the warn-level message carries the identifier, so exceed the 5s threshold. - let ev = capture(Duration::from_secs(6), &storage); + let ev = capture(Duration::from_secs(6), &storage).expect("warn emitted"); assert!(ev.fields.contains("attempt=3"), "fields: {}", ev.fields); assert!( ev.fields.contains("[TMPRL1104] run-1:12:3"), "fields: {}", ev.fields ); + // The message names the (default) threshold it exceeded. + assert!( + ev.fields.contains("exceeded 5 seconds"), + "fields: {}", + ev.fields + ); assert!( ev.fields.contains("workflow_history_size=4096"), "fields: {}", diff --git a/crates/sdk-core/src/worker/workflow/mod.rs b/crates/sdk-core/src/worker/workflow/mod.rs index 4b3340fd3..8c7fd11b1 100644 --- a/crates/sdk-core/src/worker/workflow/mod.rs +++ b/crates/sdk-core/src/worker/workflow/mod.rs @@ -549,10 +549,7 @@ impl Workflows { post_activate_hook: Option, ) -> Result<(), CompleteWfError> { let is_empty_completion = completion.is_empty(); - let task_storage_metrics = TaskStorageMetrics { - download: completion.payload_download_metrics.clone(), - upload: completion.payload_upload_metrics.clone(), - }; + let task_storage_metrics = TaskStorageMetrics::from_completion(&completion); let completion = validate_completion(completion, is_autocomplete)?; let run_id = completion.run_id().to_string(); let (tx, rx) = oneshot::channel(); @@ -1148,6 +1145,15 @@ struct TaskStorageMetrics { upload: Option, } +impl TaskStorageMetrics { + fn from_completion(completion: &WorkflowActivationCompletion) -> Self { + Self { + download: completion.payload_download_metrics.clone(), + upload: completion.payload_upload_metrics.clone(), + } + } +} + #[derive(Debug, Clone)] struct RequestEvictMsg { run_id: String, @@ -1787,6 +1793,30 @@ mod tests { protos::coresdk::workflow_activation::SignalWorkflow, }; + #[test] + fn task_storage_metrics_from_completion_keeps_directions_distinct() { + let download = ExternalStorageMetrics { + payload_count: 2, + total_size_bytes: 1024, + driver_names: vec!["s3".to_string()], + ..Default::default() + }; + let upload = ExternalStorageMetrics { + payload_count: 3, + total_size_bytes: 2048, + driver_names: vec!["gcs".to_string()], + ..Default::default() + }; + let completion = WorkflowActivationCompletion { + payload_download_metrics: Some(download.clone()), + payload_upload_metrics: Some(upload.clone()), + ..Default::default() + }; + let metrics = TaskStorageMetrics::from_completion(&completion); + assert_eq!(metrics.download, Some(download)); + assert_eq!(metrics.upload, Some(upload)); + } + #[test] fn payloads_too_large_wft_failure_is_retryable() { let violation = PayloadLimitViolation {