Add journal API endpoint

This commit is contained in:
Sipke Schoorstra 2022-06-02 11:36:56 +02:00
parent 4d1349736c
commit 8fc9fb15b6
11 changed files with 180 additions and 2 deletions

View file

@ -0,0 +1,30 @@
using Elsa.Persistence.Common.Models;
namespace Elsa.Persistence.Common.Extensions;
public static class EnumerableExtensions
{
public static Page<TTarget> Paginate<T, TTarget>(this IEnumerable<T> queryable, Func<T, TTarget> 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<T> Paginate<T>(this IEnumerable<T> 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);
}
}

View file

@ -15,6 +15,21 @@ public class MemoryStore<TEntity> where TEntity : Entity
public TEntity? Find(Func<TEntity, bool> predicate) => Entities.Values.Where(predicate).FirstOrDefault();
public IEnumerable<TEntity> FindMany(Func<TEntity, bool> predicate) => Entities.Values.Where(predicate);
public IEnumerable<TEntity> FindMany<TKey>(Func<TEntity, bool> predicate, Func<TEntity, TKey> 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<TEntity> List() => Entities.Values;
public bool Delete(string id) => Entities.Remove(id);

View file

@ -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<TDbContext, TEntity> : IStore<TDbContext, TEntity> wher
return OnLoading(dbContext, entities);
}
public async Task<Page<TEntity>> FindManyAsync<TKey>(
Expression<Func<TEntity, bool>> predicate,
Expression<Func<TEntity, TKey>> orderBy,
OrderDirection orderDirection = OrderDirection.Ascending,
PageArgs? pageArgs = default,
CancellationToken cancellationToken = default)
{
await using var dbContext = await CreateDbContextAsync(cancellationToken);
var set = dbContext.Set<TEntity>().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<bool> DeleteAsync(TEntity entity, CancellationToken cancellationToken = default)
{
await using var dbContext = await CreateDbContextAsync(cancellationToken);
@ -82,8 +106,9 @@ public class EFCoreStore<TDbContext, TEntity> : IStore<TDbContext, TEntity> 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<bool> AnyAsync(Expression<Func<TEntity, bool>> predicate, CancellationToken cancellationToken = default)

View file

@ -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<TDbContext, TEntity> where TEntity : Entity
Task SaveManyAsync(IEnumerable<TEntity> entities, CancellationToken cancellationToken = default);
Task<TEntity?> FindAsync(Expression<Func<TEntity, bool>> predicate, CancellationToken cancellationToken = default);
Task<IEnumerable<TEntity>> FindManyAsync(Expression<Func<TEntity, bool>> predicate, CancellationToken cancellationToken = default);
Task<Page<TEntity>> FindManyAsync<TKey>(
Expression<Func<TEntity, bool>> predicate,
Expression<Func<TEntity, TKey>> orderBy,
OrderDirection orderDirection = OrderDirection.Ascending,
PageArgs? pageArgs = default,
CancellationToken cancellationToken = default);
Task<bool> DeleteAsync(TEntity entity, CancellationToken cancellationToken = default);
Task<int> DeleteManyAsync(IEnumerable<TEntity> entities, CancellationToken cancellationToken = default);
Task<int> DeleteWhereAsync(Expression<Func<TEntity, bool>> predicate, CancellationToken cancellationToken = default);

View file

@ -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";
}

View file

@ -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<IActionResult> 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);
}

View file

@ -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;
}

View file

@ -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<JsonStringEnumConverter>());

View file

@ -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<WorkflowsDbContext, WorkflowExecutionLogRecord> store) => _store = store;
public async Task SaveAsync(WorkflowExecutionLogRecord record, CancellationToken cancellationToken = default) => await _store.SaveAsync(record, cancellationToken);
public async Task SaveManyAsync(IEnumerable<WorkflowExecutionLogRecord> records, CancellationToken cancellationToken = default) => await _store.SaveManyAsync(records, cancellationToken);
public async Task<Page<WorkflowExecutionLogRecord>> 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;
}
}

View file

@ -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<Page<WorkflowExecutionLogRecord>> 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);
}
}

View file

@ -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<WorkflowExecutionLogRecord> records, CancellationToken cancellationToken = default);
Task<Page<WorkflowExecutionLogRecord>> FindManyByWorkflowInstanceIdAsync(string workflowInstanceId, PageArgs? pageArgs = default, CancellationToken cancellationToken = default);
}