2017-02-15 21:05:01 +00:00
|
|
|
|
using System;
|
|
|
|
|
using System.Linq.Expressions;
|
|
|
|
|
using System.Reflection;
|
|
|
|
|
using System.Threading.Tasks;
|
|
|
|
|
using Tapeti.Config;
|
|
|
|
|
using Tapeti.Default;
|
|
|
|
|
|
|
|
|
|
namespace Tapeti.Flow.Default
|
|
|
|
|
{
|
|
|
|
|
public class FlowStarter : IFlowStarter
|
|
|
|
|
{
|
|
|
|
|
private readonly IConfig config;
|
2017-10-17 08:34:07 +00:00
|
|
|
|
private readonly ILogger logger;
|
2017-02-15 21:05:01 +00:00
|
|
|
|
|
|
|
|
|
|
2017-10-17 08:34:07 +00:00
|
|
|
|
public FlowStarter(IConfig config, ILogger logger)
|
2017-02-15 21:05:01 +00:00
|
|
|
|
{
|
|
|
|
|
this.config = config;
|
2017-10-17 08:34:07 +00:00
|
|
|
|
this.logger = logger;
|
2017-02-15 21:05:01 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public Task Start<TController>(Expression<Func<TController, Func<IYieldPoint>>> methodSelector) where TController : class
|
|
|
|
|
{
|
2017-09-22 09:19:49 +00:00
|
|
|
|
return CallControllerMethod<TController>(GetExpressionMethod(methodSelector), value => Task.FromResult((IYieldPoint)value), new object[] { });
|
2017-02-15 21:05:01 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public Task Start<TController>(Expression<Func<TController, Func<Task<IYieldPoint>>>> methodSelector) where TController : class
|
|
|
|
|
{
|
2017-09-22 09:19:49 +00:00
|
|
|
|
return CallControllerMethod<TController>(GetExpressionMethod(methodSelector), value => (Task<IYieldPoint>)value, new object[] {});
|
2017-02-15 21:05:01 +00:00
|
|
|
|
}
|
|
|
|
|
|
2017-09-22 09:19:49 +00:00
|
|
|
|
public Task Start<TController, TParameter>(Expression<Func<TController, Func<TParameter, IYieldPoint>>> methodSelector, TParameter parameter) where TController : class
|
|
|
|
|
{
|
|
|
|
|
return CallControllerMethod<TController>(GetExpressionMethod(methodSelector), value => Task.FromResult((IYieldPoint)value), new object[] {parameter});
|
|
|
|
|
}
|
2017-02-15 21:05:01 +00:00
|
|
|
|
|
2017-09-22 09:19:49 +00:00
|
|
|
|
public Task Start<TController, TParameter>(Expression<Func<TController, Func<TParameter, Task<IYieldPoint>>>> methodSelector, TParameter parameter) where TController : class
|
|
|
|
|
{
|
|
|
|
|
return CallControllerMethod<TController>(GetExpressionMethod(methodSelector), value => (Task<IYieldPoint>)value, new object[] {parameter});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private async Task CallControllerMethod<TController>(MethodInfo method, Func<object, Task<IYieldPoint>> getYieldPointResult, object[] parameters) where TController : class
|
2017-02-15 21:05:01 +00:00
|
|
|
|
{
|
|
|
|
|
var controller = config.DependencyResolver.Resolve<TController>();
|
2017-09-22 09:19:49 +00:00
|
|
|
|
var yieldPoint = await getYieldPointResult(method.Invoke(controller, parameters));
|
2017-02-15 21:05:01 +00:00
|
|
|
|
|
|
|
|
|
var context = new MessageContext
|
|
|
|
|
{
|
|
|
|
|
DependencyResolver = config.DependencyResolver,
|
|
|
|
|
Controller = controller
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
var flowHandler = config.DependencyResolver.Resolve<IFlowHandler>();
|
2017-10-17 08:34:07 +00:00
|
|
|
|
|
2017-10-17 11:29:16 +00:00
|
|
|
|
HandlingResultBuilder handlingResult = new HandlingResultBuilder
|
|
|
|
|
{
|
|
|
|
|
ConsumeResponse = ConsumeResponse.Nack,
|
|
|
|
|
};
|
2017-10-17 08:34:07 +00:00
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
await flowHandler.Execute(context, yieldPoint);
|
2017-10-17 11:29:16 +00:00
|
|
|
|
handlingResult.ConsumeResponse = ConsumeResponse.Ack;
|
2017-10-17 08:34:07 +00:00
|
|
|
|
}
|
|
|
|
|
finally
|
|
|
|
|
{
|
2017-10-17 11:29:16 +00:00
|
|
|
|
await RunCleanup(context, handlingResult.ToHandlingResult());
|
2017-10-17 08:34:07 +00:00
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2017-10-17 11:29:16 +00:00
|
|
|
|
private async Task RunCleanup(MessageContext context, HandlingResult handlingResult)
|
2017-10-17 08:34:07 +00:00
|
|
|
|
{
|
|
|
|
|
foreach (var handler in config.CleanupMiddleware)
|
|
|
|
|
{
|
|
|
|
|
try
|
|
|
|
|
{
|
2017-10-17 11:29:16 +00:00
|
|
|
|
await handler.Handle(context, handlingResult);
|
2017-10-17 08:34:07 +00:00
|
|
|
|
}
|
|
|
|
|
catch (Exception eCleanup)
|
|
|
|
|
{
|
|
|
|
|
logger.HandlerException(eCleanup);
|
|
|
|
|
}
|
|
|
|
|
}
|
2017-02-15 21:05:01 +00:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
private static MethodInfo GetExpressionMethod<TController, TResult>(Expression<Func<TController, Func<TResult>>> methodSelector)
|
|
|
|
|
{
|
|
|
|
|
var callExpression = (methodSelector.Body as UnaryExpression)?.Operand as MethodCallExpression;
|
|
|
|
|
var targetMethodExpression = callExpression?.Object as ConstantExpression;
|
|
|
|
|
|
|
|
|
|
var method = targetMethodExpression?.Value as MethodInfo;
|
|
|
|
|
if (method == null)
|
|
|
|
|
throw new ArgumentException("Unable to determine the starting method", nameof(methodSelector));
|
|
|
|
|
|
|
|
|
|
return method;
|
|
|
|
|
}
|
2017-09-22 09:19:49 +00:00
|
|
|
|
|
|
|
|
|
private static MethodInfo GetExpressionMethod<TController, TResult, TParameter>(Expression<Func<TController, Func<TParameter, TResult>>> methodSelector)
|
|
|
|
|
{
|
|
|
|
|
var callExpression = (methodSelector.Body as UnaryExpression)?.Operand as MethodCallExpression;
|
|
|
|
|
var targetMethodExpression = callExpression?.Object as ConstantExpression;
|
|
|
|
|
|
|
|
|
|
var method = targetMethodExpression?.Value as MethodInfo;
|
|
|
|
|
if (method == null)
|
|
|
|
|
throw new ArgumentException("Unable to determine the starting method", nameof(methodSelector));
|
|
|
|
|
|
|
|
|
|
return method;
|
|
|
|
|
}
|
2017-02-15 21:05:01 +00:00
|
|
|
|
}
|
|
|
|
|
}
|