diff --git a/OpenTelemetry.slnx b/OpenTelemetry.slnx index 5a807a98bae..4f779415514 100644 --- a/OpenTelemetry.slnx +++ b/OpenTelemetry.slnx @@ -155,6 +155,7 @@ + diff --git a/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterPersistentStorageTransmissionHandler.cs b/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterPersistentStorageTransmissionHandler.cs index de734be556c..32c3e6fe2fb 100644 --- a/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterPersistentStorageTransmissionHandler.cs +++ b/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterPersistentStorageTransmissionHandler.cs @@ -61,7 +61,7 @@ protected override bool OnSubmitRequestFailure(byte[] request, int contentLength protected override void OnShutdown(int timeoutMilliseconds) { - var sw = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.StartNew(); + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); try { @@ -74,16 +74,12 @@ protected override void OnShutdown(int timeoutMilliseconds) this.thread.Join(timeoutMilliseconds); - if (sw != null) + if (timestamp is { } startedAt) { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - - base.OnShutdown((int)Math.Max(timeout, 0)); - } - else - { - base.OnShutdown(timeoutMilliseconds); + timeoutMilliseconds = Stopwatch.Remaining(timeoutMilliseconds, startedAt); } + + base.OnShutdown(timeoutMilliseconds); } protected override void Dispose(bool disposing) diff --git a/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterTransmissionHandler.cs b/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterTransmissionHandler.cs index b63175ee502..b2503cd1464 100644 --- a/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterTransmissionHandler.cs +++ b/src/OpenTelemetry.Exporter.OpenTelemetryProtocol/Implementation/Transmission/OtlpExporterTransmissionHandler.cs @@ -35,12 +35,8 @@ public bool TrySubmitRequest(byte[] request, int contentLength) { var deadlineUtc = DateTime.UtcNow.AddMilliseconds(this.TimeoutMilliseconds); var response = this.ExportClient.SendExportRequest(request, contentLength, deadlineUtc); - if (response.Success) - { - return true; - } - return this.OnSubmitRequestFailure(request, contentLength, response); + return response.Success || this.OnSubmitRequestFailure(request, contentLength, response); } catch (Exception ex) { @@ -65,15 +61,13 @@ public bool Shutdown(int timeoutMilliseconds) { Guard.ThrowIfInvalidTimeout(timeoutMilliseconds); - var sw = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.StartNew(); + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); this.OnShutdown(timeoutMilliseconds); - if (sw != null) + if (timestamp is { } startedAt) { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - - return this.ExportClient.Shutdown((int)Math.Max(timeout, 0)); + timeoutMilliseconds = Stopwatch.Remaining(timeoutMilliseconds, startedAt); } return this.ExportClient.Shutdown(timeoutMilliseconds); diff --git a/src/OpenTelemetry/BatchExportProcessor.cs b/src/OpenTelemetry/BatchExportProcessor.cs index 1d26d2636e8..53605224eef 100644 --- a/src/OpenTelemetry/BatchExportProcessor.cs +++ b/src/OpenTelemetry/BatchExportProcessor.cs @@ -97,18 +97,12 @@ protected override void OnExport(T data) /// protected override bool OnForceFlush(int timeoutMilliseconds) - { - return this.worker.WaitForExport(timeoutMilliseconds); - } + => this.worker.WaitForExport(timeoutMilliseconds); /// protected override bool OnShutdown(int timeoutMilliseconds) { - Stopwatch? shutdownStopwatch = null; - if (timeoutMilliseconds > 0) - { - shutdownStopwatch = Stopwatch.StartNew(); - } + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); var result = this.worker.Shutdown(timeoutMilliseconds); @@ -124,18 +118,12 @@ protected override bool OnShutdown(int timeoutMilliseconds) return this.exporter.Shutdown(0) && result; } - int remainingTimeout = timeoutMilliseconds; - if (shutdownStopwatch != null) + if (timestamp is { } startedAt) { - shutdownStopwatch.Stop(); - remainingTimeout -= (int)shutdownStopwatch.ElapsedMilliseconds; - if (remainingTimeout < 0) - { - remainingTimeout = 0; - } + timeoutMilliseconds = Stopwatch.Remaining(timeoutMilliseconds, startedAt); } - return this.exporter.Shutdown(remainingTimeout) && result; + return this.exporter.Shutdown(timeoutMilliseconds) && result; } /// diff --git a/src/OpenTelemetry/CompositeProcessor.cs b/src/OpenTelemetry/CompositeProcessor.cs index df484950e55..83fd818cb7b 100644 --- a/src/OpenTelemetry/CompositeProcessor.cs +++ b/src/OpenTelemetry/CompositeProcessor.cs @@ -63,18 +63,18 @@ public CompositeProcessor AddProcessor(BaseProcessor processor) /// public override void OnEnd(T data) { - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - cur.Value.OnEnd(data); + current.Value.OnEnd(data); } } /// public override void OnStart(T data) { - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - cur.Value.OnStart(data); + current.Value.OnStart(data); } } @@ -82,9 +82,9 @@ internal override void SetParentProvider(BaseProvider parentProvider) { base.SetParentProvider(parentProvider); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - cur.Value.SetParentProvider(parentProvider); + current.Value.SetParentProvider(parentProvider); } } @@ -92,9 +92,9 @@ internal IReadOnlyList> ToReadOnlyList() { var list = new List>(); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - list.Add(cur.Value); + list.Add(current.Value); } return list; @@ -104,23 +104,20 @@ internal IReadOnlyList> ToReadOnlyList() protected override bool OnForceFlush(int timeoutMilliseconds) { var result = true; - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); + var initialTimeoutMilliseconds = timeoutMilliseconds; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - if (sw == null) + var currentTimeoutMilliseconds = timeoutMilliseconds; + + if (timestamp is { } startedAt) { - result = cur.Value.ForceFlush() && result; + currentTimeoutMilliseconds = Stopwatch.Remaining(initialTimeoutMilliseconds, startedAt); } - else - { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - // notify all the processors, even if we run overtime - result = cur.Value.ForceFlush((int)Math.Max(timeout, 0)) && result; - } + // Notify all the processors, even if we run overtime + result = current.Value.ForceFlush(currentTimeoutMilliseconds) && result; } return result; @@ -130,23 +127,20 @@ protected override bool OnForceFlush(int timeoutMilliseconds) protected override bool OnShutdown(int timeoutMilliseconds) { var result = true; - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); + var initialTimeoutMilliseconds = timeoutMilliseconds; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - if (sw == null) + var currentTimeoutMilliseconds = timeoutMilliseconds; + + if (timestamp is { } startedAt) { - result = cur.Value.Shutdown() && result; + currentTimeoutMilliseconds = Stopwatch.Remaining(initialTimeoutMilliseconds, startedAt); } - else - { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - // notify all the processors, even if we run overtime - result = cur.Value.Shutdown((int)Math.Max(timeout, 0)) && result; - } + // Notify all the processors, even if we run overtime + result = current.Value.Shutdown(currentTimeoutMilliseconds) && result; } return result; @@ -159,11 +153,11 @@ protected override void Dispose(bool disposing) { if (disposing) { - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { try { - cur.Value.Dispose(); + current.Value.Dispose(); } catch (Exception ex) { diff --git a/src/OpenTelemetry/Internal/BatchExportTaskWorker.cs b/src/OpenTelemetry/Internal/BatchExportTaskWorker.cs index 84e683a9857..8a06abb0e15 100644 --- a/src/OpenTelemetry/Internal/BatchExportTaskWorker.cs +++ b/src/OpenTelemetry/Internal/BatchExportTaskWorker.cs @@ -39,9 +39,7 @@ public BatchExportTaskWorker( /// public override void Start() - { - this.workerTask = Task.Run(this.ExporterProcAsync); - } + => this.workerTask = Task.Run(this.ExporterProcAsync); /// public override bool TriggerExport() @@ -97,6 +95,7 @@ public override bool WaitForExport(int timeoutMilliseconds) return true; } + // Otherwise we can just wait for the export to complete synchronously return this.WaitForExportAsync(timeoutMilliseconds, head).GetAwaiter().GetResult(); } @@ -125,12 +124,7 @@ public override bool Shutdown(int timeoutMilliseconds) return true; } - if (timeoutMilliseconds == 0) - { - return true; - } - - return this.workerTask.Wait(timeoutMilliseconds); + return timeoutMilliseconds == 0 || this.workerTask.Wait(timeoutMilliseconds); } /// @@ -152,9 +146,8 @@ protected override void Dispose(bool disposing) private async Task WaitForExportAsync(int timeoutMilliseconds, long targetHead) { - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); + var initialTimeoutMilliseconds = timeoutMilliseconds; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); // There is a chance that the export task finished processing all the data from the queue, // and signaled before we enter wait here, use polling to prevent being blocked indefinitely. @@ -163,15 +156,17 @@ private async Task WaitForExportAsync(int timeoutMilliseconds, long target while (true) { var timeout = pollingMilliseconds; - if (sw != null) + + if (timestamp is { } startedAt) { - var remaining = timeoutMilliseconds - sw.ElapsedMilliseconds; - if (remaining <= 0) + var remainingMilliseconds = Stopwatch.Remaining(initialTimeoutMilliseconds, startedAt); + + if (remainingMilliseconds <= 0) { return this.CircularBuffer.RemovedCount >= targetHead; } - timeout = Math.Min((int)remaining, pollingMilliseconds); + timeout = Math.Min(remainingMilliseconds, pollingMilliseconds); } try @@ -185,7 +180,7 @@ await Task.WhenAny( this.dataExportedNotification.Task, this.shutdownCompletionSource.Task, Task.Delay(timeout, combinedTokenSource.Token)).ConfigureAwait(false); -#if NET8_0_OR_GREATER +#if NET await combinedTokenSource.CancelAsync().ConfigureAwait(false); #else combinedTokenSource.Cancel(); @@ -229,7 +224,7 @@ private async Task ExporterProcAsync() await Task.WhenAny( this.exportTrigger.WaitAsync(waitAndDelayCts.Token), Task.Delay(this.ScheduledDelayMilliseconds, waitAndDelayCts.Token)).ConfigureAwait(false); -#if NET8_0_OR_GREATER +#if NET await waitAndDelayCts.CancelAsync().ConfigureAwait(false); #else waitAndDelayCts.Cancel(); diff --git a/src/OpenTelemetry/Internal/BatchExportThreadWorker.cs b/src/OpenTelemetry/Internal/BatchExportThreadWorker.cs index 25dcacce209..04b4ec4e94a 100644 --- a/src/OpenTelemetry/Internal/BatchExportThreadWorker.cs +++ b/src/OpenTelemetry/Internal/BatchExportThreadWorker.cs @@ -45,9 +45,7 @@ public BatchExportThreadWorker( /// public override void Start() - { - this.exporterThread.Start(); - } + => this.exporterThread.Start(); /// public override bool TriggerExport() @@ -86,9 +84,8 @@ public override bool WaitForExport(int timeoutMilliseconds) var triggers = new WaitHandle[] { this.dataExportedNotification, this.shutdownTrigger }; - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); + var initialTimeoutMilliseconds = timeoutMilliseconds; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); // There is a chance that the export thread finished processing all the data from the queue, // and signaled before we enter wait here, use polling to prevent being blocked indefinitely. @@ -96,34 +93,27 @@ public override bool WaitForExport(int timeoutMilliseconds) while (true) { - if (sw == null) - { - try - { - WaitHandle.WaitAny(triggers, pollingMilliseconds); - } - catch (ObjectDisposedException) - { - return false; - } - } - else + var timeout = pollingMilliseconds; + + if (timestamp is { } startedAt) { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; + var remainingMilliseconds = Stopwatch.Remaining(initialTimeoutMilliseconds, startedAt); - if (timeout <= 0) + if (remainingMilliseconds <= 0) { return this.CircularBuffer.RemovedCount >= head; } - try - { - WaitHandle.WaitAny(triggers, Math.Min((int)timeout, pollingMilliseconds)); - } - catch (ObjectDisposedException) - { - return false; - } + timeout = Math.Min(remainingMilliseconds, pollingMilliseconds); + } + + try + { + WaitHandle.WaitAny(triggers, timeout); + } + catch (ObjectDisposedException) + { + return false; } if (this.CircularBuffer.RemovedCount >= head) @@ -158,12 +148,7 @@ public override bool Shutdown(int timeoutMilliseconds) return true; } - if (timeoutMilliseconds == 0) - { - return true; - } - - return this.exporterThread.Join(timeoutMilliseconds); + return timeoutMilliseconds == 0 || this.exporterThread.Join(timeoutMilliseconds); } /// diff --git a/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderTaskWorker.cs b/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderTaskWorker.cs index 1b94fe2e780..cf5d4ccd450 100644 --- a/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderTaskWorker.cs +++ b/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderTaskWorker.cs @@ -33,9 +33,7 @@ public PeriodicExportingMetricReaderTaskWorker( /// public override void Start() - { - this.workerTask = Task.Run(this.ExporterProcAsync); - } + => this.workerTask = Task.Run(this.ExporterProcAsync); /// public override bool TriggerExport() @@ -84,12 +82,7 @@ public override bool Shutdown(int timeoutMilliseconds) return true; } - if (timeoutMilliseconds == 0) - { - return true; - } - - return this.workerTask.Wait(timeoutMilliseconds); + return timeoutMilliseconds == 0 || this.workerTask.Wait(timeoutMilliseconds); } /// @@ -112,13 +105,14 @@ protected override void Dispose(bool disposing) private async Task ExporterProcAsync() { var cancellationToken = this.cancellationTokenSource.Token; - var sw = Stopwatch.StartNew(); + var startedAt = Stopwatch.GetTimestamp(); try { while (!cancellationToken.IsCancellationRequested) { - var timeout = (int)(this.ExportIntervalMilliseconds - (sw.ElapsedMilliseconds % this.ExportIntervalMilliseconds)); + var elapsedMilliseconds = Stopwatch.GetElapsedTime(startedAt).Ticks / TimeSpan.TicksPerMillisecond; + var timeout = this.ExportIntervalMilliseconds - (int)(elapsedMilliseconds % this.ExportIntervalMilliseconds); Task? exportTriggerTask = null; Task? triggeredTask = null; diff --git a/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderThreadWorker.cs b/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderThreadWorker.cs index 3db3cdb722d..433ae6b8afd 100644 --- a/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderThreadWorker.cs +++ b/src/OpenTelemetry/Internal/PeriodicExportingMetricReaderThreadWorker.cs @@ -39,9 +39,7 @@ public PeriodicExportingMetricReaderThreadWorker( /// public override void Start() - { - this.exporterThread.Start(); - } + => this.exporterThread.Start(); /// public override bool TriggerExport() @@ -75,12 +73,7 @@ public override bool Shutdown(int timeoutMilliseconds) return true; } - if (timeoutMilliseconds == 0) - { - return true; - } - - return this.exporterThread.Join(timeoutMilliseconds); + return timeoutMilliseconds == 0 || this.exporterThread.Join(timeoutMilliseconds); } /// @@ -102,14 +95,14 @@ protected override void Dispose(bool disposing) private void ExporterProc() { - int index; - int timeout; var triggers = new WaitHandle[] { this.exportTrigger, this.shutdownTrigger }; - var sw = Stopwatch.StartNew(); + var startedAt = Stopwatch.GetTimestamp(); while (true) { - timeout = (int)(this.ExportIntervalMilliseconds - (sw.ElapsedMilliseconds % this.ExportIntervalMilliseconds)); + var elapsedMilliseconds = Stopwatch.GetElapsedTime(startedAt).Ticks / TimeSpan.TicksPerMillisecond; + var timeout = this.ExportIntervalMilliseconds - (int)(elapsedMilliseconds % this.ExportIntervalMilliseconds); + int index; try { diff --git a/src/OpenTelemetry/Metrics/Reader/BaseExportingMetricReader.cs b/src/OpenTelemetry/Metrics/Reader/BaseExportingMetricReader.cs index d53ada08b3b..70fcf9518a2 100644 --- a/src/OpenTelemetry/Metrics/Reader/BaseExportingMetricReader.cs +++ b/src/OpenTelemetry/Metrics/Reader/BaseExportingMetricReader.cs @@ -120,61 +120,29 @@ internal override bool OnCollectFromComposite(int timeoutMilliseconds) /// internal override bool OnShutdownFromComposite(int timeoutMilliseconds) - { - var result = true; - - if (timeoutMilliseconds == Timeout.Infinite) - { - result = this.CollectFromComposite(Timeout.Infinite) && result; - result = this.exporter.Shutdown(Timeout.Infinite) && result; - } - else - { - var sw = Stopwatch.StartNew(); - result = this.CollectFromComposite(timeoutMilliseconds) && result; - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - result = this.exporter.Shutdown((int)Math.Max(timeout, 0)) && result; - } - - return result; - } + => base.OnShutdownFromComposite(timeoutMilliseconds); /// - protected override bool OnCollect(int timeoutMilliseconds) - { - if (this.SupportedExportModes.HasFlag(ExportModes.Push)) - { - return base.OnCollect(timeoutMilliseconds); - } - - if (this.SupportedExportModes.HasFlag(ExportModes.Pull) && PullMetricScope.IsPullAllowed) - { - return base.OnCollect(timeoutMilliseconds); - } - - // TODO: add some error log - return false; - } + protected override bool OnCollect(int timeoutMilliseconds) => + this.SupportedExportModes.HasFlag(ExportModes.Push) + ? base.OnCollect(timeoutMilliseconds) + : this.SupportedExportModes.HasFlag(ExportModes.Pull) && + PullMetricScope.IsPullAllowed && + base.OnCollect(timeoutMilliseconds); /// protected override bool OnShutdown(int timeoutMilliseconds) { - var result = true; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); - if (timeoutMilliseconds == Timeout.Infinite) - { - result = this.Collect(Timeout.Infinite) && result; - result = this.exporter.Shutdown(Timeout.Infinite) && result; - } - else + var result = this.Collect(timeoutMilliseconds); + + if (timestamp is { } startedAt) { - var sw = Stopwatch.StartNew(); - result = this.Collect(timeoutMilliseconds) && result; - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - result = this.exporter.Shutdown((int)Math.Max(timeout, 0)) && result; + timeoutMilliseconds = Stopwatch.Remaining(timeoutMilliseconds, startedAt); } - return result; + return this.exporter.Shutdown(timeoutMilliseconds) && result; } /// diff --git a/src/OpenTelemetry/Metrics/Reader/CompositeMetricReader.cs b/src/OpenTelemetry/Metrics/Reader/CompositeMetricReader.cs index 7c6cd043ba7..691a1e22518 100644 --- a/src/OpenTelemetry/Metrics/Reader/CompositeMetricReader.cs +++ b/src/OpenTelemetry/Metrics/Reader/CompositeMetricReader.cs @@ -65,25 +65,23 @@ internal override bool ProcessMetrics(in Batch metrics, int timeoutMilli protected override bool OnCollect(int timeoutMilliseconds) { var result = true; - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); + var initialTimeoutMilliseconds = timeoutMilliseconds; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); this.CollectObservableInstruments(); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - if (sw == null) + var currentTimeoutMilliseconds = timeoutMilliseconds; + + if (timestamp is { } startedAt) { - result = cur.Value.CollectFromComposite(Timeout.Infinite) && result; + currentTimeoutMilliseconds = Stopwatch.Remaining(initialTimeoutMilliseconds, startedAt); } - else - { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - // notify all the readers, even if we run overtime - result = cur.Value.CollectFromComposite((int)Math.Max(timeout, 0)) && result; - } + // Collect observable instruments once at the composite level, then + // let child readers process the same snapshot. + result = current.Value.CollectFromComposite(currentTimeoutMilliseconds) && result; } return result; @@ -93,25 +91,21 @@ protected override bool OnCollect(int timeoutMilliseconds) protected override bool OnShutdown(int timeoutMilliseconds) { var result = true; - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); - this.CollectObservableInstruments(); + var initialTimeoutMilliseconds = timeoutMilliseconds; + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - if (sw == null) + var currentTimeoutMilliseconds = timeoutMilliseconds; + + if (timestamp is { } startedAt) { - result = cur.Value.ShutdownFromComposite(Timeout.Infinite) && result; + currentTimeoutMilliseconds = Stopwatch.Remaining(initialTimeoutMilliseconds, startedAt); } - else - { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - // notify all the readers, even if we run overtime - result = cur.Value.ShutdownFromComposite((int)Math.Max(timeout, 0)) && result; - } + // Notify all the readers, even if we run overtime + result = current.Value.ShutdownFromComposite(currentTimeoutMilliseconds) && result; } return result; @@ -123,11 +117,11 @@ protected override void Dispose(bool disposing) { if (disposing) { - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { try { - cur.Value?.Dispose(); + current.Value?.Dispose(); } catch (Exception ex) { diff --git a/src/OpenTelemetry/Metrics/Reader/CompositeMetricReaderExt.cs b/src/OpenTelemetry/Metrics/Reader/CompositeMetricReaderExt.cs index 407cc1d4761..030a1d53cd8 100644 --- a/src/OpenTelemetry/Metrics/Reader/CompositeMetricReaderExt.cs +++ b/src/OpenTelemetry/Metrics/Reader/CompositeMetricReaderExt.cs @@ -15,9 +15,9 @@ internal override List AddMetricWithNoViews(Instrument instrument) { var metrics = new List(this.count); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - var innerMetrics = cur.Value.AddMetricWithNoViews(instrument); + var innerMetrics = current.Value.AddMetricWithNoViews(instrument); if (innerMetrics.Count > 0) { Debug.Assert(innerMetrics.Count == 1, "Multiple metrics returned without view configuration"); @@ -33,9 +33,9 @@ internal override List AddMetricWithViews(Instrument instrument, List(this.count * metricStreamConfigs.Count); - for (var cur = this.Head; cur != null; cur = cur.Next) + for (var current = this.Head; current != null; current = current.Next) { - var innerMetrics = cur.Value.AddMetricWithViews(instrument, metricStreamConfigs); + var innerMetrics = current.Value.AddMetricWithViews(instrument, metricStreamConfigs); metrics.AddRange(innerMetrics); } diff --git a/src/OpenTelemetry/Metrics/Reader/MetricReader.cs b/src/OpenTelemetry/Metrics/Reader/MetricReader.cs index 0cc0b73f63d..2a0f206778d 100644 --- a/src/OpenTelemetry/Metrics/Reader/MetricReader.cs +++ b/src/OpenTelemetry/Metrics/Reader/MetricReader.cs @@ -14,11 +14,11 @@ public abstract partial class MetricReader : IDisposable { private const MetricReaderTemporalityPreference MetricReaderTemporalityPreferenceUnspecified = 0; - private static readonly Func CumulativeTemporalityPreferenceFunc = static (_) => AggregationTemporality.Cumulative; + private static readonly Func CumulativeTemporalityPreferenceFunc = + static (_) => AggregationTemporality.Cumulative; - private static readonly Func MonotonicDeltaTemporalityPreferenceFunc = (instrumentType) => - { - return instrumentType.GetGenericTypeDefinition() switch + private static readonly Func MonotonicDeltaTemporalityPreferenceFunc = + static (instrumentType) => instrumentType.GetGenericTypeDefinition() switch { var type when type == typeof(Counter<>) => AggregationTemporality.Delta, var type when type == typeof(ObservableCounter<>) => AggregationTemporality.Delta, @@ -34,11 +34,9 @@ public abstract partial class MetricReader : IDisposable // TODO: Consider logging here because we should not fall through to this case. _ => AggregationTemporality.Delta, }; - }; - private static readonly Func LowMemoryTemporalityPreferenceFunc = (instrumentType) => - { - return instrumentType.GetGenericTypeDefinition() switch + private static readonly Func LowMemoryTemporalityPreferenceFunc = + static (instrumentType) => instrumentType.GetGenericTypeDefinition() switch { var type when type == typeof(Counter<>) => AggregationTemporality.Delta, var type when type == typeof(Histogram<>) => AggregationTemporality.Delta, @@ -49,12 +47,12 @@ public abstract partial class MetricReader : IDisposable _ => AggregationTemporality.Cumulative, }; - }; private readonly Lock newTaskLock = new(); private readonly Lock onCollectLock = new(); private readonly TaskCompletionSource shutdownTcs = new(); private Func temporalityFunc = CumulativeTemporalityPreferenceFunc; + private int suppressObservableInstrumentsCollection; private int shutdownCount; private TaskCompletionSource? collectionTcs; private BaseProvider? parentProvider; @@ -157,7 +155,14 @@ internal bool ShutdownFromComposite(int timeoutMilliseconds = Timeout.Infinite) => this.Shutdown(timeoutMilliseconds, fromComposite: true); internal void CollectObservableInstruments() - => (this.parentProvider as MeterProviderSdk)?.CollectObservableInstruments(); + { + if (this.suppressObservableInstrumentsCollection > 0) + { + return; + } + + (this.parentProvider as MeterProviderSdk)?.CollectObservableInstruments(); + } internal virtual void SetParentProvider(BaseProvider parentProvider) { @@ -200,12 +205,16 @@ internal virtual bool ProcessMetrics(in Batch metrics, int timeoutMillis internal virtual bool OnCollectFromComposite(int timeoutMilliseconds) { OpenTelemetrySdkEventSource.Log.MetricReaderEvent("MetricReader.OnCollectFromComposite called."); + this.suppressObservableInstrumentsCollection++; - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); - - return this.ProcessMetricsCollection(sw, timeoutMilliseconds); + try + { + return this.OnCollect(timeoutMilliseconds); + } + finally + { + this.suppressObservableInstrumentsCollection--; + } } /// @@ -220,7 +229,18 @@ internal virtual bool OnCollectFromComposite(int timeoutMilliseconds) /// Returns true when shutdown succeeded; otherwise, false. /// internal virtual bool OnShutdownFromComposite(int timeoutMilliseconds) - => this.OnShutdown(timeoutMilliseconds); + { + this.suppressObservableInstrumentsCollection++; + + try + { + return this.OnShutdown(timeoutMilliseconds); + } + finally + { + this.suppressObservableInstrumentsCollection--; + } + } /// /// Called by Collect. This function should block the current @@ -243,15 +263,13 @@ protected virtual bool OnCollect(int timeoutMilliseconds) { OpenTelemetrySdkEventSource.Log.MetricReaderEvent("MetricReader.OnCollect called."); - var sw = timeoutMilliseconds == Timeout.Infinite - ? null - : Stopwatch.StartNew(); + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); this.CollectObservableInstruments(); OpenTelemetrySdkEventSource.Log.MetricReaderEvent("Observable instruments collected."); - return this.ProcessMetricsCollection(sw, timeoutMilliseconds); + return this.ProcessMetricsCollection(timestamp, timeoutMilliseconds); } /// @@ -389,50 +407,35 @@ private bool Shutdown(int timeoutMilliseconds, bool fromComposite) return result; } - private bool ProcessMetricsCollection(Stopwatch? sw, int timeoutMilliseconds) + private bool ProcessMetricsCollection(long? startedAtTimestamp, int timeoutMilliseconds) { OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetricsCollection called."); var metrics = this.GetMetricsBatch(); - bool result; - if (sw == null) + if (startedAtTimestamp is { } startedAt) { - OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics called."); - result = this.ProcessMetrics(metrics, Timeout.Infinite); - if (result) - { - OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics succeeded."); - } - else - { - OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics failed."); - } + timeoutMilliseconds = Stopwatch.Remaining(timeoutMilliseconds, startedAt); - return result; - } - else - { - var timeout = timeoutMilliseconds - sw.ElapsedMilliseconds; - - if (timeout <= 0) + if (timeoutMilliseconds <= 0) { OpenTelemetrySdkEventSource.Log.MetricReaderEvent("OnCollect failed timeout period has elapsed."); return false; } + } - OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics called."); - result = this.ProcessMetrics(metrics, (int)timeout); - if (result) - { - OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics succeeded."); - } - else - { - OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics failed."); - } + OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics called."); - return result; + var result = this.ProcessMetrics(metrics, timeoutMilliseconds); + if (result) + { + OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics succeeded."); } + else + { + OpenTelemetrySdkEventSource.Log.MetricReaderEvent("ProcessMetrics failed."); + } + + return result; } } diff --git a/src/OpenTelemetry/Metrics/Reader/PeriodicExportingMetricReader.cs b/src/OpenTelemetry/Metrics/Reader/PeriodicExportingMetricReader.cs index 48e6c4eabb6..63f5119c006 100644 --- a/src/OpenTelemetry/Metrics/Reader/PeriodicExportingMetricReader.cs +++ b/src/OpenTelemetry/Metrics/Reader/PeriodicExportingMetricReader.cs @@ -52,36 +52,16 @@ public PeriodicExportingMetricReader( /// protected override bool OnShutdown(int timeoutMilliseconds) { - Stopwatch? shutdownStopwatch = null; - if (timeoutMilliseconds > 0) - { - shutdownStopwatch = Stopwatch.StartNew(); - } + long? timestamp = timeoutMilliseconds == Timeout.Infinite ? null : Stopwatch.GetTimestamp(); var result = this.worker.Shutdown(timeoutMilliseconds); - if (timeoutMilliseconds == Timeout.Infinite) - { - return this.exporter.Shutdown() && result; - } - - if (timeoutMilliseconds == 0) + if (timestamp is { } startedAt) { - return this.exporter.Shutdown(0) && result; - } - - int remainingTimeout = timeoutMilliseconds; - if (shutdownStopwatch != null) - { - shutdownStopwatch.Stop(); - remainingTimeout -= (int)shutdownStopwatch.ElapsedMilliseconds; - if (remainingTimeout < 0) - { - remainingTimeout = 0; - } + timeoutMilliseconds = Stopwatch.Remaining(timeoutMilliseconds, startedAt); } - return this.exporter.Shutdown(remainingTimeout) && result; + return this.exporter.Shutdown(timeoutMilliseconds) && result; } /// diff --git a/src/OpenTelemetry/OpenTelemetry.csproj b/src/OpenTelemetry/OpenTelemetry.csproj index 12919a2dbee..cb11daff2df 100644 --- a/src/OpenTelemetry/OpenTelemetry.csproj +++ b/src/OpenTelemetry/OpenTelemetry.csproj @@ -24,6 +24,7 @@ + diff --git a/src/Shared/StopwatchExtensions.cs b/src/Shared/StopwatchExtensions.cs new file mode 100644 index 00000000000..f7c8a34ddc3 --- /dev/null +++ b/src/Shared/StopwatchExtensions.cs @@ -0,0 +1,36 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +namespace System.Diagnostics; + +internal static class StopwatchExtensions +{ + extension(Stopwatch) + { +#if !NET + public static TimeSpan GetElapsedTime(long begin) + { + var end = Stopwatch.GetTimestamp(); + var delta = end - begin; + var ticks = (long)(Conversion.ToTicks * delta); + + return new TimeSpan(ticks); + } +#endif + + public static int Remaining(int durationMilliseconds, long begin) + { + var elapsedMilliseconds = Stopwatch.GetElapsedTime(begin).Ticks / TimeSpan.TicksPerMillisecond; + elapsedMilliseconds = elapsedMilliseconds < 0 ? 0 : elapsedMilliseconds > int.MaxValue ? int.MaxValue : elapsedMilliseconds; + + return elapsedMilliseconds >= durationMilliseconds ? 0 : durationMilliseconds - (int)elapsedMilliseconds; + } + } + +#if !NET + private static class Conversion + { + internal static readonly double ToTicks = TimeSpan.TicksPerSecond / (double)Stopwatch.Frequency; + } +#endif +} diff --git a/test/OpenTelemetry.Tests/Internal/StopwatchExtensionsTests.cs b/test/OpenTelemetry.Tests/Internal/StopwatchExtensionsTests.cs new file mode 100644 index 00000000000..846f8158795 --- /dev/null +++ b/test/OpenTelemetry.Tests/Internal/StopwatchExtensionsTests.cs @@ -0,0 +1,22 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +using System.Diagnostics; +using Xunit; + +namespace OpenTelemetry.Internal.Tests; + +public static class StopwatchExtensionsTests +{ + [Fact] + public static void Remaining_ClampsElapsedMillisecondsPastIntMaxValue() + { + var elapsed = TimeSpan.FromDays(30); + var elapsedTimestampTicks = (long)(elapsed.Ticks * (Stopwatch.Frequency / (double)TimeSpan.TicksPerSecond)); + var begin = Stopwatch.GetTimestamp() - elapsedTimestampTicks; + + var remaining = Stopwatch.Remaining(1000, begin); + + Assert.Equal(0, remaining); + } +} diff --git a/test/OpenTelemetry.Tests/Metrics/MetricApiTests.cs b/test/OpenTelemetry.Tests/Metrics/MetricApiTests.cs index cac5c1d471f..6cf371d41d9 100644 --- a/test/OpenTelemetry.Tests/Metrics/MetricApiTests.cs +++ b/test/OpenTelemetry.Tests/Metrics/MetricApiTests.cs @@ -2487,7 +2487,7 @@ private void MultithreadedCounterTest(T deltaValueUpdatedByEachCall) } argToThread.MreToEnsureAllThreadsStart.WaitOne(); - var sw = Stopwatch.StartNew(); + var startedTimestamp = Stopwatch.GetTimestamp(); argToThread.MreToBlockUpdateThread.Set(); for (var i = 0; i < NumberOfThreads; i++) @@ -2495,7 +2495,8 @@ private void MultithreadedCounterTest(T deltaValueUpdatedByEachCall) t[i].Join(); } - this.output.WriteLine($"Took {sw.ElapsedMilliseconds} msecs. Total threads: {NumberOfThreads}, each thread doing {NumberOfMetricUpdateByEachThread} recordings."); + var elapsed = Stopwatch.GetElapsedTime(startedTimestamp); + this.output.WriteLine($"Took {elapsed.TotalMilliseconds} msecs. Total threads: {NumberOfThreads}, each thread doing {NumberOfMetricUpdateByEachThread} recordings."); meterProvider.ForceFlush(); diff --git a/test/OpenTelemetry.Tests/Metrics/MultipleReadersTests.cs b/test/OpenTelemetry.Tests/Metrics/MultipleReadersTests.cs index b21b5e9cb40..4abe0b623e8 100644 --- a/test/OpenTelemetry.Tests/Metrics/MultipleReadersTests.cs +++ b/test/OpenTelemetry.Tests/Metrics/MultipleReadersTests.cs @@ -305,6 +305,33 @@ public void ObservableInstrumentCallbacksInvokedOncePerCollection(bool hasViews) AssertLongSumValueForMetric(exportedItems2[1], 20); } + [Theory] + [InlineData(false)] + [InlineData(true)] + public void CompositeMetricReader_PreservesOriginalTimeoutBudget(bool shutdown) + { + const int delayMilliseconds = 200; + const int timeoutMilliseconds = 1500; + + using var first = new TimeoutCapturingMetricReader(delayMilliseconds); + using var second = new TimeoutCapturingMetricReader(delayMilliseconds); + using var third = new TimeoutCapturingMetricReader(); + + using var composite = new CompositeMetricReader([first, second, third]); + + if (shutdown) + { + composite.Shutdown(timeoutMilliseconds); + } + else + { + composite.Collect(timeoutMilliseconds); + } + + Assert.Equal(timeoutMilliseconds, shutdown ? first.LastShutdownTimeoutMilliseconds : first.LastCollectTimeoutMilliseconds); + Assert.InRange(shutdown ? third.LastShutdownTimeoutMilliseconds : third.LastCollectTimeoutMilliseconds, 1000, timeoutMilliseconds); + } + private static void AssertLongSumValueForMetric(Metric metric, long value) { var metricPoints = metric.GetMetricPoints(); @@ -320,4 +347,32 @@ private static void AssertLongSumValueForMetric(Metric metric, long value) Assert.Equal(value, metricPointForFirstExport.GetGaugeLastValueLong()); } } + +#pragma warning disable CA2000 // BaseExportingMetricReader owns the exporter lifecycle + private sealed class TimeoutCapturingMetricReader(int delayMilliseconds = 0) : BaseExportingMetricReader(new TimeoutCapturingMetricExporter()) +#pragma warning restore CA2000 + { + public int LastCollectTimeoutMilliseconds { get; private set; } = Timeout.Infinite; + + public int LastShutdownTimeoutMilliseconds { get; private set; } = Timeout.Infinite; + + protected override bool OnCollect(int timeoutMilliseconds) + { + this.LastCollectTimeoutMilliseconds = timeoutMilliseconds; + Thread.Sleep(delayMilliseconds); + return true; + } + + protected override bool OnShutdown(int timeoutMilliseconds) + { + this.LastShutdownTimeoutMilliseconds = timeoutMilliseconds; + Thread.Sleep(delayMilliseconds); + return true; + } + } + + private sealed class TimeoutCapturingMetricExporter : BaseExporter + { + public override ExportResult Export(in Batch batch) => ExportResult.Success; + } } diff --git a/test/OpenTelemetry.Tests/Trace/CompositeActivityProcessorTests.cs b/test/OpenTelemetry.Tests/Trace/CompositeActivityProcessorTests.cs index a287f897714..963d253c1c3 100644 --- a/test/OpenTelemetry.Tests/Trace/CompositeActivityProcessorTests.cs +++ b/test/OpenTelemetry.Tests/Trace/CompositeActivityProcessorTests.cs @@ -106,7 +106,53 @@ public void CompositeActivityProcessor_ForwardsParentProvider() Assert.Equal(provider, p2.ParentProvider); } - private sealed class TestProvider : TracerProvider + [Theory] + [InlineData(false)] + [InlineData(true)] + public void CompositeActivityProcessor_PreservesOriginalTimeoutBudget(bool shutdown) + { + const int delayMilliseconds = 200; + const int timeoutMilliseconds = 1500; + + using var first = new TimeoutCapturingActivityProcessor(delayMilliseconds); + using var second = new TimeoutCapturingActivityProcessor(delayMilliseconds); + using var third = new TimeoutCapturingActivityProcessor(); + + using var composite = new CompositeProcessor([first, second, third]); + + if (shutdown) + { + composite.Shutdown(timeoutMilliseconds); + } + else + { + composite.ForceFlush(timeoutMilliseconds); + } + + Assert.Equal(timeoutMilliseconds, shutdown ? first.LastShutdownTimeoutMilliseconds : first.LastForceFlushTimeoutMilliseconds); + Assert.InRange(shutdown ? third.LastShutdownTimeoutMilliseconds : third.LastForceFlushTimeoutMilliseconds, 1000, timeoutMilliseconds); + } + + private sealed class TestProvider : TracerProvider; + + private sealed class TimeoutCapturingActivityProcessor(int delayMilliseconds = 0) : BaseProcessor { + public int LastForceFlushTimeoutMilliseconds { get; private set; } = Timeout.Infinite; + + public int LastShutdownTimeoutMilliseconds { get; private set; } = Timeout.Infinite; + + protected override bool OnForceFlush(int timeoutMilliseconds) + { + this.LastForceFlushTimeoutMilliseconds = timeoutMilliseconds; + Thread.Sleep(delayMilliseconds); + return true; + } + + protected override bool OnShutdown(int timeoutMilliseconds) + { + this.LastShutdownTimeoutMilliseconds = timeoutMilliseconds; + Thread.Sleep(delayMilliseconds); + return true; + } } }