#![recursion_limit = "1024"]
#![allow(unused_imports)]
pub mod domain;
pub mod infrastructure;
pub mod application;
pub mod presentation;
pub mod seeders;
pub mod exports;
pub use application::service::IntegrationsWriteService;
pub use application::service::{TargetPort, MapRequest, MapOutcome, MapRejected, MappedRef};
pub use application::service::{InboundEvent, ReceiveOutcome, NewConnector, FailedEvent, IntegrationError};
pub use application::service::{IntegrationEvent, IntegrationEventMapped, IntegrationEventSink, LoggingSink};
pub use application::service::{OAuthCredentialFailure, OAuthCredentialStore, PURPOSE_OAUTH_TOKEN, TokenBundle, TokenMetadata};
pub use application::service::{
AccountStatus, AuthorizeRequest, AuthorizeResponse, CompleteOutcome, CompleteRequest,
IntegrationsOauthConfig, IntegrationsOauthService, OauthError, RefreshSummary,
STATE_TTL_SECONDS,
};
pub use infrastructure::http::{OAuthTransport, ReqwestOAuthTransport};
pub use infrastructure::persistence::*;
pub use application::service::IntegrationConnectorService;
pub use application::service::IntegrationAccountService;
pub use application::service::IntegrationEventService;
pub use application::workflows::*;
use std::sync::Arc;
use axum::Router;
use sqlx::PgPool;
pub struct IntegrationsModule {
pub(crate) integration_connector_service: Arc<IntegrationConnectorService>,
pub(crate) integration_account_service: Arc<IntegrationAccountService>,
pub(crate) integration_event_service: Arc<IntegrationEventService>,
pub integrations_write_service: Arc<IntegrationsWriteService>,
pub integrations_oauth: Option<Arc<IntegrationsOauthService>>,
}
impl IntegrationsModule {
pub fn builder() -> IntegrationsModuleBuilder {
IntegrationsModuleBuilder::new()
}
pub fn all_crud_routes(&self) -> Router {
use presentation::http::{
create_integration_connector_routes,
create_integration_event_routes,
};
Router::new()
.merge(create_integration_connector_routes(self.integration_connector_service.clone()))
.merge(create_integration_event_routes(self.integration_event_service.clone()))
}
#[deprecated(note = "mounts unvalidated generic CRUD; prefer readonly_routes() + validated writes, or all_crud_routes() for the full/unguarded surface")]
pub fn routes(&self) -> Router {
self.all_crud_routes()
}
pub fn readonly_routes(&self) -> Router {
use presentation::http::{
create_integration_connector_read_routes,
create_integration_event_read_routes,
};
Router::new()
.merge(create_integration_connector_read_routes(self.integration_connector_service.clone()))
.merge(create_integration_event_read_routes(self.integration_event_service.clone()))
}
pub fn guarded_routes(&self) -> Router {
use presentation::http::{
create_integration_connector_routes,
create_integration_event_read_routes,
};
Router::new()
.merge(create_integration_connector_routes(self.integration_connector_service.clone()))
.merge(create_integration_event_read_routes(self.integration_event_service.clone()))
}
pub fn oauth_routes(&self) -> Router {
use presentation::http::create_oauth_routes;
match &self.integrations_oauth {
Some(service) => create_oauth_routes(service.clone()),
None => Router::new(),
}
}
}
pub struct IntegrationsModuleBuilder {
db_pool: Option<PgPool>,
oauth_config: Option<IntegrationsOauthConfig>,
oauth_transport: Option<Arc<dyn OAuthTransport>>,
oauth_store: Option<Arc<dyn OAuthCredentialStore>>,
}
impl IntegrationsModuleBuilder {
pub fn new() -> Self {
Self {
db_pool: None,
oauth_config: None,
oauth_transport: None,
oauth_store: None,
}
}
pub fn with_database(mut self, pool: PgPool) -> Self {
self.db_pool = Some(pool);
self
}
pub fn with_oauth(mut self, config: IntegrationsOauthConfig) -> Self {
self.oauth_config = Some(config);
self
}
pub fn with_oauth_transport(mut self, transport: Arc<dyn OAuthTransport>) -> Self {
self.oauth_transport = Some(transport);
self
}
pub fn with_oauth_store(mut self, store: Arc<dyn OAuthCredentialStore>) -> Self {
self.oauth_store = Some(store);
self
}
pub fn build(self) -> anyhow::Result<IntegrationsModule> {
let db_pool = self.db_pool
.ok_or_else(|| anyhow::anyhow!("Database pool not configured"))?;
let integration_connector_repository = Arc::new(IntegrationConnectorRepository::new(db_pool.clone()));
let integration_connector_service = Arc::new(IntegrationConnectorService::with_repository(integration_connector_repository.clone()));
let integration_account_repository = Arc::new(IntegrationAccountRepository::new(db_pool.clone()));
let integration_account_service = Arc::new(IntegrationAccountService::with_repository(integration_account_repository.clone()));
let integration_event_repository = Arc::new(IntegrationEventRepository::new(db_pool.clone()));
let integration_event_service = Arc::new(IntegrationEventService::with_repository(integration_event_repository.clone()));
let integrations_write_service = Arc::new(IntegrationsWriteService::new(db_pool.clone()));
let integrations_oauth = match (self.oauth_config, self.oauth_transport, self.oauth_store) {
(Some(config), transport, Some(store)) => {
let transport = match transport {
Some(t) => t,
None => Arc::new(ReqwestOAuthTransport::new()
.map_err(|e| anyhow::anyhow!("building the default OAuth transport: {e}"))?),
};
Some(Arc::new(IntegrationsOauthService::build(
db_pool.clone(),
config,
transport,
store,
)?))
}
(Some(_), _, None) => {
return Err(anyhow::anyhow!(
"with_oauth needs the credential store bound through the port: with_oauth_store(..)"
))
}
(None, _, _) => None,
};
Ok(IntegrationsModule {
integration_connector_service,
integration_account_service,
integration_event_service,
integrations_write_service,
integrations_oauth,
})
}
}
impl Default for IntegrationsModuleBuilder {
fn default() -> Self {
Self::new()
}
}