97 lines
3.9 KiB
C#
97 lines
3.9 KiB
C#
using System;
|
|
using Elsa.Helpers;
|
|
using Elsa.Persistence.Models;
|
|
using Elsa.Runtime.Contracts;
|
|
using Elsa.Runtime.Extensions;
|
|
using Elsa.Runtime.ProtoActor.Grains;
|
|
using Elsa.Runtime.ProtoActor.HostedServices;
|
|
using Elsa.Runtime.ProtoActor.Services;
|
|
using Elsa.Runtime.Protos;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Proto;
|
|
using Proto.Cluster;
|
|
using Proto.Cluster.Partition;
|
|
using Proto.Cluster.Testing;
|
|
using Proto.DependencyInjection;
|
|
using Proto.Remote;
|
|
using Proto.Remote.GrpcCore;
|
|
|
|
namespace Elsa.Runtime.ProtoActor.Extensions;
|
|
|
|
public static class ServiceCollectionExtensions
|
|
{
|
|
public static IServiceCollection AddProtoActorWorkflowHost(this IServiceCollection services)
|
|
{
|
|
var systemConfig = GetSystemConfig();
|
|
|
|
// Actor System.
|
|
services.AddSingleton(sp =>
|
|
{
|
|
var system = new ActorSystem(systemConfig).WithServiceProvider(sp);
|
|
var remoteConfig = GetRemoteConfig();
|
|
var clusterConfig = GetClusterConfig(system, "my-cluster");
|
|
|
|
system
|
|
.WithRemote(remoteConfig)
|
|
.WithCluster(clusterConfig);
|
|
|
|
return system;
|
|
});
|
|
|
|
// Cluster.
|
|
services.AddSingleton(sp => sp.GetRequiredService<ActorSystem>().Cluster());
|
|
|
|
// Actors.
|
|
services
|
|
.AddSingleton(sp => new WorkflowDefinitionGrainActor((context, _) => ActivatorUtilities.CreateInstance<WorkflowDefinitionGrain>(sp, context)))
|
|
.AddSingleton(sp => new WorkflowInstanceGrainActor((context, _) => ActivatorUtilities.CreateInstance<WorkflowInstanceGrain>(sp, context)));
|
|
|
|
// Client factory.
|
|
services.AddSingleton<GrainClientFactory>();
|
|
|
|
// Configure runtime with ProtoActor workflow invoker.
|
|
services.ConfigureWorkflowRuntime(options => options.WorkflowInvokerFactory = sp => ActivatorUtilities.CreateInstance<ProtoActorWorkflowInvoker>(sp));
|
|
|
|
return services
|
|
.AddHostedService<WorkflowServerHost>();
|
|
}
|
|
|
|
private static ActorSystemConfig GetSystemConfig() =>
|
|
ActorSystemConfig
|
|
.Setup()
|
|
.WithDeadLetterThrottleCount(3)
|
|
.WithDeadLetterThrottleInterval(TimeSpan.FromSeconds(10000))
|
|
.WithDeveloperSupervisionLogging(true)
|
|
.WithDeadLetterRequestLogging(true);
|
|
|
|
private static GrpcCoreRemoteConfig GetRemoteConfig() => GrpcCoreRemoteConfig
|
|
.BindToLocalhost()
|
|
.WithProtoMessages(Protos.MessagesReflection.Descriptor);
|
|
|
|
private static ClusterConfig GetClusterConfig(ActorSystem system, string clusterName)
|
|
{
|
|
//var clusterProvider = new ConsulProvider(new ConsulProviderConfig{});
|
|
var clusterProvider = new TestProvider(new TestProviderOptions(), new InMemAgent());
|
|
|
|
var workflowDefinitionProps = system.DI().PropsFor<WorkflowDefinitionGrainActor>();
|
|
var workflowInstanceProps = system.DI().PropsFor<WorkflowInstanceGrainActor>();
|
|
|
|
var clusterConfig =
|
|
ClusterConfig
|
|
// .Setup("MyCluster", clusterProvider, new IdentityStorageLookup(GetIdentityLookup(clusterName)))
|
|
.Setup(clusterName, clusterProvider, new PartitionIdentityLookup())
|
|
.WithTimeout(TimeSpan.FromHours(1))
|
|
.WithActorRequestTimeout(TimeSpan.FromHours(1))
|
|
.WithActorActivationTimeout(TimeSpan.FromHours(1))
|
|
.WithActorSpawnTimeout(TimeSpan.FromHours(1))
|
|
.WithClusterKind(WorkflowDefinitionGrainActor.Kind, workflowDefinitionProps)
|
|
.WithClusterKind(WorkflowInstanceGrainActor.Kind, workflowInstanceProps)
|
|
;
|
|
return clusterConfig;
|
|
}
|
|
|
|
// private static IIdentityStorage GetIdentityLookup(string clusterName) =>
|
|
// new RedisIdentityStorage(clusterName, ConnectionMultiplexer
|
|
// .Connect("localhost:6379" /* use proper config */)
|
|
// );
|
|
} |