1use std::path::Path;
29use std::time::Duration;
30
31use anyhow::{Context as _, Result};
32use jiff::Timestamp;
33
34use crate::agent::{self, Invocation};
35use crate::git;
36use crate::land::seat_of;
37use crate::prompt;
38use crate::run::{QuotaLoss, RebaseFixRecord, RunState};
39
40#[derive(Debug, Clone, PartialEq, Eq)]
42pub enum Rebased {
43 Applied,
46 Stopped(String),
50}
51
52const HUNK_PER_FILE: usize = 4_000;
54const HUNK_TOTAL: usize = 16_000;
55const SUBJECTS: usize = 20;
57const PATHS_IN_REASON: usize = 8;
59
60pub async fn rebase_with_fixer(
66 state: &mut RunState,
67 scratch: &Path,
68 branch: &str,
69 onto: &str,
70) -> Result<Rebased> {
71 let repo = state.repo.clone();
72 let cap = state.config.graph.review_rounds;
73 let orig = git::rev_parse(&repo, &format!("refs/heads/{branch}")).await?;
74
75 let mut said = String::new();
76 if !git::rebase_in_progress(scratch).await {
77 match git::rebase_start(&repo, scratch, branch, onto).await? {
78 git::RebaseStart::Applied => return Ok(Rebased::Applied),
79 git::RebaseStart::Failed(why) => return Ok(Rebased::Stopped(why)),
80 git::RebaseStart::Conflicted(why) => said = why,
81 }
82 }
83 let onto_sha = git::rev_parse(&repo, onto).await?;
84 let mut touched: Vec<String> = Vec::new();
87
88 loop {
89 if !git::rebase_in_progress(scratch).await {
90 return finish(state, scratch, branch, &orig, &onto_sha, &touched, &said).await;
91 }
92 let paths = git::unmerged_paths(scratch).await.unwrap_or_default();
93 for p in &paths {
94 if !touched.contains(p) {
95 touched.push(p.clone());
96 }
97 }
98 let spent = state.rebase_fixes.len();
99 if spent >= cap {
100 let why = reason(spent, cap, &paths, &said, "the rounds are spent");
101 abandon(&repo, scratch, branch, &orig).await;
102 return Ok(Rebased::Stopped(why));
103 }
104
105 let winner = state
106 .winner()
107 .cloned()
108 .context("resolving a rebase conflict needs a winning candidate")?;
109 let roles = state
110 .config
111 .resolve_roles()
112 .context("resolve the roster for the rebase fix")?;
113 let (spec, seat_key) = match &roles.fixer {
114 Some(f) if f.id != winner.agent => (f.clone(), "fix".to_owned()),
115 _ => (
116 state
117 .config
118 .agent(&winner.agent)
119 .cloned()
120 .unwrap_or_else(|_| roles.implementers[winner.index].clone()),
121 format!("impl-{}", winner.label),
122 ),
123 };
124
125 let round = spent + 1;
126 let branch_subjects = subjects(scratch, &format!("{onto}..{branch}")).await;
127 let onto_subjects = subjects(scratch, &format!("{branch}..{onto}")).await;
128 let hunks = hunks(scratch, &paths);
129 let prompt_text = prompt::rebase_conflict(&prompt::RebaseConflict {
130 instruction: &state.instruction,
131 worktree: scratch,
132 branch,
133 onto,
134 paths: &paths,
135 branch_subjects: &branch_subjects,
136 onto_subjects: &onto_subjects,
137 hunks: &hunks,
138 round,
139 cap,
140 language: &state.config.graph.language,
141 });
142 let prompt_text = if state.config.cache_dir().is_some() {
143 format!("{prompt_text}\n\n{}", prompt::build_cache_note("fix", true))
144 } else {
145 prompt_text
146 };
147
148 state.rebase_fixes.push(RebaseFixRecord {
151 agent: spec.id.clone(),
152 paths: paths.clone(),
153 from: Some(orig.clone()),
154 finished: false,
155 error: None,
156 });
157 state.event(
158 "rebase",
159 format!(
160 "{branch} conflicts with {onto} ({} path(s)); fixer round {round} of {cap}",
161 paths.len()
162 ),
163 );
164 state.save()?;
165
166 let mut seat = seat_of(state, &seat_key, &spec.id);
167 let artifacts = agent::artifacts_dir(&state.dir());
168 let out = agent::invoke(
169 &spec,
170 &mut seat,
171 &Invocation {
172 cwd: scratch,
173 prompt: &prompt_text,
174 timeout: Duration::from_secs(state.config.graph.timeout_fix),
175 allow_write: true,
176 sessions: state.config.graph.sessions,
177 artifacts: &artifacts,
178 stem: &format!("rebase-fix-{round}"),
179 run: &state.id,
180 node: "rebase",
181 cache_dir: state.config.cache_dir().as_deref(),
182 attachments: &[],
183 writable: &[],
184 },
185 )
186 .await;
187 let seat_name = seat.key.clone();
188 state.seats.insert(seat.key.clone(), seat);
189
190 let mut error = None;
191 let mut quota = false;
192 match out {
193 Ok(o) if o.quota_exhausted() => {
194 state.quota.push(QuotaLoss {
195 seat: seat_name,
196 node: "rebase".to_owned(),
197 at: Timestamp::now(),
198 reset: o.quota.as_ref().and_then(|q| q.reset.clone()),
199 });
200 error = Some("rate limited (quota); the fixer could not run".to_owned());
201 quota = true;
202 }
203 Ok(o) if !o.usable() => {
204 error = Some(format!(
205 "the fixer produced nothing usable (exit {:?}, timed out: {})",
206 o.exit_code, o.timed_out
207 ));
208 }
209 Ok(_) => {}
210 Err(e) => error = Some(format!("{e:#}")),
211 }
212
213 let finished = !git::rebase_in_progress(scratch).await;
214 if let Some(r) = state.rebase_fixes.last_mut() {
215 r.finished = finished;
216 r.error = error.clone();
217 }
218 state.save()?;
219
220 if quota {
221 let paths = git::unmerged_paths(scratch).await.unwrap_or_default();
223 let why = reason(
224 state.rebase_fixes.len(),
225 cap,
226 &paths,
227 &said,
228 "the fixer hit its rate limit",
229 );
230 abandon(&repo, scratch, branch, &orig).await;
231 return Ok(Rebased::Stopped(why));
232 }
233 }
234}
235
236async fn finish(
238 state: &mut RunState,
239 scratch: &Path,
240 branch: &str,
241 orig: &str,
242 onto_sha: &str,
243 touched: &[String],
244 said: &str,
245) -> Result<Rebased> {
246 let repo = state.repo.clone();
247 let spent = state.rebase_fixes.len();
248 let cap = state.config.graph.review_rounds;
249 let unmerged = git::unmerged_paths(scratch).await.unwrap_or_default();
250 let head = git::rev_parse(scratch, "HEAD").await.unwrap_or_default();
251 let mut candidates: Vec<String> = touched.to_vec();
256 if let Ok(changed) = git::git(scratch, &["diff", "--name-only", onto_sha, "HEAD"]).await {
257 for p in changed.lines().map(str::trim).filter(|l| !l.is_empty()) {
258 if !candidates.iter().any(|c| c == p) {
259 candidates.push(p.to_owned());
260 }
261 }
262 }
263 let marked: Vec<String> = candidates
264 .into_iter()
265 .filter(|p| has_markers(scratch, p))
266 .collect();
267
268 let emptied = head == onto_sha
273 && git::cherry(&repo, onto_sha, orig)
274 .await
275 .map_or(true, |(unmatched, _)| !unmatched.is_empty());
276
277 let dropped = if unmerged.is_empty() && marked.is_empty() && !emptied && !head.is_empty() {
282 dropped_commits(&repo, onto_sha, orig, &head).await
283 } else {
284 Vec::new()
285 };
286
287 let problem = if !unmerged.is_empty() {
288 Some(("paths are still unmerged", unmerged))
289 } else if !marked.is_empty() {
290 Some(("conflict markers were left in the tree", marked))
291 } else if emptied {
292 Some((
293 "the rebase ended with none of the branch's commits applied (all skipped)",
294 touched.to_vec(),
295 ))
296 } else if !dropped.is_empty() {
297 Some((
298 "the rebase dropped some of the branch's commits (skipped?)",
299 dropped,
300 ))
301 } else if head.is_empty() || !git::is_ancestor(&repo, onto_sha, &head).await {
302 Some((
303 "the rebase ended without the base in the result (abandoned or skipped)",
304 touched.to_vec(),
305 ))
306 } else {
307 None
308 };
309 match problem {
310 None => {
311 git::worktree_remove(&repo, scratch).await.ok();
312 state.event(
313 "rebase",
314 format!("{branch} rebased after {spent} fixer round(s)"),
315 );
316 state.save()?;
317 Ok(Rebased::Applied)
318 }
319 Some((what, paths)) => {
320 let why = reason(spent, cap, &paths, said, what);
321 abandon(&repo, scratch, branch, orig).await;
322 Ok(Rebased::Stopped(why))
323 }
324 }
325}
326
327async fn dropped_commits(repo: &Path, onto_sha: &str, orig: &str, head: &str) -> Vec<String> {
336 let Ok((unmatched, _)) = git::cherry(repo, onto_sha, orig).await else {
337 return Vec::new();
338 };
339 if unmatched.is_empty() {
340 return Vec::new();
341 }
342 let (Ok(origin), Ok(result)) = (
343 git::commit_keys(repo, &format!("{onto_sha}..{orig}")).await,
344 git::commit_keys(repo, &format!("{onto_sha}..{head}")).await,
345 ) else {
346 return Vec::new();
347 };
348 let expected: Vec<git::CommitKey> = origin
349 .into_iter()
350 .filter(|c| unmatched.contains(&c.sha))
351 .collect();
352 let have: Vec<String> = result.iter().map(|c| c.key.clone()).collect();
353 let mut taken = vec![false; result.len()];
354 let mut lost = Vec::new();
355 let mut unverified = Vec::new();
356 for c in missing_commits(&expected, &have) {
361 if !already_in_result(repo, &c.sha, orig, head).await {
366 unverified.push(c);
367 continue;
368 }
369 let mine = change_lines(repo, &c.sha).await;
370 for (i, r) in result.iter().enumerate() {
371 if !taken[i] && r.key == c.key && change_lines(repo, &r.sha).await == mine {
372 taken[i] = true;
373 break;
374 }
375 }
376 }
377 for c in unverified {
381 let mine = commit_paths(repo, &c.sha).await;
382 let mut found = false;
383 for (i, r) in result.iter().enumerate() {
384 if taken[i] || r.key != c.key {
385 continue;
386 }
387 let theirs = commit_paths(repo, &r.sha).await;
388 if theirs.iter().any(|p| mine.contains(p)) {
389 taken[i] = true;
390 found = true;
391 break;
392 }
393 }
394 if !found {
395 lost.push(c.key.rsplit('\u{1f}').next().unwrap_or(&c.key).to_owned());
396 }
397 }
398 lost
399}
400
401async fn change_lines(repo: &Path, sha: &str) -> Vec<String> {
405 git::git(
406 repo,
407 &["diff-tree", "-p", "-U0", "--no-commit-id", "--root", sha],
408 )
409 .await
410 .map(|o| {
411 o.lines()
412 .filter(|l| {
413 !(l.starts_with("diff ")
414 || l.starts_with("index ")
415 || l.starts_with("@@")
416 || l.starts_with("--- ")
417 || l.starts_with("+++ "))
418 })
419 .map(str::to_owned)
420 .collect()
421 })
422 .unwrap_or_default()
423}
424
425async fn commit_paths(repo: &Path, sha: &str) -> Vec<String> {
428 git::git(
429 repo,
430 &[
431 "diff-tree",
432 "--no-commit-id",
433 "--name-only",
434 "-r",
435 "--root",
436 "-z",
437 sha,
438 ],
439 )
440 .await
441 .map(|o| {
442 o.split('\0')
443 .filter(|l| !l.is_empty())
444 .map(str::to_owned)
445 .collect()
446 })
447 .unwrap_or_default()
448}
449
450async fn entry(repo: &Path, rev: &str, path: &str) -> Option<String> {
453 let out = git::git(repo, &["ls-tree", rev, "--", path]).await.ok()?;
454 out.split('\t')
455 .next()
456 .filter(|e| !e.is_empty())
457 .map(str::to_owned)
458}
459
460async fn already_in_result(repo: &Path, sha: &str, orig: &str, head: &str) -> bool {
469 if replay_is_noop(repo, sha, head).await {
470 return true;
471 }
472 for p in &commit_paths(repo, sha).await {
473 let got = entry(repo, head, p).await;
474 if got != entry(repo, sha, p).await && got != entry(repo, orig, p).await {
475 return false;
476 }
477 }
478 true
479}
480
481async fn replay_is_noop(repo: &Path, sha: &str, head: &str) -> bool {
484 let base = format!("{sha}^");
485 let Ok(out) = git::git_raw(
486 repo,
487 &[
488 "merge-tree",
489 "--write-tree",
490 &format!("--merge-base={base}"),
491 head,
492 sha,
493 ],
494 )
495 .await
496 else {
497 return false;
498 };
499 if !out.ok() {
500 return false;
501 }
502 let merged = out.stdout.lines().next().unwrap_or("").trim();
503 match git::tree_of(repo, head).await {
504 Ok(t) => !merged.is_empty() && merged == t,
505 Err(_) => false,
506 }
507}
508
509fn missing_commits(expected: &[git::CommitKey], have: &[String]) -> Vec<git::CommitKey> {
517 expected
518 .iter()
519 .filter(|c| {
520 let want = expected.iter().filter(|e| e.key == c.key).count();
521 let got = have.iter().filter(|k| **k == c.key).count();
522 got < want
523 })
524 .cloned()
525 .collect()
526}
527
528async fn abandon(repo: &Path, scratch: &Path, branch: &str, orig: &str) {
531 git::rebase_abort(repo, scratch).await;
532 let full = format!("refs/heads/{branch}");
533 if git::rev_parse(repo, &full).await.ok().as_deref() != Some(orig) {
534 git::git_raw(repo, &["update-ref", &full, orig]).await.ok();
535 }
536}
537
538fn reason(spent: usize, cap: usize, paths: &[String], said: &str, what: &str) -> String {
541 let shown: Vec<&str> = paths
542 .iter()
543 .take(PATHS_IN_REASON)
544 .map(String::as_str)
545 .collect();
546 let mut list = shown.join(", ");
547 if paths.len() > shown.len() {
548 list.push_str(&format!(" and {} more", paths.len() - shown.len()));
549 }
550 if list.is_empty() {
551 list.push_str("none recorded");
552 }
553 let mut s = format!(
554 "conflict not resolved after {spent} of {cap} fixer round(s) ({what}); remaining \
555 conflicted path(s): {list}"
556 );
557 let said = said.trim();
558 if !said.is_empty() {
559 s.push_str("; git said: ");
560 s.extend(said.chars().take(250));
561 }
562 s
563}
564
565async fn subjects(worktree: &Path, range: &str) -> Vec<String> {
567 let n = format!("-n{SUBJECTS}");
568 git::git(worktree, &["log", "--format=%s", &n, range])
569 .await
570 .map(|o| o.lines().map(str::to_owned).collect())
571 .unwrap_or_default()
572}
573
574fn has_markers(worktree: &Path, path: &str) -> bool {
575 std::fs::read_to_string(worktree.join(path)).is_ok_and(|t| {
576 t.lines().any(|l| l.starts_with("<<<<<<< ")) && t.lines().any(|l| l.starts_with(">>>>>>> "))
577 })
578}
579
580fn hunks(worktree: &Path, paths: &[String]) -> String {
583 let mut out = String::new();
584 for p in paths {
585 if out.len() >= HUNK_TOTAL {
586 out.push_str("\n(more conflicted files omitted)\n");
587 break;
588 }
589 out.push_str(&format!("=== {p} ===\n"));
590 let Ok(text) = std::fs::read_to_string(worktree.join(p)) else {
591 out.push_str("(not readable as text; use git to inspect it)\n");
592 continue;
593 };
594 let mut file = String::new();
595 let mut inside = false;
596 for line in text.lines() {
597 if line.starts_with("<<<<<<< ") {
598 inside = true;
599 }
600 if inside {
601 file.push_str(line);
602 file.push('\n');
603 }
604 if line.starts_with(">>>>>>> ") {
605 inside = false;
606 }
607 }
608 if file.len() > HUNK_PER_FILE {
609 file = file.chars().take(HUNK_PER_FILE).collect();
610 file.push_str("\n(truncated)\n");
611 }
612 out.push_str(&file);
613 }
614 out
615}
616
617#[cfg(test)]
618mod tests {
619 use super::*;
620 use crate::proc::Quiet as _;
621
622 fn ck(sha: &str, subject: &str) -> git::CommitKey {
623 git::CommitKey {
624 sha: sha.to_owned(),
625 key: format!("n\u{1f}e\u{1f}1 +0000\u{1f}{subject}"),
626 }
627 }
628
629 #[test]
630 fn nothing_is_missing_when_every_key_is_present() {
631 let exp = [ck("a", "one"), ck("b", "two")];
632 let have = vec![exp[1].key.clone(), exp[0].key.clone()];
633 assert!(missing_commits(&exp, &have).is_empty());
634 }
635
636 #[test]
637 fn a_dropped_commit_is_named_by_subject() {
638 let exp = [ck("a", "one"), ck("b", "two")];
639 let have = vec![exp[1].key.clone()];
640 assert_eq!(missing_commits(&exp, &have), vec![exp[0].clone()]);
641 }
642
643 #[test]
644 fn duplicate_keys_are_counted_not_collapsed() {
645 let exp = [ck("a", "same"), ck("b", "same")];
646 let have = vec![exp[0].key.clone()];
647 assert_eq!(missing_commits(&exp, &have), exp.to_vec());
648 }
649
650 #[test]
651 fn a_commit_without_a_key_match_is_a_candidate_even_if_a_twin_exists() {
652 let exp = [ck("a", "same"), ck("b", "same")];
653 let have = vec![exp[1].key.clone()];
654 assert_eq!(missing_commits(&exp, &have).len(), 2);
655 }
656
657 fn sh(dir: &Path, args: &[&str]) {
658 let o = std::process::Command::new("git")
659 .quiet()
660 .args(args)
661 .current_dir(dir)
662 .output()
663 .unwrap();
664 assert!(
665 o.status.success(),
666 "{args:?}: {}",
667 String::from_utf8_lossy(&o.stderr)
668 );
669 }
670
671 #[tokio::test]
672 async fn a_lost_mode_change_is_not_already_in_the_result() {
673 let t = tempfile::tempdir().unwrap();
674 let d = t.path();
675 sh(d, &["init", "-q", "-b", "main"]);
676 sh(d, &["config", "user.name", "t"]);
677 sh(d, &["config", "user.email", "t@example.com"]);
678 sh(d, &["config", "core.fileMode", "true"]);
679 std::fs::write(d.join("script.sh"), "echo\n").unwrap();
680 sh(d, &["add", "-A"]);
681 sh(d, &["commit", "-q", "-m", "base"]);
682 sh(d, &["update-index", "--chmod=+x", "script.sh"]);
683 sh(d, &["commit", "-q", "-m", "chmod"]);
684 let sha = git::git(d, &["rev-parse", "HEAD"]).await.unwrap();
685 let base = git::git(d, &["rev-parse", "HEAD~1"]).await.unwrap();
686 assert!(!already_in_result(d, &sha, &sha, &base).await);
688 assert!(already_in_result(d, &sha, &sha, &sha).await);
689 }
690
691 async fn commit(d: &Path, msg: &str) -> String {
692 sh(d, &["add", "-A"]);
693 sh(d, &["commit", "-q", "-m", msg]);
694 git::git(d, &["rev-parse", "HEAD"]).await.unwrap()
695 }
696
697 #[tokio::test]
698 async fn a_change_inside_a_larger_upstream_edit_is_in_the_result() {
699 let t = tempfile::tempdir().unwrap();
700 let d = t.path();
701 sh(d, &["init", "-q", "-b", "main"]);
702 sh(d, &["config", "user.name", "t"]);
703 sh(d, &["config", "user.email", "t@example.com"]);
704 let body = "old\n1\n2\n3\n4\n5\n6\n7\n8\n9\nend\n";
705 std::fs::write(d.join("f.txt"), body).unwrap();
706 commit(d, "base").await;
707 std::fs::write(d.join("f.txt"), body.replacen("old", "new", 1)).unwrap();
708 let c = commit(d, "c1").await;
709 sh(d, &["checkout", "-q", "-b", "up", "HEAD~1"]);
710 let up = body.replacen("old", "new", 1).replace("end", "end plus");
711 std::fs::write(d.join("f.txt"), up).unwrap();
712 let head = commit(d, "upstream").await;
713 assert!(already_in_result(d, &c, &c, &head).await);
714 }
715
716 #[tokio::test]
717 async fn a_dropped_commit_on_a_non_ascii_path_is_still_noticed() {
718 let t = tempfile::tempdir().unwrap();
719 let d = t.path();
720 sh(d, &["init", "-q", "-b", "main"]);
721 sh(d, &["config", "user.name", "t"]);
722 sh(d, &["config", "user.email", "t@example.com"]);
723 std::fs::write(d.join("a.txt"), "a\n").unwrap();
724 let base = commit(d, "base").await;
725 std::fs::write(d.join("日本語.txt"), "x\n").unwrap();
726 let c = commit(d, "c").await;
727 assert!(!already_in_result(d, &c, &c, &base).await);
728 }
729
730 async fn commit_dated(d: &Path, msg: &str) -> String {
731 sh(d, &["add", "-A"]);
732 let st = std::process::Command::new("git")
733 .quiet()
734 .current_dir(d)
735 .env("GIT_AUTHOR_DATE", "2020-01-01T00:00:00+0000")
736 .env("GIT_COMMITTER_DATE", "2020-01-01T00:00:00+0000")
737 .args(["commit", "-q", "-m", msg])
738 .status()
739 .unwrap();
740 assert!(st.success());
741 git::git(d, &["rev-parse", "HEAD"]).await.unwrap()
742 }
743
744 #[tokio::test]
745 async fn a_surviving_resolved_commit_sharing_a_key_with_an_absorbed_one_is_not_lost() {
746 let t = tempfile::tempdir().unwrap();
747 let d = t.path();
748 sh(d, &["init", "-q", "-b", "main"]);
749 sh(d, &["config", "user.name", "t"]);
750 sh(d, &["config", "user.email", "t@example.com"]);
751 std::fs::write(d.join("z.txt"), "z\n").unwrap();
752 let base = commit(d, "base").await;
753 std::fs::write(d.join("a.txt"), "a\n").unwrap();
754 commit_dated(d, "same").await;
755 std::fs::write(d.join("b.txt"), "b\n").unwrap();
756 let orig = commit_dated(d, "same").await;
757 sh(d, &["checkout", "-q", "-b", "up", &base]);
760 std::fs::write(d.join("a.txt"), "a\n").unwrap();
761 std::fs::write(d.join("y.txt"), "y\n").unwrap();
762 std::fs::write(d.join("b.txt"), "other\n").unwrap();
763 let onto = commit(d, "upstream").await;
764 std::fs::write(d.join("b.txt"), "other\nb\n").unwrap();
766 let head = commit_dated(d, "same").await;
767 assert!(dropped_commits(d, &onto, &orig, &head).await.is_empty());
768 }
769
770 #[tokio::test]
771 async fn a_skipped_commit_is_not_masked_by_a_surviving_one_in_the_same_file() {
772 let t = tempfile::tempdir().unwrap();
773 let d = t.path();
774 sh(d, &["init", "-q", "-b", "main"]);
775 sh(d, &["config", "user.name", "t"]);
776 sh(d, &["config", "user.email", "t@example.com"]);
777 let body = "first\n1\n2\n3\n4\n5\n6\n7\n8\n9\nlast\n";
778 std::fs::write(d.join("f.txt"), body).unwrap();
779 let base = commit(d, "base").await;
780 std::fs::write(d.join("f.txt"), body.replacen("first", "mine", 1)).unwrap();
781 commit_dated(d, "same").await;
782 let two = body
783 .replacen("first", "mine", 1)
784 .replacen("last", "tail", 1);
785 std::fs::write(d.join("f.txt"), &two).unwrap();
786 let orig = commit_dated(d, "same").await;
787 sh(d, &["checkout", "-q", "-b", "up", &base]);
788 std::fs::write(d.join("f.txt"), body.replacen("first", "theirs", 1)).unwrap();
789 let onto = commit(d, "upstream").await;
790 std::fs::write(
792 d.join("f.txt"),
793 body.replacen("first", "theirs", 1)
794 .replacen("last", "tail", 1),
795 )
796 .unwrap();
797 let head = commit_dated(d, "same").await;
798 assert_eq!(dropped_commits(d, &onto, &orig, &head).await, vec!["same"]);
799 }
800}