use crate::{
domain::{
entities::{SchemaEnforcement, Tenant, TenantQuotas, TenantUsage, UsageMeter},
value_objects::TenantId,
},
error::Result,
};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
#[async_trait]
pub trait TenantRepository: Send + Sync {
async fn create(&self, id: TenantId, name: String, quotas: TenantQuotas) -> Result<Tenant>;
async fn create_initialized(&self, _tenant: Tenant) -> Result<Tenant> {
Err(crate::error::AllSourceError::InternalError(
"Atomic tenant initialization is unavailable".into(),
))
}
async fn save(&self, tenant: &Tenant) -> Result<()>;
async fn find_by_id(&self, id: &TenantId) -> Result<Option<Tenant>>;
async fn find_by_name(&self, name: &str) -> Result<Option<Tenant>>;
async fn find_all(&self, limit: usize, offset: usize) -> Result<Vec<Tenant>>;
async fn find_active(&self, limit: usize, offset: usize) -> Result<Vec<Tenant>>;
async fn count(&self) -> Result<usize>;
async fn count_active(&self) -> Result<usize>;
async fn delete(&self, id: &TenantId) -> Result<bool>;
async fn update_quotas(&self, id: &TenantId, quotas: TenantQuotas) -> Result<bool>;
async fn update_schema_enforcement(
&self,
id: &TenantId,
mode: SchemaEnforcement,
) -> Result<bool>;
async fn update_usage(&self, id: &TenantId, usage: TenantUsage) -> Result<bool>;
async fn increment_usage(
&self,
id: &TenantId,
meter: UsageMeter,
count: u64,
) -> Result<Option<u64>>;
async fn admit_query_usage(
&self,
_id: &TenantId,
_request: crate::domain::entities::query_usage::QueryUsageRequest,
) -> Result<Option<crate::domain::entities::query_usage::QueryUsageDecision>> {
Err(crate::error::AllSourceError::InternalError(
"Durable query admission is unavailable".into(),
))
}
async fn get_query_usage(
&self,
_id: &TenantId,
) -> Result<Option<crate::domain::entities::query_usage::QueryUsageSnapshot>> {
Err(crate::error::AllSourceError::InternalError(
"Durable query admission is unavailable".into(),
))
}
async fn reset_query_usage(
&self,
_id: &TenantId,
_request: crate::domain::entities::query_usage::QueryUsageReset,
) -> Result<Option<crate::domain::entities::query_usage::QueryUsageResetDecision>> {
Err(crate::error::AllSourceError::InternalError(
"Durable query reset is unavailable".into(),
))
}
async fn activate(&self, id: &TenantId) -> Result<bool>;
async fn deactivate(&self, id: &TenantId) -> Result<bool>;
async fn exists(&self, id: &TenantId) -> Result<bool> {
Ok(self.find_by_id(id).await?.is_some())
}
async fn is_active(&self, id: &TenantId) -> Result<bool> {
match self.find_by_id(id).await? {
Some(tenant) => Ok(tenant.is_active()),
None => Ok(false),
}
}
async fn merge_metadata(
&self,
id: &TenantId,
partial: serde_json::Value,
) -> Result<Option<serde_json::Value>> {
let Some(mut tenant) = self.find_by_id(id).await? else {
return Ok(None);
};
let mut metadata = tenant.metadata().clone();
deep_merge_metadata(&mut metadata, partial);
tenant.update_metadata(metadata.clone());
self.save(&tenant).await?;
Ok(Some(metadata))
}
}
pub fn deep_merge_metadata(target: &mut serde_json::Value, patch: serde_json::Value) {
match patch {
serde_json::Value::Object(patch_obj) => {
if !target.is_object() {
*target = serde_json::Value::Object(serde_json::Map::new());
}
let target_obj = target
.as_object_mut()
.expect("target coerced to object above");
for (key, value) in patch_obj {
deep_merge_metadata(
target_obj.entry(key).or_insert(serde_json::Value::Null),
value,
);
}
}
other => *target = other,
}
}
#[derive(Debug, Clone, Default)]
pub struct TenantQuery {
pub active_only: bool,
pub name_contains: Option<String>,
pub created_after: Option<DateTime<Utc>>,
pub created_before: Option<DateTime<Utc>>,
pub limit: Option<usize>,
pub offset: Option<usize>,
}
impl TenantQuery {
pub fn new() -> Self {
Self::default()
}
pub fn active_only(mut self) -> Self {
self.active_only = true;
self
}
pub fn with_name_filter(mut self, name: String) -> Self {
self.name_contains = Some(name);
self
}
pub fn created_after(mut self, date: DateTime<Utc>) -> Self {
self.created_after = Some(date);
self
}
pub fn created_before(mut self, date: DateTime<Utc>) -> Self {
self.created_before = Some(date);
self
}
pub fn with_pagination(mut self, limit: usize, offset: usize) -> Self {
self.limit = Some(limit);
self.offset = Some(offset);
self
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_tenant_query_builder() {
let query = TenantQuery::new()
.active_only()
.with_name_filter("acme".to_string())
.with_pagination(10, 0);
assert!(query.active_only);
assert_eq!(query.name_contains, Some("acme".to_string()));
assert_eq!(query.limit, Some(10));
assert_eq!(query.offset, Some(0));
}
#[test]
fn test_tenant_query_with_dates() {
let now = Utc::now();
let yesterday = now - chrono::Duration::days(1);
let query = TenantQuery::new()
.created_after(yesterday)
.created_before(now);
assert!(query.created_after.is_some());
assert!(query.created_before.is_some());
}
#[test]
fn deep_merge_preserves_siblings_and_recurses() {
let mut target = serde_json::json!({
"quotas": { "events_used": 142, "queries_used": 7 },
"subscription": { "tier": "studio" }
});
deep_merge_metadata(
&mut target,
serde_json::json!({ "projections": { "enabled": ["event-count"] } }),
);
assert_eq!(target["quotas"]["events_used"], 142);
assert_eq!(target["subscription"]["tier"], "studio");
assert_eq!(target["projections"]["enabled"][0], "event-count");
deep_merge_metadata(
&mut target,
serde_json::json!({ "quotas": { "events_used": 150 } }),
);
assert_eq!(target["quotas"]["events_used"], 150);
assert_eq!(target["quotas"]["queries_used"], 7);
}
#[test]
fn deep_merge_arrays_and_scalars_replace() {
let mut target = serde_json::json!({ "enabled": ["a", "b"], "n": 1 });
deep_merge_metadata(&mut target, serde_json::json!({ "enabled": ["c"], "n": 2 }));
assert_eq!(target["enabled"], serde_json::json!(["c"]));
assert_eq!(target["n"], 2);
}
#[test]
fn deep_merge_object_over_scalar_coerces() {
let mut target = serde_json::json!({ "projections": "stale" });
deep_merge_metadata(
&mut target,
serde_json::json!({ "projections": { "enabled": [] } }),
);
assert!(target["projections"].is_object());
assert_eq!(target["projections"]["enabled"], serde_json::json!([]));
}
}