Implemented rollover strategy

This commit is contained in:
Gürkan Güran 2023-01-09 16:36:58 +01:00
parent 00962a0b90
commit ce2434f9fd
22 changed files with 201 additions and 88 deletions

View file

@ -5,5 +5,5 @@ namespace Elsa.Elasticsearch.Common;
public abstract class ElasticConfiguration<T> : IElasticConfiguration
{
public abstract void Apply(ConnectionSettings connectionSettings, IDictionary<string, string> aliasConfig);
public abstract void Apply(ConnectionSettings connectionSettings, IDictionary<Type, string> indexConfig);
}

View file

@ -1,6 +1,9 @@
using Elasticsearch.Net;
using Elsa.Elasticsearch.Extensions;
using Elsa.Elasticsearch.HostedServices;
using Elsa.Elasticsearch.Models;
using Elsa.Elasticsearch.Options;
using Elsa.Elasticsearch.Scheduling;
using Elsa.Elasticsearch.Services;
using Elsa.Features.Abstractions;
using Elsa.Features.Services;
@ -16,14 +19,25 @@ public abstract class ElasticFeatureBase : FeatureBase
}
internal ElasticsearchOptions Options { get; set; } = new();
internal IDictionary<string, string> AliasConfig { get; set; } = IElasticConfiguration.GetDefaultAliasConfig();
internal IDictionary<Type, string> IndexConfig { get; set; } = IElasticConfiguration.GetDefaultIndexConfig();
internal IndexRolloverStrategy? IndexRolloverStrategy { get; set; }
public override void ConfigureHostedServices()
{
Module.ConfigureHostedService<ConfigureElasticsearchHostedService>(-1);
}
public override void Apply()
{
if (Services.Any(x => x.ServiceType == typeof(ElasticClient))) return;
var elasticClient = new ElasticClient(GetSettings());
elasticClient.ConfigureIndicesAndAliases(AliasConfig);
if (IndexRolloverStrategy != null)
{
elasticClient.ApplyRolloverStrategy(IndexConfig, IndexRolloverStrategy!);
}
Services.AddSingleton(elasticClient);
}
@ -31,7 +45,7 @@ public abstract class ElasticFeatureBase : FeatureBase
{
return new ConnectionSettings(new Uri(Options.Endpoint))
.ConfigureAuthentication(Options)
.ConfigureMapping(AliasConfig);
.ConfigureMapping(IndexConfig);
}
protected void AddStore<TModel, TStore>() where TModel : class where TStore : class

View file

@ -20,7 +20,7 @@ public class ElasticStore<T> where T : class
{
var response = await _elasticClient.GetAsync(DocumentPath<T>.Id(id), ct: cancellationToken);
if (response.IsValid) return response.Source;
if (response.ApiCall.Success) return response.Source;
_logger.LogError("Failed to fetch data from Elasticsearch: {message}", response.ServerError?.ToString());
return null;
@ -36,7 +36,8 @@ public class ElasticStore<T> where T : class
var response = await _elasticClient.SearchAsync<T>(search, cancellationToken);
if (response.IsValid) return new Page<T>(response.Hits.Select(hit => hit.Source).ToList(), response.Total);
if (response.ApiCall.Success)
return new Page<T>(response.Hits.Select(hit => hit.Source).ToList(), response.Total);
_logger.LogError("Failed to search data in Elasticsearch: {message}", response.ServerError?.ToString());
return new Page<T>(new Collection<T>(), 0);
@ -51,7 +52,8 @@ public class ElasticStore<T> where T : class
var response = await _elasticClient.SearchAsync<T>(search, cancellationToken);
if (response.IsValid) return new Page<T>(response.Hits.Select(hit => hit.Source).ToList(), response.Total);
if (response.ApiCall.Success)
return new Page<T>(response.Hits.Select(hit => hit.Source).ToList(), response.Total);
_logger.LogError("Failed to search data in Elasticsearch: {message}", response.ServerError?.ToString());
return new Page<T>(new Collection<T>(), 0);
@ -61,7 +63,7 @@ public class ElasticStore<T> where T : class
{
var response = await _elasticClient.IndexAsync(model, descriptor => descriptor, cancellationToken);
if (response.IsValid) return true;
if (response.ApiCall.Success) return true;
_logger.LogError("Failed to save data in Elasticsearch: {message}", response.ServerError?.ToString());
return false;
@ -71,7 +73,7 @@ public class ElasticStore<T> where T : class
{
var response = await _elasticClient.IndexManyAsync(documents, cancellationToken: cancellationToken);
if (response.IsValid) return true;
if (response.ApiCall.Success) return true;
_logger.LogError("Failed to save data in Elasticsearch: {message}", response.ServerError?.ToString());
return false;
@ -81,7 +83,7 @@ public class ElasticStore<T> where T : class
{
var response = await _elasticClient.DeleteAsync(DocumentPath<T>.Id(id), ct: cancellationToken);
if (response.IsValid) return true;
if (response.ApiCall.Success) return true;
_logger.LogError("Failed to delete data in Elasticsearch: {message}", response.ServerError?.ToString());
return false;
@ -91,7 +93,7 @@ public class ElasticStore<T> where T : class
{
var response = await _elasticClient.DeleteManyAsync(list, cancellationToken: cancellationToken);
if (response.IsValid) return 0;
if (response.ApiCall.Success) return 0;
_logger.LogError("Failed to delete data in Elasticsearch: {message}", response.ServerError?.ToString());
return response.Items.Count;
@ -102,7 +104,7 @@ public class ElasticStore<T> where T : class
var response = await _elasticClient.DeleteByQueryAsync<T>(q => q
.Query(query), cancellationToken);
if (response.IsValid) return true;
if (response.ApiCall.Success) return true;
_logger.LogError("Failed to delete data in Elasticsearch: {message}", response.ServerError?.ToString());
return false;
@ -113,7 +115,7 @@ public class ElasticStore<T> where T : class
var response = await _elasticClient.CountAsync<T>(s =>
s.Query(query), cancellationToken);
if (response.IsValid) return response.Count;
if (response.ApiCall.Success) return response.Count;
_logger.LogError("Failed to count data in Elasticsearch: {message}", response.ServerError?.ToString());
return 0;

View file

@ -6,10 +6,14 @@ public static class Utils
{
public static string GenerateIndexName(string aliasName)
{
var month = DateTime.Now.ToString("MM");
var year = DateTime.Now.Year;
return aliasName + "-" + year + "-" + month;
var now = DateTime.Now;
var month = now.ToString("MM");
var year = now.Year;
var day = now.Day;
var hour = now.Hour;
var minute = now.Minute;
return aliasName + "-" + year + "-" + month + "-" + day + hour + minute;
}
public static IEnumerable<Type> GetElasticConfigurationTypes() =>
@ -22,4 +26,18 @@ public static class Utils
public static IEnumerable<Type> GetElasticDocumentTypes() =>
GetElasticConfigurationTypes()
.Select(t => t.BaseType!.GenericTypeArguments.First()).ToList();
public static IDictionary<Type, string> ResolveAliasConfig(
IDictionary<Type, string> defaultConfig,
IDictionary<string, string>? option1,
IDictionary<string, string>? option2)
{
var types = GetElasticDocumentTypes();
return option1?.Select(kvp => new KeyValuePair<Type, string>(types.First(t => t.Name == kvp.Key), kvp.Value))
.ToDictionary(x => x.Key, x => x.Value) ??
option2?.Select(kvp => new KeyValuePair<Type, string>(types.First(t => t.Name == kvp.Key), kvp.Value))
.ToDictionary(x => x.Key, x => x.Value) ??
defaultConfig;
}
}

View file

@ -13,6 +13,7 @@
<ItemGroup>
<ProjectReference Include="..\..\common\Elsa.Features\Elsa.Features.csproj" />
<ProjectReference Include="..\Elsa.Jobs\Elsa.Jobs.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>

View file

@ -1,8 +1,11 @@
using Elasticsearch.Net;
using Elsa.Elasticsearch.Common;
using Elsa.Elasticsearch.Implementations.RolloverStrategies;
using Elsa.Elasticsearch.Models;
using Elsa.Elasticsearch.Options;
using Elsa.Elasticsearch.Services;
using Nest;
using Index = Nest.Index;
namespace Elsa.Elasticsearch.Extensions;
@ -22,34 +25,20 @@ public static class ElasticExtensions
return settings;
}
public static ConnectionSettings ConfigureMapping(this ConnectionSettings settings, IDictionary<string,string> aliasConfig)
public static ConnectionSettings ConfigureMapping(this ConnectionSettings settings, IDictionary<Type,string> indexConfig)
{
foreach (var config in Utils.GetElasticConfigurationTypes())
{
var configInstance = (IElasticConfiguration)Activator.CreateInstance(config)!;
configInstance.Apply(settings, aliasConfig);
configInstance.Apply(settings, indexConfig);
}
return settings;
}
public static void ConfigureIndicesAndAliases(this ElasticClient client, IDictionary<string,string> aliasConfig)
public static void ApplyRolloverStrategy(this ElasticClient client, IDictionary<Type,string> aliasConfig, IndexRolloverStrategy strategy)
{
foreach (var type in Utils.GetElasticDocumentTypes())
{
var aliasName = aliasConfig[type.Name];
var indexName = Utils.GenerateIndexName(aliasName);
var indexExists = client.Indices.Exists(indexName).Exists;
if (indexExists) continue;
var response = client.Indices.Create(indexName, s => s
.Aliases(a => a.Alias(aliasName))
.Map(m => m.AutoMap(type)));
if (response.IsValid) continue;
throw response.OriginalException;
}
var strategyInstance = (IRolloverStrategy)Activator.CreateInstance(strategy.Value, args: client)!;
strategyInstance.Apply(Utils.GetElasticDocumentTypes(), aliasConfig);
}
}

View file

@ -0,0 +1,51 @@
using Elsa.Elasticsearch.Scheduling;
using Elsa.Jobs.Schedules;
using Elsa.Jobs.Services;
using Elsa.Workflows.Management.Entities;
using Microsoft.Extensions.Hosting;
using Nest;
namespace Elsa.Elasticsearch.HostedServices;
public class ConfigureElasticsearchHostedService : IHostedService
{
private readonly ElasticClient _elasticClient;
private readonly IJobScheduler _jobScheduler;
public ConfigureElasticsearchHostedService(ElasticClient elasticClient, IJobScheduler jobScheduler)
{
_elasticClient = elasticClient;
_jobScheduler = jobScheduler;
}
public async Task StartAsync(CancellationToken cancellationToken)
{
await FlattenProperties(cancellationToken);
await ScheduleIndexAndAliasConfiguration(cancellationToken);
}
private async Task ScheduleIndexAndAliasConfiguration(CancellationToken cancellationToken)
{
var job = new ConfigureElasticIndicesJob();
var schedule = new CronSchedule
{
// At the start of every month
CronExpression = "*/5 * * * *"
};
await _jobScheduler.ScheduleAsync(job, GetType().Name, schedule, cancellationToken: cancellationToken);
}
private async Task FlattenProperties(CancellationToken cancellationToken)
{
await _elasticClient.Indices.PutMappingAsync<WorkflowInstance>(
descriptor => descriptor
.Properties(p => p
.Flattened(d => d
.Name(p => p.WorkflowState.Properties))),
cancellationToken);
}
public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask;
}

View file

@ -0,0 +1,35 @@
using Elsa.Elasticsearch.Common;
using Elsa.Elasticsearch.Services;
using Nest;
namespace Elsa.Elasticsearch.Implementations.RolloverStrategies;
public class RolloverOnMonthlyBasis : IRolloverStrategy
{
private readonly ElasticClient _client;
public RolloverOnMonthlyBasis(ElasticClient client)
{
_client = client;
}
public void Apply(IEnumerable<Type> types, IDictionary<Type, string> aliasConfig)
{
foreach (var type in types)
{
var aliasName = aliasConfig[type];
var indexName = Utils.GenerateIndexName(aliasName);
var indexExists = _client.Indices.Exists(indexName).Exists;
if (indexExists) continue;
var response = _client.Indices.Create(indexName, s => s
.Aliases(a => a.Alias(aliasName))
.Map(m => m.AutoMap(type)));
if (response.IsValid) continue;
throw response.OriginalException;
}
}
}

View file

@ -0,0 +1,17 @@
using Elsa.Elasticsearch.Implementations.RolloverStrategies;
namespace Elsa.Elasticsearch.Models;
public class IndexRolloverStrategy
{
private IndexRolloverStrategy(Type value) { Value = value; }
public Type Value { get; private set; }
public static IndexRolloverStrategy RolloverOnMonthlyBasis => new (typeof(RolloverOnMonthlyBasis));
public override string ToString()
{
return Value.Name;
}
}

View file

@ -1,3 +1,5 @@
using Elsa.Elasticsearch.Common;
using Elsa.Elasticsearch.Models;
using Elsa.Elasticsearch.Options;
using Elsa.Workflows.Management.Features;
@ -8,12 +10,18 @@ public static class Extensions
/// <summary>
/// Configures the <see cref="WorkflowInstanceFeature"/> to use the <see cref="ElasticWorkflowInstanceFeature"/>.
/// </summary>
public static WorkflowInstanceFeature UseElasticsearch(this WorkflowInstanceFeature feature, ElasticsearchOptions options, IDictionary<string,string>? aliasConfig = default, Action<ElasticWorkflowInstanceFeature>? configure = default)
public static WorkflowInstanceFeature UseElasticsearch(
this WorkflowInstanceFeature feature,
ElasticsearchOptions options,
IndexRolloverStrategy? rolloverStrategy = default,
IDictionary<string,string>? indexConfig = default,
Action<ElasticWorkflowInstanceFeature>? configure = default)
{
configure += f =>
{
f.Options = options;
f.AliasConfig = aliasConfig ?? options.Aliases ?? f.AliasConfig;
f.IndexRolloverStrategy = rolloverStrategy;
f.IndexConfig = Utils.ResolveAliasConfig(f.IndexConfig, options.IndexConfig, indexConfig);
};
feature.Module.Configure(configure);

View file

@ -6,11 +6,11 @@ namespace Elsa.Elasticsearch.Modules.Management;
public class WorkflowInstanceConfiguration : ElasticConfiguration<WorkflowInstance>
{
public override void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> aliasConfig)
public override void Apply(ConnectionSettings connectionSettings, IDictionary<Type,string> indexConfig)
{
connectionSettings
.DefaultMappingFor<WorkflowInstance>(m => m
.IndexName(aliasConfig[nameof(WorkflowInstance)])
.IndexName(indexConfig[typeof(WorkflowInstance)])
.Ignore(p => p.WorkflowState));
}
}

View file

@ -6,9 +6,9 @@ namespace Elsa.Elasticsearch.Modules.Runtime;
public class ExecutionLogConfiguration : ElasticConfiguration<WorkflowExecutionLogRecord>
{
public override void Apply(ConnectionSettings connectionSettings, IDictionary<string, string> aliasConfig)
public override void Apply(ConnectionSettings connectionSettings, IDictionary<Type, string> indexConfig)
{
connectionSettings.DefaultMappingFor<WorkflowExecutionLogRecord>(m =>
m.IndexName(aliasConfig[nameof(WorkflowExecutionLogRecord)]));
m.IndexName(indexConfig[typeof(WorkflowExecutionLogRecord)]));
}
}

View file

@ -1,3 +1,5 @@
using Elsa.Elasticsearch.Common;
using Elsa.Elasticsearch.Models;
using Elsa.Elasticsearch.Options;
using Elsa.Workflows.Runtime.Features;
@ -8,12 +10,18 @@ public static class Extensions
/// <summary>
/// Configures the <see cref="ExecutionLogRecordFeature"/> to use the <see cref="ElasticExecutionLogRecordFeature"/>.
/// </summary>
public static ExecutionLogRecordFeature UseElasticsearch(this ExecutionLogRecordFeature feature, ElasticsearchOptions options, IDictionary<string,string>? aliasConfig = default, Action<ElasticExecutionLogRecordFeature>? configure = default)
public static ExecutionLogRecordFeature UseElasticsearch(
this ExecutionLogRecordFeature feature,
ElasticsearchOptions options,
IndexRolloverStrategy? rolloverStrategy = default,
IDictionary<string,string>? indexConfig = default,
Action<ElasticExecutionLogRecordFeature>? configure = default)
{
configure += f =>
{
f.Options = options;
f.AliasConfig = aliasConfig ?? options.Aliases ?? f.AliasConfig;
f.IndexRolloverStrategy = rolloverStrategy;
f.IndexConfig = Utils.ResolveAliasConfig(f.IndexConfig, options.IndexConfig, indexConfig);
};
feature.Module.Configure(configure);

View file

@ -4,7 +4,7 @@ public class ElasticsearchOptions
{
public const string Elasticsearch = "Elasticsearch";
public Dictionary<string, string>? Aliases { get; set; }
public Dictionary<string, string>? IndexConfig { get; set; }
public string Endpoint { get; set; }
public string Username { get; set; }
public string Password { get; set; }

View file

@ -4,7 +4,7 @@ using Elsa.Jobs.Models;
using Nest;
using Job = Elsa.Jobs.Abstractions.Job;
namespace Elsa.Scheduling.Jobs;
namespace Elsa.Elasticsearch.Scheduling;
public class ConfigureElasticIndicesJob : Job
{

View file

@ -5,10 +5,10 @@ namespace Elsa.Elasticsearch.Services;
public interface IElasticConfiguration
{
void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> aliasConfig);
void Apply(ConnectionSettings connectionSettings, IDictionary<Type,string> indexConfig);
public static IDictionary<string, string> GetDefaultAliasConfig()
public static IDictionary<Type, string> GetDefaultIndexConfig()
{
return Utils.GetElasticDocumentTypes().ToDictionary(type => type.Name, type => type.Name.ToLower());
return Utils.GetElasticDocumentTypes().ToDictionary(type => type, type => type.Name.ToLower());
}
}

View file

@ -0,0 +1,6 @@
namespace Elsa.Elasticsearch.Services;
public interface IRolloverStrategy
{
void Apply(IEnumerable<Type> types, IDictionary<Type,string> aliasConfig);
}

View file

@ -1,3 +1,4 @@
using Elsa.Elasticsearch.Scheduling;
using Elsa.Quartz.Jobs;
using Elsa.Scheduling.Jobs;
using Quartz;

View file

@ -27,7 +27,6 @@ public class SchedulingFeature : FeatureBase
Services
.AddSingleton<IWorkflowTriggerScheduler, WorkflowTriggerScheduler>()
.AddSingleton<IWorkflowBookmarkScheduler, WorkflowBookmarkScheduler>()
.AddSingleton<IElasticCongurationScheduler, ElasticCongurationScheduler>()
.AddNotificationHandlersFrom<ScheduleWorkflows>();
Module.Configure<WorkflowManagementFeature>(management => management.AddActivitiesFrom<SchedulingFeature>());

View file

@ -1,28 +0,0 @@
using Elsa.Jobs.Schedules;
using Elsa.Jobs.Services;
using Elsa.Scheduling.Jobs;
using Elsa.Scheduling.Services;
namespace Elsa.Scheduling.Implementations;
public class ElasticCongurationScheduler : IElasticCongurationScheduler
{
private readonly IJobScheduler _jobScheduler;
public ElasticCongurationScheduler(IJobScheduler jobScheduler)
{
_jobScheduler = jobScheduler;
}
public async Task ScheduleAsync(CancellationToken cancellationToken = default)
{
var job = new ConfigureElasticIndicesJob();
var schedule = new CronSchedule
{
//Last day of every month
CronExpression = "0 18 L * ?"
};
await _jobScheduler.ScheduleAsync(job, GetType().Name, schedule, cancellationToken: cancellationToken);
}
}

View file

@ -1,8 +0,0 @@
using Elsa.Workflows.Core.Models;
namespace Elsa.Scheduling.Services;
public interface IElasticCongurationScheduler
{
Task ScheduleAsync(CancellationToken cancellationToken = default);
}

View file

@ -7,7 +7,7 @@ namespace Elsa.Workflows.Core.State;
/// <summary>
/// A simplified, serializable model representing an exception.
/// </summary>
public record ExceptionState(Type Type, string Message, string? StackTrace, IDictionary Data, ExceptionState? InnerException = default)
public record ExceptionState(Type Type, string Message, string? StackTrace, ExceptionState? InnerException = default)
{
// /// <summary>
// /// Constructor
@ -26,6 +26,6 @@ public record ExceptionState(Type Type, string Message, string? StackTrace, IDic
if (ex == null)
return null;
return new ExceptionState(ex.GetType(), ex.Message, ex.StackTrace, ex.Data, FromException(ex.InnerException));
return new ExceptionState(ex.GetType(), ex.Message, ex.StackTrace, FromException(ex.InnerException));
}
}