use nodedb_types::{DatabaseId, TenantId};
use crate::bridge::envelope::{ErrorCode, PhysicalPlan, Response, Status};
use crate::control::maintenance::clone_materializer::{dispatch_local, read_all_source_rows};
use crate::control::state::SharedState;
use nodedb_physical::physical_plan::DocumentOp;
use nodedb_physical::physical_plan::document::merge_types::MergeClauseOp;
use super::resolve_arms::decode_resolve;
use crate::control::target_identity::{
assign_target_surrogate, bare_collection_name, resolve_target_pk,
};
const MAX_MERGE_RETRIES: u32 = 10;
pub struct MergeArgs<'a> {
pub tenant_id: TenantId,
pub database_id: DatabaseId,
pub target_collection: &'a str,
pub source_collection: &'a str,
pub source_alias: &'a str,
pub target_join_col: &'a str,
pub source_join_col: &'a str,
pub clauses: &'a [MergeClauseOp],
}
pub async fn run_merge(state: &SharedState, args: MergeArgs<'_>) -> crate::Result<Response> {
let catalog = state.credentials.catalog();
let target_bare = bare_collection_name(args.database_id, args.target_collection);
let target = catalog
.get_collection(args.database_id, args.tenant_id.as_u64(), &target_bare)?
.ok_or_else(|| crate::Error::CollectionNotFound {
tenant_id: args.tenant_id,
collection: args.target_collection.to_string(),
})?;
let target_pk = resolve_target_pk(&target, "MERGE")?;
let mut attempt: u32 = 0;
loop {
let source_rows = read_all_source_rows(
state,
args.tenant_id,
args.database_id,
args.source_collection,
None,
)
.await?;
let resolve_plan = merge_plan(&args, true, None, Some(source_rows.clone()));
let resolve_resp = dispatch_local(
state,
args.tenant_id,
args.database_id,
args.target_collection,
resolve_plan,
None,
)
.await?;
if resolve_resp.status != Status::Ok {
return Ok(resolve_resp);
}
let insert_rows = decode_resolve(&resolve_resp.payload)?.inserts;
let mut resolved: Vec<(String, u32)> = Vec::with_capacity(insert_rows.len());
for (join_key, body) in &insert_rows {
let surrogate = assign_target_surrogate(
state,
args.database_id,
args.tenant_id,
args.target_collection,
&target_pk,
body,
)?;
resolved.push((join_key.clone(), surrogate.as_u32()));
}
let apply_plan = merge_plan(&args, false, Some(resolved), Some(source_rows));
let apply_resp = dispatch_local(
state,
args.tenant_id,
args.database_id,
args.target_collection,
apply_plan,
None,
)
.await?;
if apply_resp.error_code.as_deref() == Some(&ErrorCode::OllpRetryRequired) {
attempt += 1;
if attempt > MAX_MERGE_RETRIES {
return Err(crate::Error::OllpExhausted {
retries: MAX_MERGE_RETRIES.min(u8::MAX as u32) as u8,
});
}
continue;
}
crate::control::server::wal_dispatch::mint_dispatch_local_redo(
&state.wal,
args.tenant_id,
args.database_id,
args.target_collection,
&apply_resp,
)?;
return Ok(apply_resp);
}
}
fn merge_plan(
args: &MergeArgs<'_>,
resolve_only: bool,
resolved_inserts: Option<Vec<(String, u32)>>,
source_rows: Option<Vec<(String, Vec<u8>)>>,
) -> PhysicalPlan {
PhysicalPlan::Document(DocumentOp::Merge {
target_collection: args.target_collection.to_string(),
source_collection: args.source_collection.to_string(),
source_alias: args.source_alias.to_string(),
target_join_col: args.target_join_col.to_string(),
source_join_col: args.source_join_col.to_string(),
clauses: args.clauses.to_vec(),
returning: None,
resolve_only,
resolved_inserts,
source_rows,
})
}