在 ASP.NET Core (c#) 后端集成 WebHooks
WebHooks
webhook 是一种应用程序为三部分应用程序提供回调的方法。
当事件发生时,源应用程序通常会向预先配置的外部 URL 发出 http POST 请求,并将触发事件的数据封装在请求的有效负载中。
这种方法允许外部应用程序通过标准 WebAPI 接口响应事件。
大多数大公司都采用这种模式,我将以Github (监听新拉取请求或创建的问题)为模板,并将此功能集成到 NetCore (c#) 后端应用程序中。
演示要求
- 能够创建/更新和删除 webhook。
- 每个 WebHook 可以监听一个或多个事件。
- WebHook 处理必须与用户请求异步,并且包含一些基本的重试策略,以防交付失败。
- 您可以查看已触发的 WebHook 的历史记录。
- 您可以查看系统历史记录(审计日志),了解 webhook 的创建/删除时间或参数的更新时间。
注
:您可以在开源 GitHub 仓库中找到该应用程序的完整源代码,包括身份验证、分布式日志记录、跟踪和监控功能:
ED图
ED图(实体关系模型)展示了基本的WebHook数据库表及其之间的关系。这正是我们需要在数据库中构建的内容。
左侧是系统事件模型,用于存储诸如 WebHook 创建/更新或删除之类的信息。这代表了一些系统审计(跟踪)日志,用于监控应用程序中的 WebHook 历史记录。
右侧是WebHook数据库模型,非常简单。每个WebHook都有多个头部信息和多个触发历史记录。
WebHook结构:
- ID - 数据库的主键
- WebHookUrl - 当有人监听时将触发的 URL。
- 密钥 - 这是 API 密钥或授权令牌,由三方应用程序用于保护 API,并仅接受具有有效密钥的请求。密钥会作为 HTTP 标头“X-Trouble-Secret”随每个请求一起传输,您可以稍后根据需要重命名该标头。
- ContentType - 表示 HTTP 消息的 ContentType。
- IsActive - 是一个标志,用于启用/禁用发送事件的钩子。
- HookEvents - 钩子应该监听的事件集合。
- 标头 - 您希望随 webhook 一起发送的其他标头列表。
- 记录 - webhook 触发日志的历史记录。
- LastTrigger - 上次触发 webhook 的时间。
WebHook历史记录结构:
- ID - 数据库的主键
- WebHookID - 关联的 webhook 主键
- Guid - 备用全局唯一键
- HookType - 触发钩子的事件类型
- 结果 - webhook 运行的结果对象
- StatusCode - 结果状态码
- ResponseBody - 从触发的端点 URL 接收到的正文
- 请求体 - 请求的内容
- 请求头 - 随消息一起发送的请求头
- 异常情况 - 有关任何异常情况的信息
- 时间戳 - 触发事件的日期和时间
WebHookHeader 结构:
- ID - 数据库的主键
- WebHookID - 关联的 webhook 主键
- 名称 - 标题的名称
- 值 - 标题的值
为了将此图转换为真正的数据库表,我使用了“代码优先”方法,结合EFCore ORM(对象关系映射器)和一些设计包,帮助我们将 C# 代码转换为自动生成的数据库迁移。
堆
这是本次演示中最重要的一部分:
Netcore作为后端框架。EF-Core作为 ORM,将我们的 C# 对象映射到数据库实体并EF-Tools创建数据库迁移。MediatR作为内存消息传递来处理诸如 CreateWebHook/ProcessWebHook 等命令。Hangfire用于处理队列中的钩子和基本重试策略。PostgreSQL作为数据库使用。(您也可以使用内存进行测试)。
注意:最终的代码源代码更加复杂,本文仅关注 WebHook 集成!
领域级
领域层是实体对象和相关业务逻辑的集合,我们称之为domain models。
WebHook 域模型
由于我们采用的code-first是创建数据库表的方法,因此我们需要定义我们的 WebHook 域:
/Src/APIServer/Domain/Hooks/WebHook.cs
namespace APIServer.Domain.Core.Models.WebHooks {
public class WebHook {
public WebHook() {
this.Headers = new HashSet<WebHookHeader>();
this.HookEvents = new HookEventType[0];
this.Records = new List<WebHookRecord>();
}
/// <summary>
/// Hook DB Id
/// </summary>
public long ID { get; set; }
/// <summary>
/// <summary>
/// Webhook endpoint
/// </summary>
public string WebHookUrl { get; set; }
/// <summary>
/// Webhook secret
/// </summary>
#nullable enable
public string? Secret { get; set; }
#nullable disable
/// <summary>
/// Content Type
/// </summary>
public string ContentType { get; set; }
/// <summary>
/// Is active / NotActiv
/// </summary>
public bool IsActive { get; set; }
/// <summary>
/// Hook Events context
/// </summary>
public HookEventType[] HookEvents { get; set; }
/// <summary>
/// Additional HTTP headers. Will be sent with hook.
/// </summary>
virtual public HashSet<WebHookHeader> Headers { get; set; }
/// <summary>
/// Hook call records history
/// </summary>
virtual public ICollection<WebHookRecord> Records { get; set; }
/// <summary>
/// Timestamp of last hook trigger
/// </summary>
/// <value></value>
public DateTime? LastTrigger { get; set; }
}
}
每个钩子包含多个WebHookRecord。
/Src/APIServer/Domain/Hooks/WebHookRecord.cs
namespace APIServer.Domain.Core.Models.WebHooks {
public class WebHookRecord {
public WebHookRecord() { }
/// <summary>
/// Hook record DB Id
/// </summary>
public long ID { get; set; }
/// <summary>
/// Linked Webhook Id
/// </summary>
public long WebHookID { get; set; }
/// <summary>
/// Linked Webhook
/// </summary>
public WebHook WebHook { get; set; }
/// <summary>
/// Unique GUID
/// </summary>
public string Guid { get; set; }
/// <summary>
/// WebHookType
/// </summary>
public HookEventType HookType { get; set; }
/// <summary>
/// Hook result enum
/// </summary>
public RecordResult Result { get; set; }
/// <summary>
/// Response
/// </summary>
public int StatusCode { get; set; }
/// <summary>
/// Response json
/// </summary>
public string ResponseBody { get; set; }
/// <summary>
/// Request json
/// </summary>
public string RequestBody { get; set; }
/// <summary>
/// Request Headers json
/// </summary>
public string RequestHeaders { get; set; }
/// <summary>
/// Exception
/// </summary>
public string Exception { get; set; }
/// <summary>
/// Hook Call Timestamp
/// </summary>
public DateTime Timestamp { get; set; }
}
public enum RecordResult {
undefined = 0,
ok,
parameter_error,
http_error,
dataQueryError
}
}
每个挂钩包含多个WebHookHeader
/Src/APIServer/Domain/Hooks/WebHookHeader.cs
namespace APIServer.Domain.Core.Models.WebHooks {
public class WebHookHeader {
/// <summary>
/// Hook header ID
/// </summary>
public long ID { get; set; }
/// <summary>
/// Linked Webhook Id
/// </summary>
public long WebHookID { get; set; }
/// <summary>
/// Linked Webhook
/// </summary>
public WebHook WebHook { get; set; }
/// <summary>
/// Header Name
/// </summary>
public string Name { get; set; }
/// <summary>
/// Header content
/// </summary>
public string Value { get; set; }
/// <summary>
/// Header created time
/// </summary>
public DateTime CreatedTimestamp { get; set; }
}
}
每个 webhook 可以由多个事件触发:
/Src/APIServer/Domain/Hooks/WebHookEvents.cs
namespace APIServer.Domain.Core.Models.WebHooks {
public enum HookEventType {
// Take this as an example, you can implement any event source you like.
hook, //(Hook created, Hook deleted ...)
file, // (Some file uploaded, file deleted)
note, // (Note posted, note updated)
project, // (Some project created, project dissabled)
milestone // (Milestone created, milestone is done etc...)
//etc etc.. You can define your custom events types....
}
}
每种事件类型都包含不同的相关操作:
/Src/APIServer/Domain/Hooks/HookActions.cs
namespace APIServer.Domain.Core.Models.WebHooks {
// Actions of HookEventType
public enum HookResourceAction {
undefined,
hook_created,
hook_removed,
hook_updated,
// etc...
}
// Actions of ProjectEventType
public enum projectAction {
undefined,
project_created,
project_renamed,
project_archived,
//etc...
}
}
事件域模型
/Src/APIServer/Domain/Events/HookActions.cs
namespace APIServer.Domain.Core.Models.Events {
public class DomainEvent {
public long ID { get; set; }
#nullable enable
public Guid? ActorID { get; set; }
#nullable disable
public DateTime TimeStamp { get; set; }
public EventType EventType { get; set; }
}
}
在哪里EventType:
/Src/APIServer/Domain/Events/EventType.cs
namespace APIServer.Domain.Core.Models.Events {
public enum EventType {
WebHook,
System,
Project,
}
}
然后,我们可以按如下方式定义系统事件:
/Src/APIServer/Domain/Events/EventType.cs
/// <summary>
/// WebHookCreated
/// </summary>
public class WebHookCreated : DomainEvent {
public WebHookCreated() { }
public long WebHookId {get;set;}
// Add any custom props hire...
}
// Equal for
持久层
该层提供辅助函数来维护数据库中的应用程序状态。一般来说,它定义了DBContext数据库configuration files和migrations。
请确保您已在您的系统中安装了所需的软件包APIServer.Persistence.csproj。
<PackageReference Include="Microsoft.EntityFrameworkCore" Version="5.0.9" />
<PackageReference Include="Microsoft.EntityFrameworkCore.Tools" Version="5.0.9">
根据我们的领域模型,DBContext 必须包含以下集合:
WebHooks- 已存储的 WebHook。WebHooksHistory- WebHook 历史日志EventsWebHook 系统日志
数据库上下文:
/Src/APIServer/Persistence/ApiDbContext.cs
namespace APIServer.Persistence {
public class ApiDbContext : DbContext {
public DbSet<WebHook> WebHooks { get; set; }
public DbSet<WebHookRecord> WebHooksHistory { get; set; }
public DbSet<DomainEvent> Events { get; set; }
public ApiDbContext(
DbContextOptions<ApiDbContext> options)
: base(options) { }
protected override void OnModelCreating(ModelBuilder modelBuilder) {
modelBuilder.ApplyConfigurationsFromAssembly(typeof(ApiDbContext).Assembly);
modelBuilder.Entity<WebHookCreated>().ToTable("WebHookCreatedEvent");
modelBuilder.Entity<WebHookRemoved>().ToTable("WebHookRemovedEvent");
modelBuilder.Entity<WebHookUpdated>().ToTable("WebHookUpdatedEvent");
modelBuilder.Entity<WebHook>().HasData(
new WebHook() {
ID = 1,
WebHookUrl = "https://localhost:5015/hookloopback",
IsActive = true,
ContentType = "application/json",
HookEvents=new HookEventType[]{HookEventType.hook}
});
base.OnModelCreating(modelBuilder);
}
}
}
通过重写该OnModelCreating函数,我们可以配置数据库模型。
ApplyConfigurationsFromAssembly- 指定加载所有配置文件。这些配置文件在/Src/APIServer/Persistence/Configuration文件夹中手动定义为单独的配置类。- 定义单位与表格的关联(TPT)
- 初始数据库数据
实体配置:(使用“ApplyConfigurationsFromAssembly”加载)
/Src/APIServer/Persistence/Configuration/WebHook.cs
namespace APIServer.Presistence {
public class WebHookonfiguration : IEntityTypeConfiguration<WebHook> {
public void Configure(EntityTypeBuilder<WebHook> builder) {
builder.HasKey(e => e.ID);
builder.HasMany(e => e.Headers)
.WithOne(e => e.WebHook)
.HasForeignKey(e => e.WebHookID);
builder.HasMany(e => e.Records)
.WithOne(e => e.WebHook)
.HasForeignKey(e => e.WebHookID);
}
}
}
要根据我们的 DBContext 定义和 DB 配置文件生成 DB 迁移,EFCore Tools需要使用 。
您需要cli运行以下命令来安装它:
dotnet tool install --global dotnet-ef
dotnet ef您可以通过在终端中运行以下命令来验证安装是否成功,您应该会看到以下内容:
_/\__
---==/ \\
___ ___ |. \|\
| __|| __| | ) \\\
| _| | _| \_/ | //|\\
|___||_| / \\\/\\
Entity Framework Core .NET Command-line Tools 5.0.9
要生成迁移文件,您可以运行以下命令Src/APIServer/Persistence
dotnet ef migrations add Init
这将创建Migrations包含相应迁移文件的文件夹。要将这些迁移应用到数据库,您可以运行以下命令:
dotnet ef database update
注意
:请确保配置中定义的连接字符串MigrationConfig.cs有效。迁移工具需要它来连接到您的数据库并创建表
。
由于我们希望保持主项目的初始状态清晰,因此我们还定义了一个服务集合扩展。这有助于我们通过单个命令注册 DBContext services.AddApiDbContext(...);,并将所有配置都从Startup.cs.
/Src/APIServer/Persistence/Extensions/AddApiDbContext.cs
namespace APIServer.Persistence.Extensions {
public static partial class ServiceExtension {
public static IServiceCollection AddApiDbContext(
this IServiceCollection serviceCollection,
IConfiguration Configuration, IWebHostEnvironment Environment) {
serviceCollection.AddPooledDbContextFactory<ApiDbContext>(
(s, o) => o
.UseNpgsql(Configuration["ConnectionStrings:ApiDbContext"], option => {
option.EnableRetryOnFailure();
if (Environment.IsDevelopment()) {
o.EnableDetailedErrors();
o.EnableSensitiveDataLogging();
}
}).UseLoggerFactory(s.GetRequiredService<ILoggerFactory>()));
return serviceCollection;
}
}
}
该ConnectionStrings:ApiDbContext定义在Src/APIServer/API/appsettings.json以下文献中:
{
"ConnectionStrings": {
"ApiDbContext": "Host=localhost;Port=5432;Database=APIServer;Username=postgres;Password=postgres",
},
}
正如定义所示,这里AddPooledDbContextFactory使用了并行处理。这是因为 GraphQL 并行解析加载的数据,因此会同时向数据库发出多个查询。传统的Scoped数据库上下文无法实现这一点,因为它不是线程安全的,无法执行并行处理!
应用层
该层包含MediatR (CQRS) commands,queries用于完成特定的应用程序任务。
它还包含一个位于 MediatR 应用逻辑之上的 GraphQL 层。遗憾的是,由于 GraphQLGraphQL和MediatRMediatR 提供的概念相似,领域驱动设计 (DDD) 模式有时无法完全遵循所有已知规则,而是在这两个部分之间进行分割。理想情况下,它们应该完全分离,但由于使用了高级错误处理或游标分页,一些 DDD 规则被忽略了。
这是命令的分布式跟踪记录CreateWebHook。您可以看到我们将要实现的每个部分。
WebHook 命令和查询
让我们来看一个示例CreateWebHook命令:
Src/APIServer/Aplication/Public/Hooks/Create_WebHook.cs
namespace APIServer.Aplication.Commands.WebHooks {
/// <summary>
/// Command for creating webhook
/// </summary>
[Authorize] // <-- Activate Auth check for command
// [Authorize(FieldPolicy = true)] <-- Uncommend to activate Field Auth check for command
public class CreateWebHook : IRequest<CreateWebHookPayload> {
public CreateWebHook() {
this.HookEvents = new HashSet<HookEventType>();
}
/// <summary> Url </summary>
public string WebHookUrl { get; set; }
/// <summary> Secret </summary>
public string? Secret { get; set; }
/// <summary> IsActive </summary>
public bool IsActive { get; set; }
/// <summary> HookEvents </summary>
public HashSet<HookEventType> HookEvents { get; set; }
}
/// <summary>
/// CreateWebHook Validator
/// </summary>
public class CreateWebHookValidator : AbstractValidator<CreateWebHook> {
private readonly IDbContextFactory<ApiDbContext> _factory;
public CreateWebHookValidator(IDbContextFactory<ApiDbContext> factory){
_factory = factory;
RuleFor(e => e.WebHookUrl)
.NotEmpty()
.NotNull();
RuleFor(e => e.WebHookUrl)
.Matches(Common.URI_REGEX)
.WithMessage("Does not match URI expression");
RuleFor(e => e.WebHookUrl)
.MustAsync(BeUniqueByURL)
.WithMessage("Hook endpoint allready exist");
RuleFor(e => e.WebHookUrl)
.MustAsync(CheckMaxAllowedHooksCount)
.WithMessage("Max allowed hooks count detected");
RuleFor(e => e.HookEvents)
.NotNull();
}
public async Task<bool> BeUniqueByURL(string url, CancellationToken cancellationToken) {
await using ApiDbContext dbContext =
_factory.CreateDbContext();
return await dbContext.WebHooks.AllAsync(e => e.WebHookUrl != url);
}
public async Task<bool> CheckMaxAllowedHooksCount(string url, CancellationToken cancellationToken) {
await using ApiDbContext dbContext =
_factory.CreateDbContext();
const long MAX_HOOK_COUNT = 3;
return (await dbContext.WebHooks.CountAsync()) <= MAX_HOOK_COUNT;
}
}
/// <summary>
/// Authorization validators for CreateWebHook
/// </summary>
public class CreateWebHookAuthorizationValidator : AuthorizationValidator<CreateWebHook> {
private readonly IDbContextFactory<ApiDbContext> _factory;
public CreateWebHookAuthorizationValidator(IDbContextFactory<ApiDbContext> factory) {
_factory = factory;
// Add Field authorization cehcks.. (use [Authorize(FieldPolicy = true)] to activate)
}
}
/// <summary>
/// ICreateWebHookError
/// </summary>
public interface ICreateWebHookError { }
/// <summary>
/// CreateWebHookPayload
/// </summary>
public class CreateWebHookPayload : BasePayload<CreateWebHookPayload, ICreateWebHookError> {
/// <summary>
/// Created WebHook
/// </summary>
public WebHook hook { get; set; }
}
/// <summary>Handler for <c>CreateWebHook</c> command </summary>
public class CreateWebHookHandler : IRequestHandler<CreateWebHook, CreateWebHookPayload> {
/// <summary>
/// Injected <c>IDbContextFactory<ApiDbContext></c>
/// </summary>
private readonly IDbContextFactory<ApiDbContext> _factory;
/// <summary>
/// Injected <c>IPublisher</c>
/// </summary>
private readonly APIServer.Extensions.IPublisher _publisher;
/// <summary>
/// Injected <c>ICurrentUser</c>
/// </summary>
private readonly ICurrentUser _current;
/// <summary>
/// Main constructor
/// </summary>
public CreateWebHookHandler(
IDbContextFactory<ApiDbContext> factory,
APIServer.Extensions.IPublisher publisher,
ICurrentUser currentuser) {
_factory = factory;
_publisher = publisher;
_current = currentuser;
}
/// <summary>
/// Command handler for <c>CreateWebHook</c>
/// </summary>
public async Task<CreateWebHookPayload> Handle(CreateWebHook request, CancellationToken cancellationToken) {
await using ApiDbContext dbContext =
_factory.CreateDbContext();
WebHook hook = new WebHook {
WebHookUrl = request.WebHookUrl,
Secret = request.Secret,
ContentType = "application/json",
IsActive = request.IsActive,
HookEvents = request.HookEvents != null ? request.HookEvents.Distinct().ToArray() : new HookEventType[0]
};
dbContext.WebHooks.Add(hook);
await dbContext.SaveChangesAsync(cancellationToken);
try {
await _publisher.Publish(new WebHookCreatedNotifi() {
ActivityId = Activity.Current.Id
}, PublishStrategy.ParallelNoWait, default(CancellationToken));
} catch { }
var response = CreateWebHookPayload.Success();
response.hook = hook;
return response;
}
}
}
如您所见,命令的逻辑(作为 GraphQL 的底层逻辑)已经准备好请求/响应对象,以便很好地适应 GraphQL 层。
这意味着:
- 根据 MediatR 行为,验证和授权检查会自动转换为验证和授权错误。
- 处理程序返回一个
payload基础对象,该对象随后由 GraphQL 使用。 - 每个命令还定义了一个自定义标记接口,例如
ICreateWebHookError。
这在最后一部分“GraphQL 变更错误”(又名 6a)中有更详细的解释,但重要的是要记住,所有处理程序都会返回一个有效负载对象!
在 GraphQL 层,我们需要定义WebHookMutations如下:
Src/APIServer/Aplication/Graphql/Mutations/WebHook.cs
namespace APIServer.Aplication.GraphQL.Mutation {
/// <summary>
/// WebHooks Mutations
/// </summary>
[ExtendObjectType(OperationTypeNames.Mutation)]
public class WebHookMutations {
/// <summary>
/// Crate new webhook
/// </summary>
public class CreateWebHookInput {
/// <summary> Url </summary>
public string WebHookUrl { get; set; }
/// <summary> Secret </summary>
public string? Secret { get; set; }
/// <summary> IsActive </summary>
public bool IsActive { get; set; }
/// <summary> HookEvents </summary>
public HookEventType[] HookEvents { get; set; }
}
/// <summary>
/// Create new webhook
/// </summary>
/// <returns></returns>
public async Task<CreateWebHookPayload> CreateWebHook(
CreateWebHookInput request,
[Service] IMediator _mediator) {
return await _mediator.Send(new CreateWebHook() {
WebHookUrl = request.WebHookUrl,
Secret = request.Secret,
IsActive = request.IsActive,
HookEvents = request.HookEvents != null ? new HashSet<HookEventType>(request.HookEvents) : new HashSet<HookEventType>(new HookEventType[0]),
});
}
}
}
如您所见,mutation == commandGraphQL 层只是简单地将输入对象传递给MediatRMediatR 并接收其返回的有效负载响应。这与通过 API 控制器调用 MediatR 的方法类似,区别在于 GraphQL 可以在处理前后执行一些额外的“魔法”(语法糖)。
查询方面则略有不同。如果我们想使用内置的分页和筛选功能,需要将 DBContext 传递给我们的处理程序,或者将逻辑保留在 GraphQL 层。为了简化演示,我们将直接从 GraphQL 层查询数据:
[UseApiDbContextAttribute]
[UsePaging(typeof(WebHookType))]
[UseFiltering]
public IQueryable<GQL_WebHook> Webhooks(
[Service] ICurrentUser current,
[ScopedService] ApiDbContext context) {
if (!current.Exist) {
return null;
}
return context.WebHooks
.AsNoTracking()
.Select(e=> new GQL_WebHook {
ID = e.ID,
WebHookUrl = e.WebHookUrl,
ContentType = e.ContentType,
IsActive = e.IsActive,
LastTrigger = e.LastTrigger,
ListeningEvents = e.HookEvents
});
}
让我们快速看一下 GraphQL 层下 mutation 类型的定义:
namespace APIServer.Aplication.GraphQL.Types {
public class CreateWebHookPayloadType : ObjectType<CreateWebHookPayload> {
protected override void Configure(IObjectTypeDescriptor<CreateWebHookPayload> descriptor) {
descriptor.Field(e => e.hook).Type<WebHookType>().Resolve(context => {
WebHook e = context.Parent<CreateWebHookPayload>().hook;
if (e == null) {
return null;
}
return new GQL_WebHook {
ID = e.ID,
WebHookUrl = e.WebHookUrl,
ContentType = e.ContentType,
IsActive = e.IsActive,
LastTrigger = e.LastTrigger,
ListeningEvents = e.HookEvents
};
});
}
}
public class CreateWebHookErrorUnion : UnionType<ICreateWebHookError> {
protected override void Configure(IUnionTypeDescriptor descriptor) {
descriptor.Type<ValidationErrorType>();
descriptor.Type<UnAuthorisedType>();
descriptor.Type<InternalServerErrorType>();
}
}
}
从定义中我们可以看出,MediatR 的这种变更会返回CreateWebHookPayload一个字段hook。我们可以使用 GraphQL 扩展或覆盖任何值,这里是修改响应或定义如何解析附加数据的合适位置。
另一点需要注意的是,CreateWebHookErrorUnion该命令可以返回 3 种不同的错误类型(以联合体形式)。这些错误类型位于 payload 下。GraphQL mutation errors(又名 6a)Error[]部分对此也有更详细的解释。
GraphQL 查询的情况也类似。让我们来看一下WebHookType:
namespace APIServer.Aplication.GraphQL.Types {
/// <summary> Graphql WebHookType </summary>
public class WebHookType : ObjectType<GQL_WebHook> {
protected override void Configure(IObjectTypeDescriptor<GQL_WebHook> descriptor) {
descriptor.AsNode().IdField(t => t.ID).NodeResolver((ctx, id) =>
ctx.DataLoader<WebHookByIdDataLoader>().LoadAsync(id, ctx.RequestAborted));
descriptor.Field(t => t.ID).Type<NonNullType<IdType>>();
descriptor.Field("systemid").Type<NonNullType<LongType>>().Resolve((IResolverContext context) => {
return context.Parent<GQL_WebHook>().ID;
});
descriptor.Field(t => t.Records)
.UseDbContext<ApiDbContext>()
.Resolve(async ctx => {
ApiDbContext _context = ctx.Service<ApiDbContext>();
ICurrentUser _current = ctx.Service<ICurrentUser>();
long hook_id = ctx.Parent<GQL_WebHook>().ID;
if (!_current.Exist) {
return null;
}
if (hook_id <= 0) {
return new List<GQL_WebHookRecord>().AsQueryable();
}
return _context.WebHooksHistory
.AsNoTracking()
.Where(e => e.WebHookID == hook_id)
.Select(e => new GQL_WebHookRecord() {
ID = e.ID,
WebHookID = e.WebHookID,
WebHookSystemID = e.WebHookID,
StatusCode = e.StatusCode,
ResponseBody = e.ResponseBody,
RequestBody = e.RequestBody,
TriggerType = e.HookType,
Result = e.Result,
Guid = e.Guid,
RequestHeaders = e.RequestHeaders,
Exception = e.Exception,
Timestamp = e.Timestamp,
}!).OrderByDescending(e => e.Timestamp);
})
.UsePaging<WebHookRecordType>()
.UseFiltering();
}
}
}
它ObjectType使用 DTO 对象,并定义每个字段的解析方式。某些字段由Parent(我们的顶级查询)解析,而其他字段则包含自定义解析逻辑,例如Records字段本身。
由于 Hotchocolate GraphQL Server 遵循 Relay 规范,因此实现了 ` Relay Global Idand`Node接口。相应的定义可在 `.` 中找到WebHookType。
这样我们就可以使用该查询来查询任何实现了该功能的对象node。
descriptor.AsNode().IdField(t => t.ID).NodeResolver((ctx, id) =>
ctx.DataLoader<WebHookByIdDataLoader>().LoadAsync(id, ctx.RequestAborted));
在哪里DataLoader<WebHookByIdDataLoader>:
Src/APIServer/Aplication/Graphql/Dataloaders/WebHooks/WebHookById_DataLoader.cs
namespace APIServer.Aplication.GraphQL.DataLoaders {
public class WebHookByIdDataLoader : BatchDataLoader<long, GQL_WebHook> {
/// <summary> Injected <c>ApiDbContext</c> </summary>
private readonly IDbContextFactory<ApiDbContext> _factory;
/// <summary> Injected <c>ICurrentUser</c> </summary>
private readonly ICurrentUser _current;
private readonly SemaphoreSlim _semaphoregate = new SemaphoreSlim(1);
public WebHookByIdDataLoader(
IBatchScheduler scheduler,
IDbContextFactory<ApiDbContext> factory,
ICurrentUser current) : base(scheduler) {
_current = current;
_factory = factory;
}
protected override async Task<IReadOnlyDictionary<long, GQL_WebHook>> LoadBatchAsync(
IReadOnlyList<long> keys,
CancellationToken cancellationToken) {
if (!_current.Exist) {
return new List<GQL_WebHook>().ToDictionary(e => e.ID, null);
}
await using ApiDbContext dbContext =
_factory.CreateDbContext();
await _semaphoregate.WaitAsync();
try {
return await dbContext.WebHooks
.AsNoTracking()
.Where(s => keys.Contains(s.ID))
.Select(e => new GQL_WebHook {
ID = e.ID,
WebHookUrl = e.WebHookUrl,
ContentType = e.ContentType,
IsActive = e.IsActive,
LastTrigger = e.LastTrigger,
ListeningEvents = e.HookEvents
})
.ToDictionaryAsync(t => t.ID, cancellationToken);
} finally {
_semaphoregate.Release();
}
}
}
}
数据加载器是一种优化技术。它们并非 GraphQL 特有,也可以是 GraphQL 之外的应用程序/业务逻辑的一部分。我们使用它们来解决 N+1 问题,即同时解析具有不同 ID 的相似数据。
数据加载器相对于传统方法的优缺点及基准测试将在单独的章节中进行详细阐述,您可以在该章节中深入了解并学习如何设计方案。本章节仅对后端逻辑进行基本介绍,因此如果您对该主题不太熟悉,也无需担心。
WebHook 通知和处理
我们已经了解了如何将 mutation 传递给 MediatR 命令以及如何返回结果,现在我们可以快速轻松地设置 webhook 的处理方式。这是内部功能,此命令不应公开。
您可以在以下位置找到所有内部命令:Src/APIServer/Apliaction/Commands/Internal。
要触发内部通知,需要使用 MediatR 发布器:
await _publisher.Publish(new WebHookCreatedNotifi() {
ActivityId = Activity.Current.Id
}, PublishStrategy.ParallelNoWait, default(CancellationToken));
这些是应用户请求处理的内部通知。最好保持它们的速度,仅收集上下文数据,并将繁重的处理工作通过入队(enque)传递给调度器(后台工作进程)。
这是 MediatR 扩展:
public string Enqueue(IRequest request, string parentJobId, JobContinuationOptions continuationOption, string description = null) {
var mediatorSerializedObject = SerializeObject(request, description);
return BackgroundJob.ContinueJobWith(parentJobId, () => _commandsExecutor.ExecuteCommand(mediatorSerializedObject), continuationOption);
}
这是一个内部处理程序,用于将所有 webhook 监听器加入队列:
namespace APIServer.Aplication.Commands.Internall.Hooks {
/// <summary>
/// Command for processing WebHook event
/// </summary>
public class EnqueueRelatedWebHooks : CommandBase {
public HookEventType EventType { get; set; }
public object Event { get; set; }
}
/// <summary>
/// Command handler for <c>EnqueueRelatedWebHooks</c>
/// </summary>
public class EnqueueRelatedWebHooksHandler : IRequestHandler<EnqueueRelatedWebHooks, Unit> {
/// <summary> Injected <c>ApiDbContext</c> </summary>
private readonly IDbContextFactory<ApiDbContext> _factory;
/// <summary> Injected <c>IMediator</c> </summary>
private readonly IMediator _mediator;
/// <summary> Injected <c>IHttpClientFactory</c> </summary>
private readonly IHttpClientFactory _clientFactory;
/// <summary> Main Constructor </summary>
public EnqueueRelatedWebHooksHandler(
IDbContextFactory<ApiDbContext> factory,
IMediator mediator,
IHttpClientFactory httpClient) {
_factory = factory;
_mediator = mediator;
_clientFactory = httpClient;
}
/// <summary> Command handler for <c>EnqueueRelatedWebHooks</c> </summary>
public async Task<Unit> Handle(EnqueueRelatedWebHooks request, CancellationToken cancellationToken) {
if (request == null || request.Event == null) {
throw new ArgumentNullException();
}
await using ApiDbContext dbContext =
_factory.CreateDbContext();
List<WebHook> hooks = await dbContext.WebHooks
.AsNoTracking()
.Where(e => e.HookEvents.Contains(HookEventType.hook))
.ToListAsync(cancellationToken);
if (hooks != null) {
foreach (var hook_item in hooks) {
if (hook_item.IsActive && hook_item.ID > 0) {
try {
_mediator.Enqueue(new ProcessWebHook() {
HookId = hook_item.ID,
Event = request.Event,
EventType = request.EventType
});
} catch { }
}
}
}
return Unit.Value;
}
}
}
webhook 处理过程可能如下所示:
namespace APIServer.Aplication.Commands.Internall.Hooks {
/// <summary> Command for processing WebHook event </summary>
public class ProcessWebHook : CommandBase {
public long HookId { get; set; }
public dynamic Event { get; set; }
public HookEventType EventType { get; set; }
}
/// <summary> Command handler for <c>ProcessWebHook</c> </summary>
public class ProcessWebHookHandler : IRequestHandler<ProcessWebHook, Unit> {
/// <summary> Injected <c>ApiDbContext</c> </summary>
private readonly IDbContextFactory<ApiDbContext> _factory;
/// <summary> Injected <c>IMediator</c> </summary>
private readonly IMediator _mediator;
/// <summary> Injected <c>IHttpClientFactory</c> </summary>
private readonly IHttpClientFactory _clientFactory;
/// <summary> Main Constructor </summary>
public ProcessWebHookHandler(
IDbContextFactory<ApiDbContext> factory,
IMediator mediator,
IHttpClientFactory httpClient) {
_factory = factory;
_mediator = mediator;
_clientFactory = httpClient;
}
/// <summary> Command handler for <c>ProcessWebHook</c> </summary>
public async Task<Unit> Handle(ProcessWebHook request, CancellationToken cancellationToken) {
WebHookRecord record = new WebHookRecord() {
WebHookID = request.HookId,
Guid = Guid.NewGuid().ToString(),
HookType = request.EventType,
Timestamp = DateTime.Now
};
if (request == null || request.HookId <= 0) {
record.Result = RecordResult.parameter_error;
}
await using ApiDbContext dbContext =
_factory.CreateDbContext();
try {
WebHook hook = null;
try {
hook = await dbContext.WebHooks
.Where(e => e.ID == request.HookId)
.FirstOrDefaultAsync(cancellationToken);
} catch (Exception ex) {
record.Result = RecordResult.dataQueryError;
record.Exception = ex.ToString();
return Unit.Value;
}
if (hook != null) {
var options = new JsonSerializerOptions {
WriteIndented = true,
IncludeFields = true,
};
var serialised_request_body = new StringContent(
JsonSerializer.Serialize<dynamic>(request.Event, options),
Encoding.UTF8,
"application/json");
var httpClient = _clientFactory.CreateClient();
/* Set Headers */
httpClient.DefaultRequestHeaders.Add("X-Trouble-Delivery", record.Guid);
if (!string.IsNullOrWhiteSpace(hook.Secret)) {
httpClient.DefaultRequestHeaders.Add("X-Trouble-Secret", hook.Secret);
}
httpClient.DefaultRequestHeaders.Add("X-Trouble-Event", request.EventType.ToString().ToLowerInvariant());
record.RequestBody = await serialised_request_body.ReadAsStringAsync(cancellationToken);
var serialized_headers = new StringContent(
JsonSerializer.Serialize(httpClient.DefaultRequestHeaders.ToList(), options),
Encoding.UTF8,
"application/json");
record.RequestHeaders = await serialized_headers.ReadAsStringAsync(cancellationToken);
if (!string.IsNullOrWhiteSpace(hook.WebHookUrl)) {
try {
using var httpResponse = await httpClient.PostAsync(hook.WebHookUrl, serialised_request_body);
if (httpResponse != null) {
record.StatusCode = (int)httpResponse.StatusCode;
if (httpResponse.Content != null) {
record.ResponseBody = await httpResponse.Content.ReadAsStringAsync(cancellationToken);
}
}
record.Result = RecordResult.ok;
} catch (Exception ex) {
record.Result = RecordResult.http_error;
record.Exception = ex.ToString();
}
} else {
record.Result = RecordResult.parameter_error;
}
} else {
record.Result = RecordResult.parameter_error;
}
} finally {
try {
dbContext.WebHooksHistory.Add(record);
await dbContext.SaveChangesAsync(cancellationToken);
} catch { }
}
return Unit.Value;
}
}
}
API 层
这里配置 API 微服务。这取决于application底层persistence架构。
该类Startup.cs使用依赖注入 (DI) 设置和中间件管道配置宿主应用程序。以下是一个示例配置:
Src/APIServer/API/Startup.cs
namespace APIServer
{
public class Startup
{
public Startup(IConfiguration configuration, IWebHostEnvironment enviroment)
{
Configuration = configuration;
Environment = enviroment;
}
public IWebHostEnvironment Environment { get; }
public IConfiguration Configuration { get; }
// This method gets called by the runtime. Use this method to add services to the container.
public void ConfigureServices(IServiceCollection services)
{
services.AddControllers();
services.AddSwaggerGen(c =>{
c.SwaggerDoc("v1", new OpenApiInfo { Title = "Api", Version = "v1" });
});
services.AddAuth(Configuration);
services.AddDbContext(Configuration,Environment);
services.AddHttpClient();
services.AddHealthChecks();
services.AddGraphql(Environment);
services.AddHttpContextAccessor();
services.AddMemoryCache();
services.AddTelemerty(Configuration,Environment);
services.AddScoped<ICurrentUser, CurrentUser>();
services.AddMediatR();
services.AddScheduler(Configuration);
services.AddSingleton(Serilog.Log.Logger);
}
// This method gets called by the runtime. Use this method to configure the HTTP request pipeline.
public void Configure(
IApplicationBuilder app,
IWebHostEnvironment env,
IServiceProvider serviceProvider,
IServiceScopeFactory scopeFactory)
{
app.UseForwardedHeaders(new ForwardedHeadersOptions{
ForwardedHeaders = ForwardedHeaders.XForwardedFor | ForwardedHeaders.XForwardedProto | ForwardedHeaders.XForwardedHost,
});
app.UseEnsureApiContextCreated(serviceProvider,scopeFactory);
if (env.IsDevelopment()) {
app.UseDeveloperExceptionPage();
app.UseSwagger();
app.UseSwaggerUI(c => c.SwaggerEndpoint("/swagger/v1/swagger.json", "Api v1"));
}
app.UseHealthChecks("/health");
app.UseHttpsRedirection();
app.UseRouting();
app.UseAuthentication();
app.UseAuthorization();
app.UseHangfireServer();
if(env.IsDevelopment()){
app.UseHangfireDashboard("/scheduler");
}
app.UseEndpoints(endpoints =>
{
endpoints.MapControllers()
.RequireAuthorization("ApiCaller");
endpoints.MapGraphQL()
.WithOptions(new GraphQLServerOptions {
EnableSchemaRequests = env.IsDevelopment(),
Tool = { Enable = env.IsDevelopment() },
});
endpoints.MapControllerRoute(
name: "default",
pattern: "{controller}/{action=Index}/{id?}");
});
}
}
}
由于日志记录、追踪、监控和 PostgreSQL 错误等方面的配置已在单独的章节中进行了详细说明,因此本文不会逐日赘述,仅提供最重要的部分。如需查看所有配置,请参阅本文末尾的参考资料!
GraphQL 部分的配置包括:
namespace APIServer.Configuration {
public static partial class ServiceExtension {
public static IServiceCollection AddGraphql(
this IServiceCollection serviceCollection, IWebHostEnvironment env) {
serviceCollection.AddGraphQLServer()
.SetPagingOptions(
new PagingOptions { IncludeTotalCount = true, MaxPageSize = 200 })
.ModifyRequestOptions(requestExecutorOptions => {
if (env.IsDevelopment() ||
System.Diagnostics.Debugger.IsAttached) {
requestExecutorOptions.ExecutionTimeout = TimeSpan.FromMinutes(1);
}
requestExecutorOptions.IncludeExceptionDetails = !env.IsProduction();
})
.AddGlobalObjectIdentification()
.AddQueryFieldToMutationPayloads()
.AddFiltering()
.AddSorting()
.AddQueryType<Query>()
.AddTypeExtension<WebHookQueries>()
.AddTypeExtension<UserQueries>()
.AddTypeExtension<SystemQueries>()
.AddMutationType<Mutation>()
.AddTypeExtension<WebHookMutations>()
.BindRuntimeType<DateTime, DateTimeType>()
.BindRuntimeType<int, IntType>()
.AddType<BadRequestType>()
.AddType<InternalServerErrorType>()
.AddType<UnAuthorisedType>()
.AddType<ValidationErrorType>()
.AddType<BaseErrorType>()
.AddType<UserDeactivatedType>()
.AddType<BaseErrorInterfaceType>()
.AddType<WebHookNotFoundType>()
.AddType<WebHookRecordType>()
.AddType<WebHookType>()
.AddType<UserType>()
.AddType<UpdateWebHookUriPayloadType>()
.AddType<UpdateWebHookUriErrorUnion>()
.AddType<UpdateWebHookTriggerEventsPayloadType>()
.AddType<UpdateWebHookTriggerEventsErrorUnion>()
.AddType<UpdateWebHookSecretPayloadType>()
.AddType<UpdateWebHookSecretErrorUnion>()
.AddType<UpdateWebHookPayloadPayloadType>()
.AddType<UpdateWebHookPayloadErrorUnion>()
.AddType<UpdateWebHookActivStatePayloadType>()
.AddType<UpdateWebHookActivStateErrorUnion>()
.AddType<RemoveWebHookPayloadType>()
.AddType<RemoveWebHookErrorUnion>()
.AddType<CreateWebHookPayloadType>()
.AddType<CreateWebHookErrorUnion>()
.AddDataLoader<UserByIdDataLoader>()
.AddDataLoader<WebHookByIdDataLoader>()
.AddDataLoader<WebHookRecordByIdDataLoader>()
.UsePersistedQueryPipeline()
.UseReadPersistedQuery()
.AddReadOnlyFileSystemQueryStorage("./persisted_queries");
return serviceCollection;
}
}
}
MediatR 配置:
namespace APIServer.Configuration {
public static partial class ServiceExtension {
public static IServiceCollection AddMediatR(this IServiceCollection services) {
// Command executor
services.AddMediatR(cfg => cfg.Using<AppMediator>(), typeof(CreateWebHook).GetTypeInfo().Assembly);
services.AddTransient<APIServer.Extensions.IPublisher, APIServer.Extensions.Publisher>();
services.AddValidatorsFromAssembly(typeof(CreateWebHookValidator).GetTypeInfo().Assembly);
services.AddTransient<IRequestHandler<EnqueSaveEvent<WebHookCreated>, Unit>, EnqueSaveEventHandler<WebHookCreated>>();
services.AddTransient<IRequestHandler<EnqueSaveEvent<WebHookUpdated>, Unit>, EnqueSaveEventHandler<WebHookUpdated>>();
services.AddTransient<IRequestHandler<EnqueSaveEvent<WebHookRemoved>, Unit>, EnqueSaveEventHandler<WebHookRemoved>>();
services.AddMediatRSchedulerIntegration();
services.AddTransient(typeof(IPipelineBehavior<,>), typeof(TracingBehaviour<,>));
services.AddTransient(typeof(IPipelineBehavior<,>), typeof(CommandPerformanceBehaviour<,>));
services.AddTransient(typeof(IPipelineBehavior<,>), typeof(ValidationBehaviour<,>));
services.AddTransient(typeof(IPipelineBehavior<,>), typeof(AuthorizationBehaviour<,>));
services.AddTransient(typeof(IPipelineBehavior<,>), typeof(UnhandledExBehaviour<,>));
return services;
}
}
}
在哪里AppMediator:
namespace APIServer.Extensions {
public class AppMediator : Mediator {
private Func<IEnumerable<Func<INotification, CancellationToken, Task>>, INotification, CancellationToken, Task> _publishStrategy;
public AppMediator(
ServiceFactory serviceFactory
) : base(serviceFactory) {
}
public AppMediator(
ServiceFactory serviceFactory,
Func<IEnumerable<Func<INotification, CancellationToken, Task>>, INotification, CancellationToken, Task>? publishStrategy
) : base(serviceFactory) {
_publishStrategy = publishStrategy != null ? publishStrategy : SyncStopOnException;
}
public Task<TResponse> Send<TResponse>(ICommandBase<TResponse> request, CancellationToken cancellationToken = default) {
return base.Send<TResponse>(request, cancellationToken);
}
public Task<object?> Send(ICommandBase request, CancellationToken cancellationToken = default) {
return base.Send(request as object, cancellationToken);
}
private static async Task SyncStopOnException(IEnumerable<Func<INotification, CancellationToken, Task>> handlers, INotification notification, CancellationToken cancellationToken) {
foreach (var handler in handlers) {
await handler(notification, cancellationToken).ConfigureAwait(false);
}
}
protected override Task PublishCore(IEnumerable<Func<INotification, CancellationToken, Task>> allHandlers, INotification notification, CancellationToken cancellationToken = default) {
Activity activity = null;
if (notification is INotificationBase) {
INotificationBase I_base_notify = notification as INotificationBase;
// If any activity is in context will be set as default in case value of ActivityId == null
if (I_base_notify.ActivityId == null
&& Activity.Current != null
&& Activity.Current.Id != null) {
I_base_notify.ActivityId = Activity.Current.Id;
}
activity = Sources.DemoSource.StartActivity(
String.Format("PublishCore: Notification<{0}>", notification.GetType().FullName), ActivityKind.Producer);
// This chane activity parrent / children relation..
if (I_base_notify.ActivityId != null
&& Activity.Current != null
&& Activity.Current.ParentId == null) {
Activity.Current.SetParentId(I_base_notify.ActivityId);
}
if (Activity.Current != null
&& Activity.Current.ParentId != null) {
Activity.Current.AddTag("Parrent Id", Activity.Current.ParentId);
}
} else {
activity = Sources.DemoSource.StartActivity(
String.Format("PublishCore: Notification<{0}>", notification.GetType().FullName), ActivityKind.Producer);
}
try {
Activity.Current.AddTag("Activity Id", Activity.Current.Id);
activity.Start();
return _publishStrategy != null ? _publishStrategy(allHandlers, notification, cancellationToken) : base.PublishCore(allHandlers, notification, cancellationToken);
} finally {
activity.Stop();
activity.Dispose();
}
}
}
}
用户界面
React UI 的详细集成将在单独的资料中进行讲解。以下是当前的演示示例:
存储库
您可以在开源的Github仓库中找到该应用程序的完整源代码,包括身份验证、分布式日志记录、跟踪和监控功能:
https://github.com/damikun/trouble-training
文章来源:https://dev.to/damikun/integrate-webhook-under-net-c-backend-4f7


