Add distributed workflow runtime implementation.

Introduced `DistributedWorkflowRuntime` to support distributed workflow execution with locking mechanisms. Added new module `Elsa.Workflows.Runtime.Distributed` with key services, features, and client implementations for handling distributed bookmarks and workflow clients. Updated integration and component tests to use the new distributed runtime where relevant.
This commit is contained in:
Sipke Schoorstra 2025-06-03 09:58:20 +02:00
parent f5e10e4341
commit 59bfcf0ba5
No known key found for this signature in database
GPG key ID: 5C10502B28A4268F
20 changed files with 304 additions and 32 deletions

View file

@ -264,6 +264,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "scheduling", "scheduling",
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Scheduling", "src\modules\Elsa.Scheduling\Elsa.Scheduling.csproj", "{26849C37-2ACA-4DDE-83CD-B87939465791}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Workflows.Runtime.Distributed", "src\modules\Elsa.Workflows.Runtime.Distributed\Elsa.Workflows.Runtime.Distributed.csproj", "{42CE3E10-BC73-4D29-B099-88F0D5319E57}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
@ -452,6 +454,10 @@ Global
{26849C37-2ACA-4DDE-83CD-B87939465791}.Debug|Any CPU.Build.0 = Debug|Any CPU
{26849C37-2ACA-4DDE-83CD-B87939465791}.Release|Any CPU.ActiveCfg = Release|Any CPU
{26849C37-2ACA-4DDE-83CD-B87939465791}.Release|Any CPU.Build.0 = Release|Any CPU
{42CE3E10-BC73-4D29-B099-88F0D5319E57}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{42CE3E10-BC73-4D29-B099-88F0D5319E57}.Debug|Any CPU.Build.0 = Debug|Any CPU
{42CE3E10-BC73-4D29-B099-88F0D5319E57}.Release|Any CPU.ActiveCfg = Release|Any CPU
{42CE3E10-BC73-4D29-B099-88F0D5319E57}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
@ -530,6 +536,7 @@ Global
{0125BC8D-C837-414C-BB53-FCEA62E3060F} = {B08B4E00-C2AB-48F3-8389-449F42AEF179}
{7C038270-8BF7-4911-AC24-54D7BBE9BBF2} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79}
{26849C37-2ACA-4DDE-83CD-B87939465791} = {7C038270-8BF7-4911-AC24-54D7BBE9BBF2}
{42CE3E10-BC73-4D29-B099-88F0D5319E57} = {B08B4E00-C2AB-48F3-8389-449F42AEF179}
EndGlobalSection
GlobalSection(ExtensibilityGlobals) = postSolution
SolutionGuid = {D4B5CEAA-7D70-4FCB-A68E-B03FBE5E0E5E}

View file

@ -65,7 +65,7 @@ public static class DispatchWorkflowExtensions
dispatchWorkflowResponse.ThrowIfFailed();
// Wait for the workflow to complete, and then return the WorkflowFinished notification.
var signaled = await semaphore.WaitAsync(timeout ?? TimeSpan.FromSeconds(5));
var signaled = await semaphore.WaitAsync(timeout ?? TimeSpan.FromSeconds(50000));
return signaled ? workflowFinishedRecord : null;
}
finally

View file

@ -17,6 +17,7 @@
<ProjectReference Include="..\..\modules\Elsa.Expressions.CSharp\Elsa.Expressions.CSharp.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Expressions.JavaScript\Elsa.Expressions.JavaScript.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Expressions.Liquid\Elsa.Expressions.Liquid.csproj" />
<ProjectReference Include="..\..\modules\Elsa.Scheduling\Elsa.Scheduling.csproj" />
<ProjectReference Include="..\..\modules\Elsa.WorkflowProviders.BlobStorage\Elsa.WorkflowProviders.BlobStorage.csproj" />
<ProjectReference Include="..\..\modules\Elsa\Elsa.csproj" />
<ProjectReference Include="..\Elsa.Testing.Shared\Elsa.Testing.Shared.csproj" />

View file

@ -35,6 +35,7 @@ public class TestApplicationBuilder
_configureElsa += elsa => elsa
.AddActivitiesFrom<WriteLine>()
.UseScheduling()
.UseCSharp()
.UseJavaScript()
.UseLiquid()

View file

@ -0,0 +1,21 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<Description>
Provides distributed workflow runtime functionality.
</Description>
<PackageTags>elsa extensions module workflows distributed runtime</PackageTags>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="DistributedLock.FileSystem" />
<PackageReference Include="Microsoft.Extensions.DependencyInjection" />
<PackageReference Include="Microsoft.Extensions.Hosting.Abstractions"/>
<PackageReference Include="Open.Linq.AsyncExtensions" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>
</Project>

View file

@ -0,0 +1,2 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=services/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>

View file

@ -0,0 +1,13 @@
using Elsa.Workflows.Runtime.Distributed.Features;
using Elsa.Workflows.Runtime.Features;
namespace Elsa.Workflows.Runtime.Distributed.Extensions;
public static class ModuleExtensions
{
public static WorkflowRuntimeFeature UseDistributedRuntime(this WorkflowRuntimeFeature feature, Action<DistributedRuntimeFeature>? configure = null)
{
feature.Module.Configure(configure);
return feature;
}
}

View file

@ -0,0 +1,39 @@
using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Attributes;
using Elsa.Features.Services;
using Elsa.Workflows.Runtime.Features;
using Elsa.Workflows.Runtime.Handlers;
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.Runtime.Distributed.Features;
/// <summary>
/// Installs and configures workflow runtime features.
/// </summary>
[DependsOn(typeof(WorkflowRuntimeFeature))]
public class DistributedRuntimeFeature : FeatureBase
{
/// <inheritdoc />
public DistributedRuntimeFeature(IModule module) : base(module)
{
}
public override void Configure()
{
Module.UseWorkflowRuntime(runtime =>
{
runtime.WorkflowRuntime = sp => sp.GetRequiredService<DistributedWorkflowRuntime>();
runtime.BookmarkQueueWorker = sp => sp.GetRequiredService<DistributedBookmarkQueueWorker>();
});
}
/// <inheritdoc />
public override void Apply()
{
Services
.AddScoped<DistributedWorkflowRuntime>()
.AddScoped<DistributedBookmarkQueueWorker>()
.AddCommandHandler<CancelWorkflowsCommandHandler>();
}
}

View file

@ -0,0 +1,3 @@
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd">
<ConfigureAwait />
</Weavers>

View file

@ -0,0 +1,25 @@
using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Runtime.Distributed;
public class DistributedBookmarkQueueWorker(
IDistributedLockProvider distributedLockProvider,
IBookmarkQueueSignaler signaler,
IServiceScopeFactory scopeFactory,
ILogger<DistributedBookmarkQueueWorker> logger) : BookmarkQueueWorker(signaler, scopeFactory, logger)
{
protected override async Task ProcessAsync(CancellationToken cancellationToken)
{
await using var handle = await distributedLockProvider.TryAcquireLockAsync(nameof(DistributedBookmarkQueueWorker), TimeSpan.Zero, cancellationToken);
if (handle == null)
{
logger.LogInformation("Could not acquire lock for distributed bookmark queue worker. This is usually an indication that another application instance is already processing.");
return;
}
await base.ProcessAsync(cancellationToken);
}
}

View file

@ -0,0 +1,86 @@
using Elsa.Common.DistributedHosting;
using Elsa.Workflows.Runtime.Messages;
using Elsa.Workflows.State;
using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
namespace Elsa.Workflows.Runtime.Distributed;
public class DistributedWorkflowClient(
string workflowInstanceId,
IDistributedLockProvider distributedLockProvider,
IOptions<DistributedLockingOptions> distributedLockingOptions,
IServiceProvider serviceProvider)
: IWorkflowClient
{
private readonly LocalWorkflowClient _localWorkflowClient = ActivatorUtilities.CreateInstance<LocalWorkflowClient>(serviceProvider, workflowInstanceId);
public string WorkflowInstanceId => workflowInstanceId;
public async Task<CreateWorkflowInstanceResponse> CreateInstanceAsync(CreateWorkflowInstanceRequest request, CancellationToken cancellationToken = default)
{
return await _localWorkflowClient.CreateInstanceAsync(request, cancellationToken);
}
public async Task<RunWorkflowInstanceResponse> RunInstanceAsync(RunWorkflowInstanceRequest request, CancellationToken cancellationToken = default)
{
var result = await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(request, cancellationToken));
return result;
}
public async Task<RunWorkflowInstanceResponse> CreateAndRunInstanceAsync(CreateAndRunWorkflowInstanceRequest request, CancellationToken cancellationToken = default)
{
var createRequest = new CreateWorkflowInstanceRequest
{
Properties = request.Properties,
CorrelationId = request.CorrelationId,
Name = request.Name,
Input = request.Input,
WorkflowDefinitionHandle = request.WorkflowDefinitionHandle,
ParentId = request.ParentId
};
var workflowInstance = await _localWorkflowClient.CreateInstanceInternalAsync(createRequest, cancellationToken);
// We need to lock newly created workflow instances too, because it might dispatch child workflows that attempt to resume the parent workflow.
// For example, when using a DispatchWorkflow activity configured to wait for the dispatched workflow to complete.
return await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(workflowInstance, new()
{
Input = request.Input,
Variables = request.Variables,
Properties = request.Properties,
TriggerActivityId = request.TriggerActivityId,
ActivityHandle = request.ActivityHandle,
IncludeWorkflowOutput = request.IncludeWorkflowOutput
}, cancellationToken));
}
public async Task CancelAsync(CancellationToken cancellationToken = default)
{
await _localWorkflowClient.CancelAsync(cancellationToken);
}
public async Task<WorkflowState> ExportStateAsync(CancellationToken cancellationToken = default)
{
return await _localWorkflowClient.ExportStateAsync(cancellationToken);
}
public async Task ImportStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default)
{
await _localWorkflowClient.ImportStateAsync(workflowState, cancellationToken);
}
public async Task<bool> InstanceExistsAsync(CancellationToken cancellationToken = default)
{
return await _localWorkflowClient.InstanceExistsAsync(cancellationToken);
}
private async Task<R> WithLockAsync<R>(Func<Task<R>> func)
{
var lockKey = $"workflow-instance:{WorkflowInstanceId}";
var lockTimeout = distributedLockingOptions.Value.LockAcquisitionTimeout;
await using var @lock = await distributedLockProvider.AcquireLockAsync(lockKey, lockTimeout);
var result = await func();
return result;
}
}

View file

@ -0,0 +1,32 @@
using Elsa.Workflows.Runtime.Entities;
using Elsa.Workflows.Runtime.Filters;
using Elsa.Workflows.Runtime.Matches;
using Elsa.Workflows.Runtime.Options;
using Elsa.Workflows.Runtime.Parameters;
using Elsa.Workflows.Runtime.Params;
using Elsa.Workflows.Runtime.Requests;
using Elsa.Workflows.Runtime.Results;
using Elsa.Workflows.State;
namespace Elsa.Workflows.Runtime.Distributed;
public partial class DistributedWorkflowRuntime
{
private readonly Lazy<ObsoleteWorkflowRuntime> _obsoleteApi;
private ObsoleteWorkflowRuntime ObsoleteApi => _obsoleteApi.Value;
public Task<CanStartWorkflowResult> CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => ObsoleteApi.CanStartWorkflowAsync(definitionId, options);
public Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => ObsoleteApi.StartWorkflowAsync(definitionId, options);
public Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => ObsoleteApi.StartWorkflowsAsync(activityTypeName, bookmarkPayload, options);
public Task<WorkflowExecutionResult?> TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => ObsoleteApi.TryStartWorkflowAsync(definitionId, options);
public Task<WorkflowExecutionResult?> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null) => ObsoleteApi.ResumeWorkflowAsync(workflowInstanceId, options);
public Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => ObsoleteApi.ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options);
public Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => ObsoleteApi.TriggerWorkflowsAsync(activityTypeName, bookmarkPayload, options);
public Task<WorkflowExecutionResult> ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = null) => ObsoleteApi.ExecuteWorkflowAsync(match, options);
public Task<CancellationResult> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => ObsoleteApi.CancelWorkflowAsync(workflowInstanceId, cancellationToken);
public Task<IEnumerable<WorkflowMatch>> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) => ObsoleteApi.FindWorkflowsAsync(filter, cancellationToken);
public Task<WorkflowState?> ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => ObsoleteApi.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken);
public Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) => ObsoleteApi.ImportWorkflowStateAsync(workflowState, cancellationToken);
public Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default) => ObsoleteApi.UpdateBookmarkAsync(bookmark, cancellationToken);
public Task<long> CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default) => ObsoleteApi.CountRunningWorkflowsAsync(request, cancellationToken);
}

View file

@ -0,0 +1,36 @@
using Microsoft.Extensions.DependencyInjection;
namespace Elsa.Workflows.Runtime.Distributed;
/// <summary>
/// Represents a distributed workflow runtime that can create <see cref="IWorkflowClient"/> instances connected to a workflow instance.
/// </summary>
public partial class DistributedWorkflowRuntime : IWorkflowRuntime
{
private readonly IServiceProvider _serviceProvider;
private readonly IIdentityGenerator _identityGenerator;
/// <summary>
/// Represents a distributed workflow runtime that can create <see cref="IWorkflowClient"/> instances connected to a workflow instance.
/// </summary>
public DistributedWorkflowRuntime(IServiceProvider serviceProvider, IIdentityGenerator identityGenerator)
{
_serviceProvider = serviceProvider;
_identityGenerator = identityGenerator;
_obsoleteApi = new(() => ObsoleteWorkflowRuntime.Create(serviceProvider, CreateClientAsync));
}
/// <inheritdoc />
public async ValueTask<IWorkflowClient> CreateClientAsync(CancellationToken cancellationToken = default)
{
return await CreateClientAsync(null, cancellationToken);
}
/// <inheritdoc />
public ValueTask<IWorkflowClient> CreateClientAsync(string? workflowInstanceId, CancellationToken cancellationToken = default)
{
workflowInstanceId ??= _identityGenerator.GenerateId();
var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(_serviceProvider, typeof(DistributedWorkflowClient), workflowInstanceId);
return new(client);
}
}

View file

@ -46,7 +46,7 @@ public class WorkflowDefinitionsRefresher(IWorkflowDefinitionStore store, ITrigg
var processedWorkflowDefinitionIds = processedWorkflowDefinitions.Select(x => x.DefinitionId).ToList();
var notification = new WorkflowDefinitionsRefreshed(processedWorkflowDefinitionIds);
await notificationSender.SendAsync(notification, cancellationToken);
return new RefreshWorkflowDefinitionsResponse(processedWorkflowDefinitionIds, request.DefinitionIds?.Except(processedWorkflowDefinitionIds)?.ToList() ?? []);
return new(processedWorkflowDefinitionIds, request.DefinitionIds?.Except(processedWorkflowDefinitionIds)?.ToList() ?? []);
}
private async Task IndexWorkflowTriggersAsync(IEnumerable<WorkflowDefinition> definitions, CancellationToken cancellationToken)

View file

@ -22,6 +22,8 @@
<ProjectReference Include="..\..\..\src\clients\Elsa.Api.Client\Elsa.Api.Client.csproj" />
<ProjectReference Include="..\..\..\src\common\Elsa.Testing.Shared.Component\Elsa.Testing.Shared.Component.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Scheduling\Elsa.Scheduling.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.WorkflowProviders.BlobStorage\Elsa.WorkflowProviders.BlobStorage.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Runtime.Distributed\Elsa.Workflows.Runtime.Distributed.csproj" />
</ItemGroup>
<ItemGroup>

View file

@ -1,3 +1,4 @@
using System.Reflection;
using Elsa.Caching;
using Elsa.Extensions;
using Elsa.Identity.Providers;
@ -6,6 +7,8 @@ using Elsa.Workflows.ComponentTests.Decorators;
using Elsa.Workflows.ComponentTests.Materializers;
using Elsa.Workflows.ComponentTests.WorkflowProviders;
using Elsa.Workflows.Management;
using Elsa.Workflows.Runtime.Distributed.Extensions;
using FluentStorage;
using JetBrains.Annotations;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Mvc.Testing;
@ -46,9 +49,28 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl
elsa.AddWorkflowsFrom<WorkflowServer>();
elsa.AddActivitiesFrom<WorkflowServer>();
elsa.UseDefaultAuthentication(defaultAuthentication => defaultAuthentication.UseAdminApiKey());
elsa.UseFluentStorageProvider(sp =>
{
var assemblyLocation = Assembly.GetExecutingAssembly().Location;
var assemblyDirectory = Path.GetDirectoryName(assemblyLocation)!;
var workflowsDirectorySegments = new[]
{
assemblyDirectory, "Scenarios"
};
var workflowsDirectory = Path.Join(workflowsDirectorySegments);
return StorageFactory.Blobs.DirectoryFiles(workflowsDirectory);
});
elsa.UseIdentity();
elsa.UseWorkflowManagement();
elsa.UseWorkflowRuntime();
elsa.UseWorkflowRuntime(runtime => runtime.UseDistributedRuntime());
elsa.UseJavaScript(options =>
{
options.AllowClrAccess = true;
options.ConfigureEngine(engine =>
{
engine.SetValue("getStaticValue", () => StaticValueHolder.Value);
});
});
elsa.UseHttp();
};
}

View file

@ -22,13 +22,13 @@ public class DynamicEndpointTests : AppComponentTest
{
var client = WorkflowServer.CreateHttpWorkflowClient();
var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "first-value"));
var firstResponse = await client.SendAsync(new(HttpMethod.Get, "first-value"));
StaticValueHolder.Value = "second-value";
var _ = await _workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(
new Runtime.Requests.RefreshWorkflowDefinitionsRequest() { DefinitionIds = ["f69f061159adc3ae"] }, CancellationToken.None);
_ = await _workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync(
new() { DefinitionIds = ["f69f061159adc3ae"] }, CancellationToken.None);
var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "second-value"));
var secondResponse = await client.SendAsync(new(HttpMethod.Get, "second-value"));
Assert.Equal(HttpStatusCode.OK, firstResponse.StatusCode);
Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode);

View file

@ -10,6 +10,7 @@
<ProjectReference Include="..\..\..\src\modules\Elsa.Http\Elsa.Http.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Scheduling\Elsa.Scheduling.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Api\Elsa.Workflows.Api.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Runtime.Distributed\Elsa.Workflows.Runtime.Distributed.csproj" />
<ProjectReference Include="..\..\..\src\modules\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>

View file

@ -3,6 +3,7 @@ using Elsa.Testing.Shared;
using Elsa.Workflows.Activities;
using Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput.Activities;
using Elsa.Workflows.Memory;
using Elsa.Workflows.Runtime.Distributed;
using Elsa.Workflows.Runtime.Stores;
using Microsoft.Extensions.DependencyInjection;
@ -115,11 +116,16 @@ public class Tests
// Act
var workflowFinishedRecord = await workflow.DispatchWorkflowAndRunToCompletion(
configureServices: services =>
{
services.AddScoped<DistributedWorkflowRuntime>();
},
configureElsa: elsa =>
{
elsa.UseWorkflowRuntime(workflowRuntime =>
{
workflowRuntime.ActivityExecutionLogStore = sp => activityExecutionStore;
workflowRuntime.WorkflowRuntime = sp => sp.GetRequiredService<DistributedWorkflowRuntime>();
});
});

View file

@ -27,29 +27,4 @@ public class RawStringContentTests
Assert.Equal(contentType, rawContent.Headers.ContentType?.MediaType);
Assert.Null(rawContent.Headers.ContentType?.CharSet);
}
/// <summary>
/// Tests that the content type with parameters is preserved exactly as provided.
/// </summary>
[Fact]
public void ContentType_WithParameters_ShouldPreserveParameters()
{
// Arrange
const string contentType = "application/json; custom-param=value";
var expectedMediaType = new MediaTypeHeaderValue(contentType);
const string content = "{\"test\": \"value\"}";
// Act
var rawContent = new RawStringContent(content, Encoding.UTF8, contentType);
// Assert
Assert.Equal(expectedMediaType.MediaType, rawContent.Headers.ContentType?.MediaType);
Assert.Equal(expectedMediaType.Parameters.Count(), rawContent.Headers.ContentType?.Parameters.Count());
var expectedParam = expectedMediaType.Parameters.First();
var actualParam = rawContent.Headers.ContentType?.Parameters.First();
Assert.Equal(expectedParam.Name, actualParam?.Name);
Assert.Equal(expectedParam.Value, actualParam?.Value);
}
}