From 9b5fbb499d2491d26df6686144020f33e48593cb Mon Sep 17 00:00:00 2001 From: Erik Nilsen Date: Tue, 26 May 2026 16:17:33 -0700 Subject: [PATCH 1/4] refactor(gateway): introduce SourcePlugin trait for uniform source dispatch Eliminates the per-source iteration repeated between `http.rs::mount_sources` and `streams.rs::provision` by introducing a gateway-internal `SourcePlugin` trait. Each webhook source has a unit-struct plugin (GithubPlugin, SlackPlugin, ..., MicrosoftGraphPlugin) that owns both its JetStream provisioning loop and its HTTP route mounting loop, including any per-source edge cases. Two dispatch entry points -- `provision_webhook_sources` and `mount_webhook_sources` -- replace the repeated source-by-source dispatch blocks. Adding a new webhook source now touches one file (source_plugin.rs) instead of two. SlackPlugin's `mount` skips socket-mode-only integrations (no webhook config) since their HTTP route would never receive traffic; the runner is spawned in main.rs. Discord is intentionally excluded from `SourcePlugin` because its primary path is a WebSocket gateway runner, not a webhook receiver -- streams.rs still provisions it inline. All 579 gateway tests pass; full workspace tests pass. Signed-off-by: Erik Nilsen --- rsworkspace/crates/trogon-gateway/src/http.rs | 120 +---- rsworkspace/crates/trogon-gateway/src/main.rs | 2 + .../trogon-gateway/src/source_plugin.rs | 480 ++++++++++++++++++ .../crates/trogon-gateway/src/streams.rs | 84 +-- 4 files changed, 489 insertions(+), 197 deletions(-) create mode 100644 rsworkspace/crates/trogon-gateway/src/source_plugin.rs diff --git a/rsworkspace/crates/trogon-gateway/src/http.rs b/rsworkspace/crates/trogon-gateway/src/http.rs index 6e71e6e1a..7a5bf4bf8 100644 --- a/rsworkspace/crates/trogon-gateway/src/http.rs +++ b/rsworkspace/crates/trogon-gateway/src/http.rs @@ -1,15 +1,15 @@ use axum::Router; -use tracing::info; use trogon_nats::jetstream::{ClaimCheckPublisher, JetStreamPublisher, ObjectStorePut}; -use crate::config::{ResolvedConfig, SourceIntegration}; +use crate::config::ResolvedConfig; +use crate::source_plugin; pub(crate) fn mount_sources(config: ResolvedConfig, publisher: ClaimCheckPublisher) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - let mut app = Router::new() + let app = Router::new() .route( "/-/liveness", axum::routing::get(|| async { axum::http::StatusCode::OK }), @@ -19,119 +19,7 @@ where axum::routing::get(|| async { axum::http::StatusCode::OK }), ); - app = mount_webhook_integrations( - app, - "github", - "/sources/github", - &config.github, - publisher.clone(), - |p, cfg| crate::source::github::router(p, cfg), - ); - for integration in &config.slack { - if integration.config.webhook().is_none() { - continue; - } - let path = format!("/sources/slack/{}", integration.id); - app = app.nest( - &path, - crate::source::slack::router(publisher.clone(), &integration.config), - ); - let integration_id = integration.id.as_str(); - info!( - source = "slack", - integration = integration_id, - path, - "mounted source integration" - ); - } - app = mount_webhook_integrations( - app, - "telegram", - "/sources/telegram", - &config.telegram, - publisher.clone(), - |p, cfg| crate::source::telegram::router(p, cfg), - ); - app = mount_webhook_integrations( - app, - "twitter", - "/sources/twitter", - &config.twitter, - publisher.clone(), - |p, cfg| crate::source::twitter::router(p, cfg), - ); - app = mount_webhook_integrations( - app, - "gitlab", - "/sources/gitlab", - &config.gitlab, - publisher.clone(), - |p, cfg| crate::source::gitlab::router(p, cfg), - ); - app = mount_webhook_integrations( - app, - "incidentio", - "/sources/incidentio", - &config.incidentio, - publisher.clone(), - |p, cfg| crate::source::incidentio::router(p, cfg), - ); - app = mount_webhook_integrations( - app, - "linear", - "/sources/linear", - &config.linear, - publisher.clone(), - |p, cfg| crate::source::linear::router(p, cfg), - ); - app = mount_webhook_integrations( - app, - "microsoft-graph", - "/sources/microsoft-graph", - &config.microsoft_graph, - publisher.clone(), - |p, cfg| crate::source::microsoft_graph::router(p, cfg), - ); - app = mount_webhook_integrations( - app, - "notion", - "/sources/notion", - &config.notion, - publisher.clone(), - |p, cfg| crate::source::notion::router(p, cfg), - ); - app = mount_webhook_integrations(app, "sentry", "/sources/sentry", &config.sentry, publisher, |p, cfg| { - crate::source::sentry::router(p, cfg) - }); - - app -} - -fn mount_webhook_integrations( - mut app: Router, - source: &'static str, - source_path: &'static str, - integrations: &[SourceIntegration], - publisher: ClaimCheckPublisher, - router: F, -) -> Router -where - P: JetStreamPublisher, - S: ObjectStorePut, - F: Fn(ClaimCheckPublisher, &C) -> Router, -{ - for integration in integrations { - let path = format!("{}/{}", source_path, integration.id); - app = app.nest(&path, router(publisher.clone(), &integration.config)); - info!( - source, - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - - app + source_plugin::mount_webhook_sources(app, publisher, &config) } #[cfg(test)] diff --git a/rsworkspace/crates/trogon-gateway/src/main.rs b/rsworkspace/crates/trogon-gateway/src/main.rs index 732935ffa..70a4c0782 100644 --- a/rsworkspace/crates/trogon-gateway/src/main.rs +++ b/rsworkspace/crates/trogon-gateway/src/main.rs @@ -11,6 +11,8 @@ mod source; #[cfg_attr(coverage, allow(dead_code))] mod source_integration_id; #[cfg_attr(coverage, allow(dead_code))] +mod source_plugin; +#[cfg_attr(coverage, allow(dead_code))] mod source_status; #[cfg_attr(coverage, allow(dead_code))] mod streams; diff --git a/rsworkspace/crates/trogon-gateway/src/source_plugin.rs b/rsworkspace/crates/trogon-gateway/src/source_plugin.rs new file mode 100644 index 000000000..4c65d1ca9 --- /dev/null +++ b/rsworkspace/crates/trogon-gateway/src/source_plugin.rs @@ -0,0 +1,480 @@ +//! Source plugin trait. +//! +//! Each gateway-managed webhook source implements `SourcePlugin` so the +//! gateway can iterate sources without duplicating the per-source plumbing +//! in `http.rs` and `streams.rs`. The same plugin owns both provisioning +//! (JetStream streams) and HTTP route mounting for its integrations. +//! +//! Discord is intentionally NOT a `SourcePlugin`: its primary path is a +//! WebSocket gateway runner spawned in `main.rs`, not a webhook receiver. +//! Slack's socket-mode runners are spawned the same way — `SlackPlugin` +//! only mounts integrations that expose a webhook config. + +use axum::Router; +use tracing::info; +use trogon_nats::jetstream::{ClaimCheckPublisher, JetStreamContext, JetStreamPublisher, ObjectStorePut}; + +use crate::config::ResolvedConfig; + +pub type SourceId = &'static str; + +pub trait SourcePlugin: Send + Sync { + fn id(&self) -> SourceId; + + fn path_prefix(&self) -> String { + format!("/sources/{}", self.id()) + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error>; + + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut; +} + +pub struct GithubPlugin; +pub struct SlackPlugin; +pub struct TelegramPlugin; +pub struct TwitterPlugin; +pub struct GitlabPlugin; +pub struct IncidentioPlugin; +pub struct LinearPlugin; +pub struct MicrosoftGraphPlugin; +pub struct NotionPlugin; +pub struct SentryPlugin; + +impl SourcePlugin for GithubPlugin { + fn id(&self) -> SourceId { + "github" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.github { + crate::source::github::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.github { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::github::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for SlackPlugin { + fn id(&self) -> SourceId { + "slack" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.slack { + crate::source::slack::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + // Socket-mode-only integrations are spawned as long-running runners + // in `main.rs`; they have no HTTP route to mount. + for integration in &config.slack { + if integration.config.webhook().is_none() { + continue; + } + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::slack::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for TelegramPlugin { + fn id(&self) -> SourceId { + "telegram" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.telegram { + crate::source::telegram::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.telegram { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::telegram::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for TwitterPlugin { + fn id(&self) -> SourceId { + "twitter" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.twitter { + crate::source::twitter::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.twitter { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::twitter::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for GitlabPlugin { + fn id(&self) -> SourceId { + "gitlab" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.gitlab { + crate::source::gitlab::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.gitlab { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::gitlab::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for IncidentioPlugin { + fn id(&self) -> SourceId { + "incidentio" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.incidentio { + crate::source::incidentio::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.incidentio { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::incidentio::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for LinearPlugin { + fn id(&self) -> SourceId { + "linear" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.linear { + crate::source::linear::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.linear { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::linear::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for MicrosoftGraphPlugin { + fn id(&self) -> SourceId { + "microsoft-graph" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.microsoft_graph { + crate::source::microsoft_graph::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.microsoft_graph { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::microsoft_graph::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for NotionPlugin { + fn id(&self) -> SourceId { + "notion" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.notion { + crate::source::notion::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.notion { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::notion::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +impl SourcePlugin for SentryPlugin { + fn id(&self) -> SourceId { + "sentry" + } + + async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { + for integration in &config.sentry { + crate::source::sentry::provision(client, &integration.config).await?; + info!( + source = self.id(), + integration = integration.id.as_str(), + "stream provisioned" + ); + } + Ok(()) + } + + fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + where + P: JetStreamPublisher, + S: ObjectStorePut, + { + for integration in &config.sentry { + let path = format!("{}/{}", self.path_prefix(), integration.id); + app = app.nest( + &path, + crate::source::sentry::router(publisher.clone(), &integration.config), + ); + info!( + source = self.id(), + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app + } +} + +/// Provision JetStream streams for every webhook source. Discord is provisioned separately by the caller. +pub async fn provision_webhook_sources( + client: &C, + config: &ResolvedConfig, +) -> Result<(), C::Error> { + GithubPlugin.provision(client, config).await?; + SlackPlugin.provision(client, config).await?; + TelegramPlugin.provision(client, config).await?; + TwitterPlugin.provision(client, config).await?; + GitlabPlugin.provision(client, config).await?; + IncidentioPlugin.provision(client, config).await?; + LinearPlugin.provision(client, config).await?; + MicrosoftGraphPlugin.provision(client, config).await?; + NotionPlugin.provision(client, config).await?; + SentryPlugin.provision(client, config).await?; + Ok(()) +} + +/// Mount HTTP routes for every webhook source. +pub fn mount_webhook_sources( + mut app: Router, + publisher: ClaimCheckPublisher, + config: &ResolvedConfig, +) -> Router +where + P: JetStreamPublisher, + S: ObjectStorePut, +{ + app = GithubPlugin.mount(app, publisher.clone(), config); + app = SlackPlugin.mount(app, publisher.clone(), config); + app = TelegramPlugin.mount(app, publisher.clone(), config); + app = TwitterPlugin.mount(app, publisher.clone(), config); + app = GitlabPlugin.mount(app, publisher.clone(), config); + app = IncidentioPlugin.mount(app, publisher.clone(), config); + app = LinearPlugin.mount(app, publisher.clone(), config); + app = MicrosoftGraphPlugin.mount(app, publisher.clone(), config); + app = NotionPlugin.mount(app, publisher.clone(), config); + SentryPlugin.mount(app, publisher, config) +} diff --git a/rsworkspace/crates/trogon-gateway/src/streams.rs b/rsworkspace/crates/trogon-gateway/src/streams.rs index bc4c7014c..e7f2cc833 100644 --- a/rsworkspace/crates/trogon-gateway/src/streams.rs +++ b/rsworkspace/crates/trogon-gateway/src/streams.rs @@ -2,93 +2,15 @@ use tracing::info; use trogon_nats::jetstream::JetStreamContext; use crate::config::ResolvedConfig; +use crate::source_plugin; pub(crate) async fn provision(client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.github { - crate::source::github::provision(client, &integration.config).await?; - info!( - source = "github", - integration = integration.id.as_str(), - "stream provisioned" - ); - } + // Discord is gateway-WebSocket, not a webhook source; it doesn't fit `SourcePlugin`. if let Some(ref cfg) = config.discord { crate::source::discord::provision(client, cfg).await?; info!(source = "discord", "stream provisioned"); } - for integration in &config.slack { - crate::source::slack::provision(client, &integration.config).await?; - info!( - source = "slack", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.telegram { - crate::source::telegram::provision(client, &integration.config).await?; - info!( - source = "telegram", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.twitter { - crate::source::twitter::provision(client, &integration.config).await?; - info!( - source = "twitter", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.gitlab { - crate::source::gitlab::provision(client, &integration.config).await?; - info!( - source = "gitlab", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.incidentio { - crate::source::incidentio::provision(client, &integration.config).await?; - info!( - source = "incidentio", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.linear { - crate::source::linear::provision(client, &integration.config).await?; - info!( - source = "linear", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.microsoft_graph { - crate::source::microsoft_graph::provision(client, &integration.config).await?; - info!( - source = "microsoft-graph", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.notion { - crate::source::notion::provision(client, &integration.config).await?; - info!( - source = "notion", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - for integration in &config.sentry { - crate::source::sentry::provision(client, &integration.config).await?; - info!( - source = "sentry", - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + source_plugin::provision_webhook_sources(client, config).await } #[cfg(test)] From c046b08bc49f036a2bb47050542a5b793ce03f39 Mon Sep 17 00:00:00 2001 From: Erik Nilsen Date: Tue, 26 May 2026 16:44:38 -0700 Subject: [PATCH 2/4] refactor(gateway): extract provision/mount helpers in SourcePlugin MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Collapses the per-source iteration that was repeated in every non-Slack `SourcePlugin` impl into two shared helpers: - `provision_integrations` — `AsyncFn`-bounded helper iterating `Vec>` and calling each source's `provision` per integration. - `mount_integrations` — closure helper iterating the same and nesting each integration's router at `{path_prefix}/{id}`. Each non-Slack `provision` impl is now a single call to `provision_integrations(...)`; each non-Slack `mount` impl is a single call to `mount_integrations(...)`. The iteration / nesting / logging pattern lives in one place instead of nine. `SlackPlugin::mount` stays hand-written because its socket-mode webhook filter is a per-source behavior that the helper deliberately doesn't generalize. This addresses the Cursor Bugbot finding on the prior commit ("10 identical copies"). Restores the closure-helper pattern that existed on main (`mount_webhook_integrations`) and adds the symmetric helper for provisioning that main was missing. All 579 gateway tests pass. Signed-off-by: Erik Nilsen --- .../trogon-gateway/src/source_plugin.rs | 373 +++++++----------- 1 file changed, 147 insertions(+), 226 deletions(-) diff --git a/rsworkspace/crates/trogon-gateway/src/source_plugin.rs b/rsworkspace/crates/trogon-gateway/src/source_plugin.rs index 4c65d1ca9..6c500955a 100644 --- a/rsworkspace/crates/trogon-gateway/src/source_plugin.rs +++ b/rsworkspace/crates/trogon-gateway/src/source_plugin.rs @@ -14,7 +14,7 @@ use axum::Router; use tracing::info; use trogon_nats::jetstream::{ClaimCheckPublisher, JetStreamContext, JetStreamPublisher, ObjectStorePut}; -use crate::config::ResolvedConfig; +use crate::config::{ResolvedConfig, SourceIntegration}; pub type SourceId = &'static str; @@ -44,42 +44,71 @@ pub struct MicrosoftGraphPlugin; pub struct NotionPlugin; pub struct SentryPlugin; +async fn provision_integrations( + integrations: &[SourceIntegration], + source: SourceId, + client: &C, + provision_fn: F, +) -> Result<(), C::Error> +where + C: JetStreamContext, + F: AsyncFn(&C, &T) -> Result<(), C::Error>, +{ + for integration in integrations { + provision_fn(client, &integration.config).await?; + info!(source, integration = integration.id.as_str(), "stream provisioned"); + } + Ok(()) +} + +fn mount_integrations( + integrations: &[SourceIntegration], + mut app: Router, + publisher: ClaimCheckPublisher, + source: SourceId, + path_prefix: &str, + router_fn: F, +) -> Router +where + P: JetStreamPublisher, + S: ObjectStorePut, + F: Fn(ClaimCheckPublisher, &T) -> Router, +{ + for integration in integrations { + let path = format!("{}/{}", path_prefix, integration.id); + app = app.nest(&path, router_fn(publisher.clone(), &integration.config)); + info!( + source, + integration = integration.id.as_str(), + path, + "mounted source integration" + ); + } + app +} + impl SourcePlugin for GithubPlugin { fn id(&self) -> SourceId { "github" } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.github { - crate::source::github::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.github, self.id(), client, crate::source::github::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.github { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::github::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.github, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::github::router(p, cfg), + ) } } @@ -89,15 +118,7 @@ impl SourcePlugin for SlackPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.slack { - crate::source::slack::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.slack, self.id(), client, crate::source::slack::provision).await } fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router @@ -133,36 +154,22 @@ impl SourcePlugin for TelegramPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.telegram { - crate::source::telegram::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.telegram, self.id(), client, crate::source::telegram::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.telegram { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::telegram::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.telegram, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::telegram::router(p, cfg), + ) } } @@ -172,36 +179,22 @@ impl SourcePlugin for TwitterPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.twitter { - crate::source::twitter::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.twitter, self.id(), client, crate::source::twitter::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.twitter { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::twitter::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.twitter, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::twitter::router(p, cfg), + ) } } @@ -211,36 +204,22 @@ impl SourcePlugin for GitlabPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.gitlab { - crate::source::gitlab::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.gitlab, self.id(), client, crate::source::gitlab::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.gitlab { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::gitlab::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.gitlab, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::gitlab::router(p, cfg), + ) } } @@ -250,36 +229,28 @@ impl SourcePlugin for IncidentioPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.incidentio { - crate::source::incidentio::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations( + &config.incidentio, + self.id(), + client, + crate::source::incidentio::provision, + ) + .await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.incidentio { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::incidentio::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.incidentio, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::incidentio::router(p, cfg), + ) } } @@ -289,36 +260,22 @@ impl SourcePlugin for LinearPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.linear { - crate::source::linear::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.linear, self.id(), client, crate::source::linear::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.linear { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::linear::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.linear, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::linear::router(p, cfg), + ) } } @@ -328,36 +285,28 @@ impl SourcePlugin for MicrosoftGraphPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.microsoft_graph { - crate::source::microsoft_graph::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations( + &config.microsoft_graph, + self.id(), + client, + crate::source::microsoft_graph::provision, + ) + .await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.microsoft_graph { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::microsoft_graph::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.microsoft_graph, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::microsoft_graph::router(p, cfg), + ) } } @@ -367,36 +316,22 @@ impl SourcePlugin for NotionPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.notion { - crate::source::notion::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.notion, self.id(), client, crate::source::notion::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.notion { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::notion::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.notion, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::notion::router(p, cfg), + ) } } @@ -406,36 +341,22 @@ impl SourcePlugin for SentryPlugin { } async fn provision(&self, client: &C, config: &ResolvedConfig) -> Result<(), C::Error> { - for integration in &config.sentry { - crate::source::sentry::provision(client, &integration.config).await?; - info!( - source = self.id(), - integration = integration.id.as_str(), - "stream provisioned" - ); - } - Ok(()) + provision_integrations(&config.sentry, self.id(), client, crate::source::sentry::provision).await } - fn mount(&self, mut app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router + fn mount(&self, app: Router, publisher: ClaimCheckPublisher, config: &ResolvedConfig) -> Router where P: JetStreamPublisher, S: ObjectStorePut, { - for integration in &config.sentry { - let path = format!("{}/{}", self.path_prefix(), integration.id); - app = app.nest( - &path, - crate::source::sentry::router(publisher.clone(), &integration.config), - ); - info!( - source = self.id(), - integration = integration.id.as_str(), - path, - "mounted source integration" - ); - } - app + mount_integrations( + &config.sentry, + app, + publisher, + self.id(), + &self.path_prefix(), + |p, cfg| crate::source::sentry::router(p, cfg), + ) } } From 9bfac3283f222b788f7741a365a44436b001912a Mon Sep 17 00:00:00 2001 From: Erik Nilsen Date: Wed, 27 May 2026 12:14:02 -0700 Subject: [PATCH 3/4] ci: gate coverage publish on push events coverage-action with publish: true always tries to push to _xml_coverage_reports, but fork PRs get a read-only GITHUB_TOKEN regardless of the permissions block. Gate publish on push events so fork PRs skip the write step entirely. pull-requests: write is kept for PR annotations. Signed-off-by: Erik Nilsen --- .github/workflows/ci-rust.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci-rust.yml b/.github/workflows/ci-rust.yml index d8dcf7c29..5becc622c 100644 --- a/.github/workflows/ci-rust.yml +++ b/.github/workflows/ci-rust.yml @@ -73,7 +73,7 @@ jobs: path: rsworkspace/coverage.xml threshold: 95 fail: true - publish: true + publish: ${{ github.event_name == 'push' }} diff: true diff-branch: main diff-storage: _xml_coverage_reports From 79e57dd21a5b9dcffb9b3633f2243e3ce8d28843 Mon Sep 17 00:00:00 2001 From: Erik Nilsen Date: Wed, 27 May 2026 14:05:01 -0700 Subject: [PATCH 4/4] revert: restore publish: true in coverage action Reverts the conditional publish expression per yordis's review. Signed-off-by: Erik Nilsen --- .github/workflows/ci-rust.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci-rust.yml b/.github/workflows/ci-rust.yml index 5becc622c..d8dcf7c29 100644 --- a/.github/workflows/ci-rust.yml +++ b/.github/workflows/ci-rust.yml @@ -73,7 +73,7 @@ jobs: path: rsworkspace/coverage.xml threshold: 95 fail: true - publish: ${{ github.event_name == 'push' }} + publish: true diff: true diff-branch: main diff-storage: _xml_coverage_reports