using Elsa.Common.Entities;
using Elsa.Workflows.Runtime.Filters;
using Elsa.Workflows.Runtime.OrderDefinitions;
namespace Elsa.Workflows.Runtime;
///
public class ActivityExecutionStatsService : IActivityExecutionStatsService
{
private readonly IActivityExecutionStore _store;
///
/// Initializes a new instance of the class.
///
public ActivityExecutionStatsService(IActivityExecutionStore store)
{
_store = store;
}
///
public async Task> GetStatsAsync(string workflowInstanceId, IEnumerable activityNodeIds, CancellationToken cancellationToken = default)
{
var filter = new ActivityExecutionRecordFilter
{
WorkflowInstanceId = workflowInstanceId,
ActivityNodeIds = activityNodeIds?.ToList()
};
var order = new ActivityExecutionRecordOrder(x => x.StartedAt, OrderDirection.Ascending);
var records = (await _store.FindManySummariesAsync(filter, order, cancellationToken)).ToList();
var groupedRecords = records.GroupBy(x => x.ActivityNodeId).ToList();
var stats = groupedRecords.Select(grouping => new ActivityExecutionStats
{
ActivityNodeId = grouping.Key,
ActivityId = grouping.First().ActivityId,
StartedCount = grouping.Count(),
CompletedCount = grouping.Count(x => x.CompletedAt != null),
UncompletedCount = grouping.Count(x => x.CompletedAt == null),
IsBlocked = grouping.Any(x => x.HasBookmarks),
IsFaulted = grouping.Any(x => x.Status == ActivityStatus.Faulted),
AggregateFaultCount = grouping.Last().AggregateFaultCount
}).ToList();
return stats;
}
///
public async Task GetStatsAsync(string workflowInstanceId, string activityNodeId, CancellationToken cancellationToken = default)
{
var stats = (await GetStatsAsync(workflowInstanceId, [activityNodeId], cancellationToken)).FirstOrDefault();
return stats ?? new ActivityExecutionStats
{
ActivityNodeId = activityNodeId
};
}
}