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 3870e8958..c9297a8b2 100644 --- a/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs +++ b/crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs @@ -397,6 +397,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 0ed405b2f..fd2dd25f1 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, @@ -297,6 +298,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(); @@ -304,7 +306,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 { @@ -1573,13 +1585,230 @@ 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 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; + 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() + }; + 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() + ); +} + +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(), + ) + }) +} + +// 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 { + let mut v = names.to_vec(); + v.sort(); + v +} + #[cfg(test)] mod tests { + use super::{ + TaskStorageMetrics, log_workflow_task_duration, parse_wft_duration_warn_threshold, + }; 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) -> 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() + } + + #[test] + fn tmprl1104_warns_only_over_threshold() { + let none = TaskStorageMetrics::default(); + 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]"), + "fields: {}", + warn_ev.fields + ); + } + + #[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 { + 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, + }; + 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: {}", + 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 3ae972472..9b3657e0f 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, @@ -303,6 +304,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 @@ -323,6 +325,7 @@ impl Workflows { WorkflowActivationCompletion { run_id, status: Some(machines_err.as_failure().into()), + ..Default::default() }, true, Option::>::None, @@ -561,6 +564,7 @@ impl Workflows { post_activate_hook: Option, ) -> Result<(), CompleteWfError> { let is_empty_completion = completion.is_empty(); + 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(); @@ -630,6 +634,7 @@ impl Workflows { wft_report_status, wft_from_complete: maybe_pwft, is_autocomplete, + task_storage_metrics, }); Ok(()) @@ -1160,7 +1165,24 @@ 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, +} + +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, @@ -1800,6 +1822,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 { diff --git a/crates/sdk-core/src/worker/workflow/workflow_stream.rs b/crates/sdk-core/src/worker/workflow/workflow_stream.rs index d2b0aa588..ee47df21a 100644 --- a/crates/sdk-core/src/worker/workflow/workflow_stream.rs +++ b/crates/sdk-core/src/worker/workflow/workflow_stream.rs @@ -366,7 +366,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 { @@ -493,6 +497,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 @@ -519,7 +524,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 70a1a7d95..5bb7abddc 100644 --- a/crates/sdk-core/tests/integ_tests/worker_tests.rs +++ b/crates/sdk-core/tests/integ_tests/worker_tests.rs @@ -147,6 +147,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 @@ -158,6 +159,7 @@ async fn worker_handles_unknown_workflow_types_gracefully() { WorkflowActivationCompletion { status: Some(Status::Successful(..)), run_id, + .. } if self.unregistered_failure_seen.load(Ordering::Relaxed) && *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 ee36be0d3..4b15f91a3 100644 --- a/crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs +++ b/crates/sdk-core/tests/integ_tests/worker_versioning_tests.rs @@ -68,6 +68,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(); @@ -285,6 +286,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 50c25a6f0..a52ea3ee8 100644 --- a/crates/sdk/src/workflow_future.rs +++ b/crates/sdk/src/workflow_future.rs @@ -243,6 +243,7 @@ impl WorkflowFuture { versioning_behavior: VersioningBehavior::Unspecified.into(), }, )), + ..Default::default() }) .expect("Completion channel intact"); }