use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use chrono::{DateTime, Utc};
use crate::events::ExecutionEvent;
use crate::tui::types::WorktreeInfo;
use crate::web::remote_control_api::dto::{
AttentionState, ChangeActivity, ChangeTiming, ChangeWorktree, ParallelBlockedReason,
ParallelEligibility,
};
use crate::web::remote_control_api::projection::describe_event;
use crate::web::remote_control_api::worktrees::repository_relative_display;
const NON_ACTIVITY_EVENT_TYPES: [&str; 6] = [
"log",
"apply_output",
"archive_output",
"acceptance_output",
"resolve_output",
"progress_updated",
];
const RUN_START_EVENT_TYPES: [&str; 5] = [
"apply_started",
"acceptance_started",
"archive_started",
"resolve_started",
"push_started",
];
const RUN_END_EVENT_TYPES: [&str; 10] = [
"processing_error",
"apply_failed",
"acceptance_failed",
"archive_failed",
"change_archived",
"change_rejected",
"resolve_completed",
"resolve_failed",
"merge_completed",
"push_completed",
];
#[derive(Debug, Clone, Default)]
pub struct ChangeOperatorFacts {
pub newly_detected: bool,
pub parallel: ParallelEligibility,
pub worktree: Option<ChangeWorktree>,
pub timing: ChangeTiming,
pub latest_activity: Option<ChangeActivity>,
started_at: Option<DateTime<Utc>>,
}
impl ChangeOperatorFacts {
pub fn attention(&self, execution_marked: bool, queued: bool) -> AttentionState {
if self.newly_detected
&& !execution_marked
&& !queued
&& self.latest_activity.is_none()
&& self.timing.started_at.is_none()
{
AttentionState::New
} else {
AttentionState::None
}
}
}
#[derive(Debug, Default)]
pub struct OperatorFactsStore {
facts: HashMap<String, ChangeOperatorFacts>,
known_change_ids: HashSet<String>,
bootstrapped: bool,
process_error: Option<String>,
repo_root: Option<PathBuf>,
worktree_paths: HashMap<String, PathBuf>,
worktrees: Vec<WorktreeInfo>,
}
impl OperatorFactsStore {
pub fn new() -> Self {
Self::default()
}
pub fn set_repo_root(&mut self, repo_root: PathBuf) {
self.repo_root = Some(repo_root);
self.reproject_worktrees();
}
pub fn repo_root(&self) -> Option<PathBuf> {
self.repo_root.clone()
}
pub fn facts(&self, change_id: &str) -> ChangeOperatorFacts {
self.facts.get(change_id).cloned().unwrap_or_default()
}
pub fn process_error(&self) -> Option<String> {
self.process_error.clone()
}
pub fn observe_changes<'a>(&mut self, change_ids: impl IntoIterator<Item = &'a str>) {
let observed: HashSet<String> = change_ids.into_iter().map(str::to_string).collect();
for id in &observed {
if self.known_change_ids.insert(id.clone()) && self.bootstrapped {
self.facts.entry(id.clone()).or_default().newly_detected = true;
}
}
self.known_change_ids.retain(|id| observed.contains(id));
self.facts.retain(|id, _| observed.contains(id));
self.bootstrapped = true;
for id in observed {
self.facts.entry(id).or_default();
}
self.reproject_worktrees();
}
pub fn apply_parallel_eligibility(
&mut self,
committed_change_ids: &HashSet<String>,
uncommitted_file_change_ids: &HashSet<String>,
) {
use crate::orchestration::operator_command::ParallelEligibility as Observed;
for (id, facts) in &mut self.facts {
facts.parallel =
match Observed::observe(id, committed_change_ids, uncommitted_file_change_ids) {
Observed::Eligible => ParallelEligibility::default(),
Observed::ProposalAbsentFromHead => ParallelEligibility {
eligible: false,
blocked_reason: Some(ParallelBlockedReason::NotCommitted),
},
Observed::UncommittedProposalFiles => ParallelEligibility {
eligible: false,
blocked_reason: Some(ParallelBlockedReason::UncommittedChanges),
},
};
}
}
pub fn apply_worktree_paths(&mut self, worktree_paths: HashMap<String, PathBuf>) {
self.worktree_paths = worktree_paths;
self.reproject_worktrees();
}
pub fn apply_worktrees(&mut self, worktrees: Vec<WorktreeInfo>) {
self.worktrees = worktrees;
self.reproject_worktrees();
}
fn reproject_worktrees(&mut self) {
let Self {
facts,
repo_root,
worktree_paths,
worktrees,
..
} = self;
for (id, entry) in facts.iter_mut() {
entry.worktree = worktree_paths
.get(id)
.map(|path| Self::project_worktree(repo_root.as_deref(), path, worktrees));
}
}
fn project_worktree(
repo_root: Option<&Path>,
path: &Path,
worktrees: &[WorktreeInfo],
) -> ChangeWorktree {
let observed = worktrees.iter().find(|worktree| worktree.path == path);
ChangeWorktree {
path: match repo_root {
Some(root) => repository_relative_display(root, path),
None => path
.file_name()
.map(|name| name.to_string_lossy().to_string())
.unwrap_or_default(),
},
branch: observed
.map(|worktree| worktree.branch.clone())
.filter(|branch| !branch.is_empty()),
}
}
pub fn record_event(&mut self, event: &ExecutionEvent) {
let (event_type, change_id, _) = describe_event(event);
let detail = match event {
ExecutionEvent::ProcessingError { error, .. }
| ExecutionEvent::ApplyFailed { error, .. }
| ExecutionEvent::AcceptanceFailed { error, .. }
| ExecutionEvent::ArchiveFailed { error, .. }
| ExecutionEvent::ResolveFailed { error, .. }
| ExecutionEvent::PushFailed { error, .. }
| ExecutionEvent::RejectionReviewFailed { error, .. } => Some(error.clone()),
ExecutionEvent::ChangeRejected { reason, .. } => Some(reason.clone()),
ExecutionEvent::MergeDeferred { reason, .. } => Some(reason.clone()),
_ => None,
};
match event {
ExecutionEvent::Error { message } => {
self.process_error = Some(crate::events::sanitize_detail(message));
return;
}
ExecutionEvent::ProcessingStarted(_) => self.process_error = None,
_ => {}
}
let Some(change_id) = change_id else {
return;
};
let now = Utc::now();
let facts = self.facts.entry(change_id).or_default();
if event_type == "processing_started" {
facts.started_at = Some(now);
facts.timing = ChangeTiming {
started_at: Some(now.to_rfc3339()),
completed_at: None,
elapsed_ms: None,
};
} else if RUN_START_EVENT_TYPES.contains(&event_type) && facts.started_at.is_none() {
facts.started_at = Some(now);
facts.timing.started_at = Some(now.to_rfc3339());
facts.timing.completed_at = None;
facts.timing.elapsed_ms = None;
}
if RUN_END_EVENT_TYPES.contains(&event_type) {
if let Some(started) = facts.started_at.take() {
facts.timing.completed_at = Some(now.to_rfc3339());
facts.timing.elapsed_ms =
Some((now - started).num_milliseconds().max(0).unsigned_abs());
}
}
let repeats_current_activity = facts.latest_activity.as_ref().is_some_and(|current| {
current.event_type == event_type
&& current.detail == detail.as_deref().map(crate::events::sanitize_detail)
});
if !NON_ACTIVITY_EVENT_TYPES.contains(&event_type) && !repeats_current_activity {
facts.latest_activity = Some(ChangeActivity {
event_type: event_type.to_string(),
timestamp: now.to_rfc3339(),
detail: detail.as_deref().map(crate::events::sanitize_detail),
});
}
}
}