* Avoid null endpoint DTO metadata in tests * Enforce console logs hub read permission * Remove unused console logs hub import * Support mapped endpoint metadata in auth tests * Reduce console log capture throughput impact * Address Copilot console logs review * Refactor task scheduling to support tenant-level background work and enhance logging functionality. * Introduce ConsoleStreamHook for stdout/stderr tee and enhance logging validation. Adjust test cases and startup warnings for distributed lock provider usage. * Refactor console logging pipeline with capture optimization and new ConsoleLogsHost; update tests accordingly. * Add Ansi SGR parser for console logs and associated unit tests * Remove ANSI color renderings and parsers; integrate ConsoleLogScopeAccessor for improved logging context with workflow instance ID support. * Address console logs code quality feedback * Address PR review feedback * Preserve console logs extension points * Stabilize console logs host lifecycle * Address final automated review comments * Tighten console log capture shutdown * Address console log review feedback * Address follow-up review feedback * Cover final review feedback * Avoid recursive console provider initialization * Guard console host lease shutdown * Preserve console log scope and provider lifetime * Correlate console log scope fallback * Tighten console scope correlation * Expose host services during provider construction * Redact ANSI-normalized console lines * Add OpenTelemetry diagnostics backend foundation * Add OTLP HTTP ingestion parsing * Document OpenTelemetry diagnostics setup * Enforce OpenTelemetry hub permissions * Remove `ConsoleCaptureTee` and related services and tests * Add OpenTelemetry HTTP ingestion integration test * Use pipeline contributors for console log context * Update CShells package versions to 0.0.24-preview.132 * Add OpenTelemetry ingestion security tests * Add OpenTelemetry API authorization tests * Filter live console logs by workflow instance * Add OpenTelemetry hub tests * Add OpenTelemetry gRPC metadata hook * Assert OpenTelemetry workflow tags survive ingestion * Mark OpenTelemetry core build verified * Enhance console logging with activity execution metadata and extend test coverage. * Address console logs stream consumption comment * Wire OpenTelemetry diagnostics into core sample * Address Core diagnostics review feedback * Address Core Copilot follow-up feedback * Add OpenTelemetry metric instrument names * Address Core Copilot provider feedback * Address Core Copilot diagnostics follow-up * Address Core Copilot live feed feedback * Address Core Copilot store feedback * Integrate OpenTelemetry for logging, tracing, and metrics in ModularServer and update launch settings and docker-compose configuration. * Refactor to replace `ConsoleLogStream.Core` with `ConsoleLogStreaming.Core` across codebase and update `ConsoleStreamHook` installation. * Add diagnostics OpenTelemetry backend * Fix OpenTelemetry live hub subscription * Fix modular OpenTelemetry exporter endpoints * Add CShells logging configuration in appsettings.json * Remove obsolete unit tests and helper classes * Restore default activity exception handling * Simplify type serialization and alias management This commit refactors the internal type serialization and alias management system to reduce boilerplate, improve robustness, and simplify the developer experience: - Removed numerous explicit `ExpressionOptions` type alias registrations across various modules. - Updated `TypeJsonConverter` and polymorphic serialization to reliably handle types using assembly-qualified names when a short alias is not explicitly registered. - Streamlined `ExcludeFromHashConverter` to strictly adhere to `ExcludeFromHashAttribute` for hash calculations, removing complex `JsonIgnoreCondition` logic. - Eliminated several helper classes (`WorkflowJsonTypeResolver`, `WorkflowTypeValidator`, `IWorkflowTypeRegistry`, `WorkflowFactoryDictionary`, `JavaScriptExceptionTypeAliasRegistrar`, `WorkflowRuntimeTypeAliasRegistrar`) and their associated unit tests, simplifying the codebase. Additionally, this commit introduces a comprehensive markdown document (`product-website-feature-source.md`) outlining Elsa's core features, Studio capabilities, extension ecosystem, and architectural selling points, intended as source material for the product website. * Refine type serialization for improved robustness and alias handling This commit further enhances the type serialization and deserialization mechanisms: * Centralizes type resolution and alias management through `IWellKnownTypeRegistry` and `WorkflowJsonTypeResolver`. * Prioritizes registered type aliases when serializing type metadata in `PolymorphicObjectConverter`, resulting in more concise JSON output. * Enhances deserialization in `PolymorphicObjectConverter` and `VariableMapper` to gracefully handle unknown or non-instantiable types, providing fallbacks and logging warnings. * Simplifies `TypeJsonConverter` by delegating complex type resolution logic to the `WorkflowJsonTypeResolver`. * Adds `JsonArray` to the well-known type aliases for direct recognition. * Fix console logs packaging and workflow type resolution * Fix console log metadata and type resolution * Address Copilot review feedback * Enhance type resolution, improve console log handling, and update tests - Streamlined `WorkflowDictionaryExtensions` for better workflow registration validation. - Refined `ConsoleLogsAuthorizationTests` with the new `SetJsonRequest` helper to improve test requests handling. - Updated `OrderDefinition` to ignore JSON serialization for `KeySelector`. - Enhanced `WorkflowRuntimeFeature` for improved workflow registration and type alias configuration. - Added tests to ensure `ConsoleLogProvider` metadata filtration in various scenarios. - Improved type serialization logic in `WorkflowJsonTypeResolver`. - Updated README to fix references related to diagnostics. - Optimized `ExcludeFromHashConverter` for property serialization conditions. - Modified `TriggerIndexer` for streamlined trigger management. - Tested payload checks in `PublishEventTests`. - Adjusted `Endpoint` in `ConsoleLogs` for automatic JSON request handling. - Ensured registration of workflow type aliases in `WorkflowsFeature`. * Restore CLR workflow registration compatibility * Align JSON island serialization fixtures * Add Console Logs Services and Enhance Endpoint Handling - Introduced `ActivityExecutionsEndpointTests` to validate route exposure. - Added `ConsoleLogCaptureHostedService` for console log streaming. - Implemented `ConsoleStreamJsonConverter` for JSON conversion of console streams. - Developed `ElsaConsoleLogRecentBuffer` to handle recent log buffering. - Updated `ConsoleLogsAuthorizationTests` with new test cases for stream filter mapping. - Consolidated console log provider dependencies and registration, including recent buffering. - Enhanced `ElsaConsoleLogProvider` to use recent buffer for filtering. - Adjusted `Program.cs` for streamlined logging service setup. * Enhance type resolution and test coverage; streamline console log integration - Added `ConsoleStreamHook` for streamlined log streaming. - Updated `WorkflowJsonTypeResolverTests` to improve type resolution and test new scenarios. - Simplified type resolution by removing trusted assembly checks. * Fix CI smoke and package restore failures * Fix Docker smoke image project paths * Fix Docker Python runtime packages * Fix Docker CA smoke teardown * Refresh Elsa roadmap * Implement background processors and mediation coordination - Added `BackgroundCommandProcessor`, `BackgroundJobProcessor`, and `BackgroundNotificationProcessor` classes for handling commands, jobs, and notifications, respectively. - Introduced `MediatorBackgroundProcessingCoordinator` to coordinate the execution of all background processors. - Implemented `MediatorBackgroundTask` for wrapping `MediatorBackgroundProcessingCoordinator` in `BackgroundTask`. - Added unit tests for `MediatorBackgroundTask` to ensure proper start and stop behavior. - Refactored `BackgroundCommandSenderHostedService` to utilize `BackgroundCommandProcessor`. - Introduced 'elsa-roadmap-refresh' skill configuration for roadmap updates. * Address workflow type resolution review feedback * Address follow-up review feedback * Restore recent console logs execute path * Address Copilot follow-up review * Decouple workflow JSON aliases from expressions * Fix workflow management unit test setup * Fix console logs recent endpoint handler shape * Respect workflow JSON strict type aliases * Remove unused console log contracts reference * Address Copilot review feedback * Address Copilot follow-up comments * Synchronize ring buffer dropped count * Address background processor strategy replay * Fix diagnostics live feed regressions
305 lines
14 KiB
C#
305 lines
14 KiB
C#
using System.Security.Claims;
|
|
using Elsa.Diagnostics.OpenTelemetry.Contracts;
|
|
using Elsa.Diagnostics.OpenTelemetry.Models;
|
|
using Elsa.Diagnostics.OpenTelemetry.Options;
|
|
using Elsa.Diagnostics.OpenTelemetry.Permissions;
|
|
using Elsa.Diagnostics.OpenTelemetry.Providers.InMemory;
|
|
using Elsa.Diagnostics.OpenTelemetry.RealTime;
|
|
using FastEndpoints;
|
|
using FastEndpoints.Security;
|
|
using Microsoft.AspNetCore.Http.Features;
|
|
using Microsoft.AspNetCore.SignalR;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
using Xunit.Sdk;
|
|
using OptionsFactory = Microsoft.Extensions.Options.Options;
|
|
|
|
namespace Elsa.Diagnostics.OpenTelemetry.IntegrationTests;
|
|
|
|
public class OpenTelemetryHubTests
|
|
{
|
|
private readonly DateTimeOffset _now = new(2026, 5, 26, 10, 0, 0, TimeSpan.Zero);
|
|
|
|
[Fact]
|
|
public async Task SubscribeAsync_WhenUserLacksPermission_DeniesAccess()
|
|
{
|
|
var hub = CreateHub(new TestLiveFeed(), "write:diagnostics:opentelemetry");
|
|
|
|
await Assert.ThrowsAsync<HubException>(() => hub.SubscribeAsync(new()));
|
|
}
|
|
|
|
[Theory]
|
|
[InlineData(OpenTelemetryPermissions.Read)]
|
|
[InlineData(PermissionNames.All)]
|
|
[InlineData("read:*")]
|
|
public async Task SubscribeAsync_WhenUserCanRead_ForwardsItemsToCaller(string permission)
|
|
{
|
|
var liveFeed = new TestLiveFeed(new OpenTelemetryStreamItem { Trace = Trace("trace-1") });
|
|
var caller = new CapturingOpenTelemetryClient();
|
|
var hub = CreateHub(liveFeed, permission, caller);
|
|
|
|
await hub.SubscribeAsync(new OpenTelemetryTraceFilter { TraceId = "trace-1" });
|
|
|
|
await AssertEventuallyAsync(() =>
|
|
{
|
|
Assert.Equal("trace-1", Assert.Single(caller.Items).Trace?.TraceId);
|
|
Assert.Equal("trace-1", liveFeed.Filter?.TraceId);
|
|
});
|
|
}
|
|
|
|
[Fact]
|
|
public async Task SubscribeAsync_WhenFilterTimeRangeIsInvalid_RejectsFilter()
|
|
{
|
|
var hub = CreateHub(new TestLiveFeed(), OpenTelemetryPermissions.Read);
|
|
|
|
await Assert.ThrowsAsync<HubException>(() => hub.SubscribeAsync(new OpenTelemetryTraceFilter
|
|
{
|
|
From = _now.AddMinutes(1),
|
|
To = _now
|
|
}));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task LiveFeed_WhenTraceFilterIsSet_OnlyPublishesMatchingTraces()
|
|
{
|
|
var liveFeed = new InMemoryOpenTelemetryLiveFeed(OptionsFactory.Create(new OpenTelemetryDiagnosticsOptions()));
|
|
using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5));
|
|
await using var enumerator = liveFeed.SubscribeAsync(new OpenTelemetryTraceFilter { TraceId = "trace-keep" }, timeout.Token).GetAsyncEnumerator(timeout.Token);
|
|
var next = enumerator.MoveNextAsync().AsTask();
|
|
|
|
await liveFeed.PublishAsync(new OpenTelemetryBatch([], [Trace("trace-skip"), Trace("trace-keep")], [], [], [], []), timeout.Token);
|
|
|
|
Assert.True(await next);
|
|
Assert.Equal("trace-keep", enumerator.Current.Trace?.TraceId);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task LiveFeed_WhenServiceNameFilterIsSet_OnlyPublishesMatchingResourcesAndTraces()
|
|
{
|
|
var liveFeed = new InMemoryOpenTelemetryLiveFeed(OptionsFactory.Create(new OpenTelemetryDiagnosticsOptions()));
|
|
using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5));
|
|
await using var enumerator = liveFeed.SubscribeAsync(new OpenTelemetryTraceFilter { ServiceName = "api" }, timeout.Token).GetAsyncEnumerator(timeout.Token);
|
|
var next = enumerator.MoveNextAsync().AsTask();
|
|
|
|
await liveFeed.PublishAsync(new OpenTelemetryBatch(
|
|
[Resource("resource-skip", "worker")],
|
|
[Trace("trace-skip", ["resource-skip"])],
|
|
[], [], [], []), timeout.Token);
|
|
|
|
var completed = await Task.WhenAny(next, Task.Delay(TimeSpan.FromMilliseconds(150), timeout.Token));
|
|
Assert.NotSame(next, completed);
|
|
|
|
await liveFeed.PublishAsync(new OpenTelemetryBatch(
|
|
[Resource("resource-keep", "api")],
|
|
[Trace("trace-keep", ["resource-keep"])],
|
|
[], [], [], []), timeout.Token);
|
|
|
|
Assert.True(await next);
|
|
Assert.Equal("resource-keep", enumerator.Current.Resource?.Id);
|
|
|
|
Assert.True(await enumerator.MoveNextAsync());
|
|
Assert.Equal("trace-keep", enumerator.Current.Trace?.TraceId);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task LiveFeed_WhenResourceFilterIsSet_FiltersLogsAndMetricPoints()
|
|
{
|
|
var liveFeed = new InMemoryOpenTelemetryLiveFeed(OptionsFactory.Create(new OpenTelemetryDiagnosticsOptions()));
|
|
using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5));
|
|
await using var enumerator = liveFeed.SubscribeAsync(new OpenTelemetryTraceFilter { ResourceId = "resource-keep" }, timeout.Token).GetAsyncEnumerator(timeout.Token);
|
|
var first = enumerator.MoveNextAsync().AsTask();
|
|
|
|
await liveFeed.PublishAsync(new OpenTelemetryBatch(
|
|
[],
|
|
[],
|
|
[],
|
|
[],
|
|
[MetricPoint("point-skip", "resource-skip"), MetricPoint("point-keep", "resource-keep")],
|
|
[Log("log-skip", "resource-skip"), Log("log-keep", "resource-keep")]), timeout.Token);
|
|
|
|
Assert.True(await first);
|
|
Assert.Equal("log-keep", enumerator.Current.Log?.Id);
|
|
|
|
Assert.True(await enumerator.MoveNextAsync());
|
|
Assert.Equal("point-keep", enumerator.Current.MetricPoint?.Id);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task LiveFeed_WhenSubscriberQueueOverflows_PublishesDroppedItemSummary()
|
|
{
|
|
var liveFeed = new InMemoryOpenTelemetryLiveFeed(OptionsFactory.Create(new OpenTelemetryDiagnosticsOptions { SubscriberChannelCapacity = 1 }));
|
|
using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5));
|
|
await using var enumerator = liveFeed.SubscribeAsync(new OpenTelemetryTraceFilter(), timeout.Token).GetAsyncEnumerator(timeout.Token);
|
|
var first = enumerator.MoveNextAsync().AsTask();
|
|
|
|
await liveFeed.PublishAsync(new OpenTelemetryBatch([], [Trace("trace-1"), Trace("trace-2"), Trace("trace-3")], [], [], [], []), timeout.Token);
|
|
|
|
Assert.True(await first);
|
|
|
|
OpenTelemetryStreamItem? summary = enumerator.Current.DroppedItems != null ? enumerator.Current : null;
|
|
for (var i = 0; summary == null && i < 5 && await enumerator.MoveNextAsync(); i++)
|
|
{
|
|
if (enumerator.Current.DroppedItems != null)
|
|
summary = enumerator.Current;
|
|
}
|
|
|
|
Assert.NotNull(summary);
|
|
var nonNullSummary = summary!;
|
|
Assert.Equal(OpenTelemetrySignalType.Trace, nonNullSummary.DroppedItems!.SignalType);
|
|
Assert.Equal("SubscriberQueueFull", nonNullSummary.DroppedItems.Reason);
|
|
Assert.True(nonNullSummary.DroppedItems.Count > 0);
|
|
}
|
|
|
|
private OpenTelemetryHub CreateHub(IOpenTelemetryLiveFeed liveFeed, string permission, IOpenTelemetryClient? caller = null)
|
|
{
|
|
caller ??= new CapturingOpenTelemetryClient();
|
|
var hubContext = new TestHubContext(caller);
|
|
var subscriptionManager = new OpenTelemetrySubscriptionManager(liveFeed, hubContext, NullLogger<OpenTelemetrySubscriptionManager>.Instance);
|
|
|
|
return new OpenTelemetryHub(subscriptionManager)
|
|
{
|
|
Context = new TestHubCallerContext(CreateUser(permission)),
|
|
Clients = new TestHubCallerClients(caller)
|
|
};
|
|
}
|
|
|
|
private ClaimsPrincipal CreateUser(string permission)
|
|
{
|
|
var permissionClaimType = (string)typeof(SecurityOptions)
|
|
.GetProperty(nameof(SecurityOptions.PermissionsClaimType))!
|
|
.GetValue(new Config().Security)!;
|
|
var identity = new ClaimsIdentity([new Claim(permissionClaimType, permission)], "Test");
|
|
|
|
return new ClaimsPrincipal(identity);
|
|
}
|
|
|
|
private TelemetryResource Resource(string id, string serviceName)
|
|
{
|
|
return new(id, serviceName, null, null, new Dictionary<string, string?>(), _now, TelemetryResourceStatus.Active);
|
|
}
|
|
|
|
private TelemetryTrace Trace(string traceId, IReadOnlyCollection<string>? resourceIds = null, IReadOnlyCollection<string>? workflowInstanceIds = null)
|
|
{
|
|
return new(traceId, $"{traceId}-root", traceId, _now, _now.AddMilliseconds(10), TimeSpan.FromMilliseconds(10), SpanStatus.Ok, resourceIds ?? [], workflowInstanceIds ?? [], 1);
|
|
}
|
|
|
|
private MetricPoint MetricPoint(string id, string resourceId)
|
|
{
|
|
return new(id, $"{id}-instrument", $"{id}-instrument", resourceId, _now, 1, null, null, new Dictionary<string, string?>(), null, null);
|
|
}
|
|
|
|
private OtlpLogRecord Log(string id, string resourceId)
|
|
{
|
|
return new(id, resourceId, _now, "Information", 9, id, null, null, new Dictionary<string, string?>());
|
|
}
|
|
|
|
private class TestLiveFeed(params OpenTelemetryStreamItem[] items) : IOpenTelemetryLiveFeed
|
|
{
|
|
public OpenTelemetryTraceFilter? Filter { get; private set; }
|
|
|
|
public ValueTask PublishAsync(OpenTelemetryBatch batch, CancellationToken cancellationToken = default) => ValueTask.CompletedTask;
|
|
|
|
public async IAsyncEnumerable<OpenTelemetryStreamItem> SubscribeAsync(OpenTelemetryTraceFilter filter, [System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken cancellationToken = default)
|
|
{
|
|
Filter = filter;
|
|
|
|
foreach (var item in items)
|
|
{
|
|
cancellationToken.ThrowIfCancellationRequested();
|
|
yield return item;
|
|
await Task.Yield();
|
|
}
|
|
}
|
|
}
|
|
|
|
private class CapturingOpenTelemetryClient : IOpenTelemetryClient
|
|
{
|
|
public List<OpenTelemetryStreamItem> Items { get; } = [];
|
|
|
|
public Task ReceiveAsync(OpenTelemetryStreamItem item)
|
|
{
|
|
Items.Add(item);
|
|
return Task.CompletedTask;
|
|
}
|
|
}
|
|
|
|
private class TestHubCallerClients(IOpenTelemetryClient caller) : IHubCallerClients<IOpenTelemetryClient>
|
|
{
|
|
public IOpenTelemetryClient Caller { get; } = caller;
|
|
public IOpenTelemetryClient Others => throw new NotSupportedException();
|
|
public IOpenTelemetryClient All => throw new NotSupportedException();
|
|
public IOpenTelemetryClient AllExcept(IReadOnlyList<string> excludedConnectionIds) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Client(string connectionId) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Clients(IReadOnlyList<string> connectionIds) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Group(string groupName) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient GroupExcept(string groupName, IReadOnlyList<string> excludedConnectionIds) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Groups(IReadOnlyList<string> groupNames) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient OthersInGroup(string groupName) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient User(string userId) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Users(IReadOnlyList<string> userIds) => throw new NotSupportedException();
|
|
}
|
|
|
|
private class TestHubContext(IOpenTelemetryClient caller) : IHubContext<OpenTelemetryHub, IOpenTelemetryClient>
|
|
{
|
|
public IHubClients<IOpenTelemetryClient> Clients { get; } = new TestHubClients(caller);
|
|
|
|
public IGroupManager Groups { get; } = new TestGroupManager();
|
|
}
|
|
|
|
private class TestHubClients(IOpenTelemetryClient caller) : IHubClients<IOpenTelemetryClient>
|
|
{
|
|
public IOpenTelemetryClient All => throw new NotSupportedException();
|
|
public IOpenTelemetryClient AllExcept(IReadOnlyList<string> excludedConnectionIds) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Client(string connectionId) => caller;
|
|
public IOpenTelemetryClient Clients(IReadOnlyList<string> connectionIds) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Group(string groupName) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient GroupExcept(string groupName, IReadOnlyList<string> excludedConnectionIds) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Groups(IReadOnlyList<string> groupNames) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient User(string userId) => throw new NotSupportedException();
|
|
public IOpenTelemetryClient Users(IReadOnlyList<string> userIds) => throw new NotSupportedException();
|
|
}
|
|
|
|
private class TestGroupManager : IGroupManager
|
|
{
|
|
public Task AddToGroupAsync(string connectionId, string groupName, CancellationToken cancellationToken = default) => Task.CompletedTask;
|
|
|
|
public Task RemoveFromGroupAsync(string connectionId, string groupName, CancellationToken cancellationToken = default) => Task.CompletedTask;
|
|
}
|
|
|
|
private static async Task AssertEventuallyAsync(Action assertion)
|
|
{
|
|
var deadline = DateTimeOffset.UtcNow.AddSeconds(3);
|
|
Exception? lastException = null;
|
|
|
|
while (DateTimeOffset.UtcNow < deadline)
|
|
{
|
|
try
|
|
{
|
|
assertion();
|
|
return;
|
|
}
|
|
catch (XunitException e)
|
|
{
|
|
lastException = e;
|
|
await Task.Delay(25);
|
|
}
|
|
}
|
|
|
|
if (lastException != null)
|
|
throw lastException;
|
|
}
|
|
|
|
private class TestHubCallerContext(ClaimsPrincipal user) : HubCallerContext
|
|
{
|
|
public override string ConnectionId { get; } = "connection-1";
|
|
public override string? UserIdentifier { get; } = "user-1";
|
|
public override ClaimsPrincipal? User { get; } = user;
|
|
public override IDictionary<object, object?> Items { get; } = new Dictionary<object, object?>();
|
|
public override IFeatureCollection Features { get; } = new FeatureCollection();
|
|
public override CancellationToken ConnectionAborted { get; } = CancellationToken.None;
|
|
|
|
public override void Abort()
|
|
{
|
|
}
|
|
}
|
|
}
|