pub mod audit_logs;
pub mod channels;
pub mod cluster;
pub mod connectors;
pub mod helpers;
pub mod packages;
pub mod trace_dlq;
pub mod traces;
pub(crate) mod versioned;
pub mod workflows;
use std::sync::Arc;
use crate::errors::OrionError;
use audit_logs::AuditEvent;
pub struct Repositories {
pool: crate::storage::DbPool,
pub workflows: Arc<dyn workflows::WorkflowRepository>,
pub channels: Arc<dyn channels::ChannelRepository>,
pub connectors: Arc<dyn connectors::ConnectorRepository>,
pub traces: Arc<dyn traces::TraceRepository>,
pub audit_logs: Arc<dyn audit_logs::AuditLogRepository>,
pub trace_dlq: Arc<dyn trace_dlq::TraceDlqRepository>,
pub packages: Arc<dyn packages::PackageRepository>,
}
impl Repositories {
pub fn new(
pool: &crate::storage::DbPool,
storage: &crate::config::StorageConfig,
) -> Result<Self, crate::errors::OrionError> {
let cipher = if storage.connector_encryption_key.is_empty() {
None
} else {
Some(Arc::new(
crate::storage::config_encryption::ConfigCipher::from_hex(
&storage.connector_encryption_key,
)?,
))
};
Ok(Self {
pool: pool.clone(),
workflows: Arc::new(workflows::SqlWorkflowRepository::new(pool.clone())),
channels: Arc::new(channels::SqlChannelRepository::new(pool.clone())),
connectors: Arc::new(connectors::SqlConnectorRepository::with_cipher(
pool.clone(),
cipher,
)),
traces: Arc::new(traces::SqlTraceRepository::new(pool.clone())),
audit_logs: Arc::new(audit_logs::SqlAuditLogRepository::new(pool.clone())),
trace_dlq: Arc::new(trace_dlq::SqlTraceDlqRepository::new(pool.clone())),
packages: Arc::new(packages::SqlPackageRepository::new(pool.clone())),
})
}
pub async fn audited(&self, event: AuditEvent) -> Result<AuditedWrite<'_>, OrionError> {
Ok(AuditedWrite {
tx: self.pool.begin_write_tx().await?,
audit_logs: self.audit_logs.as_ref(),
event,
})
}
}
pub struct AuditedWrite<'a> {
tx: crate::storage::DbTransaction,
audit_logs: &'a dyn audit_logs::AuditLogRepository,
event: AuditEvent,
}
impl AuditedWrite<'_> {
pub fn tx(&mut self) -> &mut crate::storage::DbTransaction {
&mut self.tx
}
pub async fn commit(mut self) -> Result<(), OrionError> {
self.audit_logs.insert_tx(&mut self.tx, &self.event).await?;
self.tx.commit().await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::models::EntityStatus;
async fn repos() -> (crate::storage::DbPool, Repositories) {
let pool = crate::storage::test_sqlite_pool().await;
let repos = Repositories::new(&pool, &crate::config::StorageConfig::default())
.expect("repositories");
(pool, repos)
}
fn event(action: &str, resource_id: &str) -> AuditEvent {
AuditEvent {
principal: "tester".to_string(),
action: action.to_string(),
resource_type: "channel".to_string(),
resource_id: resource_id.to_string(),
details: None,
}
}
async fn seed_active_channel(repos: &Repositories, id: &str) {
let req = serde_json::from_value(serde_json::json!({
"channel_id": id,
"name": id,
"channel_type": "sync",
"protocol": "rest",
"route_pattern": format!("/{id}"),
"methods": ["POST"],
}))
.expect("request");
repos.channels.create(&req).await.expect("create");
repos.channels.activate(id).await.expect("activate");
}
async fn audit_rows(repos: &Repositories) -> i64 {
repos
.audit_logs
.list_paginated(&audit_logs::AuditLogFilter::default())
.await
.expect("list")
.total
}
#[tokio::test]
async fn a_committed_audited_write_lands_both_rows() {
let (_pool, repos) = repos().await;
seed_active_channel(&repos, "chan-commit").await;
assert_eq!(audit_rows(&repos).await, 0);
let mut write = repos
.audited(event("status_archived", "chan-commit"))
.await
.expect("begin");
let archived = repos
.channels
.archive_tx(write.tx(), "chan-commit")
.await
.expect("archive");
write.commit().await.expect("commit");
assert_eq!(archived.status, EntityStatus::Archived.as_str());
assert_eq!(
repos
.channels
.get_by_id("chan-commit")
.await
.expect("read back")
.status,
EntityStatus::Archived.as_str()
);
assert_eq!(
audit_rows(&repos).await,
1,
"the audit row must be committed with the change, not queued behind it"
);
}
#[tokio::test]
async fn dropping_an_audited_write_rolls_the_entity_write_back() {
let (_pool, repos) = repos().await;
seed_active_channel(&repos, "chan-rollback").await;
{
let mut write = repos
.audited(event("delete", "chan-rollback"))
.await
.expect("begin");
repos
.channels
.delete_tx(write.tx(), "chan-rollback")
.await
.expect("delete");
}
assert!(
repos.channels.get_by_id("chan-rollback").await.is_ok(),
"an audited write that was never committed must leave the channel in place"
);
assert_eq!(
audit_rows(&repos).await,
0,
"and must leave no audit row behind either"
);
}
}