mirror of https://github.com/abpframework/abp.git
committed by
GitHub
49 changed files with 1115 additions and 3 deletions
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,18 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>net6.0</TargetFramework> |
||||
|
<ImplicitUsings>enable</ImplicitUsings> |
||||
|
<Nullable>enable</Nullable> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.AspNetCore.Mvc.Dapr\Volo.Abp.AspNetCore.Mvc.Dapr.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.EventBus.Dapr\Volo.Abp.EventBus.Dapr.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,29 @@ |
|||||
|
using Microsoft.AspNetCore.Http.Json; |
||||
|
using Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.SystemTextJson; |
||||
|
using Volo.Abp.EventBus.Dapr; |
||||
|
using Volo.Abp.Json.SystemTextJson; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus; |
||||
|
|
||||
|
[DependsOn( |
||||
|
typeof(AbpAspNetCoreMvcDaprModule), |
||||
|
typeof(AbpEventBusDaprModule) |
||||
|
)] |
||||
|
public class AbpAspNetCoreMvcDaprEventBusModule : AbpModule |
||||
|
{ |
||||
|
public override void ConfigureServices(ServiceConfigurationContext context) |
||||
|
{ |
||||
|
// TODO: Add NewtonsoftJson json converter.
|
||||
|
|
||||
|
Configure<JsonOptions>(options => |
||||
|
{ |
||||
|
options.SerializerOptions.Converters.Add(new AbpAspNetCoreMvcDaprSubscriptionDefinitionConverter()); |
||||
|
}); |
||||
|
|
||||
|
Configure<AbpSystemTextJsonSerializerOptions>(options => |
||||
|
{ |
||||
|
options.JsonSerializerOptions.Converters.Add(new AbpAspNetCoreMvcDaprSubscriptionDefinitionConverter()); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprEventBusOptions |
||||
|
{ |
||||
|
public List<IAbpAspNetCoreMvcDaprPubSubProviderContributor> Contributors { get; } |
||||
|
|
||||
|
public AbpAspNetCoreMvcDaprEventBusOptions() |
||||
|
{ |
||||
|
Contributors = new List<IAbpAspNetCoreMvcDaprPubSubProviderContributor>(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,8 @@ |
|||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprPubSubConsts |
||||
|
{ |
||||
|
public const string DaprSubscribeUrl = "dapr/subscribe"; |
||||
|
|
||||
|
public const string DaprEventCallbackUrl = "api/abp/dapr/event"; |
||||
|
} |
||||
@ -0,0 +1,63 @@ |
|||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Models; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.EventBus; |
||||
|
using Volo.Abp.EventBus.Dapr; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprPubSubProvider : ITransientDependency |
||||
|
{ |
||||
|
protected IServiceProvider ServiceProvider { get; } |
||||
|
protected AbpAspNetCoreMvcDaprEventBusOptions AspNetCoreMvcDaprEventBusOptions { get; } |
||||
|
protected AbpDaprEventBusOptions DaprEventBusOptions { get; } |
||||
|
protected AbpDistributedEventBusOptions DistributedEventBusOptions { get; } |
||||
|
|
||||
|
public AbpAspNetCoreMvcDaprPubSubProvider( |
||||
|
IServiceProvider serviceProvider, |
||||
|
IOptions<AbpAspNetCoreMvcDaprEventBusOptions> aspNetCoreDaprEventBusOptions, |
||||
|
IOptions<AbpDaprEventBusOptions> daprEventBusOptions, |
||||
|
IOptions<AbpDistributedEventBusOptions> distributedEventBusOptions) |
||||
|
{ |
||||
|
ServiceProvider = serviceProvider; |
||||
|
AspNetCoreMvcDaprEventBusOptions = aspNetCoreDaprEventBusOptions.Value; |
||||
|
DaprEventBusOptions = daprEventBusOptions.Value; |
||||
|
DistributedEventBusOptions = distributedEventBusOptions.Value; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<List<AbpAspNetCoreMvcDaprSubscriptionDefinition>> GetSubscriptionsAsync() |
||||
|
{ |
||||
|
var subscriptions = new List<AbpAspNetCoreMvcDaprSubscriptionDefinition>(); |
||||
|
foreach (var handler in DistributedEventBusOptions.Handlers) |
||||
|
{ |
||||
|
foreach (var @interface in handler.GetInterfaces().Where(x => x.IsGenericType && x.GetGenericTypeDefinition() == typeof(IDistributedEventHandler<>))) |
||||
|
{ |
||||
|
var eventType = @interface.GetGenericArguments()[0]; |
||||
|
var eventName = EventNameAttribute.GetNameOrDefault(eventType); |
||||
|
|
||||
|
subscriptions.Add(new AbpAspNetCoreMvcDaprSubscriptionDefinition() |
||||
|
{ |
||||
|
PubSubName = DaprEventBusOptions.PubSubName, |
||||
|
Topic = eventName, |
||||
|
Route = AbpAspNetCoreMvcDaprPubSubConsts.DaprEventCallbackUrl |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
if (AspNetCoreMvcDaprEventBusOptions.Contributors.Any()) |
||||
|
{ |
||||
|
using (var scope = ServiceProvider.CreateScope()) |
||||
|
{ |
||||
|
var context = new AbpAspNetCoreMvcDaprPubSubProviderContributorContext(scope.ServiceProvider, subscriptions); |
||||
|
foreach (var contributor in AspNetCoreMvcDaprEventBusOptions.Contributors) |
||||
|
{ |
||||
|
await contributor.ContributeAsync(context); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
return subscriptions; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,16 @@ |
|||||
|
using Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Models; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprPubSubProviderContributorContext |
||||
|
{ |
||||
|
public IServiceProvider ServiceProvider { get; } |
||||
|
|
||||
|
public List<AbpAspNetCoreMvcDaprSubscriptionDefinition> Subscriptions { get; } |
||||
|
|
||||
|
public AbpAspNetCoreMvcDaprPubSubProviderContributorContext(IServiceProvider serviceProvider, List<AbpAspNetCoreMvcDaprSubscriptionDefinition> daprSubscriptionModels) |
||||
|
{ |
||||
|
ServiceProvider = serviceProvider; |
||||
|
Subscriptions = daprSubscriptionModels; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,36 @@ |
|||||
|
using System.Text.Json; |
||||
|
using Microsoft.AspNetCore.Mvc; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Models; |
||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.EventBus.Dapr; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Controllers; |
||||
|
|
||||
|
[Area("abp")] |
||||
|
[RemoteService(Name = "abp")] |
||||
|
public class AbpAspNetCoreMvcDaprPubSubController : AbpController |
||||
|
{ |
||||
|
[HttpGet(AbpAspNetCoreMvcDaprPubSubConsts.DaprSubscribeUrl)] |
||||
|
public virtual async Task<List<AbpAspNetCoreMvcDaprSubscriptionDefinition>> SubscribeAsync() |
||||
|
{ |
||||
|
return await HttpContext.RequestServices.GetRequiredService<AbpAspNetCoreMvcDaprPubSubProvider>().GetSubscriptionsAsync(); |
||||
|
} |
||||
|
|
||||
|
[HttpPost(AbpAspNetCoreMvcDaprPubSubConsts.DaprEventCallbackUrl)] |
||||
|
public virtual async Task<IActionResult> EventsAsync() |
||||
|
{ |
||||
|
var bodyJsonDocument = await JsonDocument.ParseAsync(HttpContext.Request.Body); |
||||
|
var request = JsonSerializer.Deserialize<AbpAspNetCoreMvcDaprSubscriptionRequest>(bodyJsonDocument.RootElement.GetRawText(), |
||||
|
HttpContext.RequestServices.GetRequiredService<IOptions<JsonOptions>>().Value.JsonSerializerOptions); |
||||
|
|
||||
|
var distributedEventBus = HttpContext.RequestServices.GetRequiredService<DaprDistributedEventBus>(); |
||||
|
var daprSerializer = HttpContext.RequestServices.GetRequiredService<IDaprSerializer>(); |
||||
|
|
||||
|
var eventData = daprSerializer.Deserialize(bodyJsonDocument.RootElement.GetProperty("data").GetRawText(), distributedEventBus.GetEventType(request.Topic)); |
||||
|
await distributedEventBus.TriggerHandlersAsync(distributedEventBus.GetEventType(request.Topic), eventData); |
||||
|
|
||||
|
return Ok(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,6 @@ |
|||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus; |
||||
|
|
||||
|
public interface IAbpAspNetCoreMvcDaprPubSubProviderContributor |
||||
|
{ |
||||
|
Task ContributeAsync(AbpAspNetCoreMvcDaprPubSubProviderContributorContext context); |
||||
|
} |
||||
@ -0,0 +1,10 @@ |
|||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Models; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprSubscriptionDefinition |
||||
|
{ |
||||
|
public string PubSubName { get; set; } |
||||
|
|
||||
|
public string Topic { get; set; } |
||||
|
|
||||
|
public string Route { get; set; } |
||||
|
} |
||||
@ -0,0 +1,8 @@ |
|||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Models; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprSubscriptionRequest |
||||
|
{ |
||||
|
public string PubSubName { get; set; } |
||||
|
|
||||
|
public string Topic { get; set; } |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
using System.Text.Json; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.SystemTextJson; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprPubSubJsonNamingPolicy : JsonNamingPolicy |
||||
|
{ |
||||
|
public override string ConvertName(string name) |
||||
|
{ |
||||
|
return name.ToLower(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,25 @@ |
|||||
|
using System.Text.Json; |
||||
|
using System.Text.Json.Serialization; |
||||
|
using Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.Models; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr.EventBus.SystemTextJson; |
||||
|
|
||||
|
public class AbpAspNetCoreMvcDaprSubscriptionDefinitionConverter : JsonConverter<AbpAspNetCoreMvcDaprSubscriptionDefinition> |
||||
|
{ |
||||
|
private JsonSerializerOptions _writeJsonSerializerOptions; |
||||
|
|
||||
|
public override AbpAspNetCoreMvcDaprSubscriptionDefinition Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options) |
||||
|
{ |
||||
|
throw new NotSupportedException(); |
||||
|
} |
||||
|
|
||||
|
public override void Write(Utf8JsonWriter writer, AbpAspNetCoreMvcDaprSubscriptionDefinition value, JsonSerializerOptions options) |
||||
|
{ |
||||
|
_writeJsonSerializerOptions ??= JsonSerializerOptionsHelper.Create(new JsonSerializerOptions(options) |
||||
|
{ |
||||
|
PropertyNamingPolicy = new AbpAspNetCoreMvcDaprPubSubJsonNamingPolicy() |
||||
|
}, x => x == this); |
||||
|
|
||||
|
JsonSerializer.Serialize(writer, value, _writeJsonSerializerOptions); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,22 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>net6.0</TargetFramework> |
||||
|
<ImplicitUsings>enable</ImplicitUsings> |
||||
|
<Nullable>enable</Nullable> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.AspNetCore.Mvc\Volo.Abp.AspNetCore.Mvc.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.Dapr\Volo.Abp.Dapr.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<PackageReference Include="Dapr.AspNetCore" Version="1.8.0" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,13 @@ |
|||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.AspNetCore.Mvc.Dapr; |
||||
|
|
||||
|
[DependsOn( |
||||
|
typeof(AbpAspNetCoreMvcModule), |
||||
|
typeof(AbpDaprModule) |
||||
|
)] |
||||
|
public class AbpAspNetCoreMvcDaprModule : AbpModule |
||||
|
{ |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,21 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>net6.0</TargetFramework> |
||||
|
<ImplicitUsings>enable</ImplicitUsings> |
||||
|
<Nullable>enable</Nullable> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.Json\Volo.Abp.Json.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<PackageReference Include="Dapr.Client" Version="1.8.0" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,48 @@ |
|||||
|
using System.Collections.Concurrent; |
||||
|
using System.Text.Json; |
||||
|
using Dapr.Client; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.Json.SystemTextJson; |
||||
|
|
||||
|
namespace Volo.Abp.Dapr; |
||||
|
|
||||
|
public class AbpDaprClientFactory : ITransientDependency |
||||
|
{ |
||||
|
protected AbpDaprOptions Options { get; } |
||||
|
protected AbpSystemTextJsonSerializerOptions SystemTextJsonSerializerOptions { get; } |
||||
|
|
||||
|
public AbpDaprClientFactory( |
||||
|
IOptions<AbpDaprOptions> options, |
||||
|
IOptions<AbpSystemTextJsonSerializerOptions> systemTextJsonSerializerOptions) |
||||
|
{ |
||||
|
Options = options.Value; |
||||
|
SystemTextJsonSerializerOptions = systemTextJsonSerializerOptions.Value; |
||||
|
} |
||||
|
|
||||
|
public virtual async Task<DaprClient> CreateAsync() |
||||
|
{ |
||||
|
var builder = new DaprClientBuilder() |
||||
|
.UseJsonSerializationOptions(await CreateJsonSerializerOptions()); |
||||
|
|
||||
|
if (!Options.HttpEndpoint.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
builder.UseHttpEndpoint(Options.HttpEndpoint); |
||||
|
} |
||||
|
|
||||
|
if (!Options.GrpcEndpoint.IsNullOrWhiteSpace()) |
||||
|
{ |
||||
|
builder.UseGrpcEndpoint(Options.GrpcEndpoint); |
||||
|
} |
||||
|
|
||||
|
return builder.Build(); |
||||
|
} |
||||
|
|
||||
|
private readonly static ConcurrentDictionary<string, JsonSerializerOptions> JsonSerializerOptionsCache = new ConcurrentDictionary<string, JsonSerializerOptions>(); |
||||
|
|
||||
|
protected virtual Task<JsonSerializerOptions> CreateJsonSerializerOptions() |
||||
|
{ |
||||
|
return Task.FromResult(JsonSerializerOptionsCache.GetOrAdd(nameof(AbpDaprClientFactory), |
||||
|
_ => new JsonSerializerOptions(SystemTextJsonSerializerOptions.JsonSerializerOptions))); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,15 @@ |
|||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Json; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.Dapr; |
||||
|
|
||||
|
[DependsOn(typeof(AbpJsonModule))] |
||||
|
public class AbpDaprModule : AbpModule |
||||
|
{ |
||||
|
public override void ConfigureServices(ServiceConfigurationContext context) |
||||
|
{ |
||||
|
var configuration = context.Services.GetConfiguration(); |
||||
|
Configure<AbpDaprOptions>(configuration.GetSection("Dapr")); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,10 @@ |
|||||
|
namespace Volo.Abp.Dapr; |
||||
|
|
||||
|
public class AbpDaprOptions |
||||
|
{ |
||||
|
public string AppId { get; set; } |
||||
|
|
||||
|
public string HttpEndpoint { get; set; } |
||||
|
|
||||
|
public string GrpcEndpoint { get; set; } |
||||
|
} |
||||
@ -0,0 +1,16 @@ |
|||||
|
namespace Volo.Abp.Dapr; |
||||
|
|
||||
|
public interface IDaprSerializer |
||||
|
{ |
||||
|
byte[] Serialize(object obj); |
||||
|
|
||||
|
object Deserialize(byte[] value, Type type); |
||||
|
|
||||
|
T Deserialize<T>(byte[] value); |
||||
|
|
||||
|
string SerializeToString(object obj); |
||||
|
|
||||
|
object Deserialize(string value, Type type); |
||||
|
|
||||
|
T Deserialize<T>(string value); |
||||
|
} |
||||
@ -0,0 +1,45 @@ |
|||||
|
using System.Text; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.Json; |
||||
|
|
||||
|
namespace Volo.Abp.Dapr; |
||||
|
|
||||
|
public class Utf8JsonDaprSerializer : IDaprSerializer, ITransientDependency |
||||
|
{ |
||||
|
private readonly IJsonSerializer _jsonSerializer; |
||||
|
|
||||
|
public Utf8JsonDaprSerializer(IJsonSerializer jsonSerializer) |
||||
|
{ |
||||
|
_jsonSerializer = jsonSerializer; |
||||
|
} |
||||
|
|
||||
|
public byte[] Serialize(object obj) |
||||
|
{ |
||||
|
return Encoding.UTF8.GetBytes(_jsonSerializer.Serialize(obj)); |
||||
|
} |
||||
|
|
||||
|
public object Deserialize(byte[] value, Type type) |
||||
|
{ |
||||
|
return _jsonSerializer.Deserialize(type, Encoding.UTF8.GetString(value)); |
||||
|
} |
||||
|
|
||||
|
public T Deserialize<T>(byte[] value) |
||||
|
{ |
||||
|
return _jsonSerializer.Deserialize<T>(Encoding.UTF8.GetString(value)); |
||||
|
} |
||||
|
|
||||
|
public string SerializeToString(object obj) |
||||
|
{ |
||||
|
return _jsonSerializer.Serialize(obj); |
||||
|
} |
||||
|
|
||||
|
public object Deserialize(string value, Type type) |
||||
|
{ |
||||
|
return _jsonSerializer.Deserialize(type, value); |
||||
|
} |
||||
|
|
||||
|
public T Deserialize<T>(string value) |
||||
|
{ |
||||
|
return _jsonSerializer.Deserialize<T>(value); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,18 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>net6.0</TargetFramework> |
||||
|
<ImplicitUsings>enable</ImplicitUsings> |
||||
|
<Nullable>enable</Nullable> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.Dapr\Volo.Abp.Dapr.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.DistributedLocking.Abstractions\Volo.Abp.DistributedLocking.Abstractions.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,13 @@ |
|||||
|
namespace Volo.Abp.DistributedLocking.Dapr; |
||||
|
|
||||
|
public class AbpDistributedLockDaprOptions |
||||
|
{ |
||||
|
public string StoreName { get; set; } |
||||
|
|
||||
|
public TimeSpan DefaultTimeout { get; set; } |
||||
|
|
||||
|
public AbpDistributedLockDaprOptions() |
||||
|
{ |
||||
|
DefaultTimeout = TimeSpan.FromSeconds(30); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,11 @@ |
|||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.DistributedLocking.Dapr; |
||||
|
|
||||
|
[DependsOn( |
||||
|
typeof(AbpDistributedLockingAbstractionsModule), |
||||
|
typeof(AbpDaprModule))] |
||||
|
public class AbpDistributedLockingDaprModule : AbpModule |
||||
|
{ |
||||
|
} |
||||
@ -0,0 +1,50 @@ |
|||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.DistributedLocking.Dapr; |
||||
|
|
||||
|
[Dependency(ReplaceServices = true)] |
||||
|
public class DaprAbpDistributedLock : IAbpDistributedLock, ITransientDependency |
||||
|
{ |
||||
|
protected AbpDaprClientFactory DaprClientFactory { get; } |
||||
|
protected AbpDistributedLockDaprOptions DistributedLockDaprOptions { get; } |
||||
|
protected AbpDaprOptions DaprOptions { get; } |
||||
|
|
||||
|
public DaprAbpDistributedLock( |
||||
|
AbpDaprClientFactory daprClientFactory, |
||||
|
IOptions<AbpDistributedLockDaprOptions> distributedLockDaprOptions, |
||||
|
IOptions<AbpDaprOptions> daprOptions) |
||||
|
{ |
||||
|
DaprClientFactory = daprClientFactory; |
||||
|
DaprOptions = daprOptions.Value; |
||||
|
DistributedLockDaprOptions = distributedLockDaprOptions.Value; |
||||
|
} |
||||
|
|
||||
|
public async Task<IAbpDistributedLockHandle> TryAcquireAsync( |
||||
|
string name, |
||||
|
TimeSpan timeout = default, |
||||
|
CancellationToken cancellationToken = default) |
||||
|
{ |
||||
|
if (timeout == default) |
||||
|
{ |
||||
|
timeout = DistributedLockDaprOptions.DefaultTimeout; |
||||
|
} |
||||
|
|
||||
|
var daprClient = await DaprClientFactory.CreateAsync(); |
||||
|
|
||||
|
var lockResponse = await daprClient.Lock( |
||||
|
DistributedLockDaprOptions.StoreName, |
||||
|
name, |
||||
|
DaprOptions.AppId, |
||||
|
(int)timeout.TotalSeconds, |
||||
|
cancellationToken); |
||||
|
|
||||
|
if (lockResponse == null || !lockResponse.Success) |
||||
|
{ |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
return new DaprAbpDistributedLockHandle(lockResponse); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,18 @@ |
|||||
|
using Dapr.Client; |
||||
|
|
||||
|
namespace Volo.Abp.DistributedLocking.Dapr; |
||||
|
|
||||
|
public class DaprAbpDistributedLockHandle : IAbpDistributedLockHandle |
||||
|
{ |
||||
|
protected TryLockResponse LockResponse { get; } |
||||
|
|
||||
|
public DaprAbpDistributedLockHandle(TryLockResponse lockResponse) |
||||
|
{ |
||||
|
LockResponse = lockResponse; |
||||
|
} |
||||
|
|
||||
|
public async ValueTask DisposeAsync() |
||||
|
{ |
||||
|
await LockResponse.DisposeAsync(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,18 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>net6.0</TargetFramework> |
||||
|
<ImplicitUsings>enable</ImplicitUsings> |
||||
|
<Nullable>enable</Nullable> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.EventBus\Volo.Abp.EventBus.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.Dapr\Volo.Abp.Dapr.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,11 @@ |
|||||
|
namespace Volo.Abp.EventBus.Dapr; |
||||
|
|
||||
|
public class AbpDaprEventBusOptions |
||||
|
{ |
||||
|
public string PubSubName { get; set; } |
||||
|
|
||||
|
public AbpDaprEventBusOptions() |
||||
|
{ |
||||
|
PubSubName = "pubsub"; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,20 @@ |
|||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus.Dapr; |
||||
|
|
||||
|
[DependsOn( |
||||
|
typeof(AbpEventBusModule), |
||||
|
typeof(AbpDaprModule) |
||||
|
)] |
||||
|
public class AbpEventBusDaprModule : AbpModule |
||||
|
{ |
||||
|
public override void OnApplicationInitialization(ApplicationInitializationContext context) |
||||
|
{ |
||||
|
context |
||||
|
.ServiceProvider |
||||
|
.GetRequiredService<DaprDistributedEventBus>() |
||||
|
.Initialize(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,221 @@ |
|||||
|
using System.Collections.Concurrent; |
||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
using Volo.Abp.EventBus.Distributed; |
||||
|
using Volo.Abp.Guids; |
||||
|
using Volo.Abp.MultiTenancy; |
||||
|
using Volo.Abp.Threading; |
||||
|
using Volo.Abp.Timing; |
||||
|
using Volo.Abp.Uow; |
||||
|
|
||||
|
namespace Volo.Abp.EventBus.Dapr; |
||||
|
|
||||
|
[Dependency(ReplaceServices = true)] |
||||
|
[ExposeServices(typeof(IDistributedEventBus), typeof(DaprDistributedEventBus))] |
||||
|
public class DaprDistributedEventBus : DistributedEventBusBase, ISingletonDependency |
||||
|
{ |
||||
|
protected IDaprSerializer Serializer { get; } |
||||
|
protected AbpDaprEventBusOptions DaprEventBusOptions { get; } |
||||
|
protected AbpDaprClientFactory DaprClientFactory { get; } |
||||
|
|
||||
|
protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; } |
||||
|
protected ConcurrentDictionary<string, Type> EventTypes { get; } |
||||
|
|
||||
|
public DaprDistributedEventBus( |
||||
|
IServiceScopeFactory serviceScopeFactory, |
||||
|
ICurrentTenant currentTenant, |
||||
|
IUnitOfWorkManager unitOfWorkManager, |
||||
|
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions, |
||||
|
IGuidGenerator guidGenerator, |
||||
|
IClock clock, |
||||
|
IEventHandlerInvoker eventHandlerInvoker, |
||||
|
IDaprSerializer serializer, |
||||
|
IOptions<AbpDaprEventBusOptions> daprEventBusOptions, |
||||
|
AbpDaprClientFactory daprClientFactory) |
||||
|
: base(serviceScopeFactory, currentTenant, unitOfWorkManager, abpDistributedEventBusOptions, guidGenerator, clock, eventHandlerInvoker) |
||||
|
{ |
||||
|
Serializer = serializer; |
||||
|
DaprEventBusOptions = daprEventBusOptions.Value; |
||||
|
DaprClientFactory = daprClientFactory; |
||||
|
|
||||
|
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>(); |
||||
|
EventTypes = new ConcurrentDictionary<string, Type>(); |
||||
|
} |
||||
|
|
||||
|
public void Initialize() |
||||
|
{ |
||||
|
SubscribeHandlers(AbpDistributedEventBusOptions.Handlers); |
||||
|
} |
||||
|
|
||||
|
public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) |
||||
|
{ |
||||
|
var handlerFactories = GetOrCreateHandlerFactories(eventType); |
||||
|
|
||||
|
if (factory.IsInFactories(handlerFactories)) |
||||
|
{ |
||||
|
return NullDisposable.Instance; |
||||
|
} |
||||
|
|
||||
|
handlerFactories.Add(factory); |
||||
|
|
||||
|
return new EventHandlerFactoryUnregistrar(this, eventType, factory); |
||||
|
} |
||||
|
|
||||
|
public override void Unsubscribe<TEvent>(Func<TEvent, Task> action) |
||||
|
{ |
||||
|
Check.NotNull(action, nameof(action)); |
||||
|
|
||||
|
GetOrCreateHandlerFactories(typeof(TEvent)) |
||||
|
.Locking(factories => |
||||
|
{ |
||||
|
factories.RemoveAll( |
||||
|
factory => |
||||
|
{ |
||||
|
var singleInstanceFactory = factory as SingleInstanceHandlerFactory; |
||||
|
if (singleInstanceFactory == null) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
var actionHandler = singleInstanceFactory.HandlerInstance as ActionEventHandler<TEvent>; |
||||
|
if (actionHandler == null) |
||||
|
{ |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
return actionHandler.Action == action; |
||||
|
}); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
public override void Unsubscribe(Type eventType, IEventHandler handler) |
||||
|
{ |
||||
|
GetOrCreateHandlerFactories(eventType) |
||||
|
.Locking(factories => |
||||
|
{ |
||||
|
factories.RemoveAll( |
||||
|
factory => |
||||
|
factory is SingleInstanceHandlerFactory && |
||||
|
(factory as SingleInstanceHandlerFactory).HandlerInstance == handler |
||||
|
); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
public override void Unsubscribe(Type eventType, IEventHandlerFactory factory) |
||||
|
{ |
||||
|
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Remove(factory)); |
||||
|
} |
||||
|
|
||||
|
public override void UnsubscribeAll(Type eventType) |
||||
|
{ |
||||
|
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear()); |
||||
|
} |
||||
|
|
||||
|
protected async override Task PublishToEventBusAsync(Type eventType, object eventData) |
||||
|
{ |
||||
|
await PublishToDaprAsync(eventType, eventData); |
||||
|
} |
||||
|
|
||||
|
protected override void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord) |
||||
|
{ |
||||
|
unitOfWork.AddOrReplaceDistributedEvent(eventRecord); |
||||
|
} |
||||
|
|
||||
|
protected override IEnumerable<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType) |
||||
|
{ |
||||
|
var handlerFactoryList = new List<EventTypeWithEventHandlerFactories>(); |
||||
|
|
||||
|
foreach (var handlerFactory in HandlerFactories.Where(hf => ShouldTriggerEventForHandler(eventType, hf.Key))) |
||||
|
{ |
||||
|
handlerFactoryList.Add(new EventTypeWithEventHandlerFactories(handlerFactory.Key, handlerFactory.Value)); |
||||
|
} |
||||
|
|
||||
|
return handlerFactoryList.ToArray(); |
||||
|
} |
||||
|
|
||||
|
public async override Task PublishFromOutboxAsync(OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig) |
||||
|
{ |
||||
|
await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName))); |
||||
|
} |
||||
|
|
||||
|
public async override Task PublishManyFromOutboxAsync(IEnumerable<OutgoingEventInfo> outgoingEvents, OutboxConfig outboxConfig) |
||||
|
{ |
||||
|
var outgoingEventArray = outgoingEvents.ToArray(); |
||||
|
|
||||
|
foreach (var outgoingEvent in outgoingEventArray) |
||||
|
{ |
||||
|
await PublishToDaprAsync(outgoingEvent.EventName, Serializer.Deserialize(outgoingEvent.EventData, GetEventType(outgoingEvent.EventName))); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public async override Task ProcessFromInboxAsync(IncomingEventInfo incomingEvent, InboxConfig inboxConfig) |
||||
|
{ |
||||
|
var eventType = EventTypes.GetOrDefault(incomingEvent.EventName); |
||||
|
if (eventType == null) |
||||
|
{ |
||||
|
return; |
||||
|
} |
||||
|
|
||||
|
var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); |
||||
|
var exceptions = new List<Exception>(); |
||||
|
await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); |
||||
|
if (exceptions.Any()) |
||||
|
{ |
||||
|
ThrowOriginalExceptions(eventType, exceptions); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
protected override byte[] Serialize(object eventData) |
||||
|
{ |
||||
|
return Serializer.Serialize(eventData); |
||||
|
} |
||||
|
|
||||
|
private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType) |
||||
|
{ |
||||
|
return HandlerFactories.GetOrAdd( |
||||
|
eventType, |
||||
|
type => |
||||
|
{ |
||||
|
var eventName = EventNameAttribute.GetNameOrDefault(type); |
||||
|
EventTypes[eventName] = type; |
||||
|
return new List<IEventHandlerFactory>(); |
||||
|
} |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
public Type GetEventType(string eventName) |
||||
|
{ |
||||
|
return EventTypes.GetOrDefault(eventName); |
||||
|
} |
||||
|
|
||||
|
protected virtual async Task PublishToDaprAsync(Type eventType, object eventData) |
||||
|
{ |
||||
|
await PublishToDaprAsync(EventNameAttribute.GetNameOrDefault(eventType), eventData); |
||||
|
} |
||||
|
|
||||
|
protected virtual async Task PublishToDaprAsync(string eventName, object eventData) |
||||
|
{ |
||||
|
var client = await DaprClientFactory.CreateAsync(); |
||||
|
await client.PublishEventAsync(pubsubName: DaprEventBusOptions.PubSubName, topicName: eventName, data: eventData); |
||||
|
} |
||||
|
|
||||
|
private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) |
||||
|
{ |
||||
|
//Should trigger same type
|
||||
|
if (handlerEventType == targetEventType) |
||||
|
{ |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
//TODO: Support inheritance? But it does not support on subscription to RabbitMq!
|
||||
|
//Should trigger for inherited types
|
||||
|
if (handlerEventType.IsAssignableFrom(targetEventType)) |
||||
|
{ |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,3 @@ |
|||||
|
<Weavers xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:noNamespaceSchemaLocation="FodyWeavers.xsd"> |
||||
|
<ConfigureAwait ContinueOnCapturedContext="false" /> |
||||
|
</Weavers> |
||||
@ -0,0 +1,30 @@ |
|||||
|
<?xml version="1.0" encoding="utf-8"?> |
||||
|
<xs:schema xmlns:xs="http://www.w3.org/2001/XMLSchema"> |
||||
|
<!-- This file was generated by Fody. Manual changes to this file will be lost when your project is rebuilt. --> |
||||
|
<xs:element name="Weavers"> |
||||
|
<xs:complexType> |
||||
|
<xs:all> |
||||
|
<xs:element name="ConfigureAwait" minOccurs="0" maxOccurs="1"> |
||||
|
<xs:complexType> |
||||
|
<xs:attribute name="ContinueOnCapturedContext" type="xs:boolean" /> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:all> |
||||
|
<xs:attribute name="VerifyAssembly" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'true' to run assembly verification (PEVerify) on the target assembly after all weavers have been executed.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="VerifyIgnoreCodes" type="xs:string"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>A comma-separated list of error codes that can be safely ignored in assembly verification.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
<xs:attribute name="GenerateXsd" type="xs:boolean"> |
||||
|
<xs:annotation> |
||||
|
<xs:documentation>'false' to turn off automatic generation of the XML Schema file.</xs:documentation> |
||||
|
</xs:annotation> |
||||
|
</xs:attribute> |
||||
|
</xs:complexType> |
||||
|
</xs:element> |
||||
|
</xs:schema> |
||||
@ -0,0 +1,18 @@ |
|||||
|
<Project Sdk="Microsoft.NET.Sdk"> |
||||
|
|
||||
|
<Import Project="..\..\..\configureawait.props" /> |
||||
|
<Import Project="..\..\..\common.props" /> |
||||
|
|
||||
|
<PropertyGroup> |
||||
|
<TargetFramework>net6.0</TargetFramework> |
||||
|
<ImplicitUsings>enable</ImplicitUsings> |
||||
|
<Nullable>enable</Nullable> |
||||
|
<RootNamespace /> |
||||
|
</PropertyGroup> |
||||
|
|
||||
|
<ItemGroup> |
||||
|
<ProjectReference Include="..\Volo.Abp.Http.Client\Volo.Abp.Http.Client.csproj" /> |
||||
|
<ProjectReference Include="..\Volo.Abp.Dapr\Volo.Abp.Dapr.csproj" /> |
||||
|
</ItemGroup> |
||||
|
|
||||
|
</Project> |
||||
@ -0,0 +1,23 @@ |
|||||
|
using Microsoft.Extensions.DependencyInjection; |
||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.Modularity; |
||||
|
|
||||
|
namespace Volo.Abp.Http.Client.Dapr; |
||||
|
|
||||
|
[DependsOn( |
||||
|
typeof(AbpHttpClientModule), |
||||
|
typeof(AbpDaprModule) |
||||
|
)] |
||||
|
public class AbpHttpClientDaprModule : AbpModule |
||||
|
{ |
||||
|
public override void PreConfigureServices(ServiceConfigurationContext context) |
||||
|
{ |
||||
|
PreConfigure<AbpHttpClientBuilderOptions>(options => |
||||
|
{ |
||||
|
options.ProxyClientBuildActions.Add((_, clientBuilder) => |
||||
|
{ |
||||
|
clientBuilder.AddHttpMessageHandler<AbpInvocationHandler>(); |
||||
|
}); |
||||
|
}); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,14 @@ |
|||||
|
using Dapr.Client; |
||||
|
using Microsoft.Extensions.Options; |
||||
|
using Volo.Abp.Dapr; |
||||
|
using Volo.Abp.DependencyInjection; |
||||
|
|
||||
|
namespace Volo.Abp.Http.Client.Dapr; |
||||
|
|
||||
|
public class AbpInvocationHandler : InvocationHandler, ITransientDependency |
||||
|
{ |
||||
|
public AbpInvocationHandler(IOptions<AbpDaprOptions> daprOptions) |
||||
|
{ |
||||
|
DaprEndpoint = daprOptions.Value.HttpEndpoint; |
||||
|
} |
||||
|
} |
||||
Loading…
Reference in new issue