ASP NET Core MVC обработка исключений в EasyNetQ RPC

Имеется два сервиса. Сервис A это обычный WebAPI проект. Сервис B запускается в виде службы Windows (Worker Service). Сервис B хранит некоторые данные в оперативной памяти, поэтому для получения данных из сервиса B использую RPC вызовы.

Проблема в том, что в случае возникновения исключения в сервисе B, я не знаю, как вернуть объект этого исключения в сервис A. На данный момент, при возникновении исключения, RPC вызов падает по таймауту через 5 секунд (указано в настройках)

System.AggregateException: One or more errors occurred. (A task was canceled.)
 ---> System.Threading.Tasks.TaskCanceledException: A task was canceled.
   at System.Threading.Tasks.Task.GetExceptions(Boolean includeTaskCanceledExceptions)
   at System.Threading.Tasks.Task`1.GetResultCore(Boolean waitCompletionNotification)
   at System.Threading.Tasks.Task`1.get_Result()
   at MediatR.Internal.RequestHandlerWrapperImpl`2.<>c.<Handle>b__0_0(Task`1 t)
   at System.Threading.Tasks.ContinuationResultTaskFromResultTask`2.InnerInvoke()
   at System.Threading.Tasks.Task.<>c.<.cctor>b__277_0(Object obj)
   at System.Threading.ExecutionContext.RunFromThreadPoolDispatchLoop(Thread threadPoolThread, ExecutionContext executionContext, ContextCallback callback, Object state)
   at System.Threading.Tasks.Task.ExecuteWithThreadLocal(Task& currentTaskSlot, Thread threadPoolThread)
   at System.Threading.Tasks.Task.ExecuteEntryUnsafe(Thread threadPoolThread)
   at System.Threading.Tasks.Task.ExecuteFromThreadPool(Thread threadPoolThread)
   at System.Threading.ThreadPoolWorkQueue.Dispatch()
--- End of stack trace from previous location ---

   --- End of inner exception stack trace ---
   at System.Threading.Tasks.Task`1.GetResultCore(Boolean waitCompletionNotification)
   at System.Threading.Tasks.Task`1.get_Result()
   at MediatR.Internal.RequestHandlerWrapperImpl`2.<>c.<Handle>b__0_0(Task`1 t)
   at System.Threading.Tasks.ContinuationResultTaskFromResultTask`2.InnerInvoke()
   at System.Threading.Tasks.Task.<>c.<.cctor>b__277_0(Object obj)
   at System.Threading.ExecutionContext.RunFromThreadPoolDispatchLoop(Thread threadPoolThread, ExecutionContext executionContext, ContextCallback callback, Object state)
--- End of stack trace from previous location ---
   at System.Threading.ExecutionContext.RunFromThreadPoolDispatchLoop(Thread threadPoolThread, ExecutionContext executionContext, ContextCallback callback, Object state)
   at System.Threading.Tasks.Task.ExecuteWithThreadLocal(Task& currentTaskSlot, Thread threadPoolThread)
--- End of stack trace from previous location ---
   at Pharmacy.Services.Drugs.Controllers.BaseController.Query(Object command) in C:\Users\_murashik\source\repos\abupharmacyproject\Pharmacy.Services.Drugs\Controllers\BaseController.cs:line 22
   at Pharmacy.Services.Drugs.Controllers.DrugsController.Search(String text) in C:\Users\_murashik\source\repos\abupharmacyproject\Pharmacy.Services.Drugs\Controllers\DrugsController.cs:line 31
   at Microsoft.AspNetCore.Mvc.Infrastructure.ActionMethodExecutor.TaskOfIActionResultExecutor.Execute(IActionResultTypeMapper mapper, ObjectMethodExecutor executor, Object controller, Object[] arguments)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ControllerActionInvoker.<InvokeActionMethodAsync>g__Awaited|12_0(ControllerActionInvoker invoker, ValueTask`1 actionResultValueTask)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ControllerActionInvoker.<InvokeNextActionFilterAsync>g__Awaited|10_0(ControllerActionInvoker invoker, Task lastTask, State next, Scope scope, Object state, Boolean isCompleted)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ControllerActionInvoker.Rethrow(ActionExecutedContextSealed context)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ControllerActionInvoker.Next(State& next, Scope& scope, Object& state, Boolean& isCompleted)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ControllerActionInvoker.<InvokeInnerFilterAsync>g__Awaited|13_0(ControllerActionInvoker invoker, Task lastTask, State next, Scope scope, Object state, Boolean isCompleted)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ResourceInvoker.<InvokeFilterPipelineAsync>g__Awaited|19_0(ResourceInvoker invoker, Task lastTask, State next, Scope scope, Object state, Boolean isCompleted)
   at Microsoft.AspNetCore.Mvc.Infrastructure.ResourceInvoker.<InvokeAsync>g__Awaited|17_0(ResourceInvoker invoker, Task task, IDisposable scope)
   at Microsoft.AspNetCore.Routing.EndpointMiddleware.<Invoke>g__AwaitRequestTask|6_0(Endpoint endpoint, Task requestTask, ILogger logger)
   at Microsoft.AspNetCore.Authorization.AuthorizationMiddleware.Invoke(HttpContext context)
   at Swashbuckle.AspNetCore.SwaggerUI.SwaggerUIMiddleware.Invoke(HttpContext httpContext)
   at Swashbuckle.AspNetCore.Swagger.SwaggerMiddleware.Invoke(HttpContext httpContext, ISwaggerProvider swaggerProvider)
   at Microsoft.AspNetCore.Diagnostics.DeveloperExceptionPageMiddleware.Invoke(HttpContext context)

Ниже я представлю код общей библиотеки, сервиса A и сервиса B:

Общая библиотека (Core.dll) Инициализация Rpc сервера происходит в сервисе B. В методе RespondAsync имитируем исключение:

public class RpcServer : IRpcServer
    {
        private IBus _bus;

        private ILogger _logger;
        private IMediator _mediator;
        private IMapper _mapper;
        private IDictionary<string, Type> RequestTypes { get; set; }

        public RpcServer(IServiceProvider serviceProvider)
        {
            RequestTypes = new Dictionary<string, Type>();

            _logger = serviceProvider.GetRequiredService<ILogger>();
            _mapper = serviceProvider.GetService<IMapper>();
            _mediator = serviceProvider.GetService<IMediator>();
            _bus = serviceProvider.GetService<IBus>();
        }

        public IRpcServer RegisterCommandHandler<TRequest, TResponse>() 
            where TRequest : class, IRequest<TResponse>
            where TResponse : class
        {
            Type requestType = typeof(TRequest);

            RequestTypes.Add(requestType.Name, requestType);

            _bus.Rpc.Respond<TRequest, object>(async r => await RespondAsync(r));
            return this;
        }

        public void StartResponding() { }

        private async Task<object> RespondAsync<TRequest>(TRequest request)
        {
            try
            {
                string requestName = typeof(TRequest).Name;

                if (!RequestTypes.ContainsKey(requestName))
                    return null;

                var requestObjectInstance = Activator.CreateInstance(RequestTypes[requestName]);
                var mappedObject = _mapper.Map<object, object>(request, requestObjectInstance);

                throw new Exception("Test");
                //return await _mediator.Send(mappedObject);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, nameof(RpcServer));
                return ex;
            }

            return null;
        }
    }

Диспетчер запросов RPC:

public class RpcRequestDispatcher : IRpcRequestDispatcher
    {
        private IBus _bus;
        public string ExchangeName { get; set; }

        public RpcRequestDispatcher(IBus bus)
        {
            _bus = bus;
        }

        public async Task<TResponse> Request<TRequest, TResponse>(TRequest request)
            => await _bus.Rpc.RequestAsync<TRequest, TResponse>(request);
    }

Кастомный сериализатор/десериализатор. EasyNetQ загружает типы из подключенных сборок, но так как в приложении присутствует понятие домен - стандартное поведение меня не устраивает.

public class CustomTypeNameSerializer : ITypeNameSerializer
    {
        public ITypeNameSerializerSettings TypeNameSerializerSettings { get; set; }

        public CustomTypeNameSerializer(ITypeNameSerializerSettings typeNameSerializerSettings)
        {
            TypeNameSerializerSettings = typeNameSerializerSettings;
        }
        /// <inheritdoc/>
        public Type DeSerialize(string typeName)
        {
            var splittedType = typeName.Split('_');

            Type resultType = null;

            switch (splittedType[0])
            {
                case "Response":
                    resultType = TypeNameSerializerSettings.ResponseTypes.FirstOrDefault(f => f.Name == splittedType[1]);
                    break;
                case "Command":
                    resultType = TypeNameSerializerSettings.CommandTypes.FirstOrDefault(f => f.Name == splittedType[1]);
                    break;
            }

            return resultType;
        }

        public string Serialize(Type type)
        {
            string typeName = null;

            if (TypeNameSerializerSettings.ResponseTypes != null && TypeNameSerializerSettings.ResponseTypes.Contains(type))
            {
                typeName = $"Response_{type.Name}";
            }
            else if (TypeNameSerializerSettings.CommandTypes != null && TypeNameSerializerSettings.CommandTypes.Contains(type))
            {
                typeName = $"Command_{type.Name}";
            }

            if (typeName == null)
            {
                typeName = RemoveAssemblyDetails(type.AssemblyQualifiedName);
            }

            if (typeName.Length > 255)
            {
                throw new EasyNetQException($"The serialized name of type '{type.Name}' exceeds the AMQP maximum short string length of 255 characters");
            }

            return typeName;
        }

        private static string RemoveAssemblyDetails(string fullyQualifiedTypeName)
        {
            var builder = new StringBuilder(fullyQualifiedTypeName.Length);

            // loop through the type name and filter out qualified assembly details from nested type names
            var writingAssemblyName = false;
            var skippingAssemblyDetails = false;
            foreach (var character in fullyQualifiedTypeName)
            {
                switch (character)
                {
                    case '[':
                        writingAssemblyName = false;
                        skippingAssemblyDetails = false;
                        builder.Append(character);
                        break;
                    case ']':
                        writingAssemblyName = false;
                        skippingAssemblyDetails = false;
                        builder.Append(character);
                        break;
                    case ',':
                        if (!writingAssemblyName)
                        {
                            writingAssemblyName = true;
                            builder.Append(character);
                        }
                        else
                        {
                            skippingAssemblyDetails = true;
                        }
                        break;
                    default:
                        if (!skippingAssemblyDetails)
                        {
                            builder.Append(character);
                        }
                        break;
                }
            }

            return builder.ToString();
        }
    }

Здесь хранятся массивы типов ответа и команд запросов:

public interface ITypeNameSerializerSettings
    {
        Type[] ResponseTypes { get; set; }
        Type[] CommandTypes { get; set; }
    }

Переопределение стандартного сериализатора JSON:

public class CustomJsonSerializer : ISerializer
    {
        private static readonly Encoding Encoding = new UTF8Encoding(false);
        private static readonly Newtonsoft.Json.JsonSerializerSettings DefaultSerializerSettings =
            new Newtonsoft.Json.JsonSerializerSettings
            {
                TypeNameHandling = Newtonsoft.Json.TypeNameHandling.Auto
            };

        private const int DefaultBufferSize = 1024;

        private readonly Newtonsoft.Json.JsonSerializer jsonSerializer;

        public CustomJsonSerializer() : this(DefaultSerializerSettings)
        {
        }

        public CustomJsonSerializer(Newtonsoft.Json.JsonSerializerSettings serializerSettings)
        {
            jsonSerializer = Newtonsoft.Json.JsonSerializer.Create(serializerSettings);
        }

        /// <inheritdoc />
        public byte[] MessageToBytes(Type messageType, object message)
        {
            using (var memoryStream = new MemoryStream(DefaultBufferSize))
            {
                using (var streamWriter = new StreamWriter(memoryStream, Encoding, DefaultBufferSize, true))
                using (var jsonWriter = new Newtonsoft.Json.JsonTextWriter(streamWriter))
                {
                    jsonWriter.Formatting = jsonSerializer.Formatting;
                    jsonSerializer.Serialize(jsonWriter, message);
                }

                return memoryStream.ToArray();
            }
            
        }

        /// <inheritdoc />
        public object BytesToMessage(Type messageType, byte[] bytes)
        {
            string message = System.Text.Encoding.UTF8.GetString(bytes);

            return Newtonsoft.Json.JsonConvert.DeserializeObject(message, messageType);

            //using (var memoryStream = new MemoryStream(bytes, false))
            //{
            //    using (var streamReader = new StreamReader(memoryStream, Encoding, false, DefaultBufferSize, true))
            //    {
            //        using (var reader = new Newtonsoft.Json.JsonTextReader(streamReader))
            //        {
            //            try
            //            {
            //                return jsonSerializer.Deserialize(reader, messageType);
            //            }
            //            catch (Exception ex)
            //            {

            //                throw;
            //            }
                       
            //        } 
            //    }  
            //}
        }
    }

Сервис A (WebAPI.dll)

public class BaseController : ControllerBase
    {
        private IMediator _mediator;

        public BaseController(IMediator mediator)
        {
            _mediator = mediator;
        }

        [ApiExplorerSettings(IgnoreApi = true)]
        public async Task<IActionResult> Query(object command)
        {
            return Ok(new
            {
                Code = 0,
                Message = "Ok",
                Data = await _mediator.Send(command)
            });
        }
    }

DrugsController.cs

    [Route("api/v{version:apiVersion}/[controller]")]
    [ApiController]
    [ApiVersion("1")]
    public class DrugsController : BaseController
    {
        private IMediator _mediator;

        public DrugsController(IMediator mediator) : base(mediator)
        { 
        }

        [HttpGet("get/{id}")]
        public async Task<IActionResult> Get(int id)
        {
            return null;
        }

        [HttpGet("search/{text}")]
        public async Task<IActionResult> Search(string text)
            => await Query(new GetDrugResultByTextSearchCommand() { Text = text, MaxCountResult = 20 });
}

GetDrugResultByTextSearchCommand.cs

public class GetDrugResultByTextSearchCommand : IRequest<TextSearchByDrugsResponse>, ICommand
    {
        public string Text { get; set; }
        public int MaxCountResult { get; set; }
    }

MediatR перенаправляет команду на данный обработчик:

public class GetDrugResultByTextSearchHandler : IRequestHandler<GetDrugResultByTextSearchCommand, TextSearchByDrugsResponse>
    {
        private IRpcRequestDispatcher _rpcRequestDispatcher;

        public GetDrugResultByTextSearchHandler(IRpcRequestDispatcher rpcRequestDispatcher)
        {
            _rpcRequestDispatcher = rpcRequestDispatcher;
        }

        public async Task<TextSearchByDrugsResponse> Handle(GetDrugResultByTextSearchCommand request, CancellationToken cancellationToken)
        {
            return await _rpcRequestDispatcher.Request<GetDrugResultByTextSearchCommand, TextSearchByDrugsResponse>(request);

        }
    }

Инициализация EasyNetQ в StartUp.cs

public class RabbitMqBuilder
    {
        private IServiceCollection _serviceCollection;
        private TypeNameSerializerSettings _typeNameSerializerSettings;
        private Assembly _currentAssembly;

        public RabbitMqBuilder(IServiceCollection services, Assembly currentAssembly)
        {
            _serviceCollection = services;
            _currentAssembly = currentAssembly;
            _typeNameSerializerSettings = new TypeNameSerializerSettings();

            _typeNameSerializerSettings.Exception = typeof(RabbitMqException);
            Create();
        }

        private void Create()
        {
            _serviceCollection.AddSingleton<IBus>(x => RabbitHutch.CreateBus("host=localhost;username=guest;password=guest;timeout=5", y =>
            {
                y.Register<ITypeNameSerializerSettings>(_typeNameSerializerSettings);
                y.Register<ITypeNameSerializer, CustomTypeNameSerializer>(EasyNetQ.DI.Lifetime.Singleton);
                y.Register<ISerializer>(new CustomJsonSerializer());
            }));
        }

        public RabbitMqBuilder RegisterReponsesAsImplementedInterface<T>() where T : class
        {
            _typeNameSerializerSettings.ResponseTypes = ResolveTypes(typeof(T));

            return this;
        }

        public RabbitMqBuilder RegisterCommandsAsImplementedInterface<T>() where T : class
        {
            _typeNameSerializerSettings.CommandTypes = ResolveTypes(typeof(T));

            return this;
        }

        private Type[] ResolveTypes(Type type)
        {
            if (type.IsGenericType)
            {
                return _currentAssembly.GetTypes()
                 .Where(mytype => mytype.GetInterfaces()
                    .Any(a => a.IsGenericType && a.GetGenericTypeDefinition() == type))
                 .ToArray();
            }

            return _currentAssembly.GetTypes()
                 .Where(mytype => mytype.GetInterfaces().Contains(type))
                 .ToArray();
        }
    }

public void ConfigureServices(IServiceCollection services)
        {
            services.AddApiVersioning(x => 
            {
                x.DefaultApiVersion = new ApiVersion(1, 0);
                x.AssumeDefaultVersionWhenUnspecified = true;
            });
            services.AddControllers();

            services.AddRabbitMq(Assembly.GetExecutingAssembly())
                    .RegisterCommandsAsImplementedInterface<ICommand>()
                    .RegisterReponsesAsImplementedInterface<IResponse>();

            services.AddSingleton<IRpcRequestDispatcher, RpcRequestDispatcher>();
            services.AddSwaggerGen(c =>
            {
                c.SwaggerDoc("v1", new OpenApiInfo { Title = "Service.Test", Version = "v1" });
            });

            services.AddMediatR(AppDomain.CurrentDomain.GetAssemblies());
        }

Ответы (0 шт):