diff --git a/src/Functions.WorkerProxy/ExtensionGrpcActivity.cs b/src/Functions.WorkerProxy/ExtensionGrpcActivity.cs new file mode 100644 index 000000000..4e64764de --- /dev/null +++ b/src/Functions.WorkerProxy/ExtensionGrpcActivity.cs @@ -0,0 +1,38 @@ +// Copyright (c) .NET Foundation. All rights reserved. +// Licensed under the MIT License. See License.txt in the project root for license information. + +using System.Diagnostics; + +namespace Azure.Functions.WorkerProxy; + +/// +/// Enriches worker-facing ASP.NET Core activities with extension RPC correlation and concurrency data. +/// +internal static class ExtensionGrpcActivity +{ + private const string TagPrefix = "azure.functions.worker_proxy.extension_rpc"; + + /// + /// Enriches the current activity after an extension call acquires a host stream. + /// + /// The logical extension call identifier. + /// The physical extension stream identifier. + /// The active-call count when the call opens. + public static void CallOpened(string callId, string streamId, int activeCallCount) + { + Activity? activity = Activity.Current; + activity?.SetTag($"{TagPrefix}.call_id", callId); + activity?.SetTag($"{TagPrefix}.stream_id", streamId); + activity?.SetTag($"{TagPrefix}.active_calls_at_open", activeCallCount); + } + + /// + /// Enriches the current activity with an extension call's final concurrency snapshot. + /// + /// The active-call count when the call completes. + public static void CallCompleted(int activeCallCount) + { + Activity? activity = Activity.Current; + activity?.SetTag($"{TagPrefix}.active_calls_at_completion", activeCallCount); + } +} diff --git a/src/Functions.WorkerProxy/ExtensionGrpcIngress.cs b/src/Functions.WorkerProxy/ExtensionGrpcIngress.cs index f633ba9ff..fcc6a1132 100644 --- a/src/Functions.WorkerProxy/ExtensionGrpcIngress.cs +++ b/src/Functions.WorkerProxy/ExtensionGrpcIngress.cs @@ -24,10 +24,12 @@ namespace Azure.Functions.WorkerProxy; /// /// The WorkerProxy listener configuration. /// The extension RPC stream coordinator. +/// The extension gRPC metrics. /// The logger used for ingress diagnostics. internal sealed partial class ExtensionGrpcIngress( WorkerProxyEndpointConfiguration endpoints, ExtensionRpcStreamCoordinator streamCoordinator, + ExtensionGrpcMetrics metrics, ILogger logger) { internal const string FunctionRpcEventStreamPath = "/AzureFunctionsRpcMessages.FunctionRpc/EventStream"; @@ -55,6 +57,8 @@ internal sealed partial class ExtensionGrpcIngress( private readonly ExtensionRpcStreamCoordinator _streamCoordinator = streamCoordinator ?? throw new ArgumentNullException(nameof(streamCoordinator)); + private readonly ExtensionGrpcMetrics _metrics = metrics ?? throw new ArgumentNullException(nameof(metrics)); + private readonly ILogger _logger = logger ?? throw new ArgumentNullException(nameof(logger)); /// @@ -153,6 +157,9 @@ await CompleteWithStatusAsync( int activeCallCountAtOpen = call.ActiveCallCount; double openDurationMilliseconds = Stopwatch.GetElapsedTime(callStart).TotalMilliseconds; + ExtensionGrpcActivity.CallOpened(call.CallId, call.StreamId, activeCallCountAtOpen); + _metrics.CallOpenDuration.Record(openDurationMilliseconds); + _metrics.ActiveCalls.Increment(); // Remove this per-call log when WorkerProxy metric exporting is wired up. Log.CallOpened(_logger, start.Method, call.CallId, activeCallCountAtOpen, openDurationMilliseconds); @@ -207,6 +214,9 @@ await CompleteWithStatusAsync( await StopRelayTasksAsync(cancellationTokenSource, requestTask, responseTask); int activeCallCountAtCompletion = call.ActiveCallCount; double callDurationMilliseconds = Stopwatch.GetElapsedTime(callStart).TotalMilliseconds; + ExtensionGrpcActivity.CallCompleted(activeCallCountAtCompletion); + _metrics.CallDuration.Record(callDurationMilliseconds); + _metrics.ActiveCalls.Decrement(); // Remove this per-call log when WorkerProxy metric exporting is wired up. Log.CallCompleted( diff --git a/src/Functions.WorkerProxy/ExtensionGrpcMetrics.cs b/src/Functions.WorkerProxy/ExtensionGrpcMetrics.cs new file mode 100644 index 000000000..6214f276d --- /dev/null +++ b/src/Functions.WorkerProxy/ExtensionGrpcMetrics.cs @@ -0,0 +1,119 @@ +// Copyright (c) .NET Foundation. All rights reserved. +// Licensed under the MIT License. See License.txt in the project root for license information. + +using System; +using System.Diagnostics.Metrics; + +namespace Azure.Functions.WorkerProxy; + +/// +/// Records transport-level measurements for worker-facing extension gRPC calls. +/// +internal sealed class ExtensionGrpcMetrics +{ + internal const string MeterName = "Microsoft.Azure.Functions.WorkerProxy.ExtensionGrpc"; + internal const string MeterVersion = "1.0.0"; + internal const string ActiveCallsInstrumentName = "azure.functions.worker_proxy.extension_rpc.calls.active"; + internal const string CallDurationInstrumentName = "azure.functions.worker_proxy.extension_rpc.call.duration"; + internal const string CallOpenDurationInstrumentName = + "azure.functions.worker_proxy.extension_rpc.call.open.duration"; + + /// + /// Initializes extension gRPC metrics using a factory-owned meter. + /// + /// The factory that owns the meter lifetime. + public ExtensionGrpcMetrics(IMeterFactory meterFactory) + { + ArgumentNullException.ThrowIfNull(meterFactory); + +#pragma warning disable CA2000 // IMeterFactory owns the meter lifetime. + Meter meter = meterFactory.Create(MeterName, MeterVersion); +#pragma warning restore CA2000 + + ActiveCalls = new(meter); + CallDuration = new(meter); + CallOpenDuration = new(meter); + } + + /// + /// Gets the active-call counter. + /// + public ActiveCallsCounter ActiveCalls { get; } + + /// + /// Gets the total call-duration histogram. + /// + public CallDurationHistogram CallDuration { get; } + + /// + /// Gets the call-open-duration histogram. + /// + public CallOpenDurationHistogram CallOpenDuration { get; } + + /// + /// Records the number of active extension gRPC calls. + /// + /// + /// Initializes the active-call counter. + /// + /// The meter used to create the counter. + internal sealed class ActiveCallsCounter(Meter meter) + { + private readonly UpDownCounter _counter = meter.CreateUpDownCounter( + ActiveCallsInstrumentName, + unit: "{call}", + description: "Number of worker-facing extension gRPC calls currently relayed by this proxy."); + + /// + /// Records that an extension gRPC call started relaying. + /// + public void Increment() => _counter.Add(1); + + /// + /// Records that an extension gRPC call stopped relaying. + /// + public void Decrement() => _counter.Add(-1); + } + + /// + /// Records extension gRPC call durations. + /// + /// + /// Initializes the call-duration histogram. + /// + /// The meter used to create the histogram. + internal sealed class CallDurationHistogram(Meter meter) + { + private readonly Histogram _histogram = meter.CreateHistogram( + CallDurationInstrumentName, + unit: "ms", + description: "Time spent relaying a worker-facing extension gRPC call, including stream assignment."); + + /// + /// Records a call duration. + /// + /// The total relay duration in milliseconds. + public void Record(double durationMilliseconds) => _histogram.Record(durationMilliseconds); + } + + /// + /// Records extension gRPC call-open durations. + /// + /// + /// Initializes the call-open-duration histogram. + /// + /// The meter used to create the histogram. + internal sealed class CallOpenDurationHistogram(Meter meter) + { + private readonly Histogram _histogram = meter.CreateHistogram( + CallOpenDurationInstrumentName, + unit: "ms", + description: "Time spent assigning a worker-facing extension gRPC call to the host stream."); + + /// + /// Records a call-open duration. + /// + /// The stream-assignment latency in milliseconds. + public void Record(double durationMilliseconds) => _histogram.Record(durationMilliseconds); + } +} diff --git a/src/Functions.WorkerProxy/WorkerProxyApplication.cs b/src/Functions.WorkerProxy/WorkerProxyApplication.cs index 867222d87..145ef9f94 100644 --- a/src/Functions.WorkerProxy/WorkerProxyApplication.cs +++ b/src/Functions.WorkerProxy/WorkerProxyApplication.cs @@ -53,6 +53,8 @@ public static WebApplication Build(string[] args) builder.Services.AddHostedService(static services => services.GetRequiredService()); builder.Services.AddSingleton(); builder.Services.AddSingleton(); + builder.Services.AddMetrics(); + builder.Services.AddSingleton(); builder.Services.AddSingleton(); ConfigureHttpForwarding(builder); diff --git a/test/Functions.WorkerProxy.Tests/ExtensionGrpcIngressTests.cs b/test/Functions.WorkerProxy.Tests/ExtensionGrpcIngressTests.cs index 8f23bd1c9..89b80dcd3 100644 --- a/test/Functions.WorkerProxy.Tests/ExtensionGrpcIngressTests.cs +++ b/test/Functions.WorkerProxy.Tests/ExtensionGrpcIngressTests.cs @@ -3,6 +3,8 @@ using System; using System.Buffers.Binary; +using System.Diagnostics; +using System.Diagnostics.Metrics; using System.IO; using System.Linq; using System.Text; @@ -23,6 +25,7 @@ namespace Azure.Functions.WorkerProxy.Tests; public class ExtensionGrpcIngressTests { private const int WorkerGrpcPort = 50052; + private static readonly ExtensionGrpcMetrics Metrics = new(new TestMeterFactory()); [Fact] public async Task HandleAsync_RelaysOpaqueGrpcFramesAndMetadata() @@ -42,6 +45,7 @@ public async Task HandleAsync_RelaysOpaqueGrpcFramesAndMetadata() context.Features.Set(trailersFeature); ExtensionGrpcIngress ingress = CreateIngress(streamCoordinator); + using var activity = new Activity("extension-rpc-test").Start(); Task ingressTask = ingress.HandleAsync(context); ExtensionRpcMessage start = await lease.Stream.Outbound.ReadAsync(); @@ -121,6 +125,12 @@ await lease.Stream.HandleInboundAsync(new ExtensionRpcMessage Assert.Equal("CQo=", context.Response.Headers["response-bin"]); Assert.Equal("Cww=", trailersFeature.Trailers["trailer-bin"]); Assert.Equal("0", trailersFeature.Trailers["grpc-status"]); + Assert.Equal(start.CallId, activity.GetTagItem("azure.functions.worker_proxy.extension_rpc.call_id")); + Assert.Equal(hello.ShardId, activity.GetTagItem("azure.functions.worker_proxy.extension_rpc.stream_id")); + Assert.Equal(1, activity.GetTagItem("azure.functions.worker_proxy.extension_rpc.active_calls_at_open")); + Assert.Equal(0, activity.GetTagItem("azure.functions.worker_proxy.extension_rpc.active_calls_at_completion")); + Assert.Null(activity.GetTagItem("azure.functions.worker_proxy.extension_rpc.call.open.duration_ms")); + Assert.Null(activity.GetTagItem("azure.functions.worker_proxy.extension_rpc.call.duration_ms")); } [Fact] @@ -371,6 +381,7 @@ private static ExtensionGrpcIngress CreateIngress(ExtensionRpcStreamCoordinator return new ExtensionGrpcIngress( endpoints, streamCoordinator, + Metrics, NullLogger.Instance); } @@ -417,4 +428,13 @@ private sealed class TestResponseTrailersFeature : IHttpResponseTrailersFeature { public IHeaderDictionary Trailers { get; set; } = new HeaderDictionary(); } + + private sealed class TestMeterFactory : IMeterFactory + { + public Meter Create(MeterOptions options) => new(options); + + public void Dispose() + { + } + } } diff --git a/test/Functions.WorkerProxy.Tests/ExtensionGrpcMetricsTests.cs b/test/Functions.WorkerProxy.Tests/ExtensionGrpcMetricsTests.cs new file mode 100644 index 000000000..81381fbc4 --- /dev/null +++ b/test/Functions.WorkerProxy.Tests/ExtensionGrpcMetricsTests.cs @@ -0,0 +1,93 @@ +// Copyright (c) .NET Foundation. All rights reserved. +// Licensed under the MIT License. See License.txt in the project root for license information. + +using System; +using System.Collections.Concurrent; +using System.Diagnostics.Metrics; +using Xunit; + +namespace Azure.Functions.WorkerProxy.Tests; + +public class ExtensionGrpcMetricsTests +{ + [Fact] + public void RecordsCallMeasurements() + { + using var meterFactory = new TestMeterFactory(); + var metrics = new ExtensionGrpcMetrics(meterFactory); + var measurements = new ConcurrentQueue<(string Name, double Value)>(); + using var listener = new MeterListener + { + InstrumentPublished = (instrument, meterListener) => + { + if (string.Equals( + instrument.Meter.Name, + ExtensionGrpcMetrics.MeterName, + StringComparison.Ordinal)) + { + meterListener.EnableMeasurementEvents(instrument); + } + }, + }; + listener.SetMeasurementEventCallback( + (instrument, value, _, _) => measurements.Enqueue((instrument.Name, value))); + listener.SetMeasurementEventCallback( + (instrument, value, _, _) => measurements.Enqueue((instrument.Name, value))); + listener.Start(); + + metrics.CallOpenDuration.Record(12.5); + metrics.ActiveCalls.Increment(); + metrics.CallDuration.Record(25.5); + metrics.ActiveCalls.Decrement(); + + Assert.Contains( + measurements, + measurement => string.Equals( + measurement.Name, + ExtensionGrpcMetrics.CallOpenDurationInstrumentName, + StringComparison.Ordinal) + && measurement.Value == 12.5); + Assert.Contains( + measurements, + measurement => string.Equals( + measurement.Name, + ExtensionGrpcMetrics.CallDurationInstrumentName, + StringComparison.Ordinal) + && measurement.Value == 25.5); + Assert.Contains( + measurements, + measurement => string.Equals( + measurement.Name, + ExtensionGrpcMetrics.ActiveCallsInstrumentName, + StringComparison.Ordinal) + && measurement.Value == 1); + Assert.Contains( + measurements, + measurement => string.Equals( + measurement.Name, + ExtensionGrpcMetrics.ActiveCallsInstrumentName, + StringComparison.Ordinal) + && measurement.Value == -1); + } + + private sealed class TestMeterFactory : IMeterFactory + { + private readonly ConcurrentBag _meters = []; + + public Meter Create(MeterOptions options) + { + var meter = new Meter(options); + _meters.Add(meter); + + return meter; + } + + public void Dispose() + { + foreach (Meter meter in _meters) + { + meter.Dispose(); + } + } + } +}