use super::{
Any, BTreeMap, BlockNode, EDITABLE_MARKS, HashMap, HashSet, Map, MergeError, NodeId, ObservedValue, ParseError,
ProjectedNode, SourceNode, State, TextMerge, TreePlan, Value, WorkBudget, collect_child_ids, failure, flatten,
flatten_source, get_string, insert_block_tree, merge_properties, merge_text, observe, patch_children, semantic_equal,
};
pub(super) fn execute(
baseline: &State,
current: &State,
before: &[ProjectedNode],
incoming: &[ProjectedNode],
source: &[SourceNode],
plan: &TreePlan,
budget: &mut WorkBudget,
) -> Result<HashMap<NodeId, String>, ParseError> {
let mut old_nodes = Vec::new();
flatten(before, &mut old_nodes);
let old: HashMap<_, _> = old_nodes.iter().map(|n| (n.id.as_deref().unwrap(), *n)).collect();
let mut targets = Vec::new();
flatten(incoming, &mut targets);
let mut source_nodes = Vec::new();
flatten_source(source, &mut source_nodes);
let snapshots = snapshot_blocks(¤t.pool, budget)?;
let mut allowed: HashMap<String, HashSet<String>> = HashMap::new();
let mut ids = HashMap::new();
let mut blocks = current.doc.get_map("blocks")?;
let live: HashSet<_> = plan.children().values().flatten().cloned().collect();
for (index, id) in plan.source_ids().iter().enumerate() {
if let NodeId::Existing(id) = id {
ids.insert(NodeId::Existing(id.clone()), id.clone());
}
if !live.contains(id) {
if let NodeId::Existing(id) = id
&& !semantic_equal(old[id.as_str()], targets[index])
{
return Err(failure(id, "blocks", MergeError::ConcurrentEdit));
}
continue;
}
match id {
NodeId::New(_) => {
let spec = source_nodes[index]
.spec
.clone()
.ok_or_else(|| failure("", "opaque", MergeError::InvalidProjection))?;
let id_value = insert_block_tree(
¤t.doc,
&mut blocks,
&BlockNode {
id: None,
spec,
children: Vec::new(),
},
)?;
ids.insert(id.clone(), id_value);
}
NodeId::Existing(id) => {
let previous = old
.get(id.as_str())
.ok_or_else(|| failure(id, "identity", MergeError::InvalidProjection))?;
let target = targets[index];
let base = &baseline.pool[id];
let mut block = current
.pool
.get(id)
.cloned()
.ok_or_else(|| failure(id, "blocks", MergeError::ConcurrentEdit))?;
if previous.opaque.is_some() {
if previous.opaque != target.opaque {
return Err(failure(id, "opaque", MergeError::ConcurrentEdit));
}
continue;
}
if previous.properties == target.properties && previous.text == target.text {
continue;
}
if base.identity()? != block.identity()? {
return Err(failure(id, "blocks", MergeError::InsufficientEvidence));
}
if get_string(base, "sys:flavour") != get_string(&block, "sys:flavour") {
return Err(failure(id, "sys:flavour", MergeError::ConcurrentEdit));
}
let base_values = scalar_properties(base);
let current_values = scalar_properties(&block);
let changes = merge_properties(&base_values, ¤t_values, &previous.properties, &target.properties)
.map_err(|e| failure(id, "properties", e))?;
validate_dependencies(id, &target.kind, &base_values, ¤t_values, &changes)?;
for (key, value) in changes {
allowed.entry(id.clone()).or_default().insert(key.clone());
if let Some(value) = value {
block.insert(key, value)?;
} else {
block.remove(&key);
}
}
if previous.text != target.text {
let base_text = base
.get("prop:text")
.and_then(|v| v.to_text())
.ok_or_else(|| failure(id, "prop:text", MergeError::InvalidProjection))?;
let mut text = block
.get("prop:text")
.and_then(|v| v.to_text())
.ok_or_else(|| failure(id, "prop:text", MergeError::ConcurrentEdit))?;
let text_plan = merge_text(
TextMerge {
baseline: &base_text,
current: &text,
source_before: &previous.text,
source_after: &target.text,
editable_marks: EDITABLE_MARKS,
},
budget,
)
.map_err(|e| failure(id, "prop:text", e))?;
if !text_plan.is_empty() {
text_plan.apply(&mut text).map_err(|e| failure(id, "prop:text", e))?;
allowed.entry(id.clone()).or_default().insert("prop:text".into());
}
}
}
}
}
for id in plan.removed() {
let Some(base) = baseline.pool.get(id) else {
return Err(ParseError::InvalidBinary);
};
if old.get(id.as_str()).is_some_and(|n| n.opaque.is_some()) {
return Err(failure(id, "opaque", MergeError::ConcurrentEdit));
}
if let Some(now) = current.pool.get(id) {
if observe(Value::Map(base.clone()), budget).map_err(|e| failure(id, "blocks", e))?
!= observe(Value::Map(now.clone()), budget).map_err(|e| failure(id, "blocks", e))?
{
return Err(failure(id, "blocks", MergeError::ConcurrentEdit));
}
}
}
for (parent, children) in plan.children() {
let parent = match parent {
Some(NodeId::New(_)) => ids[parent.as_ref().unwrap()].clone(),
Some(NodeId::Existing(id)) => id.clone(),
None => current.scope.clone(),
};
let children = children
.iter()
.map(|id| match id {
NodeId::Existing(id) => id.clone(),
NodeId::New(_) => ids[id].clone(),
})
.collect::<Vec<_>>();
let mut block = blocks
.get(&parent)
.and_then(|v| v.to_map())
.ok_or(ParseError::InvalidBinary)?;
if old.get(parent.as_str()).is_some_and(|n| n.opaque.is_some()) {
continue;
}
if collect_child_ids(&block) != children {
patch_children(¤t.doc, &mut block, &children)?;
allowed.entry(parent).or_default().insert("sys:children".into());
}
}
for id in plan.removed() {
blocks.remove(id);
}
let after_pool = blocks
.iter()
.map(|(id, v)| Ok((id.to_owned(), v.to_map().ok_or(ParseError::InvalidBinary)?)))
.collect::<Result<HashMap<_, _>, ParseError>>()?;
let after = snapshot_blocks(&after_pool, budget)?;
let removed: HashSet<_> = plan.removed().iter().collect();
for (id, mut old_value) in snapshots {
if removed.contains(&id) {
continue;
}
let mut new_value = after
.get(&id)
.cloned()
.ok_or_else(|| failure(&id, "retention", MergeError::InvalidProjection))?;
if let Some(keys) = allowed.get(&id) {
if let (ObservedValue::Map(_, old_fields), ObservedValue::Map(_, new_fields)) = (&mut old_value, &mut new_value) {
for key in keys {
old_fields.remove(key);
new_fields.remove(key);
}
}
}
if old_value != new_value {
return Err(failure(&id, "retention", MergeError::InvalidProjection));
}
}
Ok(ids)
}
fn scalar_properties(block: &Map) -> BTreeMap<String, Any> {
block
.iter()
.filter_map(|(key, value)| value.to_any().map(|value| (key.to_owned(), value)))
.collect()
}
fn snapshot_blocks(
pool: &HashMap<String, Map>,
budget: &mut WorkBudget,
) -> Result<BTreeMap<String, ObservedValue>, ParseError> {
pool
.iter()
.map(|(id, map)| {
Ok((
id.clone(),
observe(Value::Map(map.clone()), budget).map_err(|e| failure(id, "retention", e))?,
))
})
.collect()
}
fn validate_dependencies(
id: &str,
kind: &str,
baseline: &BTreeMap<String, Any>,
current: &BTreeMap<String, Any>,
changes: &BTreeMap<String, Option<Any>>,
) -> Result<(), ParseError> {
let (identity, dependents): (&str, &[&str]) = match kind {
"affine:list" => ("prop:type", &["prop:checked", "prop:order"]),
"affine:image" => ("prop:sourceId", &["prop:width", "prop:height", "prop:caption"]),
"affine:bookmark" => (
"prop:url",
&[
"prop:title",
"prop:description",
"prop:icon",
"prop:image",
"prop:caption",
],
),
"affine:embed-youtube" => ("prop:videoId", &["prop:title", "prop:description", "prop:thumbnail"]),
_ => return Ok(()),
};
if changes.contains_key(identity) {
for key in dependents {
if baseline.get(*key) != current.get(*key) && !changes.contains_key(*key) {
return Err(failure(id, key, MergeError::ConcurrentEdit));
}
if kind != "affine:list" && current.contains_key(*key) && !changes.contains_key(*key) {
return Err(failure(id, key, MergeError::UnsupportedMetadataEffect((*key).into())));
}
}
} else if dependents.iter().any(|key| changes.contains_key(*key)) && baseline.get(identity) != current.get(identity) {
return Err(failure(id, identity, MergeError::ConcurrentEdit));
}
Ok(())
}