using System.Collections.Concurrent; using ClaudeDo.Data; using ClaudeDo.Data.Git; using ClaudeDo.Data.Models; using ClaudeDo.Worker.Hub; using ClaudeDo.Worker.Lifecycle; using ClaudeDo.Worker.State; using Microsoft.EntityFrameworkCore; using ModelContextProtocol; using TaskStatus = ClaudeDo.Data.Models.TaskStatus; namespace ClaudeDo.Worker.Planning; /// A unit-merge conflict currently paused, driven by an MCP session rather than the UI. public sealed record ExternalPlanningMergeConflict(string PlanningTaskId, string SubtaskId); /// Outcome of a unit-merge drive (StartAsync/ContinueAsync). Status mirrors /// 's status strings — StatusMerged/StatusConflict on the two /// "everything worked (so far)" paths, or the real blocked/verify_failed/untracked_collision /// status with its Reason when a child merge (or the final approve) failed outright. public sealed record PlanningMergeResult(string Status, string? Reason) { public static readonly PlanningMergeResult Merged = new(TaskMergeService.StatusMerged, null); public static readonly PlanningMergeResult Conflict = new(TaskMergeService.StatusConflict, null); } public sealed class PlanningMergeOrchestrator : IActiveMergeState { private readonly IDbContextFactory _dbFactory; private readonly TaskMergeService _merge; private readonly PlanningAggregator _aggregator; private readonly HubBroadcaster _broadcaster; private readonly GitService _git; private readonly ITaskStateService _state; private readonly ILogger _logger; private sealed class State { public required string TargetBranch { get; init; } public required Queue RemainingSubtaskIds { get; init; } public required bool IsPlanning { get; init; } public required string WorkingDir { get; init; } /// True when this unit merge was started by an MCP tool call (e.g. review_task /// from a running Claude session) rather than a direct UI action. Conflicts on such a /// merge must not auto-pop the in-app resolver — the driving session owns resolution. public required bool ExternallyDriven { get; init; } public string? CurrentSubtaskId { get; set; } /// True from the moment the last child has merged until FinalizeParentDoneAsync /// returns. CurrentSubtaskId is already null in this window (no subtask left to merge), /// so HasActiveMerge needs this separate flag — otherwise a Cancel racing the finalize call /// would slip past TaskStateService.CancelAsync's guard while the parent is still being /// flipped to Done. public bool IsFinalizing { get; set; } } private readonly ConcurrentDictionary _states = new(); public PlanningMergeOrchestrator( IDbContextFactory dbFactory, TaskMergeService merge, PlanningAggregator aggregator, HubBroadcaster broadcaster, GitService git, ITaskStateService state, ILogger logger) { _dbFactory = dbFactory; _merge = merge; _aggregator = aggregator; _broadcaster = broadcaster; _git = git; _state = state; _logger = logger; } public async Task StartAsync( string parentTaskId, string targetBranch, CancellationToken ct, bool externallyDriven = false, IProgress? progress = null) { string workingDir; List children; bool isPlanning; bool parentHasWorktree; TaskStatus parentStatus; using (var ctx = _dbFactory.CreateDbContext()) { var parent = await ctx.Tasks .Include(t => t.List) .Include(t => t.Worktree) .Include(t => t.Children).ThenInclude(c => c.Worktree) .SingleOrDefaultAsync(t => t.Id == parentTaskId, ct) ?? throw new KeyNotFoundException($"Planning task '{parentTaskId}' not found."); workingDir = parent.List.WorkingDir ?? throw new InvalidOperationException("List has no working directory."); children = parent.Children.OrderBy(c => c.SortOrder).ToList(); isPlanning = parent.PlanningPhase != PlanningPhase.None; parentHasWorktree = parent.Worktree is { State: WorktreeState.Active }; parentStatus = parent.Status; } if (isPlanning) { foreach (var c in children) { if (c.Status != TaskStatus.Done) throw new InvalidOperationException($"subtask {c.Id} is not Done (status {c.Status})"); if (c.Worktree is null) throw new InvalidOperationException($"subtask {c.Id} has no worktree"); if (c.Worktree.State != WorktreeState.Active && c.Worktree.State != WorktreeState.Merged) throw new InvalidOperationException( $"subtask {c.Id} worktree state is {c.Worktree.State}"); } } // Applies to planning AND improvement parents alike -- a stale UI click or a second // caller after the parent already left WaitingForReview (e.g. cancelled, or a previous // Approve already drove it to Done) must not kick off a partial re-merge of its children. if (parentStatus != TaskStatus.WaitingForReview) throw new InvalidOperationException( $"planning task '{parentTaskId}' is not WaitingForReview (status: {parentStatus})"); if (await _git.IsMidMergeAsync(workingDir, ct)) throw new InvalidOperationException( "repo is mid-merge; use AbortPlanningMerge to reset the repository, then Approve again"); if (await _git.HasChangesAsync(workingDir, includeUntracked: false, ct)) throw new InvalidOperationException("working tree has uncommitted changes"); var idsToMerge = new List(); if (!isPlanning && parentHasWorktree) idsToMerge.Add(parentTaskId); idsToMerge.AddRange( children .Where(c => c.Status == TaskStatus.Done && c.Worktree is { State: WorktreeState.Active }) .Select(c => c.Id)); var queue = new Queue(idsToMerge); var state = new State { TargetBranch = targetBranch, RemainingSubtaskIds = queue, IsPlanning = isPlanning, WorkingDir = workingDir, ExternallyDriven = externallyDriven, }; if (!_states.TryAdd(parentTaskId, state)) throw new InvalidOperationException($"Merge already in progress for {parentTaskId}."); await _broadcaster.PlanningMergeStarted(parentTaskId, targetBranch); return await DrainAsync(parentTaskId, ct, progress); } /// True when a unit merge for this parent is paused on a conflict, or is in the /// window between the last child merging and FinalizeParentDoneAsync completing. public bool HasActiveMerge(string parentTaskId) => _states.TryGetValue(parentTaskId, out var s) && (s.CurrentSubtaskId is not null || s.IsFinalizing); /// Externally-driven unit merges currently paused on a conflict, so the UI can /// recover this on reconnect instead of relying solely on the one-shot broadcast. Checked /// against rather than trusting the in-memory flag /// alone, so a stale entry (e.g. the repo was already resolved through another path) never /// reports an external merge that no longer exists. public async Task> GetActiveExternalConflictsAsync(CancellationToken ct) { var result = new List(); foreach (var (planningTaskId, state) in _states) { if (!state.ExternallyDriven || state.CurrentSubtaskId is null) continue; if (!await _git.IsMidMergeAsync(state.WorkingDir, ct)) continue; result.Add(new ExternalPlanningMergeConflict(planningTaskId, state.CurrentSubtaskId)); } return result; } public async Task ContinueAsync( string planningTaskId, CancellationToken ct, IProgress? progress = null) { if (!_states.TryGetValue(planningTaskId, out var state) || state.CurrentSubtaskId is null) throw new InvalidOperationException( "no in-progress merge to continue; if the worker was restarted during a conflict, use AbortPlanningMerge to reset the repository"); var current = state.CurrentSubtaskId; var result = await _merge.ContinueMergeAsync(current, ct, progress); if (result.Status == TaskMergeService.StatusConflict) { await _broadcaster.PlanningMergeConflict(planningTaskId, current, result.ConflictFiles, state.ExternallyDriven); return PlanningMergeResult.Conflict; } if (result.Status != TaskMergeService.StatusMerged) { _logger.LogWarning( "Planning continue blocked on subtask {Subtask}: {Msg}", current, result.ErrorMessage); _states.TryRemove(planningTaskId, out _); await _broadcaster.PlanningMergeAborted(planningTaskId, result.ErrorMessage); return new PlanningMergeResult(result.Status, result.ErrorMessage); } await _broadcaster.PlanningSubtaskMerged(planningTaskId, current); state.CurrentSubtaskId = null; return await DrainAsync(planningTaskId, ct, progress); } public async Task AbortAsync(string planningTaskId, CancellationToken ct) { if (!_states.TryGetValue(planningTaskId, out var state) || state.CurrentSubtaskId is null) { // No in-memory state — worker may have been restarted while a conflict was paused. // Check whether the list repo is still mid-merge and abort it directly. await AbortStatelessAsync(planningTaskId, ct); return; } await _merge.AbortMergeAsync(state.CurrentSubtaskId, ct); _states.TryRemove(planningTaskId, out _); await _broadcaster.PlanningMergeAborted(planningTaskId, "Merge aborted."); } private async Task AbortStatelessAsync(string planningTaskId, CancellationToken ct) { string? workingDir; await using (var ctx = _dbFactory.CreateDbContext()) { workingDir = await ctx.Tasks .Where(t => t.Id == planningTaskId) .Select(t => t.List.WorkingDir) .FirstOrDefaultAsync(ct); } if (string.IsNullOrWhiteSpace(workingDir) || !await _git.IsMidMergeAsync(workingDir, ct)) throw new InvalidOperationException("no in-progress merge to abort"); await _git.MergeAbortAsync(workingDir, ct); _logger.LogInformation( "Stateless abort of mid-merge for planning task {ParentId} (post-restart recovery)", planningTaskId); await _broadcaster.PlanningMergeAborted( planningTaskId, "Merge aborted after a worker restart. Approve again to restart the merge."); // Parent remains WaitingForReview — Approve will restart the unit merge from scratch. } private async Task DrainAsync( string planningTaskId, CancellationToken ct, IProgress? progress = null) { if (!_states.TryGetValue(planningTaskId, out var state)) return new PlanningMergeResult(TaskMergeService.StatusBlocked, "no merge state found for this planning task"); var keepState = false; try { while (state.RemainingSubtaskIds.TryDequeue(out var subtaskId)) { state.CurrentSubtaskId = subtaskId; var result = await _merge.MergeAsync( subtaskId, state.TargetBranch, removeWorktree: true, commitMessage: "", // blank -> TaskMergeService builds the conventional default leaveConflictsInTree: true, ct, progress); if (result.Status == TaskMergeService.StatusConflict) { await _broadcaster.PlanningMergeConflict( planningTaskId, subtaskId, result.ConflictFiles, state.ExternallyDriven); keepState = true; return PlanningMergeResult.Conflict; } if (result.Status != TaskMergeService.StatusMerged) { _logger.LogWarning( "Planning merge blocked on subtask {Subtask}: {Msg}", subtaskId, result.ErrorMessage); await _broadcaster.PlanningMergeAborted(planningTaskId, result.ErrorMessage); return new PlanningMergeResult(result.Status, result.ErrorMessage); // keepState stays false → finally removes the state entry } await _broadcaster.PlanningSubtaskMerged(planningTaskId, subtaskId); } // No subtask left to merge, but the parent isn't Done yet -- HasActiveMerge must keep // reporting true through this window (CurrentSubtaskId is already null) so a Cancel // racing FinalizeParentDoneAsync's ApproveReviewAsync call is still refused. state.CurrentSubtaskId = null; state.IsFinalizing = true; var (finalized, reason) = await FinalizeParentDoneAsync(planningTaskId, state.IsPlanning, ct); if (finalized) { await _broadcaster.PlanningCompleted(planningTaskId); return PlanningMergeResult.Merged; } return new PlanningMergeResult(TaskMergeService.StatusBlocked, reason); } finally { if (!keepState) _states.TryRemove(planningTaskId, out _); } } private async Task<(bool Ok, string? Reason)> FinalizeParentDoneAsync(string parentTaskId, bool isPlanning, CancellationToken ct) { var result = await _state.ApproveReviewAsync(parentTaskId, ct); if (!result.Ok) { // ApproveReviewAsync requires WaitingForReview. For improvement parents whose own // worktree is in the merge queue, TaskMergeService.ApproveIfWaitingForReviewAsync // already approved the parent during the drain — check for that expected path. await using var ctx = _dbFactory.CreateDbContext(); var current = await ctx.Tasks .Where(t => t.Id == parentTaskId) .Select(t => (TaskStatus?)t.Status) .FirstOrDefaultAsync(ct); if (current != TaskStatus.Done) { // Parent was cancelled or moved to an unexpected state during the merge drain. // Do not overwrite — the external transition takes precedence. _logger.LogWarning( "Unit-merge drain completed but parent {ParentTaskId} could not be finalized (status: {Status}): {Reason}", parentTaskId, current, result.Reason); return (false, result.Reason); } } // Only planning builds an integration branch via the aggregator; skip cleanup otherwise. if (isPlanning) { try { await _aggregator.CleanupIntegrationBranchAsync(parentTaskId, ct); } catch (Exception ex) { _logger.LogWarning(ex, "integration branch cleanup failed"); } } return (true, null); } }