Implement Blob Storage Workflow Provider

This commit is contained in:
Sipke Schoorstra 2021-06-15 20:51:53 +02:00
parent ebaeabb934
commit 618f4e837c
7 changed files with 107 additions and 13 deletions

View file

@ -102,7 +102,8 @@ namespace Elsa.Activities.Http
)]
public ICollection<int>? SupportedStatusCodes { get; set; } = new HashSet<int>(new[] { 200 });
[ActivityOutput] public HttpResponseModel? Output { get; set; }
[ActivityOutput] public HttpResponseModel? Response { get; set; }
[ActivityOutput] public object? ResponseContent { get; set; }
protected override async ValueTask<IActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext context)
{
@ -123,7 +124,7 @@ namespace Elsa.Activities.Http
if (hasContent && ReadContent)
{
var formatter = SelectContentParser(contentType);
responseModel.Content = await formatter.ReadAsync(response, cancellationToken);
ResponseContent = await formatter.ReadAsync(response, cancellationToken);
}
var statusCode = (int) response.StatusCode;
@ -135,7 +136,7 @@ namespace Elsa.Activities.Http
if (!isSupportedStatusCode)
outcomes.Add("Unsupported Status Code");
Output = responseModel;
Response = responseModel;
return Outcomes(outcomes);
}

View file

@ -7,6 +7,5 @@ namespace Elsa.Activities.Http.Models
{
public HttpStatusCode StatusCode { get; set; }
public Dictionary<string, string[]> Headers { get; set; } = new();
public object? Content { get; set; }
}
}

View file

@ -196,14 +196,15 @@ namespace Microsoft.Extensions.DependencyInjection
// Workflow providers.
services
.AddWorkflowProvider<ProgrammaticWorkflowProvider>()
.AddWorkflowProvider<StorageWorkflowProvider>()
.AddWorkflowProvider<BlobStorageWorkflowProvider>()
.AddWorkflowProvider<DatabaseWorkflowProvider>();
// Workflow Storage Providers.
services
.AddSingleton<IWorkflowStorageService, WorkflowStorageService>()
.AddWorkflowStorageProvider<TransientWorkflowStorageProvider>()
.AddWorkflowStorageProvider<WorkflowInstanceWorkflowStorageProvider>();
.AddWorkflowStorageProvider<WorkflowInstanceWorkflowStorageProvider>()
.AddWorkflowStorageProvider<BlobStorageWorkflowStorageProvider>();
// Metadata.
services

View file

@ -10,5 +10,6 @@ namespace Elsa
{
public static ISetupActivity<T> WithTransientStorageFor<T, TProperty>(this ISetupActivity<T> builder, Expression<Func<T, TProperty?>> propertyAccessor) where T : IActivity => builder.WithStorageFor(propertyAccessor, TransientWorkflowStorageProvider.ProviderName);
public static ISetupActivity<T> WithWorkflowInstanceStorageFor<T, TProperty>(this ISetupActivity<T> builder, Expression<Func<T, TProperty?>> propertyAccessor) where T : IActivity => builder.WithStorageFor(propertyAccessor, WorkflowInstanceWorkflowStorageProvider.ProviderName);
public static ISetupActivity<T> WithBlobStorageFor<T, TProperty>(this ISetupActivity<T> builder, Expression<Func<T, TProperty?>> propertyAccessor) where T : IActivity => builder.WithStorageFor(propertyAccessor, BlobStorageWorkflowStorageProvider.ProviderName);
}
}

View file

@ -0,0 +1,96 @@
using System.IO;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Serialization;
using Newtonsoft.Json;
using Storage.Net.Blobs;
namespace Elsa.Providers.WorkflowStorage
{
public class BlobStorageWorkflowStorageProvider : WorkflowStorageProvider
{
public const string ProviderName = "BlobStorage";
private readonly IBlobStorage _blobStorage;
private readonly JsonSerializerSettings _serializerSettings;
public BlobStorageWorkflowStorageProvider(IBlobStorage blobStorage)
{
_blobStorage = blobStorage;
_serializerSettings = DefaultContentSerializer.CreateDefaultJsonSerializationSettings();
_serializerSettings.TypeNameHandling = TypeNameHandling.All;
}
public override string DisplayName => "Blob Storage";
public override async ValueTask SaveAsync(WorkflowStorageContext context, string key, object? value, CancellationToken cancellationToken = default)
{
if (value == null)
return;
var path = GetFullPath(context, key);
if(value is Stream stream)
{
await _blobStorage.WriteAsync(path, stream, cancellationToken: cancellationToken);
var blob = await _blobStorage.GetBlobAsync(path, cancellationToken);
blob.Metadata["ContentType"] = "Binary";
}
else if (value is byte[] bytes)
{
await _blobStorage.WriteAsync(path, bytes, cancellationToken: cancellationToken);
var blob = await _blobStorage.GetBlobAsync(path, cancellationToken);
blob.Metadata["ContentType"] = "Binary";
}
else
{
var json = JsonConvert.SerializeObject(value, _serializerSettings);
var jsonBytes = Encoding.UTF8.GetBytes(json);
await _blobStorage.WriteAsync(path, jsonBytes, cancellationToken: cancellationToken);
var blob = await _blobStorage.GetBlobAsync(path, cancellationToken);
blob.Metadata["ContentType"] = "Json";
}
}
public override async ValueTask<object?> LoadAsync(WorkflowStorageContext context, string key, CancellationToken cancellationToken = default)
{
var path = GetFullPath(context, key);
if (!await _blobStorage.ExistsAsync(path, cancellationToken))
return null;
var blob = await _blobStorage.GetBlobAsync(path, cancellationToken);
var contentType = blob.Metadata.GetItem("ContentType") ?? "Json";
if (contentType == "Json")
{
var json = await _blobStorage.ReadTextAsync(path, cancellationToken: cancellationToken);
return JsonConvert.DeserializeObject(json, _serializerSettings);
}
return await _blobStorage.ReadBytesAsync(path, cancellationToken);
}
public override async ValueTask DeleteAsync(WorkflowStorageContext context, string key, CancellationToken cancellationToken = default)
{
var path = GetFullPath(context, key);
await _blobStorage.DeleteAsync(path, cancellationToken);
}
public override async ValueTask DeleteAsync(WorkflowStorageContext context, CancellationToken cancellationToken = default)
{
var path = GetContainerPath(context);
await _blobStorage.DeleteAsync(path, cancellationToken);
}
private string GetFullPath(WorkflowStorageContext context, string key)
{
var containerPath = GetContainerPath(context);
var activityId = context.ActivityId;
return $"${containerPath}/{activityId}/{key}.dat";
}
private string GetContainerPath(WorkflowStorageContext context) => context.WorkflowInstance.Id;
}
}

View file

@ -38,11 +38,7 @@ namespace Elsa.Providers.WorkflowStorage
return new ValueTask();
}
private IDictionary<string, object> GetData(WorkflowStorageContext context)
{
return context.WorkflowInstance.ActivityData.GetItem(context.ActivityId, () => new Dictionary<string, object>());
}
private IDictionary<string, object> GetData(WorkflowStorageContext context) => context.WorkflowInstance.ActivityData.GetItem(context.ActivityId, () => new Dictionary<string, object>());
private void SetState(WorkflowStorageContext context, string propertyName, object? value) => GetData(context)!.SetState(propertyName, value);
public object? GetState(WorkflowStorageContext context, string propertyName) => GetData(context)!.GetState(propertyName);
}

View file

@ -13,14 +13,14 @@ using Storage.Net.Blobs;
namespace Elsa.Providers.Workflows
{
public class StorageWorkflowProvider : WorkflowProvider
public class BlobStorageWorkflowProvider : WorkflowProvider
{
private readonly IBlobStorage _storage;
private readonly IWorkflowBlueprintMaterializer _workflowBlueprintMaterializer;
private readonly IContentSerializer _contentSerializer;
private readonly ILogger _logger;
public StorageWorkflowProvider(IBlobStorage storage, IWorkflowBlueprintMaterializer workflowBlueprintMaterializer, IContentSerializer contentSerializer, ILogger<StorageWorkflowProvider> logger)
public BlobStorageWorkflowProvider(IBlobStorage storage, IWorkflowBlueprintMaterializer workflowBlueprintMaterializer, IContentSerializer contentSerializer, ILogger<BlobStorageWorkflowProvider> logger)
{
_storage = storage;
_workflowBlueprintMaterializer = workflowBlueprintMaterializer;