Remove MongoDB provider
This commit is contained in:
parent
bd46f80401
commit
b9d0f1e2ba
|
|
@ -1,36 +0,0 @@
|
|||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>netstandard2.0</TargetFramework>
|
||||
<LangVersion>latest</LangVersion>
|
||||
<PackageVersion>1.0.0</PackageVersion>
|
||||
<Authors>Elsa Contributors</Authors>
|
||||
<Description>
|
||||
Elsa is a set of workflow libraries and tools that enable super-fast workflowing capabilities in any .NET Core application.
|
||||
This package provides a MongoDB persistence provider.
|
||||
</Description>
|
||||
<Copyright>2019</Copyright>
|
||||
<PackageProjectUrl>https://github.com/elsa-workflows/elsa-core</PackageProjectUrl>
|
||||
<RepositoryUrl>https://github.com/elsa-workflows/elsa-core</RepositoryUrl>
|
||||
<RepositoryType>GitHub</RepositoryType>
|
||||
<PackageTags>elsa, workflows, mongodb</PackageTags>
|
||||
<PackageIcon>icon.png</PackageIcon>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<None Include="icon.png">
|
||||
<Pack>True</Pack>
|
||||
<PackagePath />
|
||||
</None>
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\core\Elsa.Core\Elsa.Core.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="MongoDb.Bson.NodaTime" Version="2.1.0" />
|
||||
<PackageReference Include="MongoDB.Driver" Version="2.11.0" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
@ -1,111 +0,0 @@
|
|||
using System;
|
||||
using Elsa;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence.MongoDb.Serialization;
|
||||
using Elsa.Persistence.MongoDb.Services;
|
||||
using MongoDB.Bson;
|
||||
using MongoDb.Bson.NodaTime;
|
||||
using MongoDB.Bson.Serialization;
|
||||
using MongoDB.Bson.Serialization.Conventions;
|
||||
using MongoDB.Driver;
|
||||
|
||||
// ReSharper disable once CheckNamespace
|
||||
namespace Microsoft.Extensions.DependencyInjection
|
||||
{
|
||||
public static class ServiceCollectionExtensions
|
||||
{
|
||||
public static ElsaOptions UseMongoDbWorkflowStores(
|
||||
this ElsaOptions options,
|
||||
string databaseName,
|
||||
string connectionString)
|
||||
{
|
||||
return options
|
||||
.AddMongoDbProvider(databaseName, connectionString)
|
||||
.UseMongoDbWorkflowDefinitionStore(databaseName, connectionString)
|
||||
.UseMongoDbWorkflowInstanceStore(databaseName, connectionString);
|
||||
}
|
||||
|
||||
public static ElsaOptions UseMongoDbWorkflowInstanceStore(
|
||||
this ElsaOptions options,
|
||||
string databaseName,
|
||||
string connectionString)
|
||||
{
|
||||
options
|
||||
.AddMongoDbProvider(databaseName, connectionString)
|
||||
.UseWorkflowInstanceStore(sp => sp.GetRequiredService<MongoWorkflowInstanceStore>())
|
||||
.Services
|
||||
.AddMongoDbCollection<WorkflowInstance>("WorkflowInstances")
|
||||
.AddScoped<MongoWorkflowInstanceStore>();
|
||||
|
||||
return options;
|
||||
}
|
||||
|
||||
public static ElsaOptions UseMongoDbWorkflowDefinitionStore(
|
||||
this ElsaOptions options,
|
||||
string databaseName,
|
||||
string connectionString)
|
||||
{
|
||||
options
|
||||
.AddMongoDbProvider(databaseName, connectionString)
|
||||
.UseWorkflowDefinitionStore(sp => sp.GetRequiredService<MongoWorkflowDefinitionStore>())
|
||||
.Services
|
||||
.AddMongoDbCollection<WorkflowDefinitionVersion>("WorkflowDefinitions")
|
||||
.AddScoped<MongoWorkflowDefinitionStore>();
|
||||
|
||||
return options;
|
||||
}
|
||||
|
||||
public static IServiceCollection AddMongoDbCollection<T>(
|
||||
this IServiceCollection services,
|
||||
string collectionName)
|
||||
{
|
||||
return services.AddSingleton(sp => CreateCollection<T>(sp, collectionName));
|
||||
}
|
||||
|
||||
private static ElsaOptions AddMongoDbProvider(
|
||||
this ElsaOptions options,
|
||||
string databaseName,
|
||||
string connectionString
|
||||
)
|
||||
{
|
||||
if (options.Services.HasService<IMongoClient>())
|
||||
return options;
|
||||
|
||||
options.Services
|
||||
.AddTransient<JObjectSerializer>()
|
||||
.AddTransient<VariableSerializer>()
|
||||
.AddSingleton(sp =>
|
||||
{
|
||||
NodaTimeSerializers.Register();
|
||||
RegisterEnumAsStringConvention();
|
||||
BsonSerializer.RegisterSerializer(sp.GetRequiredService<JObjectSerializer>());
|
||||
BsonSerializer.RegisterSerializer(sp.GetRequiredService<VariableSerializer>());
|
||||
return CreateDbClient(connectionString);
|
||||
})
|
||||
.AddSingleton(sp => CreateDatabase(sp, databaseName));
|
||||
|
||||
return options;
|
||||
}
|
||||
|
||||
private static IMongoCollection<T> CreateCollection<T>(IServiceProvider serviceProvider, string collectionName)
|
||||
{
|
||||
var database = serviceProvider.GetRequiredService<IMongoDatabase>();
|
||||
return database.GetCollection<T>(collectionName);
|
||||
}
|
||||
|
||||
private static IMongoDatabase CreateDatabase(IServiceProvider serviceProvider, string databaseName)
|
||||
{
|
||||
var client = serviceProvider.GetRequiredService<IMongoClient>();
|
||||
return client.GetDatabase(databaseName);
|
||||
}
|
||||
|
||||
private static IMongoClient CreateDbClient(string connectionString) => new MongoClient(connectionString);
|
||||
|
||||
private static void RegisterEnumAsStringConvention()
|
||||
{
|
||||
var pack = new ConventionPack { new EnumRepresentationConvention(BsonType.String) };
|
||||
|
||||
ConventionRegistry.Register("EnumStringConvention", pack, _ => true);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,12 +0,0 @@
|
|||
using Elsa.Serialization;
|
||||
using Newtonsoft.Json.Linq;
|
||||
|
||||
namespace Elsa.Persistence.MongoDb.Serialization
|
||||
{
|
||||
public class JObjectSerializer : JsonSerializerBase<JObject>
|
||||
{
|
||||
public JObjectSerializer(ITokenSerializer serializer) : base(serializer)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,33 +0,0 @@
|
|||
using Elsa.Serialization;
|
||||
using MongoDB.Bson;
|
||||
using MongoDB.Bson.Serialization;
|
||||
using MongoDB.Bson.Serialization.Serializers;
|
||||
using Newtonsoft.Json;
|
||||
|
||||
namespace Elsa.Persistence.MongoDb.Serialization
|
||||
{
|
||||
public abstract class JsonSerializerBase<T> : SerializerBase<T>
|
||||
{
|
||||
private readonly ITokenSerializer serializer;
|
||||
|
||||
protected JsonSerializerBase(ITokenSerializer serializer)
|
||||
{
|
||||
this.serializer = serializer;
|
||||
}
|
||||
|
||||
public override T Deserialize(BsonDeserializationContext context, BsonDeserializationArgs args)
|
||||
{
|
||||
var document = BsonDocumentSerializer.Instance.Deserialize(context);
|
||||
var value = serializer.Deserialize<T>(document.ToString());
|
||||
|
||||
return value;
|
||||
}
|
||||
|
||||
public override void Serialize(BsonSerializationContext context, BsonSerializationArgs args, T value)
|
||||
{
|
||||
var json = value != null ? serializer.Serialize(value) : null;
|
||||
var document = json != null ? BsonDocument.Parse(json.ToString(Formatting.None)) : new BsonDocument();
|
||||
BsonDocumentSerializer.Instance.Serialize(context, document);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,12 +0,0 @@
|
|||
using Elsa.Models;
|
||||
using Elsa.Serialization;
|
||||
|
||||
namespace Elsa.Persistence.MongoDb.Serialization
|
||||
{
|
||||
public class VariableSerializer : JsonSerializerBase<Variable>
|
||||
{
|
||||
public VariableSerializer(ITokenSerializer serializer) : base(serializer)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,91 +0,0 @@
|
|||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Extensions;
|
||||
using Elsa.Models;
|
||||
using MongoDB.Driver;
|
||||
using MongoDB.Driver.Linq;
|
||||
|
||||
namespace Elsa.Persistence.MongoDb.Services
|
||||
{
|
||||
public class MongoWorkflowDefinitionStore : IWorkflowDefinitionStore
|
||||
{
|
||||
private readonly IMongoCollection<WorkflowDefinitionVersion> workflowDefinitionCollection;
|
||||
private readonly IMongoCollection<WorkflowInstance> workflowInstanceCollection;
|
||||
|
||||
public MongoWorkflowDefinitionStore(
|
||||
IMongoCollection<WorkflowDefinitionVersion> workflowDefinitionCollection,
|
||||
IMongoCollection<WorkflowInstance> workflowInstanceCollection)
|
||||
{
|
||||
this.workflowDefinitionCollection = workflowDefinitionCollection;
|
||||
this.workflowInstanceCollection = workflowInstanceCollection;
|
||||
}
|
||||
|
||||
public async Task<WorkflowDefinitionVersion> SaveAsync(
|
||||
WorkflowDefinitionVersion definition,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
await workflowDefinitionCollection.ReplaceOneAsync(
|
||||
x => x.Id == definition.Id && x.Version == definition.Version,
|
||||
definition,
|
||||
new ReplaceOptions { IsUpsert = true },
|
||||
cancellationToken
|
||||
);
|
||||
|
||||
return definition;
|
||||
}
|
||||
|
||||
public async Task<WorkflowDefinitionVersion> AddAsync(WorkflowDefinitionVersion definition, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await workflowDefinitionCollection.InsertOneAsync(definition, new InsertOneOptions(), cancellationToken);
|
||||
return definition;
|
||||
}
|
||||
|
||||
public Task<WorkflowDefinitionVersion> GetByIdAsync(string id, CancellationToken cancellationToken = default)
|
||||
{
|
||||
return workflowDefinitionCollection.AsQueryable().FirstOrDefaultAsync(x => x.Id == id, cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<WorkflowDefinitionVersion> GetByIdAsync(
|
||||
string definitionId,
|
||||
VersionOptions version,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
var query = (IMongoQueryable<WorkflowDefinitionVersion>)workflowDefinitionCollection.AsQueryable()
|
||||
.Where(x => x.DefinitionId == definitionId)
|
||||
.WithVersion(version);
|
||||
|
||||
return await query.FirstOrDefaultAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowDefinitionVersion>> ListAsync(
|
||||
VersionOptions version,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
var query = workflowDefinitionCollection.AsQueryable();
|
||||
var results = await query.ToListAsync(cancellationToken);
|
||||
return results.WithVersion(version);
|
||||
}
|
||||
|
||||
public async Task<WorkflowDefinitionVersion> UpdateAsync(WorkflowDefinitionVersion definition,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
await workflowDefinitionCollection.ReplaceOneAsync(
|
||||
x => x.Id == definition.Id && x.Version == definition.Version,
|
||||
definition,
|
||||
new ReplaceOptions { IsUpsert = false },
|
||||
cancellationToken
|
||||
);
|
||||
|
||||
return definition;
|
||||
}
|
||||
|
||||
public async Task<int> DeleteAsync(string id, CancellationToken cancellationToken = default)
|
||||
{
|
||||
await workflowInstanceCollection.DeleteManyAsync(x => x.DefinitionId == id, cancellationToken);
|
||||
var result = await workflowDefinitionCollection.DeleteManyAsync(x => x.DefinitionId == id, cancellationToken);
|
||||
|
||||
return (int) result.DeletedCount;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,136 +0,0 @@
|
|||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Extensions;
|
||||
using Elsa.Models;
|
||||
using MongoDB.Driver;
|
||||
using MongoDB.Driver.Linq;
|
||||
|
||||
namespace Elsa.Persistence.MongoDb.Services
|
||||
{
|
||||
public class MongoWorkflowInstanceStore : IWorkflowInstanceStore
|
||||
{
|
||||
private readonly IMongoCollection<WorkflowInstance> collection;
|
||||
|
||||
public MongoWorkflowInstanceStore(IMongoCollection<WorkflowInstance> collection)
|
||||
{
|
||||
this.collection = collection;
|
||||
}
|
||||
|
||||
public async Task<WorkflowInstance> SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken)
|
||||
{
|
||||
await collection.ReplaceOneAsync(
|
||||
x => x.Id == instance.Id,
|
||||
instance,
|
||||
new ReplaceOptions { IsUpsert = true },
|
||||
cancellationToken);
|
||||
|
||||
return instance;
|
||||
}
|
||||
|
||||
public async Task<WorkflowInstance> GetByIdAsync(
|
||||
string id,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return await collection.AsQueryable().Where(x => x.Id == id).FirstOrDefaultAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<WorkflowInstance> GetByCorrelationIdAsync(
|
||||
string correlationId,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
return await collection.AsQueryable()
|
||||
.Where(x => x.CorrelationId == correlationId)
|
||||
.FirstOrDefaultAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowInstance>> ListByDefinitionAsync(
|
||||
string definitionId,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return await collection.AsQueryable()
|
||||
.Where(x => x.DefinitionId == definitionId)
|
||||
.OrderByDescending(x => x.CreatedAt)
|
||||
.ToListAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowInstance>> ListAllAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
return await collection
|
||||
.AsQueryable()
|
||||
.OrderByDescending(x => x.CreatedAt)
|
||||
.ToListAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityTagAsync(
|
||||
string activityType,
|
||||
string tag,
|
||||
string correlationId = null,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
var query = collection.AsQueryable();
|
||||
|
||||
query = query.Where(x => x.Status == WorkflowStatus.Suspended);
|
||||
|
||||
if (!string.IsNullOrWhiteSpace(correlationId))
|
||||
query = query.Where(x => x.CorrelationId == correlationId);
|
||||
|
||||
query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType && y.Tag == tag));
|
||||
query = query.OrderByDescending(x => x.CreatedAt);
|
||||
|
||||
var instances = await query.ToListAsync(cancellationToken);
|
||||
|
||||
return instances.GetBlockingActivities(activityType);
|
||||
}
|
||||
public async Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityAsync(
|
||||
string activityType,
|
||||
string correlationId = default,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
var query = collection.AsQueryable();
|
||||
|
||||
query = query.Where(x => x.Status == WorkflowStatus.Suspended);
|
||||
|
||||
if (!string.IsNullOrWhiteSpace(correlationId))
|
||||
query = query.Where(x => x.CorrelationId == correlationId);
|
||||
|
||||
query = query.Where(x => x.BlockingActivities.Any(y => y.ActivityType == activityType));
|
||||
query = query.OrderByDescending(x => x.CreatedAt);
|
||||
|
||||
var instances = await query.ToListAsync(cancellationToken);
|
||||
|
||||
return instances.GetBlockingActivities(activityType);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowInstance>> ListByStatusAsync(
|
||||
string definitionId,
|
||||
WorkflowStatus status,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return await collection
|
||||
.AsQueryable()
|
||||
.Where(x => x.DefinitionId == definitionId && x.Status == status)
|
||||
.OrderByDescending(x => x.CreatedAt)
|
||||
.ToListAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task<IEnumerable<WorkflowInstance>> ListByStatusAsync(
|
||||
WorkflowStatus status,
|
||||
CancellationToken cancellationToken)
|
||||
{
|
||||
return await collection
|
||||
.AsQueryable()
|
||||
.Where(x => x.Status == status)
|
||||
.OrderByDescending(x => x.CreatedAt)
|
||||
.ToListAsync(cancellationToken);
|
||||
}
|
||||
|
||||
public async Task DeleteAsync(
|
||||
string id,
|
||||
CancellationToken cancellationToken = default)
|
||||
{
|
||||
await collection.DeleteOneAsync(x => x.Id == id, cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
Binary file not shown.
|
Before Width: | Height: | Size: 16 KiB |
Loading…
Reference in a new issue