using System.Collections.Concurrent; using ClaudeDo.Data; using ClaudeDo.Data.Git; using ClaudeDo.Data.Models; using ClaudeDo.Data.Repositories; using ClaudeDo.Worker.Hub; using ClaudeDo.Worker.State; using Microsoft.EntityFrameworkCore; using TaskStatus = ClaudeDo.Data.Models.TaskStatus; namespace ClaudeDo.Worker.Lifecycle; public sealed record MergeResult( string Status, IReadOnlyList ConflictFiles, string? ErrorMessage); public sealed record MergeTargets( string DefaultBranch, IReadOnlyList LocalBranches); public sealed record MergePreviewResult( string Status, IReadOnlyList ConflictFiles, int ChangedFileCount); public sealed record ConflictDocuments( string TaskId, IReadOnlyList Files); public sealed record ConflictDocumentContent( string Path, bool IsBinary, IReadOnlyList Segments); public sealed record RevertResult( string Status, string? RevertCommit, IReadOnlyList ConflictFiles, string? ErrorMessage); public sealed class TaskMergeService { public const string StatusMerged = "merged"; public const string StatusConflict = "conflict"; public const string StatusBlocked = "blocked"; public const string StatusAborted = "aborted"; public const string StatusVerifyFailed = "verify_failed"; public const string StatusReverted = "reverted"; public const string StatusConflictAborted = "conflict_aborted"; public const string PreviewClean = "clean"; public const string PreviewConflict = "conflict"; public const string PreviewUnavailable = "unavailable"; // The verify command is a trusted, list-owner-configured build/test invocation (not // per-request user input), so a generous fixed timeout is enough — no need for a // per-list configurable value on top of what the spec calls for. private static readonly TimeSpan VerifyTimeout = TimeSpan.FromMinutes(10); // Serializes merge (+ verify) against the same repo working dir: a verify command running // in list.WorkingDir must not see a second merge land mid-build. Keyed by working dir since // TaskMergeService is a process-wide singleton and merges across different lists are independent. private static readonly ConcurrentDictionary MergeGates = new(StringComparer.OrdinalIgnoreCase); private static SemaphoreSlim GetMergeGate(string workingDir) => MergeGates.GetOrAdd(workingDir, static _ => new SemaphoreSlim(1, 1)); private readonly IDbContextFactory _dbFactory; private readonly GitService _git; private readonly HubBroadcaster _broadcaster; private readonly ITaskStateService _state; private readonly IVerifyCommandRunner _verify; private readonly ILogger _logger; public TaskMergeService( IDbContextFactory dbFactory, GitService git, HubBroadcaster broadcaster, ITaskStateService state, IVerifyCommandRunner verify, ILogger logger) { _dbFactory = dbFactory; _git = git; _broadcaster = broadcaster; _state = state; _verify = verify; _logger = logger; } private async Task<(TaskEntity Task, ListEntity List, WorktreeEntity? Worktree, string? VerifyCommand)> LoadMergeContextAsync( string taskId, CancellationToken ct) { using var ctx = _dbFactory.CreateDbContext(); var task = await new TaskRepository(ctx).GetByIdAsync(taskId, ct) ?? throw new KeyNotFoundException($"Task '{taskId}' not found."); var listRepo = new ListRepository(ctx); var list = await listRepo.GetByIdAsync(task.ListId, ct) ?? throw new InvalidOperationException("List not found."); var wt = await new WorktreeRepository(ctx).GetByTaskIdAsync(taskId, ct); var config = await listRepo.GetConfigAsync(task.ListId, ct); return (task, list, wt, config?.VerifyCommand); } /// /// Runs the list's configured verify command (if any) in after /// a successful merge. Returns null when there is nothing to gate on (identical to today's /// behavior); otherwise returns the terminal to report instead of /// merged (the merge itself is left in place either way — see the design notes in Worker's /// CLAUDE.md — only the Done transition is withheld). /// private async Task RunVerifyGateAsync( string? verifyCommand, string workingDir, CancellationToken ct) { if (string.IsNullOrWhiteSpace(verifyCommand)) return null; VerifyCommandResult result; try { result = await _verify.RunAsync(workingDir, verifyCommand, VerifyTimeout, ct); } catch (Exception ex) { _logger.LogWarning(ex, "verify command failed to start: {Command}", verifyCommand); return new MergeResult(StatusVerifyFailed, Array.Empty(), $"verify command failed to start: {ex.Message}"); } if (!result.TimedOut && result.ExitCode == 0) return null; var reason = result.TimedOut ? $"verify command timed out after {VerifyTimeout.TotalMinutes:0} min: {verifyCommand}" : $"verify command failed (exit {result.ExitCode}): {verifyCommand}"; return new MergeResult(StatusVerifyFailed, Array.Empty(), $"{reason}\n{TailOutput(result.Output)}"); } private static string TailOutput(string output, int maxChars = 4000) { var trimmed = output.Trim(); return trimmed.Length <= maxChars ? trimmed : trimmed[^maxChars..]; } private async Task MarkWorktreeMergedAsync(string taskId, string mergeCommitSha, CancellationToken ct) { using (var ctx = _dbFactory.CreateDbContext()) { await new WorktreeRepository(ctx).SetMergedAsync(taskId, mergeCommitSha, ct); } await _broadcaster.WorktreeUpdated(taskId); } private async Task ApproveIfWaitingForReviewAsync(TaskEntity task, CancellationToken ct) { // A merged worktree means the work is integrated, so the task must reach Done. // MarkWorktreeMergedAsync only flips the worktree state; transition the task // itself when it was still awaiting review (a Done task is already terminal). if (task.Status == TaskStatus.WaitingForReview) await _state.ApproveReviewAsync(task.Id, ct); } public async Task MergeAsync( string taskId, string targetBranch, bool removeWorktree, string commitMessage, bool leaveConflictsInTree, CancellationToken ct) { var (task, list, wt, verifyCommand) = await LoadMergeContextAsync(taskId, ct); if (task.Status == TaskStatus.Running) return Blocked("task is running"); if (wt is null) return Blocked("task has no worktree"); if (wt.State != WorktreeState.Active) return Blocked($"worktree state is {wt.State}"); if (string.IsNullOrWhiteSpace(list.WorkingDir)) return Blocked("list has no working directory"); var gate = GetMergeGate(list.WorkingDir); await gate.WaitAsync(ct); try { if (!await _git.IsGitRepoAsync(list.WorkingDir, ct)) return Blocked("working directory is not a git repository"); if (await _git.IsMidMergeAsync(list.WorkingDir, ct)) return Blocked("target working directory is mid-merge"); if (await _git.HasChangesAsync(list.WorkingDir, includeUntracked: false, ct)) return Blocked("target working tree has uncommitted changes"); var currentBranch = await _git.GetCurrentBranchAsync(list.WorkingDir, ct); if (!string.Equals(currentBranch, targetBranch, StringComparison.Ordinal)) { try { await _git.CheckoutBranchAsync(list.WorkingDir, targetBranch, ct); } catch (Exception ex) { return Blocked($"failed to switch target branch: {ex.Message}"); } } var (exitCode, stderr) = await _git.MergeNoFfAsync(list.WorkingDir, wt.BranchName, commitMessage, ct); if (exitCode != 0) { List files; try { files = await _git.ListConflictedFilesAsync(list.WorkingDir, ct); } catch { files = new(); } if (leaveConflictsInTree && files.Count > 0) { return new MergeResult(StatusConflict, files, null); } // If abort fails the repo is left mid-merge; the caller must resolve manually. // Return Blocked (not conflict) so the UI does not offer a stale conflict list. try { await _git.MergeAbortAsync(list.WorkingDir, ct); } catch (Exception ex) { _logger.LogError(ex, "git merge --abort failed after conflict — repo is mid-merge"); return Blocked($"merge conflict and abort failed: {ex.Message} — repo is mid-merge, resolve manually"); } if (files.Count == 0) { // Non-conflict failure (e.g. unrelated histories). return new MergeResult(StatusBlocked, Array.Empty(), $"merge failed: {stderr}"); } return new MergeResult(StatusConflict, files, null); } var mergeSha = await _git.RevParseHeadAsync(list.WorkingDir, ct); string? cleanupWarning = null; if (removeWorktree) { try { await _git.WorktreeRemoveAsync(list.WorkingDir, wt.Path, force: false, ct); try { await _git.BranchDeleteAsync(list.WorkingDir, wt.BranchName, force: false, ct); } catch (Exception ex) { _logger.LogWarning(ex, "branch delete failed for {Branch}", wt.BranchName); cleanupWarning = $"worktree removed, branch delete failed: {ex.Message}"; } } catch (Exception ex) { _logger.LogWarning(ex, "worktree remove failed for {Path}", wt.Path); cleanupWarning = $"worktree remove failed: {ex.Message}"; } } await MarkWorktreeMergedAsync(taskId, mergeSha, ct); var verifyFailure = await RunVerifyGateAsync(verifyCommand, list.WorkingDir, ct); if (verifyFailure is not null) { _logger.LogWarning("Verify command failed after merging task {TaskId}: {Reason}", taskId, verifyFailure.ErrorMessage); await _broadcaster.WorkerLog($"Verify failed for \"{task.Title}\" after merge into {targetBranch}", WorkerLogLevel.Warn, DateTime.UtcNow); return verifyFailure; } await ApproveIfWaitingForReviewAsync(task, ct); _logger.LogInformation( "Merged task {TaskId} branch {Branch} into {Target} (remove worktree: {Remove})", taskId, wt.BranchName, targetBranch, removeWorktree); await _broadcaster.WorkerLog($"Merged \"{task.Title}\" into {targetBranch}", WorkerLogLevel.Success, DateTime.UtcNow); return new MergeResult(StatusMerged, Array.Empty(), cleanupWarning); } finally { gate.Release(); } } public Task MergeAsync( string taskId, string targetBranch, bool removeWorktree, string commitMessage, CancellationToken ct) => MergeAsync(taskId, targetBranch, removeWorktree, commitMessage, leaveConflictsInTree: false, ct); public async Task ContinueMergeAsync(string taskId, CancellationToken ct) { var (task, list, wt, verifyCommand) = await LoadMergeContextAsync(taskId, ct); if (wt is null) return Blocked("task has no worktree"); if (wt.State != WorktreeState.Active) return Blocked($"worktree state is {wt.State}"); if (string.IsNullOrWhiteSpace(list.WorkingDir)) return Blocked("list has no working directory"); var gate = GetMergeGate(list.WorkingDir); await gate.WaitAsync(ct); try { if (!await _git.IsMidMergeAsync(list.WorkingDir, ct)) return Blocked("repo is not mid-merge"); // Validate BEFORE staging: `git add` marks a conflicted path resolved regardless of // its content, so an unresolved file with markers still in it would otherwise get // staged (and committed) as-is. Check text content for markers first; binary files // can't carry markers, so they're left to the post-stage index check below. var unresolved = await _git.ListConflictedFilesAsync(list.WorkingDir, ct); var stillConflicted = new List(); foreach (var path in unresolved) { var full = Path.Combine(list.WorkingDir, path.Replace('/', Path.DirectorySeparatorChar)); string text; try { text = await File.ReadAllTextAsync(full, ct); } catch { continue; } if (!LooksBinary(text) && ConflictMarkerParser.HasConflicts(text)) stillConflicted.Add(path); } if (stillConflicted.Count > 0) return new MergeResult(StatusConflict, stillConflicted, "conflicts not fully resolved"); // Stage exactly the resolved conflict paths — never `git add -A`, which would sweep // untracked/unrelated changes left by other sessions into this merge commit (the // target working dir is shared). foreach (var path in unresolved) await _git.AddPathAsync(list.WorkingDir, path, ct); var remaining = await _git.ListConflictedFilesAsync(list.WorkingDir, ct); if (remaining.Count > 0) return new MergeResult(StatusConflict, remaining, "conflicts not fully resolved"); try { await _git.CommitAsync(list.WorkingDir, $"Merge branch '{wt.BranchName}'", ct); } catch (Exception ex) { return Blocked($"commit failed: {ex.Message}"); } var mergeSha = await _git.RevParseHeadAsync(list.WorkingDir, ct); await MarkWorktreeMergedAsync(taskId, mergeSha, ct); var verifyFailure = await RunVerifyGateAsync(verifyCommand, list.WorkingDir, ct); if (verifyFailure is not null) { _logger.LogWarning("Verify command failed after continuing merge of task {TaskId}: {Reason}", taskId, verifyFailure.ErrorMessage); await _broadcaster.WorkerLog($"Verify failed for \"{task.Title}\" after merge", WorkerLogLevel.Warn, DateTime.UtcNow); return verifyFailure; } await ApproveIfWaitingForReviewAsync(task, ct); _logger.LogInformation("Continued merge of task {TaskId} branch {Branch}", taskId, wt.BranchName); return new MergeResult(StatusMerged, Array.Empty(), null); } finally { gate.Release(); } } public async Task AbortMergeAsync(string taskId, CancellationToken ct) { var (_, list, wt, _) = await LoadMergeContextAsync(taskId, ct); if (wt is null) return Blocked("task has no worktree"); if (wt.State != WorktreeState.Active) return Blocked($"worktree state is {wt.State}"); if (string.IsNullOrWhiteSpace(list.WorkingDir)) return Blocked("list has no working directory"); if (!await _git.IsMidMergeAsync(list.WorkingDir, ct)) return Blocked("repo is not mid-merge"); try { await _git.MergeAbortAsync(list.WorkingDir, ct); } catch (Exception ex) { return Blocked($"abort failed: {ex.Message}"); } _logger.LogInformation("Aborted merge of task {TaskId}", taskId); return new MergeResult(StatusAborted, Array.Empty(), null); } /// /// Reverts the merge commit recorded for this task () /// via `git revert -m 1`, a new commit that undoes the merge without rewriting history — the /// target working directory is shared with other sessions, so a reset/rebase is never an option. /// On success the task returns to WaitingForReview so it can be reconsidered, and the worktree /// state moves to Kept: Merged/Discarded are swept by WorktreeMaintenanceService, and by the time /// a merge can be reverted its worktree directory and branch are typically already gone (removed /// during the original merge cleanup), so Active — which implies a live, resumable worktree — /// would be misleading. A conflicting revert is aborted immediately (`git revert --abort`); no /// partial/half-resolved state is ever left in the tree. /// public async Task RevertMergeAsync(string taskId, string targetBranch, CancellationToken ct) { var (task, list, wt, _) = await LoadMergeContextAsync(taskId, ct); if (task.Status != TaskStatus.Done) return RevertBlocked("task is not Done; only a merged task's revert can be undone"); if (wt is null) return RevertBlocked("task has no worktree"); if (wt.State != WorktreeState.Merged) return RevertBlocked($"worktree state is {wt.State}, expected Merged"); if (string.IsNullOrWhiteSpace(wt.MergeCommit)) return RevertBlocked("no merge commit recorded for this task; cannot revert"); if (string.IsNullOrWhiteSpace(list.WorkingDir)) return RevertBlocked("list has no working directory"); if (!await _git.IsGitRepoAsync(list.WorkingDir, ct)) return RevertBlocked("working directory is not a git repository"); if (await _git.IsMidMergeAsync(list.WorkingDir, ct)) return RevertBlocked("target working directory is mid-merge"); if (await _git.IsMidRevertAsync(list.WorkingDir, ct)) return RevertBlocked("target working directory is mid-revert"); if (await _git.HasChangesAsync(list.WorkingDir, includeUntracked: false, ct)) return RevertBlocked("target working tree has uncommitted changes"); var currentBranch = await _git.GetCurrentBranchAsync(list.WorkingDir, ct); if (!string.Equals(currentBranch, targetBranch, StringComparison.Ordinal)) { try { await _git.CheckoutBranchAsync(list.WorkingDir, targetBranch, ct); } catch (Exception ex) { return RevertBlocked($"failed to switch target branch: {ex.Message}"); } } var (exitCode, stderr) = await _git.RevertMergeCommitAsync(list.WorkingDir, wt.MergeCommit!, ct); if (exitCode != 0) { List files; try { files = await _git.ListConflictedFilesAsync(list.WorkingDir, ct); } catch { files = new(); } try { await _git.RevertAbortAsync(list.WorkingDir, ct); } catch (Exception ex) { _logger.LogError(ex, "git revert --abort failed after conflict — repo is mid-revert"); return RevertBlocked($"revert conflict and abort failed: {ex.Message} — repo is mid-revert, resolve manually"); } if (files.Count == 0) return RevertBlocked($"revert failed: {stderr}"); return new RevertResult(StatusConflictAborted, null, files, "revert conflicted; aborted cleanly, no changes made"); } var revertSha = await _git.RevParseHeadAsync(list.WorkingDir, ct); using (var ctx = _dbFactory.CreateDbContext()) { await new WorktreeRepository(ctx).SetStateAsync(taskId, WorktreeState.Kept, ct); } await _broadcaster.WorktreeUpdated(taskId); await _state.ForceSetStatusAsync(taskId, TaskStatus.WaitingForReview, ct); _logger.LogInformation( "Reverted merge of task {TaskId} (merge commit {MergeSha}) via revert commit {RevertSha}", taskId, wt.MergeCommit, revertSha); await _broadcaster.WorkerLog($"Reverted merge of \"{task.Title}\"", WorkerLogLevel.Warn, DateTime.UtcNow); return new RevertResult(StatusReverted, revertSha, Array.Empty(), null); } /// /// Reads each conflicted working-tree file and parses its conflict markers into line-level /// segments (with the diff3 merge base when present). Binary files are flagged and skipped. /// public async Task GetConflictDocumentsAsync(string taskId, CancellationToken ct) { var (_, list, _, _) = await LoadMergeContextAsync(taskId, ct); if (string.IsNullOrWhiteSpace(list.WorkingDir)) throw new InvalidOperationException("list has no working directory"); var files = await _git.ListConflictedFilesAsync(list.WorkingDir, ct); var result = new List(files.Count); foreach (var path in files) { var full = Path.Combine(list.WorkingDir, path.Replace('/', Path.DirectorySeparatorChar)); string text; try { text = await File.ReadAllTextAsync(full, ct); } catch { text = ""; } if (LooksBinary(text)) { result.Add(new ConflictDocumentContent(path, true, Array.Empty())); continue; } result.Add(new ConflictDocumentContent(path, false, ConflictMarkerParser.Parse(text))); } return new ConflictDocuments(taskId, result); } // A NUL byte in the head of the file is the conventional binary sniff. private static bool LooksBinary(string text) { var n = Math.Min(text.Length, 8000); for (var i = 0; i < n; i++) if (text[i] == '\0') return true; return false; } public async Task WriteResolutionAsync(string taskId, string path, string content, CancellationToken ct) { var (_, list, _, _) = await LoadMergeContextAsync(taskId, ct); if (string.IsNullOrWhiteSpace(list.WorkingDir)) throw new InvalidOperationException("list has no working directory"); var full = Path.Combine(list.WorkingDir, path.Replace('/', Path.DirectorySeparatorChar)); await File.WriteAllTextAsync(full, content, ct); await _git.AddPathAsync(list.WorkingDir, path, ct); } public async Task GetTargetsAsync(string taskId, CancellationToken ct) { var (_, list, _, _) = await LoadMergeContextAsync(taskId, ct); if (string.IsNullOrWhiteSpace(list.WorkingDir)) return new MergeTargets("", Array.Empty()); var current = await _git.GetCurrentBranchAsync(list.WorkingDir, ct); var branches = await _git.ListLocalBranchesAsync(list.WorkingDir, ct); return new MergeTargets(current, branches); } public async Task PreviewAsync(string taskId, string targetBranch, CancellationToken ct) { var (_, list, wt, _) = await LoadMergeContextAsync(taskId, ct); if (wt is null || wt.State != WorktreeState.Active) return new MergePreviewResult(PreviewUnavailable, Array.Empty(), 0); if (string.IsNullOrWhiteSpace(list.WorkingDir) || !await _git.IsGitRepoAsync(list.WorkingDir, ct)) return new MergePreviewResult(PreviewUnavailable, Array.Empty(), 0); var target = string.IsNullOrWhiteSpace(targetBranch) ? await _git.GetCurrentBranchAsync(list.WorkingDir, ct) : targetBranch; var preview = await _git.PreviewMergeAsync(list.WorkingDir, target, wt.BranchName, ct); if (!preview.Supported) return new MergePreviewResult(PreviewUnavailable, Array.Empty(), 0); if (!preview.Clean) return new MergePreviewResult(PreviewConflict, preview.ConflictFiles, 0); var count = await _git.CountChangedFilesAsync(list.WorkingDir, target, wt.BranchName, ct); return new MergePreviewResult(PreviewClean, Array.Empty(), count); } public Task ApproveAndMergeAsync(string taskId, string targetBranch, CancellationToken ct) => ApproveAndMergeAsync(taskId, targetBranch, leaveConflictsInTree: false, ct); public async Task ApproveAndMergeAsync( string taskId, string targetBranch, bool leaveConflictsInTree, CancellationToken ct) { var (task, list, wt, verifyCommand) = await LoadMergeContextAsync(taskId, ct); if (task.Status != TaskStatus.WaitingForReview) return Blocked("task is not waiting for review"); if (wt is null || wt.State != WorktreeState.Active) { // There is nothing left to merge -- a sandbox run, or a list-handler task that // committed straight into the list's working dir. The verify command still has to // pass before the task may reach Done: skipping it here would exempt exactly the // runs that land the most on the target branch at once. Same working dir and same // per-repo gate as the merge path, so a concurrent merge can't land mid-verify. if (!string.IsNullOrWhiteSpace(verifyCommand) && !string.IsNullOrWhiteSpace(list.WorkingDir)) { var verifyGate = GetMergeGate(list.WorkingDir!); await verifyGate.WaitAsync(ct); try { var failed = await RunVerifyGateAsync(verifyCommand, list.WorkingDir!, ct); if (failed is not null) return failed; } finally { verifyGate.Release(); } } var done = await _state.ApproveReviewAsync(taskId, ct); return done.Ok ? new MergeResult(StatusMerged, Array.Empty(), null) : Blocked(done.Reason ?? "approve failed"); } if (string.IsNullOrWhiteSpace(list.WorkingDir)) return Blocked("list has no working directory"); var target = string.IsNullOrWhiteSpace(targetBranch) ? await _git.GetCurrentBranchAsync(list.WorkingDir, ct) : targetBranch; // MergeAsync transitions the task WaitingForReview -> Done on a successful merge. // Remove the worktree on approve (matching the unit-merge path) so merged // worktrees don't pile up; the merge commit on the target branch is the record. return await MergeAsync(taskId, target, removeWorktree: true, $"Merge {wt.BranchName}", leaveConflictsInTree, ct); } private static MergeResult Blocked(string reason) => new(StatusBlocked, Array.Empty(), reason); private static RevertResult RevertBlocked(string reason) => new(StatusBlocked, null, Array.Empty(), reason); }