YesSQL persistence for Webhooks
YesSQL persistence for Webhooks
This commit is contained in:
parent
130147f5c9
commit
6feb8f4a77
|
|
@ -4,7 +4,7 @@ namespace Elsa.Webhooks.Persistence.YesSql.Documents
|
|||
{
|
||||
public class WebhookDefinitionDocument : YesSqlDocument
|
||||
{
|
||||
public string WebhookDefinitionId { get; set; } = default!;
|
||||
public string WebhookId { get; set; } = default!;
|
||||
public string? TenantId { get; set; }
|
||||
public string Name { get; set; } = default!;
|
||||
public string Path { get; set; } = default!;
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<Import Project="..\..\..\..\common.props" />
|
||||
<Import Project="..\..\..\..\configureawait.props" />
|
||||
|
|
|
|||
|
|
@ -1,10 +1,11 @@
|
|||
using System;
|
||||
using System.Data;
|
||||
using Elsa.Activities.Webhooks;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Services;
|
||||
using Elsa.Persistence.YesSql;
|
||||
using Elsa.Persistence.YesSql.Data;
|
||||
using Elsa.Persistence.YesSql.Mapping;
|
||||
using Elsa.Persistence.YesSql.Services;
|
||||
//using Elsa.Persistence.YesSql.Data;
|
||||
//using Elsa.Persistence.YesSql.Mapping;
|
||||
//using Elsa.Persistence.YesSql.Services;
|
||||
using Elsa.Runtime;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Indexes;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Stores;
|
||||
|
|
@ -12,6 +13,8 @@ using Microsoft.Extensions.DependencyInjection;
|
|||
using YesSql;
|
||||
using YesSql.Indexes;
|
||||
using YesSql.Provider.Sqlite;
|
||||
using Elsa.Persistence.YesSql.Data;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Mapping;
|
||||
|
||||
namespace Elsa.Webhooks.Persistence.YesSql.Extensions
|
||||
{
|
||||
|
|
@ -25,22 +28,28 @@ namespace Elsa.Webhooks.Persistence.YesSql.Extensions
|
|||
webhookOptions.Services
|
||||
.AddScoped<YesSqlWebhookDefinitionStore>()
|
||||
.AddSingleton(sp => CreateStore(sp, configure))
|
||||
.AddSingleton<ISessionProvider, SessionProvider>()
|
||||
.AddScoped(CreateSession)
|
||||
.AddScoped<IDataMigrationManager, DataMigrationManager>()
|
||||
//.AddSingleton<ISessionProvider, SessionProvider>()
|
||||
//.AddScoped(CreateSession)
|
||||
//.AddScoped<IDataMigrationManager, Elsa.Persistence.YesSql.Services.DataMigrationManager>()
|
||||
.AddStartupTask<DatabaseInitializer>()
|
||||
.AddStartupTask<RunMigrations>()
|
||||
//.AddStartupTask<Elsa.Persistence.YesSql.Services.RunMigrations>()
|
||||
.AddDataMigration<Migrations>()
|
||||
.AddAutoMapperProfile<AutoMapperProfile>()
|
||||
.AddIndexProvider<WebhookDefinitionIndexProvider>();
|
||||
|
||||
var webhookOptionsBuilder = new WebhookOptionsBuilder(webhookOptions.Services);
|
||||
//var webhookOptionsBuilder = new WebhookOptionsBuilder(webhookOptions.Services);
|
||||
|
||||
webhookOptionsBuilder.UseWebhookDefinitionStore(sp => sp.GetRequiredService<YesSqlWebhookDefinitionStore>());
|
||||
|
||||
webhookOptions.UseWebhookDefinitionStore(sp => sp.GetRequiredService<YesSqlWebhookDefinitionStore>());
|
||||
|
||||
return webhookOptions;
|
||||
}
|
||||
|
||||
public static IServiceCollection AddIndexProvider<T>(this IServiceCollection services) where T : class, IIndexProvider => services.AddSingleton<IIndexProvider, T>();
|
||||
public static IServiceCollection AddScopedIndexProvider<T>(this IServiceCollection services) where T : class, IIndexProvider => services.AddScoped<IScopedIndexProvider>();
|
||||
|
||||
public static IServiceCollection AddDataMigration<T>(this IServiceCollection services) where T : class, IDataMigration => services.AddScoped<IDataMigration, T>();
|
||||
|
||||
private static IStore CreateStore(
|
||||
IServiceProvider serviceProvider,
|
||||
Action<IServiceProvider, Configuration> configure)
|
||||
|
|
@ -64,7 +73,7 @@ namespace Elsa.Webhooks.Persistence.YesSql.Extensions
|
|||
|
||||
private static ISession CreateSession(IServiceProvider serviceProvider)
|
||||
{
|
||||
var provider = serviceProvider.GetRequiredService<ISessionProvider>();
|
||||
var provider = serviceProvider.GetRequiredService<Elsa.Persistence.YesSql.Services.ISessionProvider>();
|
||||
return provider.CreateSession();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,7 +6,7 @@ namespace Elsa.Webhooks.Persistence.YesSql.Indexes
|
|||
{
|
||||
public class WebhookDefinitionIndex : MapIndex
|
||||
{
|
||||
public string WebhookDefinitionId { get; set; } = default!;
|
||||
public string WebhookId { get; set; } = default!;
|
||||
public string? TenantId { get; set; }
|
||||
public bool IsEnabled { get; set; }
|
||||
}
|
||||
|
|
@ -21,7 +21,7 @@ namespace Elsa.Webhooks.Persistence.YesSql.Indexes
|
|||
.Map(
|
||||
webhookDefinition => new WebhookDefinitionIndex
|
||||
{
|
||||
WebhookDefinitionId = webhookDefinition.WebhookDefinitionId,
|
||||
WebhookId = webhookDefinition.WebhookId,
|
||||
TenantId = webhookDefinition.TenantId,
|
||||
IsEnabled = webhookDefinition.IsEnabled
|
||||
}
|
||||
|
|
|
|||
|
|
@ -9,9 +9,10 @@ namespace Elsa.Webhooks.Persistence.YesSql.Mapping
|
|||
public AutoMapperProfile()
|
||||
{
|
||||
CreateMap<WebhookDefinition, WebhookDefinitionDocument>()
|
||||
.ForMember(d => d.WebhookId, d => d.MapFrom(s => s.Id))
|
||||
.ForMember(d => d.Id, d => d.Ignore())
|
||||
.ReverseMap()
|
||||
.ForMember(d => d.Id, d => d.MapFrom(s => s.Id));
|
||||
.ForMember(d => d.Id, d => d.MapFrom(s => s.WebhookId));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ namespace Elsa.Webhooks.Persistence.YesSql
|
|||
{
|
||||
SchemaBuilder.CreateMapIndexTable<WebhookDefinitionIndex>(
|
||||
table => table
|
||||
.Column<string>(nameof(WebhookDefinitionIndex.WebhookDefinitionId))
|
||||
.Column<string?>(nameof(WebhookDefinitionIndex.WebhookId))
|
||||
.Column<string?>(nameof(WebhookDefinitionIndex.TenantId))
|
||||
.Column<bool>(nameof(WebhookDefinitionIndex.IsEnabled)),
|
||||
CollectionNames.WebhookDefinitions);
|
||||
|
|
|
|||
|
|
@ -0,0 +1,25 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Services;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Data;
|
||||
using YesSql;
|
||||
|
||||
namespace Elsa.Webhooks.Persistence.YesSql.Services
|
||||
{
|
||||
public class DatabaseInitializer : IStartupTask
|
||||
{
|
||||
private readonly IStore _store;
|
||||
|
||||
public DatabaseInitializer(IStore store)
|
||||
{
|
||||
_store = store;
|
||||
}
|
||||
|
||||
public int Order => 0;
|
||||
|
||||
public async Task ExecuteAsync(CancellationToken cancellationToken = default)
|
||||
{
|
||||
await _store.InitializeCollectionAsync(CollectionNames.WebhookDefinitions);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,6 +1,5 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Models;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Data;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Documents;
|
||||
using Elsa.Webhooks.Persistence.YesSql.Indexes;
|
||||
|
|
@ -22,13 +21,13 @@ namespace Elsa.Webhooks.Persistence.YesSql.Stores
|
|||
{
|
||||
}
|
||||
|
||||
protected override async Task<WebhookDefinitionDocument?> FindDocumentAsync(ISession session, WebhookDefinition entity, CancellationToken cancellationToken) => await Query<WebhookDefinitionIndex>(session, x => x.WebhookDefinitionId == entity.Id).FirstOrDefaultAsync();
|
||||
protected override async Task<WebhookDefinitionDocument?> FindDocumentAsync(ISession session, WebhookDefinition entity, CancellationToken cancellationToken) => await Query<WebhookDefinitionIndex>(session, x => x.WebhookId == entity.Id).FirstOrDefaultAsync();
|
||||
|
||||
protected override IQuery<WebhookDefinitionDocument> MapSpecification(ISession session, ISpecification<WebhookDefinition> specification)
|
||||
{
|
||||
return specification switch
|
||||
{
|
||||
EntityIdSpecification<WebhookDefinition> s => Query<WebhookDefinitionIndex>(session, x => x.WebhookDefinitionId == s.Id),
|
||||
EntityIdSpecification<WebhookDefinition> s => Query<WebhookDefinitionIndex>(session, x => x.WebhookId == s.Id),
|
||||
_ => AutoMapSpecification<WebhookDefinitionIndex>(session, specification)
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -41,14 +41,14 @@ namespace Elsa.Persistence.YesSql.Stores
|
|||
{
|
||||
await _semaphore.WaitAsync(cancellationToken);
|
||||
|
||||
try
|
||||
try
|
||||
{
|
||||
await using var session = SessionProvider.CreateSession();
|
||||
var existingDocument = await FindDocumentAsync(session, entity, cancellationToken);
|
||||
var document = Mapper.Map(entity, existingDocument);
|
||||
session.Save(document, CollectionName);
|
||||
await session.SaveChangesAsync();
|
||||
}
|
||||
}
|
||||
finally
|
||||
{
|
||||
_semaphore.Release();
|
||||
|
|
|
|||
Loading…
Reference in a new issue