@ -0,0 +1,405 @@ |
|||
# Introducing the Angular Service Proxy Generation |
|||
|
|||
Angular Service Proxy System **generates TypeScript services and models** to consume your backend HTTP APIs developed using the ABP Framework. So, you **don't manually create** models for your server side DTOs and perform raw HTTP calls to the server. |
|||
|
|||
ABP Framework has introduced the **new** Angular Service Proxy Generation system with the version 3.1. While this feature was available since the [v2.3](https://blog.abp.io/abp/ABP-Framework-v2_3_0-Has-Been-Released), it was not well covering some scenarios, like inheritance and generic types and had some known problems. **With the v3.1, we've re-written** it using the [Angular Schematics](https://angular.io/guide/schematics) system. Now, it is much more stable and feature rich. |
|||
|
|||
This post introduces the service proxy generation system and highlights some important features. |
|||
|
|||
## Installation |
|||
|
|||
### ABP CLI |
|||
|
|||
You need to have the [ABP CLI](https://docs.abp.io/en/abp/latest/CLI) to use the system. So, install it if you haven't installed before: |
|||
|
|||
````bash |
|||
dotnet tool install -g Volo.Abp.Cli |
|||
```` |
|||
|
|||
If you already have installed it before, you can update to the latest version: |
|||
|
|||
````shell |
|||
dotnet tool update -g Volo.Abp.Cli |
|||
```` |
|||
|
|||
### Project Configuration |
|||
|
|||
> If you've created your project with version 3.1 or later, you can skip this part since it will be already installed in your solution. |
|||
|
|||
For a solution that was created before v3.1, follow the steps below to configure the angular application: |
|||
|
|||
* Add `@abp/ng.schematics` package to the `devDependencies` of the Angular project. Run the following command in the root folder of the angular application: |
|||
|
|||
````bash |
|||
npm install @abp/ng.schematics --save-dev |
|||
```` |
|||
|
|||
- Add `rootNamespace` entry into the `apis/default` section in the `/src/environments/environment.ts`, as shown below: |
|||
|
|||
```json |
|||
apis: { |
|||
default: { |
|||
... |
|||
rootNamespace: 'Acme.BookStore' |
|||
}, |
|||
} |
|||
``` |
|||
|
|||
`Acme.BookStore` should be replaced by the root namespace of your .NET project. This ensures to not create unnecessary nested folders while creating the service proxy code. This value is `AngularProxyDemo` for the example solution explained below. |
|||
|
|||
## Basic Usage |
|||
|
|||
### Project Creation |
|||
|
|||
> If you already have a solution, you can skip this section. |
|||
|
|||
You need to [create](https://abp.io/get-started) your solution with the Angular UI. You can use the [ABP CLI](https://docs.abp.io/en/abp/latest/CLI) to create a new solution: |
|||
|
|||
````bash |
|||
abp new AngularProxyDemo -u angular |
|||
```` |
|||
|
|||
#### Run the Application |
|||
|
|||
The backend application must be up and running to be able to use the service proxy code generation system. |
|||
|
|||
> See the [getting started](https://docs.abp.io/en/abp/latest/Getting-Started?UI=NG&DB=EF&Tiered=No) guide if you don't know details of creating and running the solution. |
|||
|
|||
### Backend |
|||
|
|||
Assume that we have an `IBookAppService` interface: |
|||
|
|||
````csharp |
|||
using System.Collections.Generic; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.Application.Services; |
|||
|
|||
namespace AngularProxyDemo.Books |
|||
{ |
|||
public interface IBookAppService : IApplicationService |
|||
{ |
|||
public Task<List<BookDto>> GetListAsync(); |
|||
} |
|||
} |
|||
```` |
|||
|
|||
That uses a `BookDto` defined as shown: |
|||
|
|||
```csharp |
|||
using System; |
|||
using Volo.Abp.Application.Dtos; |
|||
|
|||
namespace AngularProxyDemo.Books |
|||
{ |
|||
public class BookDto : EntityDto<Guid> |
|||
{ |
|||
public string Name { get; set; } |
|||
|
|||
public DateTime PublishDate { get; set; } |
|||
} |
|||
} |
|||
``` |
|||
|
|||
And implemented as the following: |
|||
|
|||
```csharp |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.Application.Services; |
|||
|
|||
namespace AngularProxyDemo.Books |
|||
{ |
|||
public class BookAppService : ApplicationService, IBookAppService |
|||
{ |
|||
public async Task<List<BookDto>> GetListAsync() |
|||
{ |
|||
//TODO: get books from a database... |
|||
} |
|||
} |
|||
} |
|||
``` |
|||
|
|||
It simply returns a list of books. You probably want to get the books from a database, but it doesn't matter for this article. |
|||
|
|||
### HTTP API |
|||
|
|||
Thanks to the [auto API controllers](https://docs.abp.io/en/abp/latest/API/Auto-API-Controllers) system of the ABP Framework, we don't have to develop API controllers manually. Just **run the backend (*HttpApi.Host*) application** that shows the [Swagger UI](https://swagger.io/tools/swagger-ui/) by default. You will see the **GET** API for the books: |
|||
|
|||
 |
|||
|
|||
### Service Proxy Generation |
|||
|
|||
Open a **command line** in the **root folder of the Angular application** and execute the following command: |
|||
|
|||
````bash |
|||
abp generate-proxy |
|||
```` |
|||
|
|||
It should produce an output like the following: |
|||
|
|||
````bash |
|||
CREATE src/app/shared/models/books/index.ts (142 bytes) |
|||
CREATE src/app/shared/services/books/book.service.ts (437 bytes) |
|||
... |
|||
```` |
|||
|
|||
> `generate-proxy` command can take some some optional parameters for advanced scenarios (like [modular development](https://docs.abp.io/en/abp/latest/Module-Development-Basics)). You can take a look at the [documentation](https://docs.abp.io/en/abp/latest/UI/Angular/Service-Proxies). |
|||
|
|||
It basically creates two files; |
|||
|
|||
src/app/shared/services/books/**book.service.ts**: This is the service that can be injected and used to get the list of books; |
|||
|
|||
````typescript |
|||
import { RestService } from '@abp/ng.core'; |
|||
import { Injectable } from '@angular/core'; |
|||
import type { BookDto } from '../../models/book'; |
|||
|
|||
@Injectable({ |
|||
providedIn: 'root', |
|||
}) |
|||
export class BookService { |
|||
apiName = 'Default'; |
|||
|
|||
getList = () => |
|||
this.restService.request<any, BookDto[]>({ |
|||
method: 'GET', |
|||
url: `/api/app/book`, |
|||
}, |
|||
{ apiName: this.apiName }); |
|||
|
|||
constructor(private restService: RestService) {} |
|||
} |
|||
|
|||
```` |
|||
|
|||
src/app/shared/models/books/**index.ts**: This file contains the modal classes corresponding to the DTOs defined in the server side; |
|||
|
|||
````typescript |
|||
import type { EntityDto } from '@abp/ng.core'; |
|||
|
|||
export interface BookDto extends EntityDto<string> { |
|||
name: string; |
|||
publishDate: string; |
|||
} |
|||
```` |
|||
|
|||
You can now inject the `BookService` into any Angular component and use the `getList()` method to get the list of books. |
|||
|
|||
### About the Generated Code |
|||
|
|||
The generated code is; |
|||
|
|||
* **Simple**: It is almost identical to the code if you've written it yourself. |
|||
* **Splitted**: Instead of a single, large file; |
|||
* It creates a separate `.ts` file for every backend **service**. **Model** (DTO) classes are also grouped per service. |
|||
* It understands the [modularity](https://docs.abp.io/en/abp/latest/Module-Development-Basics), so creates the services for your own **module** (or the module you've specified). |
|||
* **Object oriented**; |
|||
* Supports **inheritance** of server side DTOs and generates the code respecting to the inheritance structure. |
|||
* Supports **generic types**. |
|||
* Supports **re-using type definitions** across services and doesn't generate the same DTO multiple times. |
|||
* **Well-aligned to the backend**; |
|||
* Service **method signatures** match exactly with the services on the backend services. This is achieved by a special endpoint exposed by the ABP Framework that well defines the backend contracts. |
|||
* **Namespaces** are exactly matches to the backend services and DTOs. |
|||
* **Well-aligned with the ABP Framework**; |
|||
* Recognizes the **standard ABP Framework DTO types** (like `EntityDto`, `ListResultDto`... etc) and doesn't repeat these classes in the application code, but uses from the `@abp/ng.core` package. |
|||
* Uses the `RestService` defined by the `@abp/ng.core` package which simplifies the generated code, keeps it short and re-uses all the logics implemented by the `RestService` (including error handling, authorization token injection, using multiple server endpoints... etc). |
|||
|
|||
These are the main motivations behind the decision of creating a service proxy generation system, instead of using a pre-built tool like [NSWAG](https://github.com/RicoSuter/NSwag). |
|||
|
|||
## Other Examples |
|||
|
|||
Let me show you a few more examples. |
|||
|
|||
### Updating an Entity |
|||
|
|||
Assume that you added a new method to the server side application service, to update a book: |
|||
|
|||
```csharp |
|||
public Task<BookDto> UpdateAsync(Guid id, BookUpdateDto input); |
|||
``` |
|||
|
|||
`BookUpdateDto` is a simple class defined shown below: |
|||
|
|||
```csharp |
|||
using System; |
|||
|
|||
namespace AngularProxyDemo.Books |
|||
{ |
|||
public class BookUpdateDto |
|||
{ |
|||
public string Name { get; set; } |
|||
|
|||
public DateTime PublishDate { get; set; } |
|||
} |
|||
} |
|||
``` |
|||
|
|||
Let's re-run the `generate-proxy` command to see the result: |
|||
|
|||
```` |
|||
abp generate-proxy |
|||
```` |
|||
|
|||
The output of this command will be like the following: |
|||
|
|||
````bash |
|||
UPDATE src/app/shared/services/books/book.service.ts (660 bytes) |
|||
UPDATE src/app/shared/models/books/index.ts (217 bytes) |
|||
```` |
|||
|
|||
It tells us two files have been updated. Let's see the changes; |
|||
|
|||
**book.service.ts** |
|||
|
|||
````typescript |
|||
import { RestService } from '@abp/ng.core'; |
|||
import { Injectable } from '@angular/core'; |
|||
import type { BookDto, BookUpdateDto } from '../../models/books'; |
|||
|
|||
@Injectable({ |
|||
providedIn: 'root', |
|||
}) |
|||
export class BookService { |
|||
apiName = 'Default'; |
|||
|
|||
getList = () => |
|||
this.restService.request<any, BookDto[]>({ |
|||
method: 'GET', |
|||
url: `/api/app/book`, |
|||
}, |
|||
{ apiName: this.apiName }); |
|||
|
|||
update = (id: string, input: BookUpdateDto) => |
|||
this.restService.request<any, BookDto>({ |
|||
method: 'PUT', |
|||
url: `/api/app/book/${id}`, |
|||
body: input, |
|||
}, |
|||
{ apiName: this.apiName }); |
|||
|
|||
constructor(private restService: RestService) {} |
|||
} |
|||
```` |
|||
|
|||
`update` function has been added to the `BookService` that gets an `id` and a `BookUpdateDto` as the parameters. |
|||
|
|||
**index.ts** |
|||
|
|||
````typescript |
|||
import type { EntityDto } from '@abp/ng.core'; |
|||
|
|||
export interface BookDto extends EntityDto<string> { |
|||
name: string; |
|||
publishDate: string; |
|||
} |
|||
|
|||
export interface BookUpdateDto { |
|||
name: string; |
|||
publishDate: string; |
|||
} |
|||
```` |
|||
|
|||
Added a new DTO class: `BookUpdateDto`. |
|||
|
|||
### Advanced Example |
|||
|
|||
In this example, I want to show a DTO structure using inheritance, generics, arrays and dictionaries. |
|||
|
|||
I've created an `IOrderAppService` as shown below: |
|||
|
|||
````csharp |
|||
using System.Threading.Tasks; |
|||
using Volo.Abp.Application.Services; |
|||
|
|||
namespace AngularProxyDemo.Orders |
|||
{ |
|||
public interface IOrderAppService : IApplicationService |
|||
{ |
|||
public Task CreateAsync(OrderCreateDto input); |
|||
} |
|||
} |
|||
```` |
|||
|
|||
`OrderCreateDto` and the related DTOs are as the followings; |
|||
|
|||
````csharp |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using Volo.Abp.Data; |
|||
|
|||
namespace AngularProxyDemo.Orders |
|||
{ |
|||
public class OrderCreateDto : IHasExtraProperties |
|||
{ |
|||
public Guid CustomerId { get; set; } |
|||
|
|||
public DateTime CreationTime { get; set; } |
|||
|
|||
//ARRAY of DTOs |
|||
public OrderDetailDto[] Details { get; set; } |
|||
|
|||
//DICTIONARY |
|||
public Dictionary<string, object> ExtraProperties { get; set; } |
|||
} |
|||
|
|||
public class OrderDetailDto : GenericDetailDto<int> //INHERIT from GENERIC |
|||
{ |
|||
public string Note { get; set; } |
|||
} |
|||
|
|||
//GENERIC class |
|||
public abstract class GenericDetailDto<TCount> |
|||
{ |
|||
public Guid ProductId { get; set; } |
|||
|
|||
public TCount Count { get; set; } |
|||
} |
|||
} |
|||
```` |
|||
|
|||
When I run the `abp generate-proxy` command again, I see two new files have been created. |
|||
|
|||
src/app/shared/services/orders/**order.service.ts** |
|||
|
|||
````typescript |
|||
import { RestService } from '@abp/ng.core'; |
|||
import { Injectable } from '@angular/core'; |
|||
import type { OrderCreateDto } from '../../models/orders'; |
|||
|
|||
@Injectable({ |
|||
providedIn: 'root', |
|||
}) |
|||
export class OrderService { |
|||
apiName = 'Default'; |
|||
|
|||
create = (input: OrderCreateDto) => |
|||
this.restService.request<any, void>({ |
|||
method: 'POST', |
|||
url: `/api/app/order`, |
|||
body: input, |
|||
}, |
|||
{ apiName: this.apiName }); |
|||
|
|||
constructor(private restService: RestService) {} |
|||
} |
|||
```` |
|||
|
|||
src/app/shared/models/orders/**index.ts** |
|||
|
|||
````typescript |
|||
import type { GenericDetailDto } from '../../models/orders'; |
|||
|
|||
export interface OrderCreateDto { |
|||
customerId: string; |
|||
creationTime: string; |
|||
details: OrderDetailDto[]; |
|||
extraProperties: string | object; |
|||
} |
|||
|
|||
export interface OrderDetailDto extends GenericDetailDto<number> { |
|||
note: string; |
|||
} |
|||
```` |
|||
|
|||
NOTE: 3.1.0-rc2 was generating the code above, which is wrong. It will be fixed in the next RC versions. |
|||
|
After Width: | Height: | Size: 42 KiB |
|
After Width: | Height: | Size: 9.3 KiB |
|
Before Width: | Height: | Size: 123 KiB After Width: | Height: | Size: 123 KiB |
|
Before Width: | Height: | Size: 11 KiB After Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 8.4 KiB After Width: | Height: | Size: 8.4 KiB |
|
Before Width: | Height: | Size: 44 KiB After Width: | Height: | Size: 44 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 10 KiB |
|
Before Width: | Height: | Size: 106 KiB After Width: | Height: | Size: 106 KiB |
|
Before Width: | Height: | Size: 7.9 KiB After Width: | Height: | Size: 7.9 KiB |
|
Before Width: | Height: | Size: 6.2 KiB After Width: | Height: | Size: 6.2 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 10 KiB |
|
Before Width: | Height: | Size: 11 KiB After Width: | Height: | Size: 11 KiB |
|
Before Width: | Height: | Size: 8.5 KiB After Width: | Height: | Size: 8.5 KiB |
|
Before Width: | Height: | Size: 68 KiB After Width: | Height: | Size: 68 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 10 KiB |
|
Before Width: | Height: | Size: 4.5 KiB After Width: | Height: | Size: 4.5 KiB |
@ -0,0 +1,13 @@ |
|||
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
|||
using System.Collections.Generic; |
|||
|
|||
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.CropperJs |
|||
{ |
|||
public class CropperJsScriptContributor : BundleContributor |
|||
{ |
|||
public override void ConfigureBundle(BundleConfigurationContext context) |
|||
{ |
|||
context.Files.AddIfNotContains("/libs/cropperjs/js/cropper.min.js"); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,13 @@ |
|||
using System.Collections.Generic; |
|||
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
|||
|
|||
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.CropperJs |
|||
{ |
|||
public class CropperJsStyleContributor : BundleContributor |
|||
{ |
|||
public override void ConfigureBundle(BundleConfigurationContext context) |
|||
{ |
|||
context.Files.AddIfNotContains("/libs/cropperjs/css/cropper.min.css"); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,16 @@ |
|||
using System.Collections.Generic; |
|||
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
|||
using Volo.Abp.AspNetCore.Mvc.UI.Packages.JQuery; |
|||
using Volo.Abp.Modularity; |
|||
|
|||
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.StarRatingSvg |
|||
{ |
|||
[DependsOn(typeof(JQueryScriptContributor))] |
|||
public class StarRatingSvgScriptContributor : BundleContributor |
|||
{ |
|||
public override void ConfigureBundle(BundleConfigurationContext context) |
|||
{ |
|||
context.Files.AddIfNotContains("/libs/star-rating-svg/js/jquery.star-rating-svg.min.js"); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,13 @@ |
|||
using System.Collections.Generic; |
|||
using Volo.Abp.AspNetCore.Mvc.UI.Bundling; |
|||
|
|||
namespace Volo.Abp.AspNetCore.Mvc.UI.Packages.StarRatingSvg |
|||
{ |
|||
public class StarRatingSvgStyleContributor : BundleContributor |
|||
{ |
|||
public override void ConfigureBundle(BundleConfigurationContext context) |
|||
{ |
|||
context.Files.AddIfNotContains("/libs/star-rating-svg/css/star-rating-svg.css"); |
|||
} |
|||
} |
|||
} |
|||
@ -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,16 @@ |
|||
<Project Sdk="Microsoft.NET.Sdk"> |
|||
|
|||
<Import Project="..\..\..\configureawait.props" /> |
|||
<Import Project="..\..\..\common.props" /> |
|||
|
|||
<PropertyGroup> |
|||
<TargetFramework>netstandard2.0</TargetFramework> |
|||
<RootNamespace /> |
|||
</PropertyGroup> |
|||
|
|||
<ItemGroup> |
|||
<ProjectReference Include="..\Volo.Abp.EventBus\Volo.Abp.EventBus.csproj" /> |
|||
<ProjectReference Include="..\Volo.Abp.Kafka\Volo.Abp.Kafka.csproj" /> |
|||
</ItemGroup> |
|||
|
|||
</Project> |
|||
@ -0,0 +1,27 @@ |
|||
using Microsoft.Extensions.DependencyInjection; |
|||
using Volo.Abp.Kafka; |
|||
using Volo.Abp.Modularity; |
|||
|
|||
namespace Volo.Abp.EventBus.Kafka |
|||
{ |
|||
[DependsOn( |
|||
typeof(AbpEventBusModule), |
|||
typeof(AbpKafkaModule))] |
|||
public class AbpEventBusKafkaModule : AbpModule |
|||
{ |
|||
public override void ConfigureServices(ServiceConfigurationContext context) |
|||
{ |
|||
var configuration = context.Services.GetConfiguration(); |
|||
|
|||
Configure<AbpKafkaEventBusOptions>(configuration.GetSection("Kafka:EventBus")); |
|||
} |
|||
|
|||
public override void OnApplicationInitialization(ApplicationInitializationContext context) |
|||
{ |
|||
context |
|||
.ServiceProvider |
|||
.GetRequiredService<KafkaDistributedEventBus>() |
|||
.Initialize(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,12 @@ |
|||
namespace Volo.Abp.EventBus.Kafka |
|||
{ |
|||
public class AbpKafkaEventBusOptions |
|||
{ |
|||
|
|||
public string ConnectionName { get; set; } |
|||
|
|||
public string TopicName { get; set; } |
|||
|
|||
public string GroupId { get; set; } |
|||
} |
|||
} |
|||
@ -0,0 +1,208 @@ |
|||
using System; |
|||
using System.Collections.Concurrent; |
|||
using System.Collections.Generic; |
|||
using System.Linq; |
|||
using System.Threading.Tasks; |
|||
using Confluent.Kafka; |
|||
using Microsoft.Extensions.DependencyInjection; |
|||
using Microsoft.Extensions.Options; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.EventBus.Distributed; |
|||
using Volo.Abp.Kafka; |
|||
using Volo.Abp.MultiTenancy; |
|||
using Volo.Abp.Threading; |
|||
|
|||
namespace Volo.Abp.EventBus.Kafka |
|||
{ |
|||
[Dependency(ReplaceServices = true)] |
|||
[ExposeServices(typeof(IDistributedEventBus), typeof(KafkaDistributedEventBus))] |
|||
public class KafkaDistributedEventBus : EventBusBase, IDistributedEventBus, ISingletonDependency |
|||
{ |
|||
protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; } |
|||
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } |
|||
protected IKafkaMessageConsumerFactory MessageConsumerFactory { get; } |
|||
protected IKafkaSerializer Serializer { get; } |
|||
protected IProducerPool ProducerPool { get; } |
|||
protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; } |
|||
protected ConcurrentDictionary<string, Type> EventTypes { get; } |
|||
protected IKafkaMessageConsumer Consumer { get; private set; } |
|||
|
|||
public KafkaDistributedEventBus( |
|||
IServiceScopeFactory serviceScopeFactory, |
|||
ICurrentTenant currentTenant, |
|||
IOptions<AbpKafkaEventBusOptions> abpKafkaEventBusOptions, |
|||
IKafkaMessageConsumerFactory messageConsumerFactory, |
|||
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions, |
|||
IKafkaSerializer serializer, |
|||
IProducerPool producerPool) |
|||
: base(serviceScopeFactory, currentTenant) |
|||
{ |
|||
AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; |
|||
AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; |
|||
MessageConsumerFactory = messageConsumerFactory; |
|||
Serializer = serializer; |
|||
ProducerPool = producerPool; |
|||
|
|||
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>(); |
|||
EventTypes = new ConcurrentDictionary<string, Type>(); |
|||
} |
|||
|
|||
public void Initialize() |
|||
{ |
|||
Consumer = MessageConsumerFactory.Create( |
|||
AbpKafkaEventBusOptions.TopicName, |
|||
AbpKafkaEventBusOptions.GroupId, |
|||
AbpKafkaEventBusOptions.ConnectionName); |
|||
|
|||
Consumer.OnMessageReceived(ProcessEventAsync); |
|||
|
|||
SubscribeHandlers(AbpDistributedEventBusOptions.Handlers); |
|||
} |
|||
|
|||
private async Task ProcessEventAsync(Message<string, byte[]> message) |
|||
{ |
|||
var eventName = message.Key; |
|||
var eventType = EventTypes.GetOrDefault(eventName); |
|||
if (eventType == null) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
var eventData = Serializer.Deserialize(message.Value, eventType); |
|||
|
|||
await TriggerHandlersAsync(eventType, eventData); |
|||
} |
|||
|
|||
public IDisposable Subscribe<TEvent>(IDistributedEventHandler<TEvent> handler) where TEvent : class |
|||
{ |
|||
return Subscribe(typeof(TEvent), handler); |
|||
} |
|||
|
|||
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); |
|||
} |
|||
|
|||
/// <inheritdoc/>
|
|||
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; |
|||
}); |
|||
}); |
|||
} |
|||
|
|||
/// <inheritdoc/>
|
|||
public override void Unsubscribe(Type eventType, IEventHandler handler) |
|||
{ |
|||
GetOrCreateHandlerFactories(eventType) |
|||
.Locking(factories => |
|||
{ |
|||
factories.RemoveAll( |
|||
factory => |
|||
factory is SingleInstanceHandlerFactory handlerFactory && |
|||
handlerFactory.HandlerInstance == handler |
|||
); |
|||
}); |
|||
} |
|||
|
|||
/// <inheritdoc/>
|
|||
public override void Unsubscribe(Type eventType, IEventHandlerFactory factory) |
|||
{ |
|||
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Remove(factory)); |
|||
} |
|||
|
|||
/// <inheritdoc/>
|
|||
public override void UnsubscribeAll(Type eventType) |
|||
{ |
|||
GetOrCreateHandlerFactories(eventType).Locking(factories => factories.Clear()); |
|||
} |
|||
|
|||
public override async Task PublishAsync(Type eventType, object eventData) |
|||
{ |
|||
var eventName = EventNameAttribute.GetNameOrDefault(eventType); |
|||
var body = Serializer.Serialize(eventData); |
|||
|
|||
var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); |
|||
|
|||
await producer.ProduceAsync( |
|||
AbpKafkaEventBusOptions.TopicName, |
|||
new Message<string, byte[]> |
|||
{ |
|||
Key = eventName, Value = body |
|||
}); |
|||
} |
|||
|
|||
private List<IEventHandlerFactory> GetOrCreateHandlerFactories(Type eventType) |
|||
{ |
|||
return HandlerFactories.GetOrAdd( |
|||
eventType, |
|||
type => |
|||
{ |
|||
var eventName = EventNameAttribute.GetNameOrDefault(type); |
|||
EventTypes[eventName] = type; |
|||
return new List<IEventHandlerFactory>(); |
|||
} |
|||
); |
|||
} |
|||
|
|||
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(); |
|||
} |
|||
|
|||
private static bool ShouldTriggerEventForHandler(Type targetEventType, Type handlerEventType) |
|||
{ |
|||
//Should trigger same type
|
|||
if (handlerEventType == targetEventType) |
|||
{ |
|||
return true; |
|||
} |
|||
|
|||
//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,17 @@ |
|||
<Project Sdk="Microsoft.NET.Sdk"> |
|||
|
|||
<Import Project="..\..\..\configureawait.props" /> |
|||
<Import Project="..\..\..\common.props" /> |
|||
|
|||
<PropertyGroup> |
|||
<TargetFramework>netstandard2.0</TargetFramework> |
|||
<RootNamespace /> |
|||
</PropertyGroup> |
|||
|
|||
<ItemGroup> |
|||
<PackageReference Include="Confluent.Kafka" Version="1.5.0" /> |
|||
<ProjectReference Include="..\Volo.Abp.Json\Volo.Abp.Json.csproj" /> |
|||
<ProjectReference Include="..\Volo.Abp.Threading\Volo.Abp.Threading.csproj" /> |
|||
</ItemGroup> |
|||
|
|||
</Project> |
|||
@ -0,0 +1,31 @@ |
|||
using Microsoft.Extensions.DependencyInjection; |
|||
using Volo.Abp.Json; |
|||
using Volo.Abp.Modularity; |
|||
using Volo.Abp.Threading; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
[DependsOn( |
|||
typeof(AbpJsonModule), |
|||
typeof(AbpThreadingModule) |
|||
)] |
|||
public class AbpKafkaModule : AbpModule |
|||
{ |
|||
public override void ConfigureServices(ServiceConfigurationContext context) |
|||
{ |
|||
var configuration = context.Services.GetConfiguration(); |
|||
Configure<AbpKafkaOptions>(configuration.GetSection("Kafka")); |
|||
} |
|||
|
|||
public override void OnApplicationShutdown(ApplicationShutdownContext context) |
|||
{ |
|||
context.ServiceProvider |
|||
.GetRequiredService<IConsumerPool>() |
|||
.Dispose(); |
|||
|
|||
context.ServiceProvider |
|||
.GetRequiredService<IProducerPool>() |
|||
.Dispose(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,22 @@ |
|||
using System; |
|||
using Confluent.Kafka; |
|||
using Confluent.Kafka.Admin; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public class AbpKafkaOptions |
|||
{ |
|||
public KafkaConnections Connections { get; } |
|||
|
|||
public Action<ProducerConfig> ConfigureProducer { get; set; } |
|||
|
|||
public Action<ConsumerConfig> ConfigureConsumer { get; set; } |
|||
|
|||
public Action<TopicSpecification> ConfigureTopic { get; set; } |
|||
|
|||
public AbpKafkaOptions() |
|||
{ |
|||
Connections = new KafkaConnections(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,109 @@ |
|||
using System; |
|||
using System.Collections.Concurrent; |
|||
using System.Collections.Generic; |
|||
using System.Diagnostics; |
|||
using System.Linq; |
|||
using Confluent.Kafka; |
|||
using Microsoft.Extensions.Logging; |
|||
using Microsoft.Extensions.Logging.Abstractions; |
|||
using Microsoft.Extensions.Options; |
|||
using Volo.Abp.DependencyInjection; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public class ConsumerPool : IConsumerPool, ISingletonDependency |
|||
{ |
|||
protected AbpKafkaOptions Options { get; } |
|||
|
|||
protected ConcurrentDictionary<string, IConsumer<string, byte[]>> Consumers { get; } |
|||
|
|||
protected TimeSpan TotalDisposeWaitDuration { get; set; } = TimeSpan.FromSeconds(10); |
|||
|
|||
public ILogger<ConsumerPool> Logger { get; set; } |
|||
|
|||
private bool _isDisposed; |
|||
|
|||
public ConsumerPool(IOptions<AbpKafkaOptions> options) |
|||
{ |
|||
Options = options.Value; |
|||
|
|||
Consumers = new ConcurrentDictionary<string, IConsumer<string, byte[]>>(); |
|||
Logger = new NullLogger<ConsumerPool>(); |
|||
} |
|||
|
|||
public virtual IConsumer<string, byte[]> Get(string groupId, string connectionName = null) |
|||
{ |
|||
connectionName ??= KafkaConnections.DefaultConnectionName; |
|||
|
|||
return Consumers.GetOrAdd( |
|||
connectionName, connection => |
|||
{ |
|||
var config = new ConsumerConfig(Options.Connections.GetOrDefault(connection)) |
|||
{ |
|||
GroupId = groupId, |
|||
EnableAutoCommit = false |
|||
}; |
|||
|
|||
Options.ConfigureConsumer?.Invoke(config); |
|||
|
|||
return new ConsumerBuilder<string, byte[]>(config).Build(); |
|||
} |
|||
); |
|||
} |
|||
|
|||
public void Dispose() |
|||
{ |
|||
if (_isDisposed) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
_isDisposed = true; |
|||
|
|||
if (!Consumers.Any()) |
|||
{ |
|||
Logger.LogDebug($"Disposed consumer pool with no consumers in the pool."); |
|||
return; |
|||
} |
|||
|
|||
var poolDisposeStopwatch = Stopwatch.StartNew(); |
|||
|
|||
Logger.LogInformation($"Disposing consumer pool ({Consumers.Count} consumers)."); |
|||
|
|||
var remainingWaitDuration = TotalDisposeWaitDuration; |
|||
|
|||
foreach (var consumer in Consumers.Values) |
|||
{ |
|||
var poolItemDisposeStopwatch = Stopwatch.StartNew(); |
|||
|
|||
try |
|||
{ |
|||
consumer.Close(); |
|||
consumer.Dispose(); |
|||
} |
|||
catch |
|||
{ |
|||
} |
|||
|
|||
poolItemDisposeStopwatch.Stop(); |
|||
|
|||
remainingWaitDuration = remainingWaitDuration > poolItemDisposeStopwatch.Elapsed |
|||
? remainingWaitDuration.Subtract(poolItemDisposeStopwatch.Elapsed) |
|||
: TimeSpan.Zero; |
|||
} |
|||
|
|||
poolDisposeStopwatch.Stop(); |
|||
|
|||
Logger.LogInformation( |
|||
$"Disposed Kafka Consumer Pool ({Consumers.Count} consumers in {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms)."); |
|||
|
|||
if (poolDisposeStopwatch.Elapsed.TotalSeconds > 5.0) |
|||
{ |
|||
Logger.LogWarning( |
|||
$"Disposing Kafka Consumer Pool got time greather than expected: {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms."); |
|||
} |
|||
|
|||
Consumers.Clear(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,10 @@ |
|||
using System; |
|||
using Confluent.Kafka; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public interface IConsumerPool : IDisposable |
|||
{ |
|||
IConsumer<string, byte[]> Get(string groupId, string connectionName = null); |
|||
} |
|||
} |
|||
@ -0,0 +1,11 @@ |
|||
using System; |
|||
using System.Threading.Tasks; |
|||
using Confluent.Kafka; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public interface IKafkaMessageConsumer |
|||
{ |
|||
void OnMessageReceived(Func<Message<string, byte[]>, Task> callback); |
|||
} |
|||
} |
|||
@ -0,0 +1,19 @@ |
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public interface IKafkaMessageConsumerFactory |
|||
{ |
|||
/// <summary>
|
|||
/// Creates a new <see cref="IKafkaMessageConsumer"/>.
|
|||
/// Avoid to create too many consumers since they are
|
|||
/// not disposed until end of the application.
|
|||
/// </summary>
|
|||
/// <param name="topicName"></param>
|
|||
/// <param name="groupId"></param>
|
|||
/// <param name="connectionName"></param>
|
|||
/// <returns></returns>
|
|||
IKafkaMessageConsumer Create( |
|||
string topicName, |
|||
string groupId, |
|||
string connectionName = null); |
|||
} |
|||
} |
|||
@ -0,0 +1,11 @@ |
|||
using System; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public interface IKafkaSerializer |
|||
{ |
|||
byte[] Serialize(object obj); |
|||
|
|||
object Deserialize(byte[] value, Type type); |
|||
} |
|||
} |
|||
@ -0,0 +1,10 @@ |
|||
using System; |
|||
using Confluent.Kafka; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public interface IProducerPool : IDisposable |
|||
{ |
|||
IProducer<string, byte[]> Get(string connectionName = null); |
|||
} |
|||
} |
|||
@ -0,0 +1,35 @@ |
|||
using System; |
|||
using System.Collections.Generic; |
|||
using Confluent.Kafka; |
|||
using JetBrains.Annotations; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
[Serializable] |
|||
public class KafkaConnections : Dictionary<string, ClientConfig> |
|||
{ |
|||
public const string DefaultConnectionName = "Default"; |
|||
|
|||
[NotNull] |
|||
public ClientConfig Default |
|||
{ |
|||
get => this[DefaultConnectionName]; |
|||
set => this[DefaultConnectionName] = Check.NotNull(value, nameof(value)); |
|||
} |
|||
|
|||
public KafkaConnections() |
|||
{ |
|||
Default = new ClientConfig(); |
|||
} |
|||
|
|||
public ClientConfig GetOrDefault(string connectionName) |
|||
{ |
|||
if (TryGetValue(connectionName, out var connectionFactory)) |
|||
{ |
|||
return connectionFactory; |
|||
} |
|||
|
|||
return Default; |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,156 @@ |
|||
using System; |
|||
using System.Collections.Concurrent; |
|||
using System.Linq; |
|||
using System.Threading.Tasks; |
|||
using Confluent.Kafka; |
|||
using Confluent.Kafka.Admin; |
|||
using JetBrains.Annotations; |
|||
using Microsoft.Extensions.Logging; |
|||
using Microsoft.Extensions.Logging.Abstractions; |
|||
using Microsoft.Extensions.Options; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.ExceptionHandling; |
|||
using Volo.Abp.Threading; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public class KafkaMessageConsumer : IKafkaMessageConsumer, ITransientDependency, IDisposable |
|||
{ |
|||
public ILogger<KafkaMessageConsumer> Logger { get; set; } |
|||
|
|||
protected IConsumerPool ConsumerPool { get; } |
|||
|
|||
protected IExceptionNotifier ExceptionNotifier { get; } |
|||
|
|||
protected AbpKafkaOptions Options { get; } |
|||
|
|||
protected ConcurrentBag<Func<Message<string, byte[]>, Task>> Callbacks { get; } |
|||
|
|||
protected IConsumer<string, byte[]> Consumer { get; private set; } |
|||
|
|||
protected string ConnectionName { get; private set; } |
|||
|
|||
protected string GroupId { get; private set; } |
|||
|
|||
protected string TopicName { get; private set; } |
|||
|
|||
public KafkaMessageConsumer( |
|||
IConsumerPool consumerPool, |
|||
IExceptionNotifier exceptionNotifier, |
|||
IOptions<AbpKafkaOptions> options) |
|||
{ |
|||
ConsumerPool = consumerPool; |
|||
ExceptionNotifier = exceptionNotifier; |
|||
Options = options.Value; |
|||
Logger = NullLogger<KafkaMessageConsumer>.Instance; |
|||
|
|||
Callbacks = new ConcurrentBag<Func<Message<string, byte[]>, Task>>(); |
|||
} |
|||
|
|||
public virtual void Initialize( |
|||
[NotNull] string topicName, |
|||
[NotNull] string groupId, |
|||
string connectionName = null) |
|||
{ |
|||
Check.NotNull(topicName, nameof(topicName)); |
|||
Check.NotNull(groupId, nameof(groupId)); |
|||
TopicName = topicName; |
|||
ConnectionName = connectionName ?? KafkaConnections.DefaultConnectionName; |
|||
GroupId = groupId; |
|||
|
|||
AsyncHelper.RunSync(CreateTopicAsync); |
|||
Consume(); |
|||
} |
|||
|
|||
public virtual void OnMessageReceived(Func<Message<string, byte[]>, Task> callback) |
|||
{ |
|||
Callbacks.Add(callback); |
|||
} |
|||
|
|||
protected virtual async Task CreateTopicAsync() |
|||
{ |
|||
using (var adminClient = new AdminClientBuilder(Options.Connections.GetOrDefault(ConnectionName)).Build()) |
|||
{ |
|||
var topic = new TopicSpecification |
|||
{ |
|||
Name = TopicName, |
|||
NumPartitions = 1, |
|||
ReplicationFactor = 1 |
|||
}; |
|||
|
|||
Options.ConfigureTopic?.Invoke(topic); |
|||
|
|||
try |
|||
{ |
|||
await adminClient.CreateTopicsAsync(new[] {topic}); |
|||
} |
|||
catch (CreateTopicsException e) |
|||
{ |
|||
if (!e.Error.Reason.Contains($"Topic '{TopicName}' already exists")) |
|||
{ |
|||
throw; |
|||
} |
|||
} |
|||
} |
|||
} |
|||
|
|||
protected virtual void Consume() |
|||
{ |
|||
Consumer = ConsumerPool.Get(GroupId, ConnectionName); |
|||
|
|||
Task.Factory.StartNew(async () => |
|||
{ |
|||
Consumer.Subscribe(TopicName); |
|||
|
|||
while (true) |
|||
{ |
|||
try |
|||
{ |
|||
var consumeResult = Consumer.Consume(); |
|||
|
|||
if (consumeResult.IsPartitionEOF) |
|||
{ |
|||
continue; |
|||
} |
|||
|
|||
await HandleIncomingMessage(consumeResult); |
|||
} |
|||
catch (ConsumeException ex) |
|||
{ |
|||
Logger.LogException(ex, LogLevel.Warning); |
|||
AsyncHelper.RunSync(() => ExceptionNotifier.NotifyAsync(ex, logLevel: LogLevel.Warning)); |
|||
} |
|||
} |
|||
}); |
|||
} |
|||
|
|||
protected virtual async Task HandleIncomingMessage(ConsumeResult<string, byte[]> consumeResult) |
|||
{ |
|||
try |
|||
{ |
|||
foreach (var callback in Callbacks) |
|||
{ |
|||
await callback(consumeResult.Message); |
|||
} |
|||
|
|||
Consumer.Commit(consumeResult); |
|||
} |
|||
catch (Exception ex) |
|||
{ |
|||
Logger.LogException(ex); |
|||
await ExceptionNotifier.NotifyAsync(ex); |
|||
} |
|||
} |
|||
|
|||
public virtual void Dispose() |
|||
{ |
|||
if (Consumer == null) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
Consumer.Close(); |
|||
Consumer.Dispose(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,32 @@ |
|||
using System; |
|||
using System.Threading.Tasks; |
|||
using Microsoft.Extensions.DependencyInjection; |
|||
using Volo.Abp.DependencyInjection; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public class KafkaMessageConsumerFactory : IKafkaMessageConsumerFactory, ISingletonDependency, IDisposable |
|||
{ |
|||
protected IServiceScope ServiceScope { get; } |
|||
|
|||
public KafkaMessageConsumerFactory(IServiceScopeFactory serviceScopeFactory) |
|||
{ |
|||
ServiceScope = serviceScopeFactory.CreateScope(); |
|||
} |
|||
|
|||
public IKafkaMessageConsumer Create( |
|||
string topicName, |
|||
string groupId, |
|||
string connectionName = null) |
|||
{ |
|||
var consumer = ServiceScope.ServiceProvider.GetRequiredService<KafkaMessageConsumer>(); |
|||
consumer.Initialize(topicName, groupId, connectionName); |
|||
return consumer; |
|||
} |
|||
|
|||
public void Dispose() |
|||
{ |
|||
ServiceScope?.Dispose(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,102 @@ |
|||
using System; |
|||
using System.Collections.Concurrent; |
|||
using System.Diagnostics; |
|||
using System.Linq; |
|||
using Confluent.Kafka; |
|||
using Microsoft.Extensions.Logging; |
|||
using Microsoft.Extensions.Logging.Abstractions; |
|||
using Microsoft.Extensions.Options; |
|||
using Volo.Abp.DependencyInjection; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public class ProducerPool : IProducerPool, ISingletonDependency |
|||
{ |
|||
protected AbpKafkaOptions Options { get; } |
|||
|
|||
protected ConcurrentDictionary<string, IProducer<string, byte[]>> Producers { get; } |
|||
|
|||
protected TimeSpan TotalDisposeWaitDuration { get; set; } = TimeSpan.FromSeconds(10); |
|||
|
|||
public ILogger<ProducerPool> Logger { get; set; } |
|||
|
|||
private bool _isDisposed; |
|||
|
|||
public ProducerPool(IOptions<AbpKafkaOptions> options) |
|||
{ |
|||
Options = options.Value; |
|||
|
|||
Producers = new ConcurrentDictionary<string, IProducer<string, byte[]>>(); |
|||
Logger = new NullLogger<ProducerPool>(); |
|||
} |
|||
|
|||
public virtual IProducer<string, byte[]> Get(string connectionName = null) |
|||
{ |
|||
connectionName ??= KafkaConnections.DefaultConnectionName; |
|||
|
|||
return Producers.GetOrAdd( |
|||
connectionName, connection => |
|||
{ |
|||
var config = Options.Connections.GetOrDefault(connection); |
|||
|
|||
Options.ConfigureProducer?.Invoke(new ProducerConfig(config)); |
|||
|
|||
return new ProducerBuilder<string, byte[]>(config).Build(); |
|||
}); |
|||
} |
|||
|
|||
public void Dispose() |
|||
{ |
|||
if (_isDisposed) |
|||
{ |
|||
return; |
|||
} |
|||
|
|||
_isDisposed = true; |
|||
|
|||
if (!Producers.Any()) |
|||
{ |
|||
Logger.LogDebug($"Disposed producer pool with no producers in the pool."); |
|||
return; |
|||
} |
|||
|
|||
var poolDisposeStopwatch = Stopwatch.StartNew(); |
|||
|
|||
Logger.LogInformation($"Disposing producer pool ({Producers.Count} producers)."); |
|||
|
|||
var remainingWaitDuration = TotalDisposeWaitDuration; |
|||
|
|||
foreach (var producer in Producers.Values) |
|||
{ |
|||
var poolItemDisposeStopwatch = Stopwatch.StartNew(); |
|||
|
|||
try |
|||
{ |
|||
producer.Dispose(); |
|||
} |
|||
catch |
|||
{ |
|||
} |
|||
|
|||
poolItemDisposeStopwatch.Stop(); |
|||
|
|||
remainingWaitDuration = remainingWaitDuration > poolItemDisposeStopwatch.Elapsed |
|||
? remainingWaitDuration.Subtract(poolItemDisposeStopwatch.Elapsed) |
|||
: TimeSpan.Zero; |
|||
} |
|||
|
|||
poolDisposeStopwatch.Stop(); |
|||
|
|||
Logger.LogInformation( |
|||
$"Disposed Kafka Producer Pool ({Producers.Count} producers in {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms)."); |
|||
|
|||
if (poolDisposeStopwatch.Elapsed.TotalSeconds > 5.0) |
|||
{ |
|||
Logger.LogWarning( |
|||
$"Disposing Kafka Producer Pool got time greather than expected: {poolDisposeStopwatch.Elapsed.TotalMilliseconds:0.00} ms."); |
|||
} |
|||
|
|||
Producers.Clear(); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,27 @@ |
|||
using System; |
|||
using System.Text; |
|||
using Volo.Abp.DependencyInjection; |
|||
using Volo.Abp.Json; |
|||
|
|||
namespace Volo.Abp.Kafka |
|||
{ |
|||
public class Utf8JsonKafkaSerializer : IKafkaSerializer, ITransientDependency |
|||
{ |
|||
private readonly IJsonSerializer _jsonSerializer; |
|||
|
|||
public Utf8JsonKafkaSerializer(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)); |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,46 @@ |
|||
var abp = abp || {}; |
|||
(function () { |
|||
|
|||
if (!luxon) { |
|||
throw "abp/luxon library requires the luxon library included to the page!"; |
|||
} |
|||
|
|||
/* TIMING *************************************************/ |
|||
|
|||
abp.timing = abp.timing || {}; |
|||
|
|||
var setObjectValue = function (obj, property, value) { |
|||
if (typeof property === "string") { |
|||
property = property.split('.'); |
|||
} |
|||
|
|||
if (property.length > 1) { |
|||
var p = property.shift(); |
|||
setObjectValue(obj[p], property, value); |
|||
} else { |
|||
obj[property[0]] = value; |
|||
} |
|||
} |
|||
|
|||
var getObjectValue = function (obj, property) { |
|||
return property.split('.').reduce((a, v) => a[v], obj) |
|||
} |
|||
|
|||
abp.timing.convertFieldsToIsoDate = function (form, fields) { |
|||
for (var field of fields) { |
|||
var dateTime = luxon.DateTime |
|||
.fromFormat( |
|||
getObjectValue(form, field), |
|||
abp.localization.currentCulture.dateTimeFormat.shortDatePattern, |
|||
{locale: abp.localization.currentCulture.cultureName} |
|||
); |
|||
|
|||
if (!dateTime.invalid) { |
|||
setObjectValue(form, field, dateTime.toFormat("yyyy-MM-dd HH:mm:ss")) |
|||
} |
|||
} |
|||
|
|||
return form; |
|||
} |
|||
|
|||
})(jQuery); |
|||
@ -1,25 +1,29 @@ |
|||
@inject ICurrentUser CurrentUser |
|||
@using Volo.Abp.Users |
|||
@model Volo.CmsKit.Public.Web.Pages.CmsKit.Shared.Components.ReactionSelection.ReactionSelectionViewComponent.ReactionSelectionViewModel |
|||
<span class="cms-reaction-area" data-entity-type="@Model.EntityType" data-entity-id="@Model.EntityId"> |
|||
<div class="text-right"> |
|||
<div class="px-2 py-1 my-3 card border-0 shadow-sm d-inline-block"> |
|||
<span class="cms-reaction-area" data-entity-type="@Model.EntityType" data-entity-id="@Model.EntityId"> |
|||
|
|||
@if (CurrentUser.IsAuthenticated) |
|||
{ |
|||
<a class="cms-reaction-select-icon" tabindex="0"><i class="fa fa-smile-o"></i></a> |
|||
<div class="cms-reaction-selection-popover-content" style="display: none"> |
|||
@foreach (var reaction in Model.Reactions) |
|||
{ |
|||
<span class="mr-1 cms-reaction-icon @(reaction.IsSelectedByCurrentUser ? "cms-reaction-icon-selected" : "")" data-reaction-name="@reaction.Name"> |
|||
<img src="@reaction.Icon" width="18" height="18"/> |
|||
</span> |
|||
} |
|||
@if (CurrentUser.IsAuthenticated) |
|||
{ |
|||
<a class="cms-reaction-select-icon" href="#"><i class="fa fa-smile-o text-muted"></i></a> |
|||
<div class="cms-reaction-selection-popover-content" style="display: none"> |
|||
@foreach (var reaction in Model.Reactions) |
|||
{ |
|||
<span class="m-2 p-2 w-25 d-inline-block text-center cms-reaction-icon @(reaction.IsSelectedByCurrentUser ? "shadow-sm bg-light rounded cms-reaction-icon-selected" : "")" data-reaction-name="@reaction.Name"> |
|||
<i class="@reaction.Icon fa-2x"></i> |
|||
</span> |
|||
} |
|||
</div> |
|||
} |
|||
@foreach (var reaction in Model.Reactions.Where(r => r.Count > 0)) |
|||
{ |
|||
<span class="ml-3 cms-reaction-icon @(reaction.IsSelectedByCurrentUser ? "cms-reaction-icon-selected" : "")" data-reaction-name="@reaction.Name" data-click-action="@(CurrentUser.IsAuthenticated ? "true" : "false")"> |
|||
<i class="@reaction.Icon"></i> |
|||
<small class="text-muted" style="opacity: .75;">@(reaction.Count)</small> |
|||
</span> |
|||
} |
|||
</span> |
|||
</div> |
|||
} |
|||
@foreach (var reaction in Model.Reactions.Where(r => r.Count > 0)) |
|||
{ |
|||
<span class="mr-1 cms-reaction-icon @(reaction.IsSelectedByCurrentUser ? "cms-reaction-icon-selected" : "")" data-reaction-name="@reaction.Name" data-click-action="@(CurrentUser.IsAuthenticated ? "true" : "false")"> |
|||
<img src="@reaction.Icon" width="18" height="18"/> |
|||
@(reaction.Count) |
|||
</span> |
|||
} |
|||
</span> |
|||
</div> |
|||
@ -1,13 +1,8 @@ |
|||
.cms-reaction-select-icon |
|||
{ |
|||
.cms-reaction-select-icon, .cms-reaction-icon { |
|||
cursor: pointer; |
|||
} |
|||
.cms-reaction-icon |
|||
{ |
|||
cursor: pointer; |
|||
padding: 3px 5px 5px; |
|||
} |
|||
.cms-reaction-icon-selected |
|||
{ |
|||
background-color: #eef; |
|||
} |
|||
.cms-reaction-selection-popover-content i.fa-2x{ |
|||
width: 25%; |
|||
display: inline-block; |
|||
float: left; |
|||
} |
|||
@ -0,0 +1,9 @@ |
|||
namespace Volo.Abp.FeatureManagement |
|||
{ |
|||
public class FeatureProviderDto |
|||
{ |
|||
public string Name { get; set; } |
|||
|
|||
public string Key { get; set; } |
|||
} |
|||
} |
|||
@ -0,0 +1,19 @@ |
|||
using System; |
|||
using JetBrains.Annotations; |
|||
|
|||
namespace Volo.Abp.FeatureManagement |
|||
{ |
|||
[Serializable] |
|||
public class FeatureNameValueWithGrantedProvider : NameValue |
|||
{ |
|||
public FeatureValueProviderInfo Provider { get; set; } |
|||
|
|||
public FeatureNameValueWithGrantedProvider([NotNull] string name, string value) |
|||
{ |
|||
Check.NotNull(name, nameof(name)); |
|||
|
|||
Name = name; |
|||
Value = value; |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
using System; |
|||
using JetBrains.Annotations; |
|||
|
|||
namespace Volo.Abp.FeatureManagement |
|||
{ |
|||
[Serializable] |
|||
public class FeatureValueProviderInfo |
|||
{ |
|||
public string Name { get; } |
|||
|
|||
public string Key { get; } |
|||
|
|||
public FeatureValueProviderInfo([NotNull]string name, string key) |
|||
{ |
|||
Check.NotNull(name, nameof(name)); |
|||
|
|||
Name = name; |
|||
Key = key; |
|||
} |
|||
} |
|||
} |
|||