use std::time::{Duration, Instant};
use async_trait::async_trait;
use mongodb_driver::bson::{self, Document, doc};
use mongodb_driver::options::ReturnDocument;
use mongodb_driver::{Collection, Database, IndexModel};
use super::{CanonicalStore, DurabilityToken};
const OUTBOX_SEQ_ID: &str = "outbox_seq";
const LEASE_COLLECTION: &str = "udb_advisory_leases";
pub struct MongoDbCanonicalStore {
pub(super) db: Database,
pub(super) instance_name: String,
pub(super) outbox_collection: String,
}
pub(super) const PROJECTION_COLLECTION: &str = "udb_projection_tasks";
pub(super) const SAGA_COLLECTION: &str = "udb_sagas";
pub(super) const ADMIN_AUDIT_COLLECTION: &str = "udb_admin_audit_log";
pub(super) const ADMIN_AUDIT_LOCK_COLLECTION: &str = "udb_admin_audit_lock";
pub(super) const MIGRATION_RUNS_COLLECTION: &str = "udb_migration_runs";
pub(super) const MIGRATION_LEDGER_COLLECTION: &str = "udb_migration_op_ledger";
pub(super) const MIGRATION_OP_SEQ_ID: &str = "migration_op_seq";
impl MongoDbCanonicalStore {
pub fn new(
db: Database,
instance_name: impl Into<String>,
outbox_collection: impl Into<String>,
) -> Self {
Self {
db,
instance_name: instance_name.into(),
outbox_collection: outbox_collection.into(),
}
}
pub fn from_executor(
executor: &crate::runtime::executors::mongodb::MongoDbExecutor,
instance_name: impl Into<String>,
outbox_collection: impl Into<String>,
) -> Option<Self> {
executor
.native_database()
.map(|db| Self::new(db, instance_name, outbox_collection))
}
pub(super) fn db(&self) -> &Database {
&self.db
}
fn outbox(&self) -> Collection<Document> {
self.db.collection::<Document>(&self.outbox_collection)
}
fn leases(&self) -> Collection<Document> {
self.db.collection::<Document>(LEASE_COLLECTION)
}
async fn current_seq(&self) -> Result<i64, String> {
let counter = self
.outbox()
.find_one(doc! { "_id": OUTBOX_SEQ_ID })
.await
.map_err(|err| format!("mongodb outbox counter read failed: {err}"))?;
Ok(counter
.as_ref()
.and_then(|doc| doc.get_i64("seq").ok())
.unwrap_or(0))
}
pub(super) fn is_namespace_exists(err: &mongodb_driver::error::Error) -> bool {
let text = err.to_string();
text.contains("NamespaceExists") || text.contains("already exists")
}
}
#[async_trait]
impl CanonicalStore for MongoDbCanonicalStore {
fn backend_label(&self) -> &'static str {
"mongodb"
}
fn instance_name(&self) -> &str {
&self.instance_name
}
async fn current_durability_token(&self) -> Result<DurabilityToken, String> {
let seq = self.current_seq().await?;
Ok(DurabilityToken::new("mongodb", seq.to_string()))
}
async fn wait_for_token(
&self,
token: &DurabilityToken,
timeout: Duration,
) -> Result<bool, String> {
if !token.is_for("mongodb") {
return Err(format!(
"MongoDbCanonicalStore cannot wait on a '{}' token",
token.backend_label
));
}
let target: i64 = token
.value
.parse()
.map_err(|_| format!("malformed mongodb durability token '{}'", token.value))?;
let started = Instant::now();
let poll = super::durability_poll_interval(timeout, super::MONGODB_DURABILITY_POLL_MS);
loop {
if self.current_seq().await? >= target {
return Ok(true);
}
if started.elapsed() >= timeout {
return Ok(false);
}
tokio::time::sleep(poll).await;
}
}
async fn enqueue_outbox_event(
&self,
event_id: &str,
topic: &str,
partition_key: &str,
payload: &serde_json::Value,
) -> Result<i64, String> {
let counter = self
.outbox()
.find_one_and_update(
doc! { "_id": OUTBOX_SEQ_ID },
doc! { "$inc": { "seq": 1_i64 } },
)
.upsert(true)
.return_document(ReturnDocument::After)
.await
.map_err(|err| format!("mongodb outbox counter increment failed: {err}"))?;
let seq = counter
.as_ref()
.and_then(|doc| doc.get_i64("seq").ok())
.ok_or_else(|| "mongodb outbox counter returned no seq".to_string())?;
let payload_bson = bson::to_bson(payload)
.map_err(|err| format!("mongodb outbox payload encode failed: {err}"))?;
let created_at = bson::DateTime::from_millis(chrono::Utc::now().timestamp_millis());
let event = doc! {
"_id": event_id,
"event_seq": seq,
"topic": topic,
"partition_key": partition_key,
"payload": payload_bson,
"created_at": created_at,
};
self.outbox()
.insert_one(event)
.await
.map_err(|err| format!("mongodb outbox insert failed: {err}"))?;
Ok(seq)
}
async fn outbox_max_seq(&self) -> Result<i64, String> {
self.current_seq().await
}
async fn ensure_system_tables(&self) -> Result<(), String> {
match self.db.create_collection(&self.outbox_collection).await {
Ok(_) => {}
Err(err) if Self::is_namespace_exists(&err) => {}
Err(err) => {
return Err(format!(
"mongodb ensure_system_tables createCollection '{}' failed: {err}",
self.outbox_collection
));
}
}
self.outbox()
.find_one_and_update(
doc! { "_id": OUTBOX_SEQ_ID },
doc! { "$setOnInsert": { "seq": 0_i64 } },
)
.upsert(true)
.return_document(ReturnDocument::After)
.await
.map_err(|err| format!("mongodb ensure_system_tables counter init failed: {err}"))?;
Ok(())
}
async fn ensure_advisory_lease_table(&self) -> Result<(), String> {
match self.db.create_collection(LEASE_COLLECTION).await {
Ok(_) => {}
Err(err) if Self::is_namespace_exists(&err) => {}
Err(err) => {
return Err(format!(
"mongodb ensure_advisory_lease_table createCollection failed: {err}"
));
}
}
let index = IndexModel::builder()
.keys(doc! { "lease_name": 1 })
.options(
mongodb_driver::options::IndexOptions::builder()
.unique(true)
.build(),
)
.build();
self.leases()
.create_index(index)
.await
.map_err(|err| format!("mongodb advisory lease unique index failed: {err}"))?;
Ok(())
}
async fn try_acquire_advisory_lease(
&self,
lease_name: &str,
owner_id: &str,
ttl: std::time::Duration,
) -> Result<bool, String> {
let now = chrono::Utc::now();
let now_bson = bson::DateTime::from_millis(now.timestamp_millis());
let expires_at = bson::DateTime::from_millis(
(now + chrono::Duration::from_std(ttl).unwrap_or_default()).timestamp_millis(),
);
let filter = doc! {
"lease_name": lease_name,
"$or": [
{ "expires_at": { "$lte": now_bson } },
{ "owner_id": owner_id },
],
};
let update = doc! {
"$set": { "owner_id": owner_id, "expires_at": expires_at },
"$setOnInsert": { "lease_name": lease_name },
};
let outcome = self
.leases()
.find_one_and_update(filter, update)
.upsert(true)
.return_document(ReturnDocument::After)
.await;
match outcome {
Ok(Some(doc)) => Ok(doc
.get_str("owner_id")
.map(|o| o == owner_id)
.unwrap_or(false)),
Ok(None) => Ok(false),
Err(err) if Self::is_duplicate_key(&err) => {
Ok(false)
}
Err(err) => Err(format!("mongodb try_acquire_advisory_lease failed: {err}")),
}
}
async fn release_advisory_lease(&self, lease_name: &str, owner_id: &str) -> Result<(), String> {
self.leases()
.delete_one(doc! { "lease_name": lease_name, "owner_id": owner_id })
.await
.map_err(|err| format!("mongodb release_advisory_lease failed: {err}"))?;
Ok(())
}
}
impl MongoDbCanonicalStore {
pub(super) fn is_duplicate_key(err: &mongodb_driver::error::Error) -> bool {
let text = err.to_string();
text.contains("E11000") || text.contains("duplicate key")
}
}