Implemented aliases

This commit is contained in:
Gürkan Güran 2023-01-07 20:31:42 +01:00
parent 14c66f7d83
commit f995a4f35a
15 changed files with 135 additions and 29 deletions

View file

@ -16,21 +16,22 @@ public abstract class ElasticFeatureBase : FeatureBase
}
internal ElasticsearchOptions Options { get; set; } = new();
internal IDictionary<string,string> IndexConfig { get; set; }
internal IDictionary<string, string> AliasConfig { get; set; } = IElasticConfiguration.GetDefaultAliasConfig();
public override void Apply()
{
if (Services.All(x => x.ServiceType != typeof(ElasticClient)))
{
Services.AddSingleton(new ElasticClient(GetSettings()));
}
if (Services.Any(x => x.ServiceType == typeof(ElasticClient))) return;
var elasticClient = new ElasticClient(GetSettings());
elasticClient.ConfigureIndicesAndAliases(AliasConfig);
Services.AddSingleton(elasticClient);
}
private ConnectionSettings GetSettings()
{
return new ConnectionSettings(new Uri(Options.Endpoint))
.ConfigureAuthentication(Options)
.ConfigureMapping(IndexConfig);
.ConfigureMapping(AliasConfig);
}
protected void AddStore<TModel, TStore>() where TModel : class where TStore : class

View file

@ -0,0 +1,19 @@
using Elsa.Elasticsearch.Services;
namespace Elsa.Elasticsearch.Common;
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;
}
public static IEnumerable<Type> GetElasticDocumentTypes() =>
AppDomain.CurrentDomain.GetAssemblies()
.SelectMany(s => s.GetTypes())
.Where(p => typeof(IElasticConfiguration).IsAssignableFrom(p) && p.IsClass);
}

View file

@ -1,11 +1,12 @@
using Elasticsearch.Net;
using Elsa.Elasticsearch.Common;
using Elsa.Elasticsearch.Options;
using Elsa.Elasticsearch.Services;
using Nest;
namespace Elsa.Elasticsearch.Extensions;
public static class ConnectionSettingsExtensions
public static class ElasticExtensions
{
public static ConnectionSettings ConfigureAuthentication(this ConnectionSettings settings, ElasticsearchOptions options)
{
@ -21,7 +22,7 @@ public static class ConnectionSettingsExtensions
return settings;
}
public static ConnectionSettings ConfigureMapping(this ConnectionSettings settings, IDictionary<string,string> indexConfig)
public static ConnectionSettings ConfigureMapping(this ConnectionSettings settings, IDictionary<string,string> aliasConfig)
{
var configs = AppDomain.CurrentDomain.GetAssemblies()
.SelectMany(s => s.GetTypes())
@ -30,9 +31,22 @@ public static class ConnectionSettingsExtensions
foreach (var config in configs)
{
var configInstance = (IElasticConfiguration)Activator.CreateInstance(config)!;
configInstance.Apply(settings, indexConfig);
configInstance.Apply(settings, aliasConfig);
}
return settings;
}
public static void ConfigureIndicesAndAliases(this ElasticClient client, IDictionary<string,string> aliasConfig)
{
foreach (var type in Utils.GetElasticDocumentTypes())
{
var aliasName = aliasConfig[type.Name];
var indexName = Utils.GenerateIndexName(aliasName);
client.Indices.Create(indexName, s => s
.Aliases(a => a.Alias(aliasName))
.Map(m => m.AutoMap(type)));
}
}
}

View file

@ -8,12 +8,12 @@ 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>? indexConfig = default, Action<ElasticWorkflowInstanceFeature>? configure = default)
public static WorkflowInstanceFeature UseElasticsearch(this WorkflowInstanceFeature feature, ElasticsearchOptions options, IDictionary<string,string>? aliasConfig = default, Action<ElasticWorkflowInstanceFeature>? configure = default)
{
configure += f =>
{
f.Options = options;
f.IndexConfig = indexConfig ?? options.Indices ?? new Dictionary<string, string>();
f.AliasConfig = aliasConfig ?? options.Aliases ?? f.AliasConfig;
};
feature.Module.Configure(configure);

View file

@ -1,4 +1,3 @@
using Elsa.Elasticsearch.Options;
using Elsa.Elasticsearch.Services;
using Elsa.Workflows.Management.Entities;
using Nest;
@ -7,11 +6,9 @@ namespace Elsa.Elasticsearch.Modules.Management;
public class WorkflowInstanceConfiguration : IElasticConfiguration
{
private const string IndexName = "workflow-instance";
public void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> indexConfig)
public void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> aliasConfig)
{
connectionSettings.DefaultMappingFor<WorkflowInstance>(m =>
m.IndexName(IElasticConfiguration.ResolveIndexName<WorkflowInstance>(indexConfig, IndexName)));
m.IndexName(aliasConfig[nameof(WorkflowInstance)]));
}
}

View file

@ -7,11 +7,9 @@ namespace Elsa.Elasticsearch.Modules.Runtime;
public class ExecutionLogConfiguration : IElasticConfiguration
{
private const string IndexName = "workflow-execution-log";
public void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> indexConfig)
public void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> aliasConfig)
{
connectionSettings.DefaultMappingFor<WorkflowExecutionLogRecord>(m =>
m.IndexName(IElasticConfiguration.ResolveIndexName<WorkflowExecutionLogRecord>(indexConfig, IndexName)));
m.IndexName(aliasConfig[nameof(WorkflowExecutionLogRecord)]));
}
}

View file

@ -8,12 +8,12 @@ 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>? indexConfig = default, Action<ElasticExecutionLogRecordFeature>? configure = default)
public static ExecutionLogRecordFeature UseElasticsearch(this ExecutionLogRecordFeature feature, ElasticsearchOptions options, IDictionary<string,string>? aliasConfig = default, Action<ElasticExecutionLogRecordFeature>? configure = default)
{
configure += f =>
{
f.Options = options;
f.IndexConfig = indexConfig ?? options.Indices ?? new Dictionary<string, string>();
f.AliasConfig = aliasConfig ?? options.Aliases ?? f.AliasConfig;
};
feature.Module.Configure(configure);

View file

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

View file

@ -1,15 +1,15 @@
using Elsa.Elasticsearch.Options;
using Elsa.Elasticsearch.Common;
using Nest;
namespace Elsa.Elasticsearch.Services;
public interface IElasticConfiguration
{
void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> indexConfig);
public static string ResolveIndexName<T>(IDictionary<string,string> indices, string? indexName = default)
void Apply(ConnectionSettings connectionSettings, IDictionary<string,string> aliasConfig);
internal static IDictionary<string,string> GetDefaultAliasConfig()
{
var indexNameFromConfig = indices[typeof(T).Name];
return string.IsNullOrWhiteSpace(indexNameFromConfig) ? indexName : indexNameFromConfig;
var types = Utils.GetElasticDocumentTypes();
return new Dictionary<string, string>(types.Select(t => new KeyValuePair<string, string>(t.Name, t.Name)));
}
}

View file

@ -12,6 +12,7 @@ public static class DependencyInjectionExtensions
{
quartz.AddJob<RunWorkflowJob>();
quartz.AddJob<ResumeWorkflowJob>();
quartz.AddJob<ConfigureElasticIndicesJob>();
return quartz;
}

View file

@ -12,6 +12,7 @@
</PropertyGroup>
<ItemGroup>
<ProjectReference Include="..\Elsa.Elasticsearch\Elsa.Elasticsearch.csproj" />
<ProjectReference Include="..\Elsa.Jobs\Elsa.Jobs.csproj" />
<ProjectReference Include="..\Elsa.Workflows.Runtime\Elsa.Workflows.Runtime.csproj" />
</ItemGroup>

View file

@ -27,6 +27,7 @@ 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

@ -0,0 +1,28 @@
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

@ -0,0 +1,38 @@
using System.Text.Json.Serialization;
using Elsa.Elasticsearch.Common;
using Elsa.Jobs.Models;
using Nest;
using Job = Elsa.Jobs.Abstractions.Job;
namespace Elsa.Scheduling.Jobs;
public class ConfigureElasticIndicesJob : Job
{
[JsonConstructor]
public ConfigureElasticIndicesJob()
{
}
protected override async ValueTask ExecuteAsync(JobExecutionContext context)
{
var client = context.GetRequiredService<ElasticClient>();
var indexAliasGroups = await client.Indices.GetAliasAsync();
foreach (var group in indexAliasGroups.Indices)
{
// Only 1 alias exists per index in Elsa Elasticsearch configuration
var aliasPointingCurrentIndex = group.Value.Aliases.Keys.Single();
var currentIndexName = group.Key.Name;
var newIndexName = Utils.GenerateIndexName(aliasPointingCurrentIndex);
var indexExists = (await client.Indices.ExistsAsync(newIndexName)).Exists;
if (indexExists) continue;
// Point the alias to the new index
await client.Indices.BulkAliasAsync(aliases => aliases
.Remove(a => a.Alias(aliasPointingCurrentIndex).Index(currentIndexName))
.Add(a => a.Alias(aliasPointingCurrentIndex).Index(newIndexName)));
}
}
}

View file

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