diff --git a/src/common/Elsa.Persistence.Abstractions/Extensions/EnumerableExtensions.cs b/src/common/Elsa.Persistence.Abstractions/Extensions/EnumerableExtensions.cs new file mode 100644 index 000000000..204ef46b7 --- /dev/null +++ b/src/common/Elsa.Persistence.Abstractions/Extensions/EnumerableExtensions.cs @@ -0,0 +1,30 @@ +using Elsa.Persistence.Common.Models; + +namespace Elsa.Persistence.Common.Extensions; + +public static class EnumerableExtensions +{ + public static Page Paginate(this IEnumerable queryable, Func projection, PageArgs? pageArgs = default) + { + var items = queryable.ToList(); + var count = items.Count; + + if (pageArgs?.Offset != null) items = items.Skip(pageArgs.Offset.Value).ToList(); + if (pageArgs?.Limit != null) items = items.Take(pageArgs.Limit.Value).ToList(); + + var results = items.Select(projection).ToList(); + return Page.Of(results, count); + } + + public static Page Paginate(this IEnumerable queryable, PageArgs? pageArgs = default) + { + var items = queryable.ToList(); + var count = items.Count; + + if (pageArgs?.Offset != null) items = items.Skip(pageArgs.Offset.Value).ToList(); + if (pageArgs?.Limit != null) items = items.Take(pageArgs.Limit.Value).ToList(); + + var results = items.ToList(); + return Page.Of(results, count); + } +} \ No newline at end of file diff --git a/src/common/Elsa.Persistence.Abstractions/Implementations/MemoryStore.cs b/src/common/Elsa.Persistence.Abstractions/Implementations/MemoryStore.cs index 1f6199253..48039a0f8 100644 --- a/src/common/Elsa.Persistence.Abstractions/Implementations/MemoryStore.cs +++ b/src/common/Elsa.Persistence.Abstractions/Implementations/MemoryStore.cs @@ -15,6 +15,21 @@ public class MemoryStore where TEntity : Entity public TEntity? Find(Func predicate) => Entities.Values.Where(predicate).FirstOrDefault(); public IEnumerable FindMany(Func predicate) => Entities.Values.Where(predicate); + + public IEnumerable FindMany(Func predicate, Func orderBy, OrderDirection orderDirection = OrderDirection.Ascending) + { + var query = Entities.Values.Where(predicate); + + query = orderDirection switch + { + OrderDirection.Ascending => query.OrderBy(orderBy), + OrderDirection.Descending => query.OrderByDescending(orderBy), + _ => query.OrderBy(orderBy) + }; + + return query; + } + public IEnumerable List() => Entities.Values; public bool Delete(string id) => Entities.Remove(id); diff --git a/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Implementations/EFCoreStore.cs b/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Implementations/EFCoreStore.cs index b8589ae64..7dd5e864c 100644 --- a/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Implementations/EFCoreStore.cs +++ b/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Implementations/EFCoreStore.cs @@ -1,6 +1,8 @@ using System.Linq.Expressions; using EFCore.BulkExtensions; using Elsa.Persistence.Common.Entities; +using Elsa.Persistence.Common.Extensions; +using Elsa.Persistence.Common.Models; using Elsa.Persistence.EntityFrameworkCore.Common.Extensions; using Elsa.Persistence.EntityFrameworkCore.Common.Services; using Microsoft.EntityFrameworkCore; @@ -52,6 +54,28 @@ public class EFCoreStore : IStore wher return OnLoading(dbContext, entities); } + public async Task> FindManyAsync( + Expression> predicate, + Expression> orderBy, + OrderDirection orderDirection = OrderDirection.Ascending, + PageArgs? pageArgs = default, + CancellationToken cancellationToken = default) + { + await using var dbContext = await CreateDbContextAsync(cancellationToken); + var set = dbContext.Set().Where(predicate); + + set = orderDirection switch + { + OrderDirection.Ascending => set.OrderBy(orderBy), + OrderDirection.Descending => set.OrderByDescending(orderBy), + _ => set.OrderBy(orderBy) + }; + + var page = await set.PaginateAsync(pageArgs); + OnLoading(dbContext, page.Items); + return page; + } + public async Task DeleteAsync(TEntity entity, CancellationToken cancellationToken = default) { await using var dbContext = await CreateDbContextAsync(cancellationToken); @@ -82,8 +106,9 @@ public class EFCoreStore : IStore wher var queryable = query(set.AsQueryable()); queryable = query(queryable); + var entities = await queryable.ToListAsync(cancellationToken); - return Load(dbContext, queryable).ToList(); + return Load(dbContext, entities).ToList(); } public async Task AnyAsync(Expression> predicate, CancellationToken cancellationToken = default) diff --git a/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Services/IStore.cs b/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Services/IStore.cs index 63787a05c..b6829b75c 100644 --- a/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Services/IStore.cs +++ b/src/common/Elsa.Persistence.EntityFrameworkCore.Common/Services/IStore.cs @@ -1,5 +1,6 @@ using System.Linq.Expressions; using Elsa.Persistence.Common.Entities; +using Elsa.Persistence.Common.Models; namespace Elsa.Persistence.EntityFrameworkCore.Common.Services; @@ -10,6 +11,14 @@ public interface IStore where TEntity : Entity Task SaveManyAsync(IEnumerable entities, CancellationToken cancellationToken = default); Task FindAsync(Expression> predicate, CancellationToken cancellationToken = default); Task> FindManyAsync(Expression> predicate, CancellationToken cancellationToken = default); + + Task> FindManyAsync( + Expression> predicate, + Expression> orderBy, + OrderDirection orderDirection = OrderDirection.Ascending, + PageArgs? pageArgs = default, + CancellationToken cancellationToken = default); + Task DeleteAsync(TEntity entity, CancellationToken cancellationToken = default); Task DeleteManyAsync(IEnumerable entities, CancellationToken cancellationToken = default); Task DeleteWhereAsync(Expression> predicate, CancellationToken cancellationToken = default); diff --git a/src/modules/Elsa.Workflows.Api/ControllerNames.cs b/src/modules/Elsa.Workflows.Api/ControllerNames.cs index 4f3d5cd28..9f2be349d 100644 --- a/src/modules/Elsa.Workflows.Api/ControllerNames.cs +++ b/src/modules/Elsa.Workflows.Api/ControllerNames.cs @@ -8,4 +8,5 @@ public class ControllerNames public const string WorkflowInstances = "WorkflowInstances"; public const string Labels = "Labels"; public const string WorkflowDefinitionLabels = "WorkflowDefinitionLabels"; + public const string WorkflowJournal = "WorkflowJournal"; } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/Get.cs b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/Get.cs new file mode 100644 index 000000000..183e6e45b --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/WorkflowInstances/Journal/Get.cs @@ -0,0 +1,65 @@ +using System; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Elsa.AspNetCore.Attributes; +using Elsa.Persistence.Common.Models; +using Elsa.Workflows.Core.Serialization; +using Elsa.Workflows.Persistence.Services; +using Microsoft.AspNetCore.Mvc; + +// ReSharper disable NotAccessedPositionalProperty.Global + +namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal; + +[Area(AreaNames.Elsa)] +[ApiEndpoint(ControllerNames.WorkflowJournal, "Get")] +public class Get : Controller +{ + private readonly IWorkflowExecutionLogStore _store; + private readonly WorkflowSerializerOptionsProvider _serializerOptionsProvider; + + public Get(IWorkflowExecutionLogStore store, WorkflowSerializerOptionsProvider serializerOptionsProvider) + { + _store = store; + _serializerOptionsProvider = serializerOptionsProvider; + } + + [HttpGet] + public async Task HandleAsync( + [FromRoute(Name = "id")] string workflowInstanceId, + [FromQuery] int? page, + [FromQuery] int? pageSize, + CancellationToken cancellationToken) + { + var serializerOptions = _serializerOptionsProvider.CreateApiOptions(); + var pageArgs = new PageArgs(page, pageSize); + var pageOfRecords = await _store.FindManyByWorkflowInstanceIdAsync(workflowInstanceId, pageArgs, cancellationToken); + + var models = pageOfRecords.Items.Select(x => + new WorkflowExecutionLogRecordModel( + x.Id, + x.ActivityId, + x.ActivityType, + x.Timestamp, + x.EventName, + x.Message, + x.Source, + x.Payload)) + .ToList(); + + var pageOfModels = Page.Of(models, pageOfRecords.TotalCount); + + return Json(pageOfModels, serializerOptions); + } + + public record WorkflowExecutionLogRecordModel( + string Id, + string ActivityId, + string ActivityType, + DateTimeOffset Timestamp, + string? EventName, + string? Message, + string? Source, + object? Payload); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Extensions/EndpointRouteBuilderExtensions.cs b/src/modules/Elsa.Workflows.Api/Extensions/EndpointRouteBuilderExtensions.cs index 248334bbc..861c71120 100644 --- a/src/modules/Elsa.Workflows.Api/Extensions/EndpointRouteBuilderExtensions.cs +++ b/src/modules/Elsa.Workflows.Api/Extensions/EndpointRouteBuilderExtensions.cs @@ -40,6 +40,9 @@ public static class EndpointRouteBuilderExtensions Map("WorkflowInstances.Get", "workflow-instances/{id}", new { Controller = ControllerNames.WorkflowInstances, Action = "Get" }); Map("WorkflowInstances.Delete", "workflow-instances/{id}", new { Controller = ControllerNames.WorkflowInstances, Action = "Delete" }); Map("WorkflowInstances.List", "workflow-instances", new { Controller = ControllerNames.WorkflowInstances, Action = "List" }); + + // Workflow Journal. + Map("WorkflowJournal.Get", "workflow-instances/{id}/journal", new { Controller = ControllerNames.WorkflowJournal, Action = "Get" }); return endpoints; } diff --git a/src/modules/Elsa.Workflows.Core/Serialization/WorkflowSerializerOptionsProvider.cs b/src/modules/Elsa.Workflows.Core/Serialization/WorkflowSerializerOptionsProvider.cs index d4430dd84..84f96dd42 100644 --- a/src/modules/Elsa.Workflows.Core/Serialization/WorkflowSerializerOptionsProvider.cs +++ b/src/modules/Elsa.Workflows.Core/Serialization/WorkflowSerializerOptionsProvider.cs @@ -27,7 +27,8 @@ public class WorkflowSerializerOptionsProvider var options = new JsonSerializerOptions { PropertyNamingPolicy = JsonNamingPolicy.CamelCase, - PropertyNameCaseInsensitive = true + PropertyNameCaseInsensitive = true, + DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull }; options.Converters.Add(Create()); diff --git a/src/modules/Elsa.Workflows.Persistence.EntityFrameworkCore/Implementations/EFCoreWorkflowExecutionLogStore.cs b/src/modules/Elsa.Workflows.Persistence.EntityFrameworkCore/Implementations/EFCoreWorkflowExecutionLogStore.cs index 3f1509f35..25116a480 100644 --- a/src/modules/Elsa.Workflows.Persistence.EntityFrameworkCore/Implementations/EFCoreWorkflowExecutionLogStore.cs +++ b/src/modules/Elsa.Workflows.Persistence.EntityFrameworkCore/Implementations/EFCoreWorkflowExecutionLogStore.cs @@ -1,3 +1,5 @@ +using Elsa.Persistence.Common.Entities; +using Elsa.Persistence.Common.Models; using Elsa.Persistence.EntityFrameworkCore.Common.Services; using Elsa.Workflows.Persistence.Entities; using Elsa.Workflows.Persistence.Services; @@ -10,4 +12,16 @@ public class EFCoreWorkflowExecutionLogStore : IWorkflowExecutionLogStore public EFCoreWorkflowExecutionLogStore(IStore store) => _store = store; public async Task SaveAsync(WorkflowExecutionLogRecord record, CancellationToken cancellationToken = default) => await _store.SaveAsync(record, cancellationToken); public async Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default) => await _store.SaveManyAsync(records, cancellationToken); + + public async Task> FindManyByWorkflowInstanceIdAsync(string workflowInstanceId, PageArgs? pageArgs = default, CancellationToken cancellationToken = default) + { + var records = await _store.FindManyAsync( + x => x.WorkflowInstanceId == workflowInstanceId, + x => x.Timestamp, + OrderDirection.Ascending, + pageArgs, + cancellationToken); + + return records; + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Persistence/Implementations/MemoryWorkflowExecutionLogStore.cs b/src/modules/Elsa.Workflows.Persistence/Implementations/MemoryWorkflowExecutionLogStore.cs index bef2b4bbe..56e7096d4 100644 --- a/src/modules/Elsa.Workflows.Persistence/Implementations/MemoryWorkflowExecutionLogStore.cs +++ b/src/modules/Elsa.Workflows.Persistence/Implementations/MemoryWorkflowExecutionLogStore.cs @@ -1,4 +1,6 @@ +using Elsa.Persistence.Common.Extensions; using Elsa.Persistence.Common.Implementations; +using Elsa.Persistence.Common.Models; using Elsa.Workflows.Persistence.Entities; using Elsa.Workflows.Persistence.Services; @@ -24,4 +26,15 @@ public class MemoryWorkflowExecutionLogStore : IWorkflowExecutionLogStore _store.SaveMany(records); return Task.CompletedTask; } + + public Task> FindManyByWorkflowInstanceIdAsync(string workflowInstanceId, PageArgs? pageArgs = default, CancellationToken cancellationToken = default) + { + var page = _store + .FindMany( + x => x.WorkflowInstanceId == workflowInstanceId, + x => x.Timestamp) + .Paginate(); + + return Task.FromResult(page); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Persistence/Services/IWorkflowExecutionLogStore.cs b/src/modules/Elsa.Workflows.Persistence/Services/IWorkflowExecutionLogStore.cs index f2611c171..238250ca1 100644 --- a/src/modules/Elsa.Workflows.Persistence/Services/IWorkflowExecutionLogStore.cs +++ b/src/modules/Elsa.Workflows.Persistence/Services/IWorkflowExecutionLogStore.cs @@ -1,3 +1,4 @@ +using Elsa.Persistence.Common.Models; using Elsa.Workflows.Persistence.Entities; namespace Elsa.Workflows.Persistence.Services; @@ -6,4 +7,5 @@ public interface IWorkflowExecutionLogStore { Task SaveAsync(WorkflowExecutionLogRecord record, CancellationToken cancellationToken = default); Task SaveManyAsync(IEnumerable records, CancellationToken cancellationToken = default); + Task> FindManyByWorkflowInstanceIdAsync(string workflowInstanceId, PageArgs? pageArgs = default, CancellationToken cancellationToken = default); } \ No newline at end of file