1use std::borrow::Borrow;
18use std::collections::BTreeMap;
19use std::collections::HashSet;
20use std::iter::zip;
21use std::sync::Arc;
22use std::vec;
23
24use futures::AsyncReadExt as _;
25use futures::FutureExt as _;
26use futures::StreamExt as _;
27use futures::future::BoxFuture;
28use futures::future::try_join_all;
29use futures::stream::FuturesUnordered;
30use itertools::Itertools as _;
31
32use crate::backend;
33use crate::backend::BackendError;
34use crate::backend::BackendResult;
35use crate::backend::MergedTreeVal;
36use crate::backend::MergedTreeValue;
37use crate::backend::MergedTreeValueExt as _;
38use crate::backend::TreeId;
39use crate::backend::TreeValue;
40use crate::config::ConfigGetError;
41use crate::files;
42use crate::files::FileMergeHunkLevel;
43use crate::merge::Merge;
44use crate::merge::SameChange;
45use crate::merged_tree::all_merged_tree_entries;
46use crate::object_id::ObjectId as _;
47use crate::repo_path::RepoPath;
48use crate::repo_path::RepoPathBuf;
49use crate::repo_path::RepoPathComponentBuf;
50use crate::settings::UserSettings;
51use crate::store::Store;
52use crate::tree::ToTreeMergeExt as _;
53use crate::tree::Tree;
54
55#[derive(Clone, Debug)]
57pub struct MergeOptions {
58 pub hunk_level: FileMergeHunkLevel,
60 pub same_change: SameChange,
62}
63
64impl MergeOptions {
65 pub fn from_settings(settings: &UserSettings) -> Result<Self, ConfigGetError> {
67 Ok(Self {
68 hunk_level: settings.get("merge.hunk-level")?,
71 same_change: settings.get("merge.same-change")?,
72 })
73 }
74}
75
76pub async fn merge_trees(store: &Arc<Store>, merge: Merge<TreeId>) -> BackendResult<Merge<TreeId>> {
79 let merge = match merge.into_resolved() {
80 Ok(tree) => return Ok(Merge::resolved(tree)),
81 Err(merge) => merge,
82 };
83
84 let mut merger = TreeMerger {
85 store: store.clone(),
86 trees_to_resolve: BTreeMap::new(),
87 work: FuturesUnordered::new(),
88 unstarted_work: BTreeMap::new(),
89 };
90 merger.enqueue_tree_read(
91 RepoPathBuf::root(),
92 merge.map(|tree_id| Some(TreeValue::Tree(tree_id.clone()))),
93 );
94 let trees = merger.merge().await?;
95 Ok(trees.map(|tree| tree.id().clone()))
96}
97
98struct MergedTreeInput {
99 resolved: BTreeMap<RepoPathComponentBuf, TreeValue>,
100 pending_lookup: HashSet<RepoPathComponentBuf>,
103 conflicts: BTreeMap<RepoPathComponentBuf, MergedTreeValue>,
104}
105
106impl MergedTreeInput {
107 fn new(resolved: BTreeMap<RepoPathComponentBuf, TreeValue>) -> Self {
108 Self {
109 resolved,
110 pending_lookup: HashSet::new(),
111 conflicts: BTreeMap::new(),
112 }
113 }
114
115 fn mark_completed(
116 &mut self,
117 basename: RepoPathComponentBuf,
118 value: MergedTreeValue,
119 same_change: SameChange,
120 ) {
121 let was_pending = self.pending_lookup.remove(&basename);
122 assert!(was_pending, "No pending lookup for {basename:?}");
123 if let Some(resolved) = value.resolve_trivial(same_change) {
124 if let Some(resolved) = resolved.as_ref() {
125 self.resolved.insert(basename, resolved.clone());
126 }
127 } else {
128 self.conflicts.insert(basename, value);
129 }
130 }
131
132 fn into_backend_trees(self) -> Merge<backend::Tree> {
133 assert!(self.pending_lookup.is_empty());
134
135 fn by_name(
136 (name1, _): &(RepoPathComponentBuf, TreeValue),
137 (name2, _): &(RepoPathComponentBuf, TreeValue),
138 ) -> bool {
139 name1 < name2
140 }
141
142 if self.conflicts.is_empty() {
143 let all_entries = self.resolved.into_iter().collect();
144 Merge::resolved(backend::Tree::from_sorted_entries(all_entries))
145 } else {
146 let mut conflict_entries = self.conflicts.first_key_value().unwrap().1.map(|_| vec![]);
148 for (basename, value) in self.conflicts {
149 assert_eq!(value.num_sides(), conflict_entries.num_sides());
150 for (entries, value) in zip(&mut conflict_entries, value) {
151 if let Some(value) = value {
152 entries.push((basename.clone(), value));
153 }
154 }
155 }
156
157 let mut backend_trees = vec![];
158 for entries in conflict_entries {
159 let backend_tree = backend::Tree::from_sorted_entries(
160 self.resolved
161 .iter()
162 .map(|(name, value)| (name.clone(), value.clone()))
163 .merge_by(entries, by_name)
164 .collect(),
165 );
166 backend_trees.push(backend_tree);
167 }
168 Merge::from_vec(backend_trees)
169 }
170 }
171}
172
173enum TreeMergerWorkOutput {
175 ReadTrees {
177 dir: RepoPathBuf,
178 result: BackendResult<Merge<Tree>>,
179 },
180 WrittenTrees {
181 dir: RepoPathBuf,
182 result: BackendResult<Merge<Tree>>,
183 },
184 MergedFiles {
185 path: RepoPathBuf,
186 result: BackendResult<MergedTreeValue>,
187 },
188}
189
190#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
191enum TreeMergeWorkItemKey {
192 MergeFiles { path: RepoPathBuf },
195 ReadTrees { dir: RepoPathBuf },
196}
197
198struct TreeMerger {
199 store: Arc<Store>,
200 trees_to_resolve: BTreeMap<RepoPathBuf, MergedTreeInput>,
202 work: FuturesUnordered<BoxFuture<'static, TreeMergerWorkOutput>>,
204 unstarted_work: BTreeMap<TreeMergeWorkItemKey, BoxFuture<'static, TreeMergerWorkOutput>>,
206}
207
208impl TreeMerger {
209 async fn merge(mut self) -> BackendResult<Merge<Tree>> {
210 while let Some(work_item) = self.work.next().await {
211 match work_item {
212 TreeMergerWorkOutput::ReadTrees { dir, result } => {
213 let tree = result?;
214 self.process_tree(dir, tree);
215 }
216 TreeMergerWorkOutput::WrittenTrees { dir, result } => {
217 let tree = result?;
218 if dir.is_root() {
219 assert!(self.trees_to_resolve.is_empty());
220 assert!(self.work.is_empty());
221 assert!(self.unstarted_work.is_empty());
222 return Ok(tree);
223 }
224 let new_value = tree.map(|tree| {
226 (tree.id() != self.store.empty_tree_id())
227 .then(|| TreeValue::Tree(tree.id().clone()))
228 });
229 self.mark_completed(&dir, new_value);
230 }
231 TreeMergerWorkOutput::MergedFiles { path, result } => {
232 let value = result?;
233 self.mark_completed(&path, value);
234 }
235 }
236
237 while self.work.len() < self.store.concurrency() {
238 if let Some((_key, work)) = self.unstarted_work.pop_first() {
239 self.work.push(work);
240 } else {
241 break;
242 }
243 }
244 }
245
246 unreachable!("There was no work item for writing the root tree");
247 }
248
249 fn process_tree(&mut self, dir: RepoPathBuf, tree: Merge<Tree>) {
250 let same_change = self.store.merge_options().same_change;
253 let mut resolved = vec![];
254 let mut non_trivial = vec![];
255 for (basename, path_merge) in all_merged_tree_entries(&tree) {
256 if let Some(value) = path_merge.resolve_trivial(same_change) {
257 if let Some(value) = value.cloned() {
258 resolved.push((basename.to_owned(), value));
259 }
260 } else {
261 non_trivial.push((basename.to_owned(), path_merge.cloned()));
262 }
263 }
264
265 if non_trivial.is_empty() {
267 let backend_trees = Merge::resolved(backend::Tree::from_sorted_entries(resolved));
268 self.enqueue_tree_write(dir, backend_trees);
269 return;
270 }
271
272 let mut unmerged_tree = MergedTreeInput::new(resolved.into_iter().collect());
273 for (basename, value) in non_trivial {
274 let path = dir.join(&basename);
275 unmerged_tree.pending_lookup.insert(basename);
276 if value.is_tree() {
277 self.enqueue_tree_read(path, value);
278 } else {
279 self.enqueue_file_merge(path, value);
283 }
284 }
285
286 self.trees_to_resolve.insert(dir, unmerged_tree);
287 }
288
289 fn enqueue_tree_read(&mut self, dir: RepoPathBuf, value: MergedTreeValue) {
290 let key = TreeMergeWorkItemKey::ReadTrees { dir: dir.clone() };
291 let work_fut = read_trees(self.store.clone(), dir.clone(), value)
292 .map(|result| TreeMergerWorkOutput::ReadTrees { dir, result });
293 if self.work.len() < self.store.concurrency() {
294 self.work.push(Box::pin(work_fut));
295 } else {
296 self.unstarted_work.insert(key, Box::pin(work_fut));
297 }
298 }
299
300 fn enqueue_tree_write(&mut self, dir: RepoPathBuf, backend_trees: Merge<backend::Tree>) {
301 let work_fut = write_trees(self.store.clone(), dir.clone(), backend_trees)
302 .map(|result| TreeMergerWorkOutput::WrittenTrees { dir, result });
303 self.work.push(Box::pin(work_fut));
306 }
307
308 fn enqueue_file_merge(&mut self, path: RepoPathBuf, value: MergedTreeValue) {
309 let key = TreeMergeWorkItemKey::MergeFiles { path: path.clone() };
310 let work_fut = resolve_file_values_owned(self.store.clone(), path.clone(), value)
311 .map(|result| TreeMergerWorkOutput::MergedFiles { path, result });
312 if self.work.len() < self.store.concurrency() {
313 self.work.push(Box::pin(work_fut));
314 } else {
315 self.unstarted_work.insert(key, Box::pin(work_fut));
316 }
317 }
318
319 fn mark_completed(&mut self, path: &RepoPath, value: MergedTreeValue) {
320 let (dir, basename) = path.split().unwrap();
321 let tree = self.trees_to_resolve.get_mut(dir).unwrap();
322 let same_change = self.store.merge_options().same_change;
323 tree.mark_completed(basename.to_owned(), value, same_change);
324 if tree.pending_lookup.is_empty() {
327 let tree = self.trees_to_resolve.remove(dir).unwrap();
328 self.enqueue_tree_write(dir.to_owned(), tree.into_backend_trees());
329 }
330 }
331}
332
333async fn read_trees(
334 store: Arc<Store>,
335 dir: RepoPathBuf,
336 value: MergedTreeValue,
337) -> BackendResult<Merge<Tree>> {
338 let trees = value
339 .to_tree_merge(&store, &dir)
340 .await?
341 .expect("Should be tree merge");
342 Ok(trees)
343}
344
345async fn write_trees(
346 store: Arc<Store>,
347 dir: RepoPathBuf,
348 backend_trees: Merge<backend::Tree>,
349) -> BackendResult<Merge<Tree>> {
350 let trees = try_join_all(
353 backend_trees
354 .into_iter()
355 .map(|backend_tree| store.write_tree(&dir, backend_tree)),
356 )
357 .await?;
358 Ok(Merge::from_vec(trees))
359}
360
361async fn resolve_file_values_owned(
362 store: Arc<Store>,
363 path: RepoPathBuf,
364 values: MergedTreeValue,
365) -> BackendResult<MergedTreeValue> {
366 let maybe_resolved = try_resolve_file_values(&store, &path, &values).await?;
367 Ok(maybe_resolved.unwrap_or(values))
368}
369
370pub async fn resolve_file_values(
374 store: &Arc<Store>,
375 path: &RepoPath,
376 values: MergedTreeValue,
377) -> BackendResult<MergedTreeValue> {
378 let same_change = store.merge_options().same_change;
379 if let Some(resolved) = values.resolve_trivial(same_change) {
380 return Ok(Merge::resolved(resolved.clone()));
381 }
382
383 let maybe_resolved = try_resolve_file_values(store, path, &values).await?;
384 Ok(maybe_resolved.unwrap_or(values))
385}
386
387async fn try_resolve_file_values<T: Borrow<TreeValue>>(
388 store: &Arc<Store>,
389 path: &RepoPath,
390 values: &Merge<Option<T>>,
391) -> BackendResult<Option<MergedTreeValue>> {
392 let simplified = values
395 .map(|value| value.as_ref().map(Borrow::borrow))
396 .simplify();
397 if let Some(resolved) = try_resolve_file_conflict(store, path, &simplified).await? {
400 Ok(Some(Merge::normal(resolved)))
401 } else {
402 Ok(None)
404 }
405}
406
407async fn try_resolve_file_conflict(
412 store: &Store,
413 filename: &RepoPath,
414 conflict: &MergedTreeVal<'_>,
415) -> BackendResult<Option<TreeValue>> {
416 let options = store.merge_options();
417 let Ok(file_id_conflict) = conflict.try_map(|term| match term {
422 Some(TreeValue::File {
423 id,
424 executable: _,
425 copy_id: _,
426 }) => Ok(id),
427 _ => Err(()),
428 }) else {
429 return Ok(None);
430 };
431 let Ok(executable_conflict) = conflict.try_map(|term| match term {
432 Some(TreeValue::File {
433 id: _,
434 executable,
435 copy_id: _,
436 }) => Ok(executable),
437 _ => Err(()),
438 }) else {
439 return Ok(None);
440 };
441 let Ok(copy_id_conflict) = conflict.try_map(|term| match term {
442 Some(TreeValue::File {
443 id: _,
444 executable: _,
445 copy_id,
446 }) => Ok(copy_id),
447 _ => Err(()),
448 }) else {
449 return Ok(None);
450 };
451 let Some(&&executable) = executable_conflict.resolve_trivial(SameChange::Accept) else {
454 return Ok(None);
456 };
457 let Some(©_id) = copy_id_conflict.resolve_trivial(SameChange::Accept) else {
458 return Ok(None);
460 };
461 if let Some(&resolved_file_id) = file_id_conflict.resolve_trivial(options.same_change) {
462 return Ok(Some(TreeValue::File {
465 id: resolved_file_id.clone(),
466 executable,
467 copy_id: copy_id.clone(),
468 }));
469 }
470
471 let file_id_conflict = file_id_conflict.simplify();
478
479 let contents = file_id_conflict
480 .try_map_async(async |file_id| {
481 let mut content = vec![];
482 let mut reader = store.read_file(filename, file_id).await?;
483 reader
484 .read_to_end(&mut content)
485 .await
486 .map_err(|err| BackendError::ReadObject {
487 object_type: file_id.object_type(),
488 hash: file_id.hex(),
489 source: err.into(),
490 })?;
491 BackendResult::Ok(content)
492 })
493 .await?;
494 if let Some(merged_content) = files::try_merge(&contents, options) {
495 let id = store
496 .write_file(filename, &mut merged_content.as_slice())
497 .await?;
498 Ok(Some(TreeValue::File {
499 id,
500 executable,
501 copy_id: copy_id.clone(),
502 }))
503 } else {
504 Ok(None)
505 }
506}