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());
}