use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use rmcp::ErrorData;
use rmcp::model::Tool;
use serde_json::{Map, Value, json};
use crate::server::tool_trait::{
McpTool, ToolContext, ToolOutput, get_bool, get_f64, get_int, get_str, get_str_array,
require_resolved_path,
};
use crate::tool_defs::tool_def;
fn per_file_lock(path: &str) -> Arc<Mutex<()>> {
crate::core::path_locks::per_file_lock(path)
}
pub struct CtxReadTool;
impl McpTool for CtxReadTool {
fn name(&self) -> &'static str {
"ctx_read"
}
fn tool_def(&self) -> Tool {
tool_def(
"ctx_read",
"Read source files. mode REQUIRED — choose by intent (see `mode` below).\n\
To UNDERSTAND code run ctx_compose FIRST; ctx_read after it identified files.\n\
anchored → edit by reference via ctx_patch (no exact-recall).",
json!({
"type": "object",
"properties": {
"path": { "type": "string", "description": "Absolute path" },
"paths": { "type": "array", "items": { "type": "string" }, "description": "Batch read" },
"mode": {
"type": "string",
"description": "REQUIRED. full=verbatim(edit-ready) anchored=full+N:hh|anchors(edit via ctx_patch) raw=exact-bytes signatures=API map=structure auto=smart diff=git-delta lines:N-M=window (comma multi-selects: lines:5,10-20) reference=quotes task=focus"
},
"raw": { "type": "boolean", "description": "Verbatim (= mode=raw + fresh)" },
"start_line": { "type": "integer", "description": "1-based" },
"offset": { "type": "integer", "description": "start_line alias" },
"limit": { "type": "integer", "description": "Max lines" },
"fresh": { "type": "boolean", "description": "Bypass cache" },
"aggressiveness": { "type": "number", "description": "0.0–1.0 density (entropy/task)" },
"protect": { "type": "array", "items": { "type": "string" }, "description": "Symbols kept verbatim" }
},
"required": []
}),
)
}
fn handle(
&self,
args: &Map<String, Value>,
ctx: &ToolContext,
) -> Result<ToolOutput, ErrorData> {
if args
.get("paths")
.and_then(|v| v.as_array())
.is_some_and(|a| !a.is_empty())
{
return super::ctx_multi_read::batch_read(args, ctx);
}
let path = if let Some(repo) = get_str(args, "repo") {
let root = crate::core::multi_repo::resolve_repo_root(&repo).ok_or_else(|| {
let known = crate::core::multi_repo::known_aliases().join(", ");
let known = if known.is_empty() {
"none registered — use ctx_multi_repo add_root".to_string()
} else {
known
};
ErrorData::invalid_params(
format!("unknown repo alias: {repo} (known: {known})"),
None,
)
})?;
let rel = get_str(args, "path").unwrap_or_else(|| ".".to_string());
crate::core::path_resolve::resolve_tool_path(Some(&root), None, &rel)
.map_err(|e| ErrorData::invalid_params(e, None))?
} else {
require_resolved_path(ctx, args, "path")?
};
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
self.handle_inner(args, ctx, &path)
})) {
Ok(result) => result,
Err(_) => Err(ErrorData::internal_error(
format!(
"ctx_read panicked while processing '{path}'. This is a bug — please report it."
),
None,
)),
}
}
}
impl CtxReadTool {
#[allow(clippy::unused_self)]
fn handle_inner(
&self,
args: &Map<String, Value>,
ctx: &ToolContext,
path: &str,
) -> Result<ToolOutput, ErrorData> {
let session_lock = ctx
.session
.as_ref()
.ok_or_else(|| ErrorData::internal_error("session not available", None))?;
let cache_lock = ctx
.cache
.as_ref()
.ok_or_else(|| ErrorData::internal_error("cache not available", None))?;
let current_task = {
let rt = tokio::runtime::Handle::current();
let mut attempt = 0u32;
loop {
if let Ok(session) = rt.block_on(tokio::time::timeout(
std::time::Duration::from_secs(5),
session_lock.read(),
)) {
break session.task.as_ref().map(|t| t.description.clone());
}
attempt += 1;
if attempt >= 3 {
tracing::warn!(
"session read-lock timeout after {attempt} attempts in ctx_read for {path}"
);
return Err(ErrorData::internal_error(
"session lock timeout — another tool may be holding it. Retry in a moment.",
None,
));
}
tracing::debug!(
"session read-lock attempt {attempt}/3 timed out for {path}, retrying"
);
std::thread::sleep(std::time::Duration::from_millis(100 * u64::from(attempt)));
}
};
let task_ref = current_task.as_deref();
let profile = crate::core::profiles::active_profile();
let arg_raw = get_bool(args, "raw").unwrap_or(false);
let explicit_mode_arg = resolve_raw_alias(arg_raw, get_str(args, "mode"));
let explicit_mode = explicit_mode_arg.is_some();
let policy_default_mode = if explicit_mode {
None
} else {
crate::core::policy::runtime::active()
.and_then(|p| p.resolved.default_read_mode.clone())
};
let persona_default_mode = if explicit_mode || policy_default_mode.is_some() {
None
} else {
crate::core::persona::active().read_mode_override()
};
let mut mode = if let Some(m) = explicit_mode_arg {
m
} else if let Some(pd) = policy_default_mode {
pd
} else if let Some(pm) = persona_default_mode {
pm
} else if profile.read.default_mode_effective() == "auto" {
if let Ok(cache) = cache_lock.try_read() {
crate::tools::ctx_smart_read::select_mode_with_task(&cache, path, task_ref)
} else {
tracing::debug!(
"cache lock contested during auto-mode selection for {path}; \
falling back to full"
);
"full".to_string()
}
} else {
profile.read.default_mode_effective().to_string()
};
let mut fresh = get_bool(args, "fresh").unwrap_or(false);
if arg_raw {
fresh = true;
}
let cache_policy = crate::server::compaction_sync::effective_cache_policy();
if cache_policy == "off" {
fresh = true;
}
let aggressiveness =
crate::core::aggressiveness::effective(get_f64(args, "aggressiveness"));
let protect = get_str_array(args, "protect").unwrap_or_default();
if !explicit_mode && let Some(a) = aggressiveness {
mode = crate::tools::ctx_read::ReadMode::Density(
crate::core::aggressiveness::AggressivenessProfile::from_level(a).density_target,
)
.to_string();
}
apply_line_window(
&mut mode,
&mut fresh,
explicit_mode,
get_int(args, "start_line"),
get_int(args, "offset"),
get_int(args, "limit"),
);
let pressure_action = ctx.pressure_snapshot.as_ref().map(|p| &p.recommendation);
let resolved_agent_id = ctx.agent_id.as_ref().and_then(|a| match a.try_read() {
Ok(guard) => guard.clone(),
Err(_) => None,
});
let gate_result = crate::server::context_gate::pre_dispatch_read_for_agent(
path,
&mode,
task_ref,
Some(&ctx.project_root),
pressure_action,
resolved_agent_id.as_deref(),
);
if gate_result.budget_blocked {
let msg = gate_result
.budget_warning
.unwrap_or_else(|| "Agent token budget exceeded".to_string());
return Err(ErrorData::invalid_params(msg, None));
}
let budget_warning = gate_result.budget_warning.clone();
if mode != "raw"
&& let Some(overridden) = gate_result.overridden_mode
{
mode = overridden;
}
let (mut mode, degrade_warning) = if crate::tools::ctx_read::is_instruction_file(path) {
("full".to_string(), None)
} else if mode == "raw" {
("raw".to_string(), None)
} else {
auto_degrade_read_mode(&mode)
};
let mut delta_explicit_note: Option<String> = None;
if !fresh
&& explicit_mode
&& (mode == "full" || mode == "full-compact" || mode.starts_with("lines:"))
&& crate::core::config::Config::load().delta_explicit_effective()
&& let Ok(cache) = cache_lock.try_read()
{
let decision = crate::tools::ctx_read::resolve_explicit_delta_mode(
&cache,
path,
&mode,
explicit_mode,
fresh,
true,
);
mode = decision.mode;
delta_explicit_note = decision.note;
}
if mode.starts_with("lines:") {
fresh = true;
}
if crate::core::binary_detect::is_binary_file(path) {
let msg = crate::core::binary_detect::binary_file_message(path);
return Err(ErrorData::invalid_params(msg, None));
}
{
let cap = crate::core::limits::max_read_bytes() as u64;
if let Ok(meta) = std::fs::metadata(path)
&& meta.len() > cap
{
let msg = format!(
"File too large ({} bytes, limit {} bytes via LCTX_MAX_READ_BYTES). \
Use mode=\"lines:1-100\" for partial reads or increase the limit.",
meta.len(),
cap
);
return Err(ErrorData::invalid_params(msg, None));
}
}
if !fresh
&& let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir()
&& let Ok(mut cache) = cache_lock.try_write()
{
crate::server::compaction_sync::sync_if_compacted(&mut cache, &data_dir);
}
let read_timeout = std::time::Duration::from_secs(30);
let cancelled = Arc::new(AtomicBool::new(false));
let (output, resolved_mode, original, is_cache_hit, file_ref, cache_stats) = {
let crp_mode = ctx.crp_mode;
let task_ref = current_task.as_deref();
let fast_result = 'fast: {
let file_lock = per_file_lock(path);
let Some(_file_guard) = file_lock.try_lock().ok() else {
break 'fast None;
};
if !fresh
&& (mode == "full" || mode == "full-compact" || mode == "auto")
&& let Ok(cache) = cache_lock.try_read()
&& let Some(read_output) =
crate::tools::ctx_read::try_stub_hit_readonly(&cache, path)
{
let content = read_output.content;
let rmode = read_output.resolved_mode;
let orig = cache.get(path).map_or(0, |e| e.original_tokens);
let hit = content.contains(" cached ")
|| content.contains("[unchanged")
|| content.contains("[delta:");
let fref = cache.file_ref_map().get(path).cloned();
let stats = cache.get_stats();
let stats_snapshot = (stats.total_reads(), stats.cache_hits());
break 'fast Some((content, rmode, orig, hit, fref, stats_snapshot));
}
let Some(mut cache) = cache_lock.try_write().ok() else {
break 'fast None;
};
let read_output = if fresh {
crate::tools::ctx_read::handle_fresh_with_task_resolved_tuned(
&mut cache,
path,
&mode,
crp_mode,
task_ref,
aggressiveness,
&protect,
)
} else {
crate::tools::ctx_read::handle_with_task_resolved_tuned(
&mut cache,
path,
&mode,
crp_mode,
task_ref,
aggressiveness,
&protect,
)
};
let content = read_output.content;
let rmode = read_output.resolved_mode;
let orig = cache.get(path).map_or(0, |e| e.original_tokens);
let hit = content.contains(" cached ")
|| content.contains("[unchanged")
|| content.contains("[delta:");
let fref = cache.file_ref_map().get(path).cloned();
let stats = cache.get_stats();
let stats_snapshot = (stats.total_reads(), stats.cache_hits());
Some((content, rmode, orig, hit, fref, stats_snapshot))
};
if let Some(result) = fast_result {
result
} else {
let cache_lock = cache_lock.clone();
let mode = mode.clone();
let task_owned = current_task.clone();
let protect_owned = protect.clone();
let path_owned = path.to_string();
let cancel_flag = cancelled.clone();
let (tx, rx) = std::sync::mpsc::sync_channel(1);
std::thread::spawn(move || {
let file_lock = per_file_lock(&path_owned);
let _file_guard = {
let deadline =
std::time::Instant::now() + std::time::Duration::from_secs(25);
loop {
if cancel_flag.load(Ordering::Relaxed) {
return;
}
if let Ok(guard) = file_lock.try_lock() {
break guard;
}
if std::time::Instant::now() >= deadline {
tracing::error!(
"ctx_read: per-file lock timeout after 25s for {path_owned}"
);
let _ = tx.send((
format!("per-file lock contention for {path_owned} — retry in a moment"),
"error".to_string(), 0, false, None, (0, 0),
));
return;
}
std::thread::sleep(std::time::Duration::from_millis(50));
}
};
if cancel_flag.load(Ordering::Relaxed) {
return;
}
if !fresh
&& (mode == "full" || mode == "full-compact" || mode == "auto")
&& let Ok(cache) = cache_lock.try_read()
&& let Some(read_output) =
crate::tools::ctx_read::try_stub_hit_readonly(&cache, &path_owned)
{
let content = read_output.content;
let rmode = read_output.resolved_mode;
let orig = cache.get(&path_owned).map_or(0, |e| e.original_tokens);
let hit = true;
let fref = cache.file_ref_map().get(path_owned.as_str()).cloned();
let stats = cache.get_stats();
let stats_snapshot = (stats.total_reads(), stats.cache_hits());
let _ = tx.send((content, rmode, orig, hit, fref, stats_snapshot));
return;
}
let preread = crate::tools::ctx_read::read_file_lossy(&path_owned).ok();
if cancel_flag.load(Ordering::Relaxed) {
return;
}
let mut cache = {
let deadline =
std::time::Instant::now() + std::time::Duration::from_secs(25);
loop {
if cancel_flag.load(Ordering::Relaxed) {
return;
}
if let Ok(guard) = cache_lock.try_write() {
break guard;
}
if std::time::Instant::now() >= deadline {
tracing::error!(
"ctx_read: cache write-lock timeout after 25s for {path_owned}"
);
let _ = tx.send((
format!(
"cache lock contention for {path_owned} — retry in a moment"
),
"error".to_string(),
0,
false,
None,
(0, 0),
));
return;
}
std::thread::sleep(std::time::Duration::from_millis(50));
}
};
let task_ref = task_owned.as_deref();
let read_output = if let Some(content) = preread {
crate::tools::ctx_read::handle_with_preread(
&mut cache,
&path_owned,
&mode,
fresh,
crp_mode,
task_ref,
aggressiveness,
&protect_owned,
content,
)
} else if fresh {
crate::tools::ctx_read::handle_fresh_with_task_resolved_tuned(
&mut cache,
&path_owned,
&mode,
crp_mode,
task_ref,
aggressiveness,
&protect_owned,
)
} else {
crate::tools::ctx_read::handle_with_task_resolved_tuned(
&mut cache,
&path_owned,
&mode,
crp_mode,
task_ref,
aggressiveness,
&protect_owned,
)
};
let content = read_output.content;
let rmode = read_output.resolved_mode;
let orig = cache.get(&path_owned).map_or(0, |e| e.original_tokens);
let hit = content.contains(" cached ");
let fref = cache.file_ref_map().get(path_owned.as_str()).cloned();
let stats = cache.get_stats();
let stats_snapshot = (stats.total_reads(), stats.cache_hits());
let _ = tx.send((content, rmode, orig, hit, fref, stats_snapshot));
});
if let Ok(result) = rx.recv_timeout(read_timeout) {
result
} else {
cancelled.store(true, Ordering::Relaxed);
tracing::error!("ctx_read timed out after {read_timeout:?} for {path}");
let msg = format!(
"ERROR: ctx_read timed out after {}s reading {path}. \
The file may be very large or a blocking I/O issue occurred. \
Try mode=\"lines:1-100\" for a partial read.",
read_timeout.as_secs()
);
return Err(ErrorData::internal_error(msg, None));
}
} };
if resolved_mode == "error" {
return Err(ErrorData::invalid_params(output, None));
}
let output_tokens = crate::core::tokens::count_tokens(&output);
let saved = original.saturating_sub(output_tokens);
let mut ensured_root: Option<String> = None;
let mut traversal_working_set: Vec<String> = Vec::new();
let project_root_snapshot;
{
let rt = tokio::runtime::Handle::current();
let session_guard = rt.block_on(tokio::time::timeout(
std::time::Duration::from_secs(10),
session_lock.write(),
));
if let Ok(mut session) = session_guard {
session.touch_file(path, file_ref.as_deref(), &resolved_mode, original);
traversal_working_set =
crate::core::tool_lifecycle::recent_working_set(&session, path);
let file_summary = extract_file_summary(&output, path);
if !file_summary.is_empty() {
session.set_file_summary(path, &file_summary);
}
if is_cache_hit {
session.record_cache_hit();
}
if session.active_structured_intent.is_none() && session.files_touched.len() >= 2 {
let touched: Vec<String> = session
.files_touched
.iter()
.map(|f| f.path.clone())
.collect();
let inferred =
crate::core::intent_engine::StructuredIntent::from_file_patterns(&touched);
if inferred.confidence >= 0.4 {
session.active_structured_intent = Some(inferred);
}
}
if session.task.is_none() && session.stats.files_read % 5 == 0 {
session.auto_infer_task();
}
let root_missing = session
.project_root
.as_deref()
.is_none_or(|r| r.trim().is_empty());
if root_missing && let Some(root) = crate::core::protocol::detect_project_root(path)
{
session.project_root = Some(root.clone());
ensured_root = Some(root);
}
project_root_snapshot = session
.project_root
.clone()
.unwrap_or_else(|| ".".to_string());
} else {
tracing::warn!(
"session write-lock timeout (5s) in ctx_read post-update for {path}"
);
project_root_snapshot = ctx.project_root.clone();
}
}
if let Some(root) = ensured_root.as_deref() {
crate::core::index_orchestrator::ensure_all_background(root);
}
{
let path_bg = path.to_string();
let resolved_mode_bg = resolved_mode.clone();
let project_root_bg = project_root_snapshot.clone();
let (turns, hits) = cache_stats;
let ledger_cache = (crate::core::savings_ledger::ledger_family()
!= crate::core::tokens::TokenizerFamily::O200kBase)
.then(|| cache_lock.clone());
let ledger_output = ledger_cache.as_ref().map(|_| output.clone());
std::thread::spawn(move || {
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
crate::core::heatmap::record_file_access(&path_bg, original, saved);
{
use crate::core::savings_ledger as ledger;
let (lbase, lsaved) = match (&ledger_cache, &ledger_output) {
(Some(cl), Some(out)) => match cl.try_read().ok().and_then(|c| {
c.get(&path_bg)
.and_then(crate::core::cache::CacheEntry::content)
}) {
Some(raw) => {
let lo = ledger::count_for_ledger(&raw);
(lo, lo.saturating_sub(ledger::count_for_ledger(out)))
}
None => (original, saved),
},
_ => (original, saved),
};
ledger::record_read_event(lbase, lsaved);
}
if let Some(root) =
crate::core::tool_lifecycle::usable_root(Some(project_root_bg.as_str()))
{
crate::core::cooccurrence::record_focus_access(
root,
&path_bg,
&traversal_working_set,
);
}
let sig =
crate::core::mode_predictor::FileSignature::from_path(&path_bg, original);
let density = if output_tokens > 0 {
original as f64 / output_tokens as f64
} else {
1.0
};
let outcome = crate::core::mode_predictor::ModeOutcome {
mode: resolved_mode_bg,
tokens_in: original,
tokens_out: output_tokens,
density: density.min(1.0),
};
let mut predictor = crate::core::mode_predictor::ModePredictor::new();
predictor.set_project_root(&project_root_bg);
predictor.record(sig, outcome);
predictor.save();
let ext = std::path::Path::new(&path_bg)
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
.to_string();
let thresholds =
crate::core::adaptive_thresholds::thresholds_for_path(&path_bg);
let feedback_outcome = crate::core::feedback::CompressionOutcome {
session_id: format!("{}", std::process::id()),
language: ext,
entropy_threshold: thresholds.bpe_entropy,
jaccard_threshold: thresholds.jaccard,
total_turns: turns as u32,
tokens_saved: saved as u64,
tokens_original: original as u64,
cache_hits: hits as u32,
total_reads: turns as u32,
task_completed: crate::core::bounce_tracker::global()
.lock()
.ok()
.and_then(|bt| bt.bounce_rate_for_extension(&path_bg))
.is_none_or(|rate| rate < 0.30),
timestamp: chrono::Local::now().to_rfc3339(),
};
let mut store = crate::core::feedback::FeedbackStore::load();
store.project_root = Some(project_root_bg);
store.record_outcome(feedback_outcome);
}));
});
}
if let Some(aid) = resolved_agent_id.as_deref() {
crate::core::agent_budget::record_consumption(aid, output_tokens);
}
let graph_hint = if !is_cache_hit
&& !resolved_mode.starts_with("lines:")
&& crate::core::profiles::active_profile()
.output_hints
.related_hint()
{
crate::tools::ctx_read::graph_related_hint(path)
} else {
None
};
let hints_suffix = {
let graph_db =
crate::core::property_graph::graph_dir(&ctx.project_root).join("graph.db");
let edges = if graph_db.exists() {
crate::core::property_graph::CodeGraph::open(&ctx.project_root)
.map(|g| g.all_cross_source_edges())
.unwrap_or_default()
} else {
Vec::new()
};
if edges.is_empty() {
String::new()
} else {
let hints = crate::core::cross_source_hints::hints_for_file(
path,
&edges,
&ctx.project_root,
);
crate::core::cross_source_hints::format_hints(&hints)
}
};
let mut warnings = Vec::new();
if let Some(ref w) = budget_warning {
warnings.push(w.as_str());
}
if let Some(ref w) = degrade_warning {
warnings.push(w.as_str());
}
if let Some(ref w) = delta_explicit_note {
warnings.push(w.as_str());
}
let graph_suffix = graph_hint.map(|h| format!("\n{h}")).unwrap_or_default();
let final_output = if !warnings.is_empty() {
format!(
"{output}{hints_suffix}{graph_suffix}\n\n{}",
warnings.join("\n")
)
} else if hints_suffix.is_empty() && graph_suffix.is_empty() {
output
} else {
format!("{output}{hints_suffix}{graph_suffix}")
};
Ok(ToolOutput {
text: final_output,
original_tokens: original,
saved_tokens: saved,
mode: Some(resolved_mode),
path: Some(path.to_string()),
changed: false,
shell_outcome: None,
})
}
}
fn resolve_line_window(
start_line: Option<i64>,
offset: Option<i64>,
limit: Option<i64>,
) -> Option<(i64, Option<i64>)> {
let start = start_line.or(offset).map(|v| v.max(1));
let limit = limit.filter(|&l| l > 0);
match (start, limit) {
(Some(s), l) => Some((s, l)),
(None, Some(_)) => Some((1, limit)),
(None, None) => None,
}
}
fn lines_mode(start: i64, limit: Option<i64>) -> String {
match limit {
Some(l) => format!("lines:{start}-{}", start + l - 1),
None => format!("lines:{start}-999999"),
}
}
fn apply_line_window(
mode: &mut String,
fresh: &mut bool,
explicit_mode: bool,
start_line: Option<i64>,
offset: Option<i64>,
limit: Option<i64>,
) {
let Some((start, limit)) = resolve_line_window(start_line, offset, limit) else {
return;
};
if start <= 1 && limit.is_none() {
return;
}
*fresh = true;
if !explicit_mode || mode.starts_with("lines") {
*mode = lines_mode(start, limit);
}
}
fn resolve_raw_alias(arg_raw: bool, mode_arg: Option<String>) -> Option<String> {
if arg_raw {
Some("raw".to_string())
} else {
mode_arg
}
}
fn apply_verdict(
mode: &str,
verdict: crate::core::degradation_policy::DegradationVerdictV1,
) -> (String, bool) {
use crate::core::degradation_policy::DegradationVerdictV1;
match verdict {
DegradationVerdictV1::Ok => (mode.to_string(), false),
DegradationVerdictV1::Warn => match mode {
"full" => ("map".to_string(), true),
other => (other.to_string(), false),
},
DegradationVerdictV1::Throttle => match mode {
"full" | "map" => ("signatures".to_string(), true),
other => (other.to_string(), false),
},
DegradationVerdictV1::Block => {
if mode == "signatures" {
("signatures".to_string(), false)
} else {
("signatures".to_string(), true)
}
}
}
}
fn auto_degrade_read_mode(mode: &str) -> (String, Option<String>) {
if crate::core::config::Config::load().no_degrade_effective() {
return (mode.to_string(), None);
}
let profile = crate::core::profiles::active_profile();
if !profile.degradation.enforce_effective() {
return (mode.to_string(), None);
}
let policy = crate::core::degradation_policy::evaluate_v1_for_tool("ctx_read", None);
let (new_mode, degraded) = apply_verdict(mode, policy.decision.verdict);
let warning = if degraded {
Some(format!(
"⚠ Context pressure: mode={mode} was downgraded to mode={new_mode} \
(verdict: {:?}). Use start_line=1 to bypass, or run ctx_compress to free budget.",
policy.decision.verdict
))
} else {
None
};
(new_mode, warning)
}
fn extract_file_summary(output: &str, path: &str) -> String {
let hint = crate::core::auto_findings::extract_content_hint(output);
if !hint.is_empty() {
return hint;
}
let ext = std::path::Path::new(path)
.extension()
.and_then(|e| e.to_str())
.unwrap_or("");
let line_count = output.lines().count();
if line_count > 5 {
format!("{ext} file, {line_count} lines")
} else {
String::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn raw_alias_forces_raw_mode_over_explicit_mode() {
assert_eq!(
resolve_raw_alias(true, Some("signatures".to_string())),
Some("raw".to_string())
);
assert_eq!(resolve_raw_alias(true, None), Some("raw".to_string()));
}
#[test]
fn raw_alias_absent_passes_mode_through() {
assert_eq!(
resolve_raw_alias(false, Some("full".to_string())),
Some("full".to_string())
);
assert_eq!(resolve_raw_alias(false, None), None);
}
#[test]
fn per_file_lock_same_path_returns_same_mutex() {
let lock_a1 = per_file_lock("/tmp/test_same_path.txt");
let lock_a2 = per_file_lock("/tmp/test_same_path.txt");
assert!(Arc::ptr_eq(&lock_a1, &lock_a2));
}
#[test]
fn per_file_lock_different_paths_return_different_mutexes() {
let lock_a = per_file_lock("/tmp/test_path_a.txt");
let lock_b = per_file_lock("/tmp/test_path_b.txt");
assert!(!Arc::ptr_eq(&lock_a, &lock_b));
}
#[test]
fn per_file_lock_serializes_concurrent_access() {
let counter = Arc::new(AtomicUsize::new(0));
let max_concurrent = Arc::new(AtomicUsize::new(0));
let path = "/tmp/test_concurrent_serialization.txt";
let mut handles = Vec::new();
for _ in 0..5 {
let counter = counter.clone();
let max_concurrent = max_concurrent.clone();
let path = path.to_string();
handles.push(std::thread::spawn(move || {
let lock = per_file_lock(&path);
let _guard = lock.lock().unwrap();
let active = counter.fetch_add(1, Ordering::SeqCst) + 1;
max_concurrent.fetch_max(active, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(10));
counter.fetch_sub(1, Ordering::SeqCst);
}));
}
for h in handles {
h.join().unwrap();
}
assert_eq!(max_concurrent.load(Ordering::SeqCst), 1);
}
#[test]
fn per_file_lock_allows_parallel_different_paths() {
let counter = Arc::new(AtomicUsize::new(0));
let max_concurrent = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for i in 0..4 {
let counter = counter.clone();
let max_concurrent = max_concurrent.clone();
let path = format!("/tmp/test_parallel_{i}.txt");
handles.push(std::thread::spawn(move || {
let lock = per_file_lock(&path);
let _guard = lock.lock().unwrap();
let active = counter.fetch_add(1, Ordering::SeqCst) + 1;
max_concurrent.fetch_max(active, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(50));
counter.fetch_sub(1, Ordering::SeqCst);
}));
}
for h in handles {
h.join().unwrap();
}
assert!(max_concurrent.load(Ordering::SeqCst) > 1);
}
#[test]
fn zombie_thread_does_not_block_subsequent_cache_access() {
let cache: Arc<tokio::sync::RwLock<u32>> = Arc::new(tokio::sync::RwLock::new(0));
let zombie_lock = cache.clone();
let _zombie = std::thread::spawn(move || {
let _guard = zombie_lock.blocking_write();
std::thread::sleep(std::time::Duration::from_secs(2));
});
std::thread::sleep(std::time::Duration::from_millis(50));
assert!(cache.try_read().is_err());
let cancel = Arc::new(AtomicBool::new(false));
let cancel2 = cancel.clone();
let lock2 = cache.clone();
let waiter = std::thread::spawn(move || {
let start = std::time::Instant::now();
loop {
if cancel2.load(Ordering::Relaxed) {
return (false, start.elapsed());
}
if let Ok(_guard) = lock2.try_write() {
return (true, start.elapsed());
}
std::thread::sleep(std::time::Duration::from_millis(50));
}
});
std::thread::sleep(std::time::Duration::from_millis(200));
cancel.store(true, Ordering::Relaxed);
let (acquired, elapsed) = waiter.join().unwrap();
assert!(
!acquired,
"should not have acquired lock while zombie holds it"
);
assert!(
elapsed < std::time::Duration::from_secs(1),
"cancellation should have stopped the loop promptly"
);
}
fn apply_start_line(
mode: &mut String,
fresh: &mut bool,
explicit_mode: bool,
start_line: Option<i64>,
) {
super::apply_line_window(mode, fresh, explicit_mode, start_line, None, None);
}
#[test]
fn start_line_1_does_not_override_mode() {
let mut mode = "auto".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, false, Some(1));
assert_eq!(mode, "auto", "start_line=1 should not change mode");
assert!(!fresh, "start_line=1 should not force fresh=true");
}
#[test]
fn start_line_gt1_overrides_implicit_mode() {
let mut mode = "auto".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, false, Some(50));
assert_eq!(mode, "lines:50-999999");
assert!(fresh);
}
#[test]
fn start_line_gt1_does_not_override_explicit_map() {
let mut mode = "map".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, true, Some(50));
assert_eq!(
mode, "map",
"explicit mode=map must not be clobbered by start_line"
);
assert!(fresh, "start_line>1 should still force fresh");
}
#[test]
fn start_line_gt1_does_not_override_explicit_signatures() {
let mut mode = "signatures".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, true, Some(100));
assert_eq!(mode, "signatures");
assert!(fresh);
}
#[test]
fn start_line_gt1_honors_explicit_lines_mode() {
let mut mode = "lines:1-50".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, true, Some(30));
assert_eq!(
mode, "lines:30-999999",
"explicit lines mode should accept start_line override"
);
assert!(fresh);
}
#[test]
fn start_line_none_does_nothing() {
let mut mode = "map".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, true, None);
assert_eq!(mode, "map");
assert!(!fresh);
}
#[test]
fn start_line_1_with_explicit_mode_preserves_it() {
let mut mode = "map".to_string();
let mut fresh = false;
apply_start_line(&mut mode, &mut fresh, true, Some(1));
assert_eq!(mode, "map");
assert!(!fresh);
}
#[test]
fn offset_is_alias_for_start_line() {
let mut mode = "auto".to_string();
let mut fresh = false;
super::apply_line_window(&mut mode, &mut fresh, false, None, Some(40), None);
assert_eq!(mode, "lines:40-999999");
assert!(fresh);
}
#[test]
fn offset_and_limit_make_bounded_window() {
let mut mode = "auto".to_string();
let mut fresh = false;
super::apply_line_window(&mut mode, &mut fresh, false, None, Some(40), Some(20));
assert_eq!(mode, "lines:40-59", "20 inclusive lines starting at 40");
assert!(fresh);
}
#[test]
fn limit_alone_reads_from_first_line() {
let mut mode = "auto".to_string();
let mut fresh = false;
super::apply_line_window(&mut mode, &mut fresh, false, None, None, Some(25));
assert_eq!(mode, "lines:1-25");
assert!(fresh);
}
#[test]
fn start_line_wins_over_offset_when_both_present() {
assert_eq!(
super::resolve_line_window(Some(10), Some(99), None),
Some((10, None))
);
}
#[test]
fn resolve_clamps_start_and_drops_nonpositive_limit() {
assert_eq!(
super::resolve_line_window(Some(-5), None, Some(0)),
Some((1, None))
);
assert_eq!(super::resolve_line_window(None, None, Some(-3)), None);
assert_eq!(super::resolve_line_window(None, None, None), None);
}
#[test]
fn lines_mode_bounds_are_inclusive() {
assert_eq!(super::lines_mode(40, Some(20)), "lines:40-59");
assert_eq!(super::lines_mode(5, None), "lines:5-999999");
}
#[test]
fn explicit_map_not_clobbered_by_offset_limit() {
let mut mode = "map".to_string();
let mut fresh = false;
super::apply_line_window(&mut mode, &mut fresh, true, None, Some(40), Some(20));
assert_eq!(mode, "map", "explicit mode wins over offset/limit");
assert!(fresh);
}
#[test]
fn schema_advertises_line_window_aliases() {
let tool = CtxReadTool.tool_def();
let props = tool
.input_schema
.get("properties")
.and_then(|p| p.as_object())
.expect("ctx_read schema has a properties object");
for key in ["path", "mode", "start_line", "offset", "limit", "fresh"] {
assert!(props.contains_key(key), "ctx_read schema missing '{key}'");
}
}
use crate::core::degradation_policy::DegradationVerdictV1;
#[test]
fn verdict_ok_does_not_degrade() {
let (mode, degraded) = super::apply_verdict("full", DegradationVerdictV1::Ok);
assert_eq!(mode, "full");
assert!(!degraded);
}
#[test]
fn verdict_warn_degrades_full_to_map() {
let (mode, degraded) = super::apply_verdict("full", DegradationVerdictV1::Warn);
assert_eq!(mode, "map");
assert!(degraded, "full→map must be flagged as degraded");
}
#[test]
fn verdict_warn_keeps_map() {
let (mode, degraded) = super::apply_verdict("map", DegradationVerdictV1::Warn);
assert_eq!(mode, "map");
assert!(!degraded, "map is not degraded under Warn");
}
#[test]
fn verdict_warn_keeps_signatures() {
let (mode, degraded) = super::apply_verdict("signatures", DegradationVerdictV1::Warn);
assert_eq!(mode, "signatures");
assert!(!degraded);
}
#[test]
fn verdict_throttle_degrades_full_to_signatures() {
let (mode, degraded) = super::apply_verdict("full", DegradationVerdictV1::Throttle);
assert_eq!(mode, "signatures");
assert!(degraded);
}
#[test]
fn verdict_throttle_degrades_map_to_signatures() {
let (mode, degraded) = super::apply_verdict("map", DegradationVerdictV1::Throttle);
assert_eq!(mode, "signatures");
assert!(degraded);
}
#[test]
fn verdict_throttle_keeps_lines() {
let (mode, degraded) = super::apply_verdict("lines:1-50", DegradationVerdictV1::Throttle);
assert_eq!(mode, "lines:1-50");
assert!(!degraded, "lines mode bypasses degradation");
}
#[test]
fn verdict_block_degrades_full_to_signatures() {
let (mode, degraded) = super::apply_verdict("full", DegradationVerdictV1::Block);
assert_eq!(mode, "signatures");
assert!(degraded);
}
#[test]
fn verdict_block_does_not_degrade_signatures() {
let (mode, degraded) = super::apply_verdict("signatures", DegradationVerdictV1::Block);
assert_eq!(mode, "signatures");
assert!(!degraded, "already at signatures — no degradation needed");
}
#[test]
fn degrade_warning_message_contains_mode_info() {
let (new_mode, degraded) = super::apply_verdict("full", DegradationVerdictV1::Warn);
assert!(degraded);
let warning = format!(
"⚠ Context pressure: mode=full was downgraded to mode={new_mode} (verdict: {:?}).",
DegradationVerdictV1::Warn
);
assert!(warning.contains("mode=full"));
assert!(warning.contains("mode=map"));
assert!(warning.contains("Warn"));
}
#[test]
fn auto_degrade_preserves_full_when_default_config() {
if std::env::var("LCTX_NO_DEGRADE").is_ok() {
return;
}
let (mode, warning) = super::auto_degrade_read_mode("full");
assert_eq!(mode, "full");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_map_when_default_config() {
if std::env::var("LCTX_NO_DEGRADE").is_ok() {
return;
}
let (mode, warning) = super::auto_degrade_read_mode("map");
assert_eq!(mode, "map");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_signatures_when_default_config() {
if std::env::var("LCTX_NO_DEGRADE").is_ok() {
return;
}
let (mode, warning) = super::auto_degrade_read_mode("signatures");
assert_eq!(mode, "signatures");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_diff_always() {
let (mode, warning) = super::auto_degrade_read_mode("diff");
assert_eq!(mode, "diff");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_lines_mode_always() {
let (mode, warning) = super::auto_degrade_read_mode("lines:10-50");
assert_eq!(mode, "lines:10-50");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_aggressive_when_default_config() {
if std::env::var("LCTX_NO_DEGRADE").is_ok() {
return;
}
let (mode, warning) = super::auto_degrade_read_mode("aggressive");
assert_eq!(mode, "aggressive");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_entropy_when_default_config() {
if std::env::var("LCTX_NO_DEGRADE").is_ok() {
return;
}
let (mode, warning) = super::auto_degrade_read_mode("entropy");
assert_eq!(mode, "entropy");
assert!(warning.is_none());
}
#[test]
fn auto_degrade_preserves_auto_when_default_config() {
if std::env::var("LCTX_NO_DEGRADE").is_ok() {
return;
}
let (mode, warning) = super::auto_degrade_read_mode("auto");
assert_eq!(mode, "auto");
assert!(warning.is_none());
}
#[test]
fn verdict_warn_does_not_degrade_diff() {
let (mode, degraded) = super::apply_verdict("diff", DegradationVerdictV1::Warn);
assert_eq!(mode, "diff");
assert!(!degraded);
}
#[test]
fn verdict_throttle_does_not_degrade_signatures() {
let (mode, degraded) = super::apply_verdict("signatures", DegradationVerdictV1::Throttle);
assert_eq!(mode, "signatures");
assert!(!degraded);
}
#[test]
fn verdict_ok_preserves_map() {
let (mode, degraded) = super::apply_verdict("map", DegradationVerdictV1::Ok);
assert_eq!(mode, "map");
assert!(!degraded);
}
#[test]
fn verdict_ok_preserves_signatures() {
let (mode, degraded) = super::apply_verdict("signatures", DegradationVerdictV1::Ok);
assert_eq!(mode, "signatures");
assert!(!degraded);
}
#[test]
fn verdict_ok_preserves_lines() {
let (mode, degraded) = super::apply_verdict("lines:1-100", DegradationVerdictV1::Ok);
assert_eq!(mode, "lines:1-100");
assert!(!degraded);
}
#[test]
fn verdict_block_degrades_map_to_signatures() {
let (mode, degraded) = super::apply_verdict("map", DegradationVerdictV1::Block);
assert_eq!(mode, "signatures");
assert!(degraded);
}
}
#[cfg(test)]
#[path = "ctx_read_repo_param_tests.rs"]
mod repo_param_tests;