use std::collections::{BTreeSet, HashMap};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, OnceLock};
use serde::{Deserialize, Serialize};
#[cfg(feature = "web-monitoring")]
use utoipa::ToSchema;
use super::same_path;
use crate::worktree_ops::service::SafetyFact;
#[allow(unused_imports)] pub use crate::vcs::git::commands::{
check_merge_conflicts, MergeSimulation, MAX_CONFLICT_SAMPLE, MAX_OUTPUT_PREFIX_BYTES,
};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[cfg_attr(feature = "web-monitoring", derive(ToSchema))]
#[serde(rename_all = "snake_case")]
pub enum InspectionState {
Checked,
Reused,
#[default]
NotInspected,
}
impl InspectionState {
pub fn is_inspected(self) -> bool {
!matches!(self, Self::NotInspected)
}
pub fn label(self) -> &'static str {
match self {
Self::Checked => "",
Self::Reused => "cached",
Self::NotInspected => "not inspected",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ObservationRequest {
Periodic,
Listing,
Target(PathBuf),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum InspectionScope {
Periodic {
eligible_branches: BTreeSet<String>,
},
Listing,
Target(PathBuf),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Admission {
Fresh,
Cacheable,
Skip,
}
#[derive(Debug, Clone, Copy)]
pub struct InspectionCandidate<'a> {
pub path: &'a Path,
pub branch: &'a str,
pub is_main: bool,
pub has_no_branch_identity: bool,
}
impl<'a> InspectionCandidate<'a> {
pub fn new(path: &'a Path, branch: &'a str, is_main: bool, is_detached: bool) -> Self {
Self {
path,
branch,
is_main,
has_no_branch_identity: is_detached || branch.is_empty(),
}
}
}
impl InspectionScope {
pub fn resolve(repo_root: &Path, request: &ObservationRequest) -> Self {
match request {
ObservationRequest::Listing => Self::Listing,
ObservationRequest::Target(path) => Self::Target(path.clone()),
ObservationRequest::Periodic => Self::Periodic {
eligible_branches: eligible_branches(repo_root),
},
}
}
pub fn admits(&self, candidate: &InspectionCandidate<'_>) -> Admission {
if candidate.is_main || candidate.has_no_branch_identity {
return Admission::Skip;
}
match self {
Self::Listing => Admission::Skip,
Self::Target(path) => {
if same_path(path, candidate.path) {
Admission::Fresh
} else {
Admission::Skip
}
}
Self::Periodic { eligible_branches } => {
if eligible_branches.contains(candidate.branch) {
return Admission::Cacheable;
}
match crate::vcs::GitWorkspaceManager::extract_change_id_from_worktree_name(
candidate.branch,
) {
Some(change_id) if eligible_branches.contains(&change_id) => {
Admission::Cacheable
}
_ => Admission::Skip,
}
}
}
}
}
pub fn eligible_branches(repo_root: &Path) -> BTreeSet<String> {
let active = crate::openspec::list_changes_native_from(repo_root).unwrap_or_default();
let rejected =
crate::openspec::list_rejected_changes_native_from(repo_root).unwrap_or_default();
let mut branches = BTreeSet::new();
for change in active.into_iter().chain(rejected) {
branches.insert(super::service::branch_name_for_change(&change.id));
branches.insert(change.id);
}
branches
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct ObservationKey {
pub repository: PathBuf,
pub branch: String,
pub base_head: String,
pub worktree_head: String,
pub merge_base: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Observation {
pub conflict_files: Vec<String>,
pub has_commits_ahead: SafetyFact,
}
pub const MAX_CACHED_OBSERVATIONS: usize = 512;
#[derive(Debug, Default)]
pub struct ObservationCache {
entries: Mutex<HashMap<ObservationKey, Observation>>,
}
impl ObservationCache {
pub fn new() -> Self {
Self::default()
}
pub fn get(&self, key: &ObservationKey) -> Option<Observation> {
self.lock().get(key).cloned()
}
pub fn insert(&self, key: ObservationKey, observation: Observation) {
let mut entries = self.lock();
if entries.len() >= MAX_CACHED_OBSERVATIONS && !entries.contains_key(&key) {
entries.clear();
}
entries.insert(key, observation);
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn len(&self) -> usize {
self.lock().len()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn clear(&self) {
self.lock().clear();
}
fn lock(&self) -> std::sync::MutexGuard<'_, HashMap<ObservationKey, Observation>> {
self.entries
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
}
pub fn shared_cache() -> &'static ObservationCache {
static SHARED: OnceLock<ObservationCache> = OnceLock::new();
SHARED.get_or_init(ObservationCache::new)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum InspectionCommand {
MergeBase,
Conflicts,
CommitsAhead,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct InspectionRecord {
pub command: InspectionCommand,
pub worktree: PathBuf,
}
static RECORDING: AtomicBool = AtomicBool::new(false);
fn records() -> &'static Mutex<Vec<InspectionRecord>> {
static RECORDS: OnceLock<Mutex<Vec<InspectionRecord>>> = OnceLock::new();
RECORDS.get_or_init(|| Mutex::new(Vec::new()))
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn record_inspection_commands() {
RECORDING.store(true, Ordering::Relaxed);
}
pub(crate) fn record(command: InspectionCommand, worktree: &Path) {
if !RECORDING.load(Ordering::Relaxed) {
return;
}
records()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.push(InspectionRecord {
command,
worktree: worktree.to_path_buf(),
});
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn recorded_inspections_under(prefix: &Path) -> Vec<InspectionRecord> {
let prefix = std::fs::canonicalize(prefix).unwrap_or_else(|_| prefix.to_path_buf());
records()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.iter()
.filter(|record| {
let path =
std::fs::canonicalize(&record.worktree).unwrap_or_else(|_| record.worktree.clone());
path.starts_with(&prefix)
})
.cloned()
.collect()
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn recorded_inspection_count(prefix: &Path, command: InspectionCommand) -> usize {
recorded_inspections_under(prefix)
.into_iter()
.filter(|record| record.command == command)
.count()
}
#[cfg(test)]
mod tests {
use super::*;
fn periodic(eligible: &[&str]) -> InspectionScope {
InspectionScope::Periodic {
eligible_branches: eligible.iter().map(|id| (*id).to_string()).collect(),
}
}
fn candidate<'a>(path: &'a Path, branch: &'a str) -> InspectionCandidate<'a> {
InspectionCandidate::new(path, branch, false, false)
}
fn key(branch: &str, base: &str, head: &str, merge_base: &str) -> ObservationKey {
ObservationKey {
repository: PathBuf::from("/repo"),
branch: branch.to_string(),
base_head: base.to_string(),
worktree_head: head.to_string(),
merge_base: merge_base.to_string(),
}
}
fn observation() -> Observation {
Observation {
conflict_files: Vec::new(),
has_commits_ahead: SafetyFact::Yes,
}
}
#[test]
fn periodic_refresh_inspects_only_current_change_branches() {
let scope = periodic(&["live-change"]);
let live = PathBuf::from("/w/live-change");
let stale = PathBuf::from("/w/stale-change");
assert_eq!(
scope.admits(&candidate(&live, "live-change")),
Admission::Cacheable,
"a branch naming a current change is the case automatic inspection exists for"
);
assert_eq!(
scope.admits(&candidate(&stale, "stale-change")),
Admission::Skip,
"a branch that maps to no current change must not cost a merge simulation"
);
assert_eq!(
scope.admits(&candidate(&stale, "ws-session-a1b2c3")),
Admission::Skip,
"a session branch maps to no change and is periodically ineligible"
);
}
#[test]
fn legacy_prefixed_branches_resolve_to_their_change_id() {
let scope = periodic(&["live-change"]);
let path = PathBuf::from("/w/legacy");
assert_eq!(
scope.admits(&candidate(&path, "ws-live-change-a1b2c3d4")),
Admission::Cacheable,
"the legacy ws-<id>-<hex> branch form names the same change"
);
}
#[test]
fn structural_rows_never_spawn_inspection_commands() {
let scope = periodic(&["live-change"]);
let path = PathBuf::from("/w/live-change");
assert_eq!(
scope.admits(&InspectionCandidate::new(&path, "main", true, false)),
Admission::Skip,
"the main worktree is the base; there is nothing to simulate merging it into"
);
assert_eq!(
scope.admits(&InspectionCandidate::new(&path, "", false, true)),
Admission::Skip,
"a detached worktree has no branch identity to compare"
);
assert_eq!(
scope.admits(&InspectionCandidate::new(&path, "", false, false)),
Admission::Skip,
"a branchless worktree has no branch identity to compare"
);
}
#[test]
fn listing_scope_inspects_nothing_at_all() {
let path = PathBuf::from("/w/live-change");
assert_eq!(
InspectionScope::Listing.admits(&candidate(&path, "live-change")),
Admission::Skip,
"a structural listing must not pay for merge simulation even for a live change"
);
}
#[test]
fn an_operator_target_is_inspected_fresh_even_when_periodic_refresh_would_skip_it() {
let target = PathBuf::from("/w/ws-session-a1b2c3");
let other = PathBuf::from("/w/live-change");
let scope = InspectionScope::Target(target.clone());
assert_eq!(
scope.admits(&candidate(&target, "ws-session-a1b2c3")),
Admission::Fresh,
"an operator-addressed worktree is inspected from current evidence, mapped or not"
);
assert_eq!(
scope.admits(&candidate(&other, "live-change")),
Admission::Skip,
"a targeted observation pays for its target only"
);
}
#[test]
fn a_cached_observation_is_reused_only_for_an_identical_revision_tuple() {
let cache = ObservationCache::new();
let stored = key("live-change", "base1", "head1", "merge1");
cache.insert(stored.clone(), observation());
assert_eq!(cache.get(&stored), Some(observation()));
for (name, changed) in [
("branch", key("other-change", "base1", "head1", "merge1")),
("base head", key("live-change", "base2", "head1", "merge1")),
(
"worktree head",
key("live-change", "base1", "head2", "merge1"),
),
("merge base", key("live-change", "base1", "head1", "merge2")),
] {
assert_eq!(
cache.get(&changed),
None,
"a changed {name} must invalidate the observation rather than reuse it"
);
}
let mut other_repository = stored.clone();
other_repository.repository = PathBuf::from("/other-repo");
assert_eq!(
cache.get(&other_repository),
None,
"two repositories holding the same branch at the same commit are different observations"
);
}
#[test]
fn the_cache_is_bounded_and_disposable() {
let cache = ObservationCache::new();
for index in 0..(MAX_CACHED_OBSERVATIONS + 1) {
cache.insert(
key("live-change", "base", &format!("head{index}"), "merge"),
observation(),
);
}
assert!(
cache.len() <= MAX_CACHED_OBSERVATIONS,
"an unbounded cache would grow for the life of the process"
);
cache.clear();
assert!(
cache.is_empty(),
"the cache must be discardable at any time"
);
}
#[test]
fn inspection_state_separates_looked_from_found_nothing() {
assert!(InspectionState::Checked.is_inspected());
assert!(InspectionState::Reused.is_inspected());
assert!(!InspectionState::NotInspected.is_inspected());
assert_eq!(InspectionState::default(), InspectionState::NotInspected);
assert_eq!(InspectionState::Reused.label(), "cached");
assert_eq!(InspectionState::NotInspected.label(), "not inspected");
}
}