use std::collections::HashSet;
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::doc_format;
use crate::engine::document::store::doc_id_to_surrogate;
use nodedb_physical::physical_plan::document::merge_types::{
MergeActionOp, MergeClauseKind as MergeClauseKindOp,
};
use nodedb_types::Surrogate;
use super::super::merge::MergeParams;
use super::super::merge_helpers::{
build_insert_doc, build_merged, build_update_doc, find_arm, json_to_str,
};
pub(super) struct MergeUpdate {
pub(super) doc_id: String,
pub(super) surrogate: Option<Surrogate>,
pub(super) body: Vec<u8>,
}
pub(super) struct MergeDelete {
pub(super) doc_id: String,
pub(super) surrogate: Option<Surrogate>,
pub(super) body: Vec<u8>,
}
pub(super) struct MergeInsert {
pub(super) join_key: String,
pub(super) body: Vec<u8>,
}
pub(super) struct MergePlanActions {
pub(super) updates: Vec<MergeUpdate>,
pub(super) deletes: Vec<MergeDelete>,
pub(super) inserts: Vec<MergeInsert>,
}
fn encode_doc_body(doc: &serde_json::Value) -> Vec<u8> {
let value: nodedb_types::Value = doc.clone().into();
nodedb_types::value_to_msgpack(&value).unwrap_or_else(|_| doc_format::encode_to_msgpack(doc))
}
fn decode_target(
bytes: &[u8],
strict_schema: &Option<nodedb_types::columnar::StrictSchema>,
) -> Option<serde_json::Value> {
if let Some(schema) = strict_schema {
crate::data::executor::strict_format::binary_tuple_to_json(bytes, schema)
} else {
doc_format::decode_document(bytes)
}
}
impl CoreLoop {
pub(super) fn collect_merge_plan(
&self,
database_id: u64,
tid: u64,
txn_id: Option<crate::types::TxnId>,
params: &MergeParams<'_>,
) -> crate::Result<MergePlanActions> {
let source_map = self.build_merge_source_map(
database_id,
tid,
params.source_collection,
params.source_join_col,
params.source_rows,
)?;
let strict_schema = self.merge_strict_schema(database_id, tid, params.target_collection);
let target_docs =
self.collect_target_docs(database_id, tid, params.target_collection, txn_id)?;
let mut updates: Vec<MergeUpdate> = Vec::new();
let mut deletes: Vec<MergeDelete> = Vec::new();
let mut matched_source_keys: HashSet<String> = HashSet::new();
let null_source = serde_json::Value::Null;
for (doc_id, bytes) in &target_docs {
let target_doc = match decode_target(bytes, &strict_schema) {
Some(v) => v,
None => continue,
};
let join_val = target_doc
.get(params.target_join_col)
.map(json_to_str)
.unwrap_or_default();
let surrogate = doc_id_to_surrogate(doc_id);
let (arm_kind, source_doc): (MergeClauseKindOp, &serde_json::Value) =
if let Some(source_doc) = source_map.get(&join_val) {
matched_source_keys.insert(join_val.clone());
(MergeClauseKindOp::Matched, source_doc)
} else {
(MergeClauseKindOp::NotMatchedBySource, &null_source)
};
let context = if arm_kind == MergeClauseKindOp::Matched {
build_merged(&target_doc, source_doc, params.source_alias)
} else {
target_doc.clone()
};
if let Some(arm) = find_arm(params.clauses, arm_kind, &context) {
match &arm.action {
MergeActionOp::Update { updates: upd } => {
let updated =
build_update_doc(&target_doc, source_doc, params.source_alias, upd);
updates.push(MergeUpdate {
doc_id: doc_id.clone(),
surrogate,
body: encode_doc_body(&updated),
});
}
MergeActionOp::Delete => deletes.push(MergeDelete {
doc_id: doc_id.clone(),
surrogate,
body: encode_doc_body(&target_doc),
}),
MergeActionOp::Insert { .. } | MergeActionOp::DoNothing => {}
}
}
}
let mut inserts: Vec<MergeInsert> = Vec::new();
for (src_key, src_doc) in &source_map {
if matched_source_keys.contains(src_key.as_str()) {
continue;
}
if let Some(arm) = find_arm(params.clauses, MergeClauseKindOp::NotMatched, src_doc)
&& let MergeActionOp::Insert { columns, values } = &arm.action
{
let body = encode_doc_body(&build_insert_doc(
columns,
values,
src_doc,
params.source_alias,
));
inserts.push(MergeInsert {
join_key: src_key.clone(),
body,
});
}
}
Ok(MergePlanActions {
updates,
deletes,
inserts,
})
}
}