use std::collections::BTreeMap;
use std::fs;
use std::io;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender, channel};
use std::sync::{Arc, Mutex};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
use cargoless_proto::{
BuildIdentity, CheckResult, ContentHash, Diagnostic, FileState, Profile, StateEvent,
TargetTriple, TreeState,
};
use crate::lsp::LspEvent;
const CHECK_HARD_CAP: Duration = Duration::from_secs(180);
const CHECK_SETTLE: Duration = Duration::from_secs(2);
const DEFAULT_WATCH_DEBOUNCE: Duration = Duration::from_millis(150);
fn resolve_watch_debounce() -> Duration {
std::env::var("TF_DEBOUNCE_MS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.filter(|&ms| ms > 0)
.map(Duration::from_millis)
.unwrap_or(DEFAULT_WATCH_DEBOUNCE)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum VerdictProvenance {
Authoritative,
Advisory,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Verdict {
pub tree: TreeState,
pub provenance: VerdictProvenance,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum LifecycleEvent {
AnalyzerRestarting,
}
pub trait IdentityProvider: Send {
fn current_identity(&self) -> BuildIdentity;
}
impl<F> IdentityProvider for F
where
F: Fn() -> BuildIdentity + Send,
{
fn current_identity(&self) -> BuildIdentity {
self()
}
}
pub struct Model {
auth: BTreeMap<String, FileState>,
native: BTreeMap<String, FileState>,
flycheck_done: bool,
tree: TreeState,
subscribers: Vec<Sender<StateEvent>>,
advisory_subscribers: Vec<Sender<Verdict>>,
identity: Box<dyn IdentityProvider>,
diagnostics: BTreeMap<String, Vec<Diagnostic>>,
lifecycle_subscribers: Vec<Sender<LifecycleEvent>>,
procmacro_downranked: u64,
}
fn file_state_for(pd: &crate::lsp::PublishDiagnostics, downrank: bool) -> (FileState, bool) {
if downrank {
let st = if pd.has_authoritative_error() {
FileState::Red
} else {
FileState::Green
};
let suppressed = pd.has_any_severity_error() && !pd.has_authoritative_error();
(st, suppressed)
} else {
let st = if pd.has_any_severity_error() {
FileState::Red
} else {
FileState::Green
};
(st, false)
}
}
impl Model {
pub fn new<I: IdentityProvider + 'static>(identity: I) -> Self {
Self {
auth: BTreeMap::new(),
native: BTreeMap::new(),
flycheck_done: false,
tree: TreeState::Red,
subscribers: Vec::new(),
advisory_subscribers: Vec::new(),
identity: Box::new(identity),
diagnostics: BTreeMap::new(),
lifecycle_subscribers: Vec::new(),
procmacro_downranked: 0,
}
}
pub fn subscribe(&mut self) -> Receiver<StateEvent> {
let (tx, rx) = channel();
self.subscribers.push(tx);
rx
}
pub fn subscribe_advisory(&mut self) -> Receiver<Verdict> {
let (tx, rx) = channel();
self.advisory_subscribers.push(tx);
rx
}
pub fn subscribe_lifecycle(&mut self) -> Receiver<LifecycleEvent> {
let (tx, rx) = channel();
self.lifecycle_subscribers.push(tx);
rx
}
pub(crate) fn emit_lifecycle(&mut self, ev: LifecycleEvent) {
self.lifecycle_subscribers.retain(|s| s.send(ev).is_ok());
}
pub fn tree_state(&self) -> TreeState {
self.tree
}
pub fn flycheck_done(&self) -> bool {
self.flycheck_done
}
pub fn procmacro_downranked(&self) -> u64 {
self.procmacro_downranked
}
pub fn last_verdict(&self) -> Verdict {
Verdict {
tree: self.tree,
provenance: if self.flycheck_done {
VerdictProvenance::Authoritative
} else {
VerdictProvenance::Advisory
},
}
}
pub fn file_state(&self, path: &str) -> Option<FileState> {
self.auth.get(path).copied()
}
pub fn file_diagnostics(&self, path: &str) -> &[Diagnostic] {
self.diagnostics.get(path).map(Vec::as_slice).unwrap_or(&[])
}
pub fn all_diagnostics(&self) -> Vec<Diagnostic> {
let total: usize = self.diagnostics.values().map(Vec::len).sum();
let mut out = Vec::with_capacity(total);
for v in self.diagnostics.values() {
out.extend(v.iter().cloned());
}
out
}
pub fn apply_event(&mut self, ev: &LspEvent) {
match ev {
LspEvent::Diagnostics(pd) => {
let Some(path) = crate::lsp::path_from_uri(&pd.uri) else {
return;
};
let (file_state, downranked) = file_state_for(pd, crate::procmacro::enabled());
if downranked {
self.procmacro_downranked += 1;
}
let native_state = if pd.advisory_errors > 0 {
FileState::Red
} else {
FileState::Green
};
self.auth.insert(path.clone(), file_state);
self.native.insert(path.clone(), native_state);
if pd.diagnostics.is_empty() {
self.diagnostics.remove(&path);
} else {
self.diagnostics
.insert(path.clone(), pd.diagnostics.clone());
}
self.emit(StateEvent::FileVerdict {
path,
state: file_state,
});
self.reconcile();
self.emit_advisory();
}
LspEvent::FlycheckEnded => {
self.flycheck_done = true;
self.reconcile();
self.emit_advisory();
}
LspEvent::IndexingEnded => {
}
}
}
pub fn forget_file(&mut self, path: &str) {
let a = self.auth.remove(path).is_some();
let n = self.native.remove(path).is_some();
let d = self.diagnostics.remove(path).is_some();
if a || n || d {
self.reconcile();
self.emit_advisory();
}
}
fn authoritative_tree(&self) -> TreeState {
if !self.flycheck_done {
return TreeState::Red;
}
if self.auth.values().any(|s| *s == FileState::Red) {
TreeState::Red
} else {
TreeState::Green
}
}
fn reconcile(&mut self) {
let next = self.authoritative_tree();
if next == self.tree {
return;
}
self.tree = next;
match next {
TreeState::Green => {
let identity = self.identity.current_identity();
self.emit(StateEvent::BecameGreen { identity });
}
TreeState::Red => self.emit(StateEvent::BecameRed),
}
}
fn emit(&mut self, ev: StateEvent) {
self.subscribers.retain(|s| s.send(ev.clone()).is_ok());
}
fn emit_advisory(&mut self) {
let v = self.last_verdict();
self.advisory_subscribers.retain(|s| s.send(v).is_ok());
}
}
fn poisoned<T>(m: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
m.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn collect_rs_files(root: &Path, ignore: &crate::watcher::IgnoreRules) -> Vec<PathBuf> {
let mut out = Vec::new();
let mut stack = vec![root.to_path_buf()];
while let Some(dir) = stack.pop() {
let Ok(entries) = fs::read_dir(&dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
let rel = path.strip_prefix(root).unwrap_or(&path);
if ignore.is_ignored(rel) {
continue;
}
match entry.file_type() {
Ok(ft) if ft.is_dir() => stack.push(path),
Ok(ft) if ft.is_file() => {
if path.extension().is_some_and(|e| e == "rs") {
out.push(path);
}
}
_ => {}
}
}
}
out
}
pub fn placeholder_identity() -> BuildIdentity {
let sentinel = ContentHash::new("placeholder-display-only-not-a-build-key");
BuildIdentity {
source_tree: sentinel.clone(),
cargo_lock: sentinel.clone(),
rust_toolchain: sentinel.clone(),
tf_config: sentinel,
target: TargetTriple::new("wasm32-unknown-unknown"),
profile: Profile::Dev,
}
}
pub fn check_verdict(root: &Path) -> io::Result<Verdict> {
let root = fs::canonicalize(root)?;
let mut cmd = crate::analyzer::rust_analyzer_command()?;
cmd.current_dir(&root);
let mut guard = crate::analyzer::ReapOnDrop::new(cmd.spawn()?);
let (stdin, stdout) = guard
.take_stdio()
.ok_or_else(|| io::Error::other("rust-analyzer stdio unavailable"))?;
let root_str = root.to_string_lossy().into_owned();
let init_opts = crate::lsp::InitOpts::from_env_and_project(&root);
let (client, events) = crate::lsp::LspClient::initialize(stdin, stdout, &root_str, &init_opts)?;
let ignore = crate::watcher::IgnoreRules::for_root(&root);
for f in collect_rs_files(&root, &ignore) {
if let Ok(text) = fs::read_to_string(&f) {
let _ = client.did_open(&f.to_string_lossy(), &text, 1);
}
}
let cap = std::env::var("TF_CHECK_TIMEOUT_SECS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.map(Duration::from_secs)
.unwrap_or(CHECK_HARD_CAP);
let deadline = Instant::now() + cap;
let mut auth: BTreeMap<String, FileState> = BTreeMap::new();
let mut flycheck_seen = false;
let mut got_any = false;
let mut indexing_done = false;
while !flycheck_seen {
let now = Instant::now();
if now >= deadline {
break;
}
let wait = CHECK_SETTLE.min(deadline - now);
match events.recv_timeout(wait) {
Ok(LspEvent::Diagnostics(pd)) => {
got_any = true;
if let Some(p) = crate::lsp::path_from_uri(&pd.uri) {
let s = if pd.has_any_severity_error() {
FileState::Red
} else {
FileState::Green
};
auth.insert(p, s);
}
}
Ok(LspEvent::FlycheckEnded) => {
flycheck_seen = true;
}
Ok(LspEvent::IndexingEnded) => {
indexing_done = true;
}
Err(RecvTimeoutError::Timeout) => {
if got_any && indexing_done {
break;
}
}
Err(RecvTimeoutError::Disconnected) => break, }
}
drop(guard);
if flycheck_seen {
let tree = if auth.values().any(|s| *s == FileState::Red) {
TreeState::Red
} else {
TreeState::Green
};
Ok(Verdict {
tree,
provenance: VerdictProvenance::Authoritative,
})
} else {
Ok(Verdict {
tree: TreeState::Red,
provenance: VerdictProvenance::Advisory,
})
}
}
pub fn check_once(root: &Path) -> io::Result<TreeState> {
check_verdict(root).map(|v| v.tree)
}
pub fn check_once_with_diagnostics(root: &Path) -> io::Result<CheckResult> {
let root = fs::canonicalize(root)?;
let mut cmd = crate::analyzer::rust_analyzer_command()?;
cmd.current_dir(&root);
let mut guard = crate::analyzer::ReapOnDrop::new(cmd.spawn()?);
let (stdin, stdout) = guard
.take_stdio()
.ok_or_else(|| io::Error::other("rust-analyzer stdio unavailable"))?;
let root_str = root.to_string_lossy().into_owned();
let init_opts = crate::lsp::InitOpts::from_env_and_project(&root);
let (client, events) = crate::lsp::LspClient::initialize(stdin, stdout, &root_str, &init_opts)?;
let ignore = crate::watcher::IgnoreRules::for_root(&root);
for f in collect_rs_files(&root, &ignore) {
if let Ok(text) = fs::read_to_string(&f) {
let _ = client.did_open(&f.to_string_lossy(), &text, 1);
}
}
let cap = std::env::var("TF_CHECK_TIMEOUT_SECS")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.map(Duration::from_secs)
.unwrap_or(CHECK_HARD_CAP);
let deadline = Instant::now() + cap;
let mut auth: BTreeMap<String, FileState> = BTreeMap::new();
let mut diagnostics: BTreeMap<String, Vec<Diagnostic>> = BTreeMap::new();
let mut flycheck_seen = false;
let mut got_any = false;
let mut indexing_done = false;
while !flycheck_seen {
let now = Instant::now();
if now >= deadline {
break;
}
let wait = CHECK_SETTLE.min(deadline - now);
match events.recv_timeout(wait) {
Ok(LspEvent::Diagnostics(pd)) => {
got_any = true;
if let Some(p) = crate::lsp::path_from_uri(&pd.uri) {
let s = if pd.has_any_severity_error() {
FileState::Red
} else {
FileState::Green
};
auth.insert(p.clone(), s);
if pd.diagnostics.is_empty() {
diagnostics.remove(&p);
} else {
diagnostics.insert(p, pd.diagnostics.clone());
}
}
}
Ok(LspEvent::FlycheckEnded) => {
flycheck_seen = true;
}
Ok(LspEvent::IndexingEnded) => {
indexing_done = true;
}
Err(RecvTimeoutError::Timeout) => {
if got_any && indexing_done {
break; }
}
Err(RecvTimeoutError::Disconnected) => break, }
}
drop(guard);
let tree = if flycheck_seen {
if auth.values().any(|s| *s == FileState::Red) {
TreeState::Red
} else {
TreeState::Green
}
} else {
TreeState::Red
};
let total: usize = diagnostics.values().map(Vec::len).sum();
let mut flat = Vec::with_capacity(total);
for v in diagnostics.values() {
flat.extend(v.iter().cloned());
}
Ok(CheckResult {
tree,
diagnostics: flat,
})
}
pub struct ModelSession {
model: Arc<Mutex<Model>>,
stop: Arc<AtomicBool>,
supervisor: Option<crate::analyzer::Supervisor>,
watch: Option<crate::watcher::WatchHandle>,
threads: Vec<JoinHandle<()>>,
structural_counters: Arc<crate::structural::StructuralCounters>,
idle_evict_counters: Arc<crate::idle::IdleEvictCounters>,
}
impl ModelSession {
pub fn subscribe(&self) -> Receiver<StateEvent> {
poisoned(&self.model).subscribe()
}
pub fn subscribe_advisory(&self) -> Receiver<Verdict> {
poisoned(&self.model).subscribe_advisory()
}
pub fn tree_state(&self) -> TreeState {
poisoned(&self.model).tree_state()
}
pub fn last_verdict(&self) -> Verdict {
poisoned(&self.model).last_verdict()
}
pub fn current_diagnostics(&self) -> Vec<Diagnostic> {
poisoned(&self.model).all_diagnostics()
}
pub fn subscribe_lifecycle(&self) -> Receiver<LifecycleEvent> {
poisoned(&self.model).subscribe_lifecycle()
}
pub fn structural_counters(&self) -> (u64, u64) {
self.structural_counters.snapshot()
}
pub fn idle_evict_counters(&self) -> (u64, u64) {
self.idle_evict_counters.snapshot()
}
pub fn procmacro_downranked(&self) -> u64 {
poisoned(&self.model).procmacro_downranked()
}
pub fn shutdown(mut self) {
self.do_shutdown();
}
fn do_shutdown(&mut self) {
self.stop.store(true, Ordering::SeqCst);
if let Some(sup) = self.supervisor.take() {
sup.shutdown();
}
drop(self.watch.take());
for t in self.threads.drain(..) {
let _ = t.join();
}
}
}
impl Drop for ModelSession {
fn drop(&mut self) {
self.do_shutdown();
}
}
pub fn watch<I: IdentityProvider + 'static>(
root: &Path,
identity: I,
) -> io::Result<(ModelSession, Receiver<StateEvent>)> {
let root = fs::canonicalize(root)?;
let root_str = root.to_string_lossy().into_owned();
let model = Arc::new(Mutex::new(Model::new(identity)));
let events = poisoned(&model).subscribe();
let stop = Arc::new(AtomicBool::new(false));
let current: Arc<Mutex<Option<Arc<crate::lsp::LspClient>>>> = Arc::new(Mutex::new(None));
let spawn_root = root.clone();
let spawn = move || {
let mut cmd = crate::analyzer::rust_analyzer_command()?;
cmd.current_dir(&spawn_root);
cmd.spawn()
};
let hook_root = root_str.clone();
let hook_model = Arc::clone(&model);
let hook_current = Arc::clone(¤t);
let spawn_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let hook_spawn_count = Arc::clone(&spawn_count);
let on_spawn = move |child: &mut std::process::Child| {
let n = hook_spawn_count.fetch_add(1, Ordering::SeqCst);
if n > 0 {
poisoned(&hook_model).emit_lifecycle(LifecycleEvent::AnalyzerRestarting);
}
let (Some(stdin), Some(stdout)) = (child.stdin.take(), child.stdout.take()) else {
return;
};
let init_opts = crate::lsp::InitOpts::from_env_and_project(Path::new(&hook_root));
let Ok((client, events)) =
crate::lsp::LspClient::initialize(stdin, stdout, &hook_root, &init_opts)
else {
return; };
let client = Arc::new(client);
let ig = crate::watcher::IgnoreRules::for_root(Path::new(&hook_root));
for f in collect_rs_files(Path::new(&hook_root), &ig) {
if let Ok(text) = fs::read_to_string(&f) {
let _ = client.did_open(&f.to_string_lossy(), &text, 1);
}
}
*poisoned(&hook_current) = Some(Arc::clone(&client));
let m = Arc::clone(&hook_model);
let _ = thread::Builder::new()
.name("tf-model-events".into())
.spawn(move || {
while let Ok(ev) = events.recv() {
poisoned(&m).apply_event(&ev);
}
});
};
let supervisor = crate::analyzer::Supervisor::start_with_hook(spawn, on_spawn)?;
let suspend_handle = supervisor.suspend_handle();
let (watch_handle, batches) =
crate::watcher::watch(&root, resolve_watch_debounce()).map_err(io::Error::other)?;
let structural_counters = Arc::new(crate::structural::StructuralCounters::new());
let idle_counters = Arc::new(crate::idle::IdleEvictCounters::new());
let mut threads = Vec::new();
{
let model = Arc::clone(&model);
let stop = Arc::clone(&stop);
let current = Arc::clone(¤t);
let structural = Arc::clone(&structural_counters);
let idle_counters = Arc::clone(&idle_counters);
let suspend_handle = suspend_handle.clone();
threads.push(
thread::Builder::new()
.name("tf-model-fs".into())
.spawn(move || {
let mut version: i64 = 2;
let mut last_activity = std::time::Instant::now();
let mut suspended_since: Option<std::time::Instant> = None;
loop {
match batches.recv_timeout(Duration::from_millis(250)) {
Ok(batch) => {
if crate::idle::enabled() && suspend_handle.is_suspended() {
suspend_handle.resume();
let deadline =
std::time::Instant::now() + Duration::from_secs(35);
while !suspend_handle.child_alive()
&& std::time::Instant::now() < deadline
{
thread::sleep(Duration::from_millis(50));
}
if let Some(since) = suspended_since.take() {
idle_counters.add_suspended(since.elapsed());
}
}
last_activity = std::time::Instant::now();
let client = poisoned(¤t).as_ref().cloned();
if crate::structural::enabled() {
let mut files: Vec<(String, String)> = Vec::new();
for path in batch {
if path.extension().is_none_or(|e| e != "rs") {
continue;
}
let p = path.to_string_lossy().into_owned();
match fs::read_to_string(&path) {
Ok(text) => files.push((p, text)),
Err(_) => poisoned(&model).forget_file(&p),
}
}
let all_closed =
files.iter().all(|(_, t)| crate::structural::is_closed(t));
structural.record(all_closed);
for (p, text) in &files {
version += 1;
if let Some(c) = client.as_ref() {
let _ = c.did_change(p, text, version);
if all_closed {
let _ = c.did_save(p);
}
}
}
} else {
for path in batch {
if path.extension().is_none_or(|e| e != "rs") {
continue;
}
let p = path.to_string_lossy().into_owned();
match fs::read_to_string(&path) {
Ok(text) => {
version += 1;
if let Some(c) = client.as_ref() {
let _ = c.did_change(&p, &text, version);
let _ = c.did_save(&p);
}
}
Err(_) => poisoned(&model).forget_file(&p),
}
}
}
}
Err(RecvTimeoutError::Timeout) => {
if stop.load(Ordering::SeqCst) {
break;
}
if crate::idle::enabled()
&& !suspend_handle.is_suspended()
&& poisoned(&model).flycheck_done()
&& last_activity.elapsed() >= crate::idle::idle_window()
{
suspend_handle.suspend();
idle_counters.record_eviction();
suspended_since = Some(std::time::Instant::now());
}
}
Err(RecvTimeoutError::Disconnected) => break,
}
}
})
.expect("spawn tf-model-fs"),
);
}
let session = ModelSession {
model,
stop,
supervisor: Some(supervisor),
watch: Some(watch_handle),
threads,
structural_counters,
idle_evict_counters: idle_counters,
};
Ok((session, events))
}
#[cfg(test)]
mod tests {
use super::*;
fn ident() -> BuildIdentity {
BuildIdentity {
source_tree: ContentHash::new("src"),
cargo_lock: ContentHash::new("lock"),
rust_toolchain: ContentHash::new("tc"),
tf_config: ContentHash::new("cfg"),
target: TargetTriple::new("wasm32-unknown-unknown"),
profile: Profile::Dev,
}
}
fn model() -> Model {
Model::new(ident)
}
fn diag(uri: &str, auth_err: usize, adv_err: usize) -> LspEvent {
LspEvent::Diagnostics(crate::lsp::PublishDiagnostics {
uri: uri.into(),
authoritative_errors: auth_err,
advisory_errors: adv_err,
total: auth_err + adv_err,
diagnostics: Vec::new(),
})
}
fn diag_rich(uri: &str, ds: Vec<Diagnostic>) -> LspEvent {
let mut auth = 0usize;
let mut adv = 0usize;
for d in &ds {
if d.severity == cargoless_proto::Severity::Error {
if d.source.as_deref() == Some("rustc") {
auth += 1;
} else {
adv += 1;
}
}
}
LspEvent::Diagnostics(crate::lsp::PublishDiagnostics {
uri: uri.into(),
authoritative_errors: auth,
advisory_errors: adv,
total: ds.len(),
diagnostics: ds,
})
}
#[test]
fn file_state_for_downrank_vs_f8redo_default() {
fn pd(auth: usize, adv: usize) -> crate::lsp::PublishDiagnostics {
crate::lsp::PublishDiagnostics {
uri: "file:///x.rs".into(),
authoritative_errors: auth,
advisory_errors: adv,
total: auth + adv,
diagnostics: Vec::new(),
}
}
assert_eq!(file_state_for(&pd(0, 0), false), (FileState::Green, false));
assert_eq!(file_state_for(&pd(1, 0), false), (FileState::Red, false));
assert_eq!(
file_state_for(&pd(0, 1), false),
(FileState::Red, false),
"F8-redo: RA-native-only error still drives RED by default"
);
assert_eq!(file_state_for(&pd(0, 0), true), (FileState::Green, false));
assert_eq!(
file_state_for(&pd(1, 0), true),
(FileState::Red, false),
"no false-GREEN: a real cargo-check error still drives RED"
);
assert_eq!(
file_state_for(&pd(0, 1), true),
(FileState::Green, true),
"the fix: RA-native-only (proc-macro hallucination) is \
demoted — would-be false-RED suppressed + counted"
);
assert_eq!(
file_state_for(&pd(1, 1), true),
(FileState::Red, false),
"authoritative present ⇒ RED, NOT a suppression (rustc \
evidence is real; nothing was downranked away)"
);
}
fn mk_diag(
path: &str,
line: u32,
col: u32,
sev: cargoless_proto::Severity,
code: Option<&str>,
msg: &str,
source: Option<&str>,
) -> Diagnostic {
Diagnostic {
file_path: std::path::PathBuf::from(path),
line,
col,
severity: sev,
code: code.map(str::to_owned),
message: msg.to_owned(),
source: source.map(str::to_owned),
}
}
#[test]
fn lifecycle_subscribers_receive_analyzer_restarting() {
let mut m = model();
let r1 = m.subscribe_lifecycle();
{
let r2 = m.subscribe_lifecycle();
m.emit_lifecycle(LifecycleEvent::AnalyzerRestarting);
assert_eq!(
drain(&r1),
vec![LifecycleEvent::AnalyzerRestarting],
"r1 receives"
);
assert_eq!(
drain(&r2),
vec![LifecycleEvent::AnalyzerRestarting],
"r2 receives"
);
}
m.emit_lifecycle(LifecycleEvent::AnalyzerRestarting);
assert_eq!(drain(&r1).len(), 1);
}
#[test]
fn lifecycle_no_emit_before_first_restart_is_silent() {
let mut m = model();
let rx = m.subscribe_lifecycle();
assert!(drain(&rx).is_empty(), "no lifecycle events on quiet model");
m.apply_event(&LspEvent::FlycheckEnded);
assert!(
drain(&rx).is_empty(),
"flycheck-end is a verdict event, not a lifecycle event"
);
}
fn drain<T>(rx: &Receiver<T>) -> Vec<T> {
let mut v = Vec::new();
while let Ok(e) = rx.try_recv() {
v.push(e);
}
v
}
#[test]
fn starts_red_advisory() {
let m = model();
assert_eq!(m.tree_state(), TreeState::Red);
assert_eq!(
m.last_verdict(),
Verdict {
tree: TreeState::Red,
provenance: VerdictProvenance::Advisory
}
);
assert_eq!(m.file_state("x"), None);
}
#[test]
fn native_only_clean_never_green_without_flycheck() {
let mut m = model();
let rx = m.subscribe();
m.apply_event(&diag("file:///p/src/lib.rs", 0, 0));
assert_eq!(m.tree_state(), TreeState::Red, "no green without flycheck");
assert_eq!(m.last_verdict().provenance, VerdictProvenance::Advisory);
let evs = drain(&rx);
assert!(
evs.iter()
.all(|e| !matches!(e, StateEvent::BecameGreen { .. }))
);
assert!(evs.contains(&StateEvent::FileVerdict {
path: "/p/src/lib.rs".into(),
state: FileState::Green
}));
}
#[test]
fn native_error_pre_flycheck_is_red_advisory_no_event_green() {
let mut m = model();
let rx = m.subscribe();
let arx = m.subscribe_advisory();
m.apply_event(&diag("file:///p/a.rs", 0, 3)); assert_eq!(m.tree_state(), TreeState::Red);
assert_eq!(
m.last_verdict(),
Verdict {
tree: TreeState::Red,
provenance: VerdictProvenance::Advisory
}
);
assert!(
drain(&rx)
.iter()
.all(|e| !matches!(e, StateEvent::BecameGreen { .. }))
);
assert!(!drain(&arx).is_empty());
}
#[test]
fn completed_clean_flycheck_is_authoritative_green() {
let mut m = model();
let rx = m.subscribe();
m.apply_event(&diag("file:///p/src/lib.rs", 0, 0)); assert_eq!(m.tree_state(), TreeState::Red);
m.apply_event(&LspEvent::FlycheckEnded);
assert_eq!(m.tree_state(), TreeState::Green);
assert_eq!(
m.last_verdict(),
Verdict {
tree: TreeState::Green,
provenance: VerdictProvenance::Authoritative
}
);
let evs = drain(&rx);
assert!(evs.contains(&StateEvent::BecameGreen { identity: ident() }));
}
#[test]
fn ra_native_severity_error_post_flycheck_is_red_too() {
let mut m = model();
let rx = m.subscribe();
m.apply_event(&diag("file:///p/lib.rs", 0, 0));
m.apply_event(&LspEvent::FlycheckEnded);
assert_eq!(m.tree_state(), TreeState::Green);
let _ = drain(&rx);
m.apply_event(&diag("file:///p/lib.rs", 0, 1)); assert_eq!(
m.tree_state(),
TreeState::Red,
"RA-native severity:Error post-flycheck must flip tree Red"
);
assert!(drain(&rx).contains(&StateEvent::BecameRed));
assert_eq!(
m.last_verdict().provenance,
VerdictProvenance::Authoritative
);
}
#[test]
fn rustc_error_after_pass_is_authoritative_red() {
let mut m = model();
let rx = m.subscribe();
m.apply_event(&diag("file:///p/a.rs", 0, 0));
m.apply_event(&LspEvent::FlycheckEnded);
assert_eq!(m.tree_state(), TreeState::Green);
let _ = drain(&rx);
m.apply_event(&diag("file:///p/a.rs", 1, 0));
assert_eq!(m.tree_state(), TreeState::Red);
assert_eq!(
m.last_verdict().provenance,
VerdictProvenance::Authoritative
);
assert!(drain(&rx).contains(&StateEvent::BecameRed));
}
#[test]
fn empty_clean_pass_is_green() {
let mut m = model();
m.apply_event(&LspEvent::FlycheckEnded);
assert_eq!(m.tree_state(), TreeState::Green);
assert_eq!(
m.last_verdict().provenance,
VerdictProvenance::Authoritative
);
}
#[test]
fn forget_last_rustc_red_file_flips_green_post_pass() {
let mut m = model();
let rx = m.subscribe();
m.apply_event(&diag("file:///p/keep.rs", 0, 0));
m.apply_event(&diag("file:///p/scratch.rs", 1, 0)); m.apply_event(&LspEvent::FlycheckEnded);
assert_eq!(m.tree_state(), TreeState::Red);
let _ = drain(&rx);
m.forget_file("/p/scratch.rs");
assert_eq!(m.tree_state(), TreeState::Green);
assert!(drain(&rx).contains(&StateEvent::BecameGreen { identity: ident() }));
}
#[test]
fn advisory_channel_receives_and_prunes() {
let mut m = model();
let a1 = m.subscribe_advisory();
{
let a2 = m.subscribe_advisory();
m.apply_event(&diag("file:///p/a.rs", 0, 1));
assert!(!drain(&a1).is_empty());
assert!(!drain(&a2).is_empty());
}
m.apply_event(&LspEvent::FlycheckEnded);
let got = drain(&a1);
assert!(
got.iter()
.any(|v| v.provenance == VerdictProvenance::Authoritative)
);
}
#[test]
fn non_file_uri_ignored() {
let mut m = model();
m.apply_event(&diag("untitled:Untitled-1", 5, 5));
assert_eq!(m.file_state("untitled:Untitled-1"), None);
assert_eq!(m.tree_state(), TreeState::Red);
}
struct EnvGuard {
prev: Option<String>,
}
impl EnvGuard {
fn set(value: &str) -> Self {
let prev = std::env::var("TF_DEBOUNCE_MS").ok();
unsafe { std::env::set_var("TF_DEBOUNCE_MS", value) };
Self { prev }
}
}
impl Drop for EnvGuard {
fn drop(&mut self) {
unsafe {
match &self.prev {
Some(v) => std::env::set_var("TF_DEBOUNCE_MS", v),
None => std::env::remove_var("TF_DEBOUNCE_MS"),
}
}
}
}
#[test]
fn watch_debounce_defaults_when_env_unset() {
let _g = EnvGuard::set("");
unsafe { std::env::remove_var("TF_DEBOUNCE_MS") };
let d = resolve_watch_debounce();
assert_eq!(d, Duration::from_millis(150), "default = 150ms");
}
#[test]
fn watch_debounce_honors_valid_env_override() {
let _g = EnvGuard::set("300");
assert_eq!(resolve_watch_debounce(), Duration::from_millis(300));
}
#[test]
fn watch_debounce_rejects_zero_and_garbage() {
let _g0 = EnvGuard::set("0");
assert_eq!(resolve_watch_debounce(), Duration::from_millis(150));
drop(_g0);
let _g1 = EnvGuard::set("nope");
assert_eq!(resolve_watch_debounce(), Duration::from_millis(150));
}
#[test]
fn watch_debounce_accepts_large_values_for_flicker_free_refactor() {
let _g = EnvGuard::set("2500");
assert_eq!(resolve_watch_debounce(), Duration::from_millis(2500));
}
#[test]
fn placeholder_identity_is_sentinel_dev_wasm() {
let id = placeholder_identity();
assert_eq!(id.profile, Profile::Dev);
assert_eq!(id.target.as_str(), "wasm32-unknown-unknown");
assert_eq!(id.source_tree, id.cargo_lock); }
#[test]
fn diagnostics_accumulate_per_file_and_aggregate() {
let mut m = model();
m.apply_event(&diag_rich(
"file:///p/src/a.rs",
vec![mk_diag(
"/p/src/a.rs",
10,
3,
cargoless_proto::Severity::Error,
Some("E0277"),
"the trait bound not satisfied",
Some("rustc"),
)],
));
m.apply_event(&diag_rich(
"file:///p/src/b.rs",
vec![mk_diag(
"/p/src/b.rs",
1,
1,
cargoless_proto::Severity::Warning,
Some("unused_imports"),
"unused import",
Some("rust-analyzer"),
)],
));
let all = m.all_diagnostics();
assert_eq!(all.len(), 2, "two diagnostics, one per file");
assert_eq!(m.file_diagnostics("/p/src/a.rs").len(), 1);
assert_eq!(m.file_diagnostics("/p/src/b.rs").len(), 1);
let codes: Vec<&str> = all.iter().filter_map(|d| d.code.as_deref()).collect();
assert!(codes.contains(&"E0277"));
assert!(codes.contains(&"unused_imports"));
}
#[test]
fn later_publish_replaces_prior_per_file() {
let mut m = model();
m.apply_event(&diag_rich(
"file:///p/src/a.rs",
vec![mk_diag(
"/p/src/a.rs",
1,
1,
cargoless_proto::Severity::Error,
Some("E0277"),
"first",
Some("rustc"),
)],
));
assert_eq!(m.file_diagnostics("/p/src/a.rs").len(), 1);
m.apply_event(&diag_rich(
"file:///p/src/a.rs",
vec![mk_diag(
"/p/src/a.rs",
2,
2,
cargoless_proto::Severity::Error,
Some("E0308"),
"second",
Some("rustc"),
)],
));
let now = m.file_diagnostics("/p/src/a.rs");
assert_eq!(now.len(), 1);
assert_eq!(now[0].code.as_deref(), Some("E0308"));
m.apply_event(&diag_rich("file:///p/src/a.rs", vec![]));
assert!(m.file_diagnostics("/p/src/a.rs").is_empty());
}
#[test]
fn forget_file_drops_diagnostics_too() {
let mut m = model();
m.apply_event(&diag_rich(
"file:///p/src/a.rs",
vec![mk_diag(
"/p/src/a.rs",
1,
1,
cargoless_proto::Severity::Error,
Some("E0277"),
"x",
Some("rustc"),
)],
));
assert!(!m.all_diagnostics().is_empty());
m.forget_file("/p/src/a.rs");
assert!(
m.all_diagnostics().is_empty(),
"deleted file's diagnostics must be evicted from aggregate"
);
}
#[test]
fn aggregate_diagnostics_are_path_ordered() {
let mut m = model();
m.apply_event(&diag_rich(
"file:///p/z.rs",
vec![mk_diag(
"/p/z.rs",
1,
1,
cargoless_proto::Severity::Error,
None,
"z",
Some("rustc"),
)],
));
m.apply_event(&diag_rich(
"file:///p/a.rs",
vec![mk_diag(
"/p/a.rs",
1,
1,
cargoless_proto::Severity::Error,
None,
"a",
Some("rustc"),
)],
));
let all = m.all_diagnostics();
assert_eq!(all.len(), 2);
assert_eq!(all[0].file_path, std::path::PathBuf::from("/p/a.rs"));
assert_eq!(all[1].file_path, std::path::PathBuf::from("/p/z.rs"));
}
#[test]
fn collect_rs_files_skips_target_git_and_gitignored() {
let base = std::env::temp_dir().join(format!("tf-model-walk-{}", std::process::id()));
let _ = fs::remove_dir_all(&base);
let mk = |rel: &str, body: &str| {
let p = base.join(rel);
fs::create_dir_all(p.parent().unwrap()).unwrap();
fs::write(&p, body).unwrap();
};
mk(".gitignore", "ignored.rs\n");
mk("src/lib.rs", "");
mk("src/nested/m.rs", "");
mk("ignored.rs", "");
mk("target/debug/build.rs", "");
mk(".git/hooks/pre.rs", "");
mk("README.md", "");
let root = fs::canonicalize(&base).unwrap();
let ignore = crate::watcher::IgnoreRules::for_root(&root);
let mut got: Vec<String> = collect_rs_files(&root, &ignore)
.into_iter()
.map(|p| {
p.strip_prefix(&root)
.unwrap()
.to_string_lossy()
.replace('\\', "/")
})
.collect();
got.sort();
assert_eq!(
got,
vec!["src/lib.rs".to_string(), "src/nested/m.rs".to_string()]
);
let _ = fs::remove_dir_all(&base);
}
}