use crate::bridge::envelope::{PhysicalPlan, Response};
use crate::control::change_stream::ChangeOperation;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId};
use nodedb_physical::physical_plan::{ClusterArrayOp, TimeseriesOp};
use super::extract::{cluster_array_change_meta, extract_write_metadata};
fn current_timestamp_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
fn is_timeseries_cdc_enabled(
shared: &SharedState,
database_id: DatabaseId,
tenant_id: TenantId,
collection: &str,
) -> bool {
let catalog = shared.credentials.catalog();
if let Ok(Some(coll)) = catalog.get_collection(database_id, tenant_id.as_u64(), collection)
&& coll.collection_type.is_timeseries()
{
if let Some(config) = coll.get_timeseries_config()
&& let Some(cdc_val) = config.get("cdc")
{
return cdc_val.as_str() == Some("true") || cdc_val.as_bool() == Some(true);
}
return false;
}
true
}
fn publish_change_event(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
is_columnar_collection: bool,
change_meta: (String, String, ChangeOperation),
lsn: nodedb_types::Lsn,
) {
let (collection, doc_id, op) = change_meta;
let should_publish = if is_columnar_collection {
is_timeseries_cdc_enabled(shared, database_id, tenant_id, &collection)
} else {
true
};
if !should_publish {
return;
}
use crate::control::change_stream::ChangeEvent;
let event = ChangeEvent {
lsn,
tenant_id,
collection,
document_id: doc_id,
operation: op,
timestamp_ms: current_timestamp_ms(),
after: None,
};
if let (Some(transport), Some(topology)) = (&shared.cluster_transport, &shared.cluster_topology)
{
use std::sync::atomic::Ordering;
static NOTIFY_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
let seq = NOTIFY_SEQ.fetch_add(1, Ordering::Relaxed);
crate::control::change_stream::broadcast_notify_to_cluster(
&event,
shared.node_id,
seq,
transport,
topology,
);
}
shared.change_stream.publish(event);
}
pub(crate) struct WriteChangeSet {
is_columnar_collection: bool,
metas: Vec<(String, String, ChangeOperation)>,
}
pub(crate) fn extract_write_change_set(plan: &PhysicalPlan, tenant_id: TenantId) -> WriteChangeSet {
WriteChangeSet {
is_columnar_collection: matches!(
plan,
PhysicalPlan::Columnar(_)
| PhysicalPlan::Timeseries(TimeseriesOp::Ingest { .. })
| PhysicalPlan::Timeseries(TimeseriesOp::Scan { .. })
),
metas: extract_write_metadata(plan, tenant_id),
}
}
pub(crate) fn publish_change_set_with_lsn(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
change_set: WriteChangeSet,
lsn: nodedb_types::Lsn,
) {
let WriteChangeSet {
is_columnar_collection,
metas,
} = change_set;
for meta in metas {
publish_change_event(
shared,
tenant_id,
database_id,
is_columnar_collection,
meta,
lsn,
);
}
}
pub(crate) fn publish_change_set(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
change_set: WriteChangeSet,
response: &Response,
) {
publish_change_set_with_lsn(
shared,
tenant_id,
database_id,
change_set,
response.watermark_lsn,
);
}
pub(crate) fn publish_origin_change_events(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
plan: &PhysicalPlan,
response: &Response,
) {
publish_change_set(
shared,
tenant_id,
database_id,
extract_write_change_set(plan, tenant_id),
response,
);
}
pub(crate) fn publish_cluster_array_change_events(
shared: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
op: &ClusterArrayOp,
lsn: u64,
) {
let change_set = WriteChangeSet {
is_columnar_collection: false,
metas: cluster_array_change_meta(op),
};
publish_change_set_with_lsn(
shared,
tenant_id,
database_id,
change_set,
nodedb_types::Lsn::new(lsn),
);
}