Bitwarden Server Event Integrations 架构解析:双层 AMQP 消息流水线、重试机制与新集成扩展指南
【免费下载链接】serverBitwarden infrastructure/backend (API, database, Docker, etc).项目地址: https://gitcode.com/GitHub_Trending/ser/server
本文基于 Bitwarden 服务端仓库中的设计文档 src/Core/Dirt/EventIntegrations/README.md,深入讲解 Bitwarden Server 事件集成(Event Integrations)子系统的整体架构:如何借助 RabbitMQ / Azure Service Bus 的 fan-out 能力,把组织审计事件广播到 Slack、Webhook、HEC、Datadog、Teams 等多种外部集成,并通过“事件层 + 集成层”的双层 Exchange 设计实现独立重试、死信队列与模板化消息。读完本篇,你将掌握该子系统的消息流转链路、重试与缓存机制,并能够按照官方步骤独立完成一个新事件集成的接入与部署。
一、设计目标
事件集成子系统的核心目标是:在不进行大量定制开发的前提下,让新集成可以随时间轻松加入。具体包括四点(原文档“Design goals”章节的完整继承):
- Fan-out 即插即用:借助 AMQP(RabbitMQ 或 Azure Service Bus)提供的扇出(fan-out)能力,任何数量的新集成都可以挂接到现有事件系统上,无需额外广播代码。只要向现有流水线添加一个新的 listener,它就能获得一条独立的事件流。
- 健壮的失败与重试处理:通过下述双层 Exchange 方法,在服务层面内建重试支持。新集成只需关注集成自身的业务逻辑与状态上报,重试与延迟完全由消息系统托管。
- 自托管(Self-Hosted)同等支持:该能力不仅面向云端版本,也对自托管实例开放。RabbitMQ 为自托管实例提供了一个轻量级的接入方式,无需 Azure Service Bus 即可使用同一套健壮架构。
- 组织管理员的灵活控制:允许组织管理员决定哪些事件有重要意义、事件发送到哪里、以及消息中携带哪些数据。配置架构允许组织对特定集成的细节进行自定义(见下文“集成与集成配置”一节)。
二、整体架构与入口:IEventWriteService
事件集成的入口是IEventWriteService。通过把EventIntegrationEventWriteService配置为EventWriteService,所有发送到该服务的事件都会在 RabbitMQ 或 Azure Service Bus 的消息交换器(exchange)上广播。为屏蔽“发布到具体 AMQP 提供商”的细节,EventIntegrationEventWriteService中注入了一个IEventIntegrationPublisher,由它负责真正把事件发布到 RabbitMQ 或 Azure Service Bus。
2.1 源码印证:写服务的装配优先级
从 EventIntegrationsServiceCollectionExtensions.cs 中的AddEventWriteServices扩展方法可以看到,IEventWriteService的具体实现是按配置“择优”注册的,优先级依次为:
- Azure Service Bus:
EventLogging.AzureServiceBus的连接串、事件 Topic、集成 Topic 三项齐全时,注册EventIntegrationEventWriteService,并以AzureServiceBusService作为IEventIntegrationPublisher; - RabbitMQ:
EventLogging.RabbitMq的 HostName、Username、Password、EventExchangeName、IntegrationExchangeName 五项齐全时(见IsRabbitMqEnabled,同文件 #L541-L548),注册EventIntegrationEventWriteService+RabbitMqService; - Azure Queue Storage:配置了
Events.ConnectionString与Events.QueueName时,注册AzureQueueEventWriteService; - Repository(自托管):
SelfHosted为 true 时,注册RepositoryEventWriteService(事件直接落库); - Noop:以上均未配置时,注册不做任何事的
NoopEventWriteService。
这一装配顺序解释了文档中“云端存 Azure Tables、自托管存数据库”的行为差异来源:没有配置消息总线时,事件走仓储直写。
2.2 发布侧实现
EventIntegrationEventWriteService.cs 的实现非常薄:CreateAsync把单个IEvent序列化为 JSON,CreateManyAsync把事件数组序列化为 JSON(取第一个事件的OrganizationId作为路由依据),然后统一调用IEventIntegrationPublisher.PublishEventAsync。这正是文档所述“消息体是单个EventMessage或EventMessage数组的 JSON 表示”的发布端实现。
三、双层 Exchange(Two-tier Exchange)
当EventIntegrationEventWriteService发布消息时,它发送到双层消息处理方案的第一层。每一层在 AMQP 栈中对应一个独立的 exchange(RabbitMQ 术语)或 topic(Azure Service Bus 术语)。
3.1 事件层(Event Tier)
在第一层,事件通过 fan-out 广播给一组 listener。消息体是单个EventMessage或EventMessage数组的 JSON 表示,该层的 handler 负责处理每个事件或事件数组。目前有两个handler:
EventRepositoryHandler- 负责事件的长期存储。它接收所有事件,通过注入的
IEventRepository存入数据库。 - 这与事件集成关闭时的行为一致:云端存 Azure Tables,自托管存数据库。
- 负责事件的长期存储。它接收所有事件,通过注入的
EventIntegrationHandler- 一个泛型通用类,按集成的配置细节定制化,职责是:判断“该事件 / 该组织 / 该集成”是否存在配置、拉取该配置、并把事件详情解析进模板字符串。
- 它使用注入的
IOrganizationIntegrationConfigurationRepository,按事件类型、组织、集成类型取出对应的配置与模板集合。这些配置决定了:该集成是否应触发、发送所需细节、以及实际要发送的消息内容。 - 其输出是一个新的
IntegrationMessage,携带与该集成交互所需的配置详情和(已填入事件细节的)待发消息,并发布到消息总线的集成层。
源码印证:EventIntegrationHandler.cs 中HandleEventAsync的完整流程与文档描述一一对应:
- 按
OrganizationId + IntegrationType + EventType从扩展缓存(IFusionCache)读取List<OrganizationIntegrationConfigurationDetails>,缓存未命中才回源到configurationRepository.GetManyByEventTypeOrganizationIdIntegrationType(L126-L155); - 若配置带
Filters,先反序列化为IntegrationFilterGroup并用IIntegrationFilterService求值,为false则continue丢弃该事件(L35-L44); - 用
IntegrationTemplateProcessor.ReplaceTokens渲染模板(上下文按需从缓存补齐 Group / User / ActingUser / Organization 等实体,默认 TTL 30 分钟); - 组装
IntegrationMessage<T>(MessageId优先复用事件的IdempotencyId,保证幂等),调用eventIntegrationPublisher.PublishAsync发布到集成层(L46-L65)。
3.2 集成层(Integration Tier)
在集成层,消息是IIntegrationMessage的 JSON 表示,具体类型是泛型IntegrationMessage<T>的某个具体实现,其中<T>即为该集成的配置详情。这些消息代表“把某个特定事件发送到某个特定集成”所需的全部细节,包括重试与延迟处理。
集成层的 handler 直接绑定到具体集成(例如SlackIntegrationHandler、WebhookIntegrationHandler)。它们接收IntegrationMessage<T>,输出IntegrationHandlerResult,告诉 listener 集成的结果(成功/失败、是否可重试、以及最短延迟等)。这种设计让它们可以在完全脱离 AMQP 与消息系统的情况下被单独单元测试。
该层的 listener 负责:收到新消息后调用 handler,再根据结果采取行动——成功结果直接确认(ack)消息并结束;失败则要么进入死信队列(DLQ),要么在正确的延迟后重新发布以重试。
四、重试机制(Retries)
引入集成层的目标之一,就是简化并支撑针对单个事件集成的多次重试。例如:某个下游服务短暂宕机时,不希望某个 handler 阻塞其余队列的等待重试;如果某个事件的 N 个集成中只有一个失败,也不希望对其余集成全部重试,更不希望重新查询配置。通过将IntegrationMessage<T>拆分为携带配置、消息体与重试细节的独立消息,每个“事件/集成”组合都可以被独立处理、独立重试。
当IntegrationHandlerResult.Success为false(本次集成尝试失败)时,Retryable标志告诉 listener 该失败是暂时性的还是终局性的:
Retryable为false:消息立即送入 DLQ;Retryable为true:listener 调用IntegrationMessage的ApplyRetry(DateTime)方法,它同时负责递增RetryCount、按给定 DateTime 更新DelayUntilDate,并叠加基于RetryCount的指数退避与随机抖动。listener 再比较RetryCount是否超过 Global Settings 中定义的MaxRetries:超过则进 DLQ,否则安排重试。
4.1 源码印证:退避算法与参数默认值
IntegrationMessage.cs 中ApplyRetry的具体实现与文档完全吻合,并给出了可复现的退避公式:
public void ApplyRetry(DateTime? handlerDelayUntilDate) { RetryCount++; var baseTime = handlerDelayUntilDate ?? DateTime.UtcNow; var backoffSeconds = Math.Pow(2, RetryCount); // 指数退避:2^n 秒 var jitterSeconds = Random.Shared.Next(0, 3); // 0~2 秒随机抖动 DelayUntilDate = baseTime.AddSeconds(backoffSeconds + jitterSeconds); }即第 1、2、3 次重试的基础延迟分别为约 2s、4s、8s(叠加 0~2s 抖动),且 handler 可通过DelayUntilDate提供更晚的起始时间。相关默认值在 GlobalSettings.cs 中确认:EventLogging.MaxRetries默认为3(#L318),EventLogging.RabbitMq.UseDelayPlugin默认为false(#L369)。ListenerConfiguration基类把MaxRetries直接绑定到globalSettings.EventLogging.MaxRetries(见 ListenerConfiguration.cs)。
4.2 Azure Service Bus 的重试调度
Azure Service Bus 在核心能力上原生支持消息定时:重试被调度到某一具体时间点,ASB 会扣留消息并在正确时间再发布。
4.3 RabbitMQ 的重试选项(仅自托管使用)
对 RabbitMQ(仅自托管使用)有两种选项,由GlobalSettings.RabbitMqSettings中的useDelayPlugin标志决定。默认false,表示使用“重试队列 + 时间检查”方案。
选项 1:延迟插件(Delay Plugin)
- 对应 RabbitMQ 官方仓库中的
rabbitmq-delayed-message-exchange插件(可在 RabbitMQ 官方仓库搜索该名称)。 - 该插件在 RabbitMQ 中提供“延迟消息交换器”,支持按特殊头(header)中指定的时长延迟一条消息。
- 这使得方案可以完全不用重试队列,直接依赖 delay exchange:消息打上 header 后发布到 exchange,由 exchange 负责在合适时间前扣留消息(类似 ASB 的内建支持)。
- 该插件必须先安装并启用,之后才能打开此选项(因此默认关闭)。
选项 2:重试队列 + 时间检查(默认)
- 关闭延迟插件时,消息被推入一个重试队列,该队列有固定的时间长度后才把消息重新发布回主队列。
- 消息从队列取出时,检查
DelayUntilDate是否已过期:- 已过期:正常处理集成并重试请求;
- 仍在未来:把消息放回重试队列继续等待。
- 虽然这消耗额外处理开销,但在延迟插件未启用的情况下能更好地遵守延迟承诺。由于该方案仅用于自托管,短延迟、少量重试下的开销可以忽略。
源码印证:RabbitMqIntegrationListenerService.cs 的ProcessReceivedMessageAsync完整实现了上述逻辑——先反序列化消息并检查DelayUntilDate是否仍在未来,是则RepublishToRetryQueueAsync后 ack 返回;否则调用 handler,按result.Success / result.Retryable分支执行 ack、ApplyRetry后PublishToRetryAsync(当RetryCount < MaxRetries)或PublishToDeadLetterAsync(超出上限或不可重试)。RabbitMqService.cs 中则依据_useDelayPlugin决定是在 header 中写入延迟时间、还是走 retry routing key。
五、Listener / Handler 模式
为了同时支持多个 AMQP 服务(RabbitMQ 与 Azure Service Bus),监听消息流的动作与响应消息的动作被解耦:
5.1 Listeners
- 处理通信平台(RabbitMQ / Azure Service Bus)的细节;
- 每个平台、每个层级各有一个 listener,即各有一个 event listener 和一个 integration listener;
- 负责消息平台的一切装配/拆除、订阅、消息确认(ack)等,但自身不处理任何事件,而是委托给与之配对的 handler;
- 可以配置多个实例独立运行,各自拥有独立的 handler 与订阅/队列。
5.2 Handlers
- 每个队列/订阅(集成层即每个集成)对应一个 handler;
- 完全隔离于消息平台、对其一无所知,因此可以跨通信平台自由复用;
- 负责事件处理的所有环节;
- 因其隔离与解耦,具备高度可测试性。
这种组合使得 ServiceCollectionExtensions 中的配置 可以把“当前消息平台的 listener 实例”与“任意数量的 handler”配对。AddEventIntegrationServices会为每个集成注册三样东西(以AddAzureServiceBusIntegration<TConfig, TListenerConfig>/AddRabbitMqIntegration<TConfig, TListenerConfig>为例):
- 以
ListenerConfiguration.RoutingKey为 key 的EventIntegrationHandler<TConfig>键控单例; - 事件层的
*EventListenerService宿主服务; - 集成层的
*IntegrationListenerService宿主服务。
目前代码中实际注册的集成 handler 包括 Slack、Webhook、Hec(复用 Webhook handler 类型)、Datadog、Teams 五个(见 EventIntegrationsServiceCollectionExtensions.cs #L305-L316)。
六、Publishers 与 Services
Listeners(以及EventIntegrationHandler)通过IEventPublisher接口与消息系统交互,其后是由 RabbitMQ 与 ASB 专属服务支撑的实现。把消息平台的大部分细节放在服务层,可以在一处集中处理连接配置、绑定/创建特定队列等公共事务。IRabbitMqService与IAzureServiceBusService都实现了IEventPublisher接口,因此也能直接处理所有消息发布功能。
七、集成与集成配置(Integrations & Configurations)
组织可以为不同端点配置集成配置——每个 handler 映射到一个特定集成,收到事件时检查对应配置。当前已有 Slack、Webhooks 与 HTTP Event Collector(HEC)的集成/handler(代码中还扩展了 Datadog 与 Teams)。
7.1OrganizationIntegration
- 组织级启用某个集成的顶层对象;
- 包含适用于“该集成所有事件”的属性。
- 例如:Slack 把 token 存在
Configuration中,对每个事件都生效;而把 channel id 存在OrganizationIntegrationConfiguration的Configuration中。token 适用于整个 Slack 集成,但 channel 可以按事件类型不同而不同。 - 各级具体存什么,见下表。
- 例如:Slack 把 token 存在
7.2OrganizationIntegrationConfiguration
- 包含该集成针对每个
EventType的配置; Configuration包含事件级配置。- 此层级的任何属性都会覆盖
OrganizationIntegration中的Configuration; - 具体集成的例子见下表。
- 此层级的任何属性都会覆盖
Template包含模板字符串,预期用实际事件内容填充。- 字符串中的 token 用
#字符包裹,例如 UserId 写作#UserId#。 IntegrationTemplateProcessor负责把 token 替换为从给定EventMessage内省(introspect)出来的值。- 模板不强制任何结构——它既可以是发到 Slack 的自由文本,也可以是发到 webhook 的 JSON body;它只作为字符串存储和使用,以最大化灵活性。
- 字符串中的 token 用
7.3OrganizationIntegrationConfigurationDetails
- 它是
OrganizationIntegration与OrganizationIntegrationConfiguration二合一的合并对象,合并内容告诉集成的 handler 调用外部服务所需的全部细节; OrganizationIntegrationConfiguration优先于OrganizationIntegration——两者都存在的 key,取OrganizationIntegrationConfiguration的值;EventIntegrationHandler从数据库取出的正是一个OrganizationIntegrationConfigurationDetails数组,用于决定在集成层发布什么。
7.4 现有集成及其两级配置明细
下表说明各集成如何配置、Configuration属性在两级(OrganizationIntegration或OrganizationIntegrationConfiguration)中分别存什么。OrganizationIntegration列中,有效的OrganizationIntegrationStatus以粗体标出,并给出每种状态下的存储示例。
| 集成 | OrganizationIntegration | OrganizationIntegrationConfiguration |
|---|---|---|
| CloudBillingSync | 不适用(尚未使用) | 不适用(尚未使用) |
| Scim | 不适用(尚未使用) | 不适用(尚未使用) |
| Slack | Initiated:nullCompleted: { "Token": "xoxb-token-from-slack" } | { "channelId": "C123456" } |
| Webhook | null或{ "Scheme": "Bearer", "Token": "AUTH-TOKEN", "Uri": "https://example.com" } | null或{ "Scheme": "Bearer", "Token":"AUTH-TOKEN", "Uri": "https://example.com" }此层级定义的值优先 |
| Hec | { "Scheme": "Bearer", "Token": "AUTH-TOKEN", "Uri": "https://example.com" } | 恒为null |
| Datadog | { "ApiKey": "TheKey12345", "Uri": "https://api.us5.datadoghq.com/api/v1/events"} | 恒为null |
| Teams | Initiated:nullIn Progress: { "TenantID": "tenant", "Teams": ["Id": "team", DisplayName: "MyTeam"]}Completed: { "TenantID": "tenant", "Teams": ["Id": "team", DisplayName: "MyTeam"], "ServiceUrl":"https://example.com", ChannelId: "channel-1234"} | 恒为null |
八、过滤(Filtering)
除了上文所述的集成配置能力,组织管理员还可以在OrganizationIntegrationConfiguration中添加可选的Filters。过滤器完全可选,管理员可以把它做得简单或复杂。过滤器以 JSON 形式存库,并被反序列化为IntegrationFilterGroup,随后交给IntegrationFilterService求值为bool:true时集成按上述流程继续,false时忽略该事件,不路由到集成层。
8.1IntegrationFilterGroup
若干规则与其他子分组的逻辑 AND / OR 组合。
| 属性 | 说明 |
|---|---|
AndOperator | 指示Rules与Groups中全部(true)还是任意(false)为真即可。该语义同时作用于内部分组与规则列表;例如本组包含 Rule1、Rule2 以及 Group1、Group2 时:true:Rule1 && Rule2 && Group1 && Group2false:Rule1 \|\| Rule2 \|\| Group1 \|\| Group2 |
Rules | IntegrationFilterRule列表。可为 null 或空,此时返回true。 |
Groups | 嵌套的IntegrationFilterGroup列表。可为 null 或空,此时返回true。 |
8.2IntegrationFilterRule
过滤框架的核心:判断本条 EventMessage 的数据是否匹配过滤器所查找的数据。
| 属性 | 说明 |
|---|---|
Property | EventMessage上要评估的属性(例如CollectionId)。 |
Operation | 属性与Value之间执行的比较。支持的操作: • Equals:Guid等于Value• NotEquals:Equals的逻辑反向• In:Guid在Value列表中• NotIn:In的逻辑反向 |
Value | 比较值。类型取决于Operation:• Equals、NotEquals:Guid• In、NotIn:Guid列表 |
源码印证:EventIntegrationHandler.HandleEventAsync中对Filters字段的处理(反序列化 +EvaluateFilterGroup求值,为false即continue)就是这段描述的直接实现,见 EventIntegrationHandler.cs #L35-L44。
九、缓存(Caching)
为降低数据库负载、提升性能,事件集成使用自己独立的具名扩展缓存(extended cache,更多信息见 src/Core/Utilities/CACHING.md)。没有缓存时,每条进入的EventMessage都会触发一次数据库查询,去取相关的OrganizationIntegrationConfigurationDetails。
9.1EventIntegrationsCacheConstants
该常量类让代码在操作扩展缓存时能以强类型引用大量缓存细节:缓存名以及所有缓存 key 和 tag 都从EventIntegrationsCacheConstants程序化访问,而不是散落的字符串字面量。例如EventIntegrationsCacheConstants.CacheName被用于缓存装配、键控服务、依赖注入等场景,而不是在代码里写死字符串 "EventIntegrations"。
9.2OrganizationIntegrationConfigurationDetails的缓存策略
- 这是架构中被使用最频繁的部件之一:任何带有组织的关联事件都需要检查配置,判断是否要触发集成;
- 借助扩展缓存,所有读操作在触及数据库前都先命中 L1 或 L2 缓存;
- 读取返回给定 key 的
List<OrganizationIntegrationConfigurationDetails>,无匹配时返回空列表; - 这些记录的 TTL 设得很高(1 天)。原因是:管理端 API 每做任何变更,都会通知缓存删除对应 key,该删除会通过扩展缓存的 backplane 传播到事件监听代码,缓存随即失效,下次读取时拉取新值。因此高 TTL 是安全的——只在必要时才刷新。
按集成打标签(Tagging per integration)
- 缓存中的每条条目(返回
List<OrganizationIntegrationConfigurationDetails>)都会打上组织 id 与集成类型的标签; - 这让管理员在集成层面做变更时,可以一次性移除某组织某集成的全部配置详情。
- 例如:某组织的 webhook 配置了 5 个事件,管理员改了集成层的 URL,变更就必须被传播,否则缓存会继续返回过期的 URL;
- 通过对每条条目打标签,API 可以一次调用请求扩展缓存移除某个组织集成的全部条目,缓存会以高性能方式处理这些条目的丢弃/刷新;
- 代码中有两处都知道标签机制,且必须保持同步:
EventIntegrationHandler在拉取相关配置详情时必须使用该 tag——这样缓存从仓储成功加载时会带上 tag 存储该条目;CreateOrganizationIntegrationCommand、UpdateOrganizationIntegrationCommand与DeleteOrganizationIntegrationCommand命令在管理员创建、更新、删除集成时必须使用该 tag 移除所有带 tag 的条目;- 为保证两处对“如何打 tag”的认知一致,它们都调用
EventIntegrationsCacheConstants.BuildCacheTagForOrganizationIntegration构建 tag。
源码印证:EventIntegrationHandler.cs #L135-L152 中正是用BuildCacheTagForOrganizationIntegration(organizationId, integrationType)生成 tag,并以FusionCacheEntryOptions(duration: DurationForOrganizationIntegrationConfigurationDetails)的高 TTL 选项执行cache.GetOrSetAsync,与文档描述完全一致。
9.3 模板属性缓存(Template Properties)
IntegrationTemplateProcessor支持一些需要额外查询的属性。例如UserId直接来自EventMessage,但UserName意味着要把 user id 映射到真实姓名的额外查询;User(含ActingUser)、Group与Organization的属性通过扩展缓存缓存,默认 TTL 30 分钟;- 这些属性同时缓存在 L1(内存)与 L2(Redis),并在需要时自动刷新。
十、构建一个新集成(Building a New Integration)
以下是构建新集成所需的全部部件。为便于命名说明,下文假设新集成名为 “Example”。完整示例可参考 Bitwarden 官方仓库中新增 Datadog 集成的 PR #6289(可检索该 PR 号查看提交上下文)。
10.1 IntegrationType
为IntegrationType枚举添加新集成的类型。
10.2 Configuration Models(配置模型)
配置模型决定OrganizationIntegration与OrganizationIntegrationConfiguration在数据库里存什么。Configuration列就是对应对象的序列化版本,代表该集成与事件类型的配置细节:
ExampleIntegration- 整个集成的配置细节(例如 Slack 的 token);
- 对该集成定义的每一个事件类型配置生效;
- 映射到
OrganizationIntegration的Configuration中存储的 JSON 结构。
ExampleIntegrationConfiguration- 可能逐事件变化的配置细节(例如 Slack 的 channelId);
- 映射到
OrganizationIntegrationConfiguration的Configuration中存储的 JSON 结构。
ExampleIntegrationConfigurationDetails- Integration 与 IntegrationConfiguration 的合并配置;
- 即
OrganizationIntegrationConfigurationDetails中MergedConfiguration的反序列化结果。
新增集成后,应在上文“现有集成及其两级配置明细”表格中为新集成添加一行。
10.3 Request Models(请求模型)
- 在
OrganizationIntegrationRequestModel.Validate的 switch 方法中新增分支——同时在OrganizationIntegrationRequestModelTests中添加测试; - 在
OrganizationIntegrationConfigurationRequestModel.IsValidForType的 switch 方法中新增分支——同时在OrganizationIntegrationConfigurationRequestModelTests中添加/更新测试。
10.4 Response Model(响应模型)
- 在
OrganizationIntegrationResponseModel.Status的 switch 方法中新增分支——同时在OrganizationIntegrationResponseModelTests中添加/更新测试。
10.5 Integration Handler(集成处理器)
例如ExampleIntegrationHandler:
- 这里存放执行集成动作的实际代码(即发出 HTTP 请求等);
- Handler 接收
IntegrationMessage<T>,其中<T>是上面定义的ExampleIntegrationConfigurationDetails,包含 Configuration 与已渲染好待发送的模板消息; - Handler 返回
IntegrationHandlerResult,携带请求结果细节——成功/失败、是否可重试、应延迟到何时等; - Handler 的职责范围仅仅是“执行集成并上报结果”。其余一切(重试次数、何时重试、失败如何处理)都由 Listener 负责。
10.6 GlobalSettings(全局设置)
RabbitMQ
添加该集成的队列名。它们通常带有默认值,这样首次被代码访问时 RabbitMQ 会自动创建它们:
ExampleEventQueueNameExampleIntegrationQueueNameExampleIntegrationRetryQueueName
Azure Service Bus
添加 ASB 使用的订阅名。与 RabbitMQ 类似,也提供默认值,避免必须写入 secrets 才能配置,同时允许覆盖。但是与 RabbitMQ 不同,这些订阅必须在代码访问它们之前就已存在,不会按需自动创建。参见下文“部署新集成”。
ExampleEventSubscriptionNameExampleIntegrationSubscriptionName
Service Bus Emulator 本地配置
为在本地创建 ASB 资源,还需要更新 dev/servicebusemulator_config.json 加入新订阅:
- 在现有事件 topic(
event-logging)下,为新集成的事件层添加订阅(events-example-subscription); - 在现有集成 topic(
event-integrations)下,为集成层消息添加新订阅(integration-example-subscription):- 从其他集成层订阅复制 correlation filter。它应基于
IntegrationType.ToRoutingKey过滤,本例即example。
- 从其他集成层订阅复制 correlation filter。它应基于
此处添加的名称必须与 secrets 中提供的值或 Global Settings 中给出的默认值一致。必须就位(并重启本地 ASB emulator)之后,才能使用任何本地访问 ASB 资源的代码。
10.7 ListenerConfiguration(监听器配置)
新集成需要一个ListenerConfiguration的子类,且该类还需实现IIntegrationListenerConfiguration。这个类提供访问前文在GlobalSettings中配置的 RabbitMQ 队列与 ASB 订阅的途径。新的 listener 配置将用于对 listener 进行类型化,并提供访问该集成所需配置的手段。
10.8 ServiceCollectionExtensions(服务装配)
在ServiceCollectionExtensions中,把上述所有部件串联起来,在每一消息层启动带 handler 的 listener。所有事件集成装配的核心方法是AddEventIntegrationServices。两个“添加 listener”的方法都会调用它,从而保证跨消息平台的公共依赖与集成只有唯一装配点。例如SlackIntegrationHandler需要SlackService,所以AddEventIntegrationServices里有AddSlackService的调用;webhook 的具名 HttpClient 定义同理。
在AddEventIntegrationServices中:
- 创建 handler 的单例:
services.TryAddSingleton<IIntegrationHandler<ExampleIntegrationConfigurationDetails>, ExampleIntegrationHandler>();- 创建 listener 配置:
var exampleConfiguration = new ExampleListenerConfiguration(globalSettings);- 把集成加入 RabbitMQ 与 ASB 各自的声明:
services.AddRabbitMqIntegration<ExampleIntegrationConfigurationDetails, ExampleListenerConfiguration>(exampleConfiguration);以及
services.AddAzureServiceBusIntegration<ExampleIntegrationConfigurationDetails, ExampleListenerConfiguration>(exampleConfiguration);以上三步与当前仓库中 AddEventIntegrationServices 的实际写法(注册 handler 单例 → 构造各 ListenerConfiguration → 按 ASB/RabbitMQ 分别调用Add*Integration泛型扩展)完全对应。
十一、部署新集成(Deploying a New Integration)
11.1 RabbitMQ
RabbitMQ 会在队列和交换器首次被代码访问时动态创建它们。因此部署新集成时无需手工创建队列。当然也可以提前创建和配置,但不是必须。注意:一旦创建完成,若需要修改任何配置,队列或交换器必须删除后重建。
11.2 Azure Service Bus
与 RabbitMQ 相反,ASB 资源必须在代码访问之前分配,不会按需创建。这意味着新集成所需的任何订阅,都必须在部署该代码之前在 ASB 中创建好。上文在 Global Settings 与servicebusemulator_config.json中定义的两个订阅,需要在部署代码前通过 Azure 门户或 CLI 为对应环境创建:
ExampleEventSubscriptionName- 这是从主事件 topic 扇出(fan-out)的订阅;
- 因此它一经声明就会开始接收所有事件;
- 这可能在“集成专属 handler 声明并部署”之前形成积压(backlog);
- 一种规避策略是:以假过滤器(例如
1 = 0)创建该订阅:- 订阅被创建,但过滤器确保没有任何消息真正落入该订阅;
- 可以部署引用该订阅的代码,因为订阅合法存在(只是为空);
- 代码就位、准备让新集成开始接收消息时,移除过滤器即可让订阅恢复接收所有 fan-out 消息。
ExampleIntegrationSubscriptionName- 该订阅必须在新集成代码部署前创建;
- 但它不是 fan-out,而是基于
IntegrationType.ToRoutingKey的过滤器; - 因此,在组织拥有激活的配置之前,它不会开始接收消息。这意味着提前声明它不会造成积压风险。
十二、小结:关键路径速查
| 关注点 | 关键类型/文件 |
|---|---|
| 事件发布入口 | IEventWriteService / EventIntegrationEventWriteService.cs |
| 平台选择与装配 | EventIntegrationsServiceCollectionExtensions.cs |
| 事件层处理器 | EventIntegrationHandler.cs、EventRepositoryHandler.cs |
| 集成层重试 | RabbitMqIntegrationListenerService.cs、AzureServiceBusIntegrationListenerService.cs |
| 消息与退避 | IntegrationMessage.cs |
| 重试/延迟参数 | GlobalSettings.cs(EventLogging.MaxRetries、EventLogging.RabbitMq.UseDelayPlugin) |
| 集成配置管理命令 | CreateOrganizationIntegrationCommand.cs、UpdateOrganizationIntegrationCommand.cs、DeleteOrganizationIntegrationCommand.cs |
| 本地 ASB 资源 | dev/servicebusemulator_config.json |
需要说明的适用前提:本架构中的 Azure Service Bus 路径对应云端部署,RabbitMQ 路径主要面向自托管实例;两者的资源创建语义(按需创建 vs 预先创建)是部署新集成时最需要注意的差异。所有队列/订阅默认值均可通过 Global Settings 或 secrets 覆盖,但修改已创建的 RabbitMQ 队列/交换器配置需要先删除再重建。
【免费下载链接】serverBitwarden infrastructure/backend (API, database, Docker, etc).项目地址: https://gitcode.com/GitHub_Trending/ser/server
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考