diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index edd0320..c93f3e6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -40,6 +40,8 @@ jobs: run: python3 .github/scripts/test_dco_check.py - name: Test grounding qualification evidence gates run: python3 -m unittest discover -s testing/grounding -p test_qualify.py + - name: Test search adoption metric + run: python3 -m unittest discover -s testing/adoption -p test_search_adoption.py - name: Test npm launcher run: npm test - name: Inspect npm package payload diff --git a/crates/mcp-client/src/client.rs b/crates/mcp-client/src/client.rs index f344796..c6c7640 100644 --- a/crates/mcp-client/src/client.rs +++ b/crates/mcp-client/src/client.rs @@ -5160,14 +5160,24 @@ impl ContextStreamClient { return true; } - let committed_files = body - .get("indexed_files") - .and_then(serde_json::Value::as_i64) - .or_else(|| { - body.get("indexed_file_count") - .and_then(serde_json::Value::as_i64) - }) - .unwrap_or(0); + // The hosted status labels a project with no index row as + // `project_index_state: "ready"` / `status: "completed"`. A reported + // file count of zero is therefore authoritative: nothing is searchable, + // whatever the state label or a stale generation number says. State and + // generation only stand in for the count when the response omits it. + let reported_files = [ + body.get("indexed_files") + .and_then(serde_json::Value::as_i64), + body.get("indexed_file_count") + .and_then(serde_json::Value::as_i64), + ] + .into_iter() + .flatten() + .max(); + let committed_files = reported_files.unwrap_or(0); + if reported_files.is_some() && committed_files <= 0 { + return false; + } let ready_state = body .get("project_index_state") .or_else(|| body.get("status")) @@ -21107,6 +21117,41 @@ mod tests { ); } + #[test] + fn empty_project_labelled_ready_is_not_canonically_ready() { + // Shape the hosted status returns for a project that was never + // indexed: a "ready"/"completed" label with zero files and no `indexed`. + for status in [ + serde_json::json!({ + "project_index_state": "ready", + "status": "completed", + "status_detail": "no_files_indexed", + "indexed_files": 0, + "indexed_file_count": 0, + "total_files": 0, + "committed_generation": 0 + }), + serde_json::json!({ + "project_index_state": "ready", + "indexed_file_count": 0, + "committed_generation": 12 + }), + serde_json::json!({"indexed_files": 0, "status": "completed"}), + ] { + assert!( + !ContextStreamClient::project_index_status_reports_canonical_ready(&status), + "empty project must not report canonical readiness: {status}" + ); + } + // A positive count still wins over a missing label. + assert!( + ContextStreamClient::project_index_status_reports_canonical_ready(&serde_json::json!({ + "indexed_files": 0, + "indexed_file_count": 7 + })) + ); + } + #[tokio::test] async fn unroutable_checkout_status_keeps_canonical_evidence_separate_from_checkout_readiness() { diff --git a/crates/mcp-server/src/hook_handlers/pre_tool_use.rs b/crates/mcp-server/src/hook_handlers/pre_tool_use.rs index e8b17dc..d180596 100644 --- a/crates/mcp-server/src/hook_handlers/pre_tool_use.rs +++ b/crates/mcp-server/src/hook_handlers/pre_tool_use.rs @@ -65,6 +65,10 @@ const QUESTION_WORDS: &[&str] = &[ const DEFAULT_INDEX_WAIT_SECONDS: u64 = 20; const MIN_INDEX_WAIT_SECONDS: u64 = 15; const MAX_INDEX_WAIT_SECONDS: u64 = 20; +/// How long the "no index recorded for this checkout" shell-search nudge stays +/// quiet after it fires. Long enough that a session full of `rg` calls hears +/// it once, short enough to resurface if the agent keeps ignoring it. +const SHELL_SEARCH_NUDGE_COOLDOWN_SECONDS: u64 = 600; /// Check if a glob pattern is a broad discovery pattern. /// Allows targeted patterns like "src/models/*.rs", "web/src/**/*sidebar*", @@ -1174,6 +1178,27 @@ fn is_local_discovery_tool_during_index_wait(tool_lower: &str, tool_input: &Valu } } +/// Nudge for shell code search (`rg`, `grep -r`, `find -name`, `fd`) in a +/// checkout that has no recorded index. +/// +/// Telling the agent to "use search" here would send it to an empty index; one +/// empty result and most agents stop searching for the session. So the nudge is +/// honest: search has nothing yet, build the index once, then prefer search. +/// Returns `None` for commands that are not code discovery (log filtering, +/// process lists, a single targeted file), so those never get nudged. +fn unindexed_shell_search_nudge(editor: &EditorFormat, command: &str) -> Option { + let (tool_name, query_hint) = detect_bash_code_search(command)?; + let (mode, _) = recommend_search_mode(&query_hint); + let search = search_call(editor, mode, &query_hint); + let project = contextstream_tool_name(editor, "project"); + Some(format!( + "ContextStream has no index recorded for this checkout, so search may come back empty until one exists. \ + Build it once with {project}(action=\"index\") (it runs in the background and search fills in as files commit), \ + then use {search} instead of shell `{tool_name}` for code discovery. \ + Shell search is fine in the meantime." + )) +} + fn is_contextstream_read_only_operation(tool_name: &str, tool_input: &Value) -> bool { let action = first_str(tool_input, &["action"]) .unwrap_or("") @@ -2428,6 +2453,20 @@ pub async fn handle() -> Result<()> { return Ok(()); } + // Shell code search is the same discovery the indexed path redirects, + // but it used to fall through to the silent Allow below, so a checkout + // missing from the local registry never got a nudge at all. + if tool_lower == "bash" { + let command = first_str(&tool_input, &["command"]).unwrap_or("").trim(); + if let Some(msg) = unindexed_shell_search_nudge(&editor, command) { + if prompt_state::claim_shell_search_nudge(&cwd, SHELL_SEARCH_NUDGE_COOLDOWN_SECONDS) + { + emit(HookDecision::AllowWithContext(msg))?; + return Ok(()); + } + } + } + // Non-discovery tools should continue while refresh runs in background. emit(HookDecision::Allow)?; return Ok(()); @@ -3611,6 +3650,41 @@ mod tests { )); } + #[test] + fn unindexed_checkout_nudges_shell_code_search_toward_indexing() { + for command in [ + "rg -n 'canonical_index_ready' crates", + "grep -rn \"PreToolUse\" .", + "cd crates && rg -n foo", + "find . -name '*.rs'", + ] { + let nudge = unindexed_shell_search_nudge(&EditorFormat::Claude, command) + .unwrap_or_else(|| panic!("expected a nudge for: {command}")); + assert!(nudge.contains("mcp__contextstream__project(action=\"index\")")); + assert!(nudge.contains("mcp__contextstream__search")); + // Must not claim an index exists. + assert!(nudge.contains("no index recorded")); + } + } + + #[test] + fn unindexed_checkout_does_not_nudge_non_discovery_shell_commands() { + for command in [ + "ps aux | grep node", + "grep ERROR /var/log/app.log", + "grep -n foo src/main.rs", + "git status", + "cargo test", + "find . -newer Cargo.lock", + "", + ] { + assert!( + unindexed_shell_search_nudge(&EditorFormat::Claude, command).is_none(), + "unexpected nudge for: {command}" + ); + } + } + #[test] fn index_wait_never_blocks_grep() { assert!(!is_local_discovery_tool_during_index_wait( diff --git a/crates/mcp-server/src/hook_handlers/prompt_state.rs b/crates/mcp-server/src/hook_handlers/prompt_state.rs index dd512b9..4fbabc1 100644 --- a/crates/mcp-server/src/hook_handlers/prompt_state.rs +++ b/crates/mcp-server/src/hook_handlers/prompt_state.rs @@ -27,6 +27,10 @@ struct PromptStateEntry { index_wait_started_at: Option, #[serde(default)] index_wait_until: Option, + /// When the "no index recorded for this checkout" shell-search nudge last + /// fired, so it repeats at most once per cooldown instead of on every `rg`. + #[serde(default)] + shell_search_nudged_at: Option, updated_at: String, } @@ -97,6 +101,7 @@ pub fn mark_context_required(cwd: &str) { last_state_change_at: None, index_wait_started_at: None, index_wait_until: None, + shell_search_nudged_at: None, updated_at: now.clone(), }); entry.require_context = true; @@ -135,6 +140,7 @@ pub fn mark_init_required(cwd: &str) { last_state_change_at: None, index_wait_started_at: None, index_wait_until: None, + shell_search_nudged_at: None, updated_at: now.clone(), }); entry.require_init = true; @@ -189,6 +195,7 @@ pub fn mark_state_changed(cwd: &str) { last_state_change_at: None, index_wait_started_at: None, index_wait_until: None, + shell_search_nudged_at: None, updated_at: now.clone(), }); entry.last_state_change_at = Some(now.clone()); @@ -272,6 +279,7 @@ pub fn start_index_wait_window(cwd: &str, wait_seconds: u64) { last_state_change_at: None, index_wait_started_at: None, index_wait_until: None, + shell_search_nudged_at: None, updated_at: now_iso.clone(), }); @@ -322,6 +330,55 @@ pub fn index_wait_remaining_seconds(cwd: &str) -> Option { (remaining > 0).then_some(remaining as u64) } +/// Whether a nudge last shown at `last` may be shown again at `now`. +fn nudge_cooldown_elapsed( + last: Option<&str>, + now: chrono::DateTime, + cooldown_seconds: u64, +) -> bool { + match last.and_then(parse_rfc3339_utc) { + Some(last) => now.signed_duration_since(last).num_seconds() >= cooldown_seconds as i64, + None => true, + } +} + +/// Claim the right to show the shell-search "no index recorded" nudge for +/// `cwd`. Returns `true` at most once per `cooldown_seconds`; the caller emits +/// the nudge only when it does. +pub fn claim_shell_search_nudge(cwd: &str, cooldown_seconds: u64) -> bool { + if cwd.trim().is_empty() { + return false; + } + let mut state = read_state(); + let now = chrono::Utc::now(); + let now_iso = now.to_rfc3339(); + let entry = state + .workspaces + .entry(cwd.to_string()) + .or_insert(PromptStateEntry { + require_context: false, + require_init: false, + last_context_at: None, + last_state_change_at: None, + index_wait_started_at: None, + index_wait_until: None, + shell_search_nudged_at: None, + updated_at: now_iso.clone(), + }); + + if !nudge_cooldown_elapsed( + entry.shell_search_nudged_at.as_deref(), + now, + cooldown_seconds, + ) { + return false; + } + entry.shell_search_nudged_at = Some(now_iso.clone()); + entry.updated_at = now_iso; + write_state(&state); + true +} + pub fn cleanup_stale(max_age_minutes: u64) { let mut state = read_state(); let now = chrono::Utc::now(); @@ -357,6 +414,7 @@ mod tests { last_state_change_at: None, index_wait_started_at: None, index_wait_until: None, + shell_search_nudged_at: None, updated_at: chrono::Utc::now().to_rfc3339(), }, ); @@ -375,6 +433,37 @@ mod tests { assert!(parsed.workspaces["/tmp/project"].index_wait_until.is_none()); } + #[test] + fn shell_search_nudge_cooldown_only_repeats_after_it_elapses() { + let now = chrono::Utc::now(); + // Never nudged, or an unreadable timestamp: allowed. + assert!(nudge_cooldown_elapsed(None, now, 600)); + assert!(nudge_cooldown_elapsed(Some("not a timestamp"), now, 600)); + // Nudged a minute ago with a ten minute cooldown: suppressed. + let recent = (now - chrono::Duration::seconds(60)).to_rfc3339(); + assert!(!nudge_cooldown_elapsed(Some(&recent), now, 600)); + // Nudged eleven minutes ago: allowed again. + let old = (now - chrono::Duration::seconds(660)).to_rfc3339(); + assert!(nudge_cooldown_elapsed(Some(&old), now, 600)); + } + + #[test] + fn legacy_prompt_state_defaults_shell_search_nudge_to_never() { + let legacy = serde_json::json!({ + "workspaces": { + "/tmp/project": { + "require_context": false, + "updated_at": chrono::Utc::now().to_rfc3339() + } + } + }); + let parsed: PromptStateFile = + serde_json::from_value(legacy).expect("parse legacy prompt state"); + assert!(parsed.workspaces["/tmp/project"] + .shell_search_nudged_at + .is_none()); + } + #[test] fn legacy_prompt_state_defaults_require_init_false() { let legacy = serde_json::json!({ diff --git a/crates/mcp-server/src/setup/mod.rs b/crates/mcp-server/src/setup/mod.rs index 7bf28c2..ebca829 100644 --- a/crates/mcp-server/src/setup/mod.rs +++ b/crates/mcp-server/src/setup/mod.rs @@ -1983,7 +1983,12 @@ pub async fn update_rules_scoped( // 3. Resolved from the API (so the rule header shows a real UUID + name // instead of the null UUID when neither 1 nor 2 is available) // 4. Inferred from existing rule-file headers (offline fallback) - let (ws_id, ws_name): (Option, Option) = if workspace_id.is_some() { + // Only project rules carry a workspace identity (global rules are + // workspace-neutral), so a global-only refresh needs none and skips the API + // lookup entirely. + let (ws_id, ws_name): (Option, Option) = if !include_project { + (None, None) + } else if workspace_id.is_some() { ( workspace_id.map(String::from), workspace_name.map(String::from), @@ -2022,7 +2027,7 @@ pub async fn update_rules_scoped( // Update global rules if scope == "global" || scope == "all" { - match rules::write_editor_rules(editor, ws_id_ref, ws_name_ref) { + match rules::write_editor_rules(editor) { Ok(()) => updated.push("global rules"), Err(e) => { if !e.to_string().contains("Could not determine rules path") { @@ -3641,7 +3646,7 @@ pub async fn configure_editor_with_workspace( } // Generate AI rules (global) - match rules::write_editor_rules(editor, workspace_id, workspace_name) { + match rules::write_editor_rules(editor) { Ok(()) => { rules_targets.push("global"); } diff --git a/crates/mcp-server/src/setup/rules.rs b/crates/mcp-server/src/setup/rules.rs index 749a59e..f72d0de 100644 --- a/crates/mcp-server/src/setup/rules.rs +++ b/crates/mcp-server/src/setup/rules.rs @@ -1175,11 +1175,13 @@ fn apply_mcp_prefix(content: &str) -> String { } /// Write editor rules file (global). -pub fn write_editor_rules( - editor: &Editor, - workspace_id: Option<&str>, - workspace_name: Option<&str>, -) -> Result<()> { +/// +/// Global rules apply to every project the editor opens, so they never carry a +/// workspace identity: the rules tell the agent to use the ids `init(...)` and +/// `context(...)` return. Stamping the workspace of whichever directory setup or +/// repair happened to run in made the one global file name a different workspace +/// each time (and disagree with the project rules loaded beside it). +pub fn write_editor_rules(editor: &Editor) -> Result<()> { let paths = editor.all_rules_paths(None); if paths.is_empty() { return Err(anyhow::anyhow!( @@ -1188,7 +1190,7 @@ pub fn write_editor_rules( )); } - let rules = global_rules_content(editor, workspace_id, workspace_name); + let rules = global_rules_content(editor, None, None); let primary = &paths[0]; write_contextstream_block_to_path(primary, &rules, true)?; @@ -1209,7 +1211,7 @@ pub fn write_editor_rules( if *editor == Editor::Aider { if let Some(home) = dirs::home_dir() { let shared = home.join(".contextstream").join("rules.md"); - let shared_content = shared_rules_content(workspace_id, workspace_name, None); + let shared_content = shared_rules_content(None, None, None); write_contextstream_block_to_path(&shared, &shared_content, true)?; } } @@ -4272,6 +4274,52 @@ mod tests { } } + #[test] + fn rewriting_global_rules_drops_a_stale_workspace_stamp_and_keeps_user_content() { + let _guard = env_test_mutex().lock().unwrap_or_else(|e| e.into_inner()); + let temp = tempdir().expect("tempdir"); + let previous_home = std::env::var_os("HOME"); + std::env::set_var("HOME", temp.path()); + + // A global block stamped with some other directory's workspace, the way + // repeated setup/repair runs left it, with the user's own text around it. + let global_rules = temp.path().join(".codex").join("AGENTS.md"); + std::fs::create_dir_all(global_rules.parent().expect("parent")).expect("mkdirs"); + std::fs::write( + &global_rules, + format!( + "my own notes above\n\n{}\n{} 0123456789abcdef -->\n# Workspace: Stale Example Workspace\n# Workspace ID: 11111111-2222-4333-8444-555555555555\n# ContextStream Rules\nMANDATORY STARTUP:\n{}\n\nmy own notes below\n", + CONTEXTSTREAM_START, RULES_HASH_MARKER_PREFIX, CONTEXTSTREAM_END + ), + ) + .expect("seed stale global rules"); + + let result = write_editor_rules(&Editor::Codex); + let rewritten = std::fs::read_to_string(&global_rules).unwrap_or_default(); + + if let Some(value) = previous_home { + std::env::set_var("HOME", value); + } else { + std::env::remove_var("HOME"); + } + + result.expect("global rules rewrite"); + assert!( + !rewritten.contains("Stale Example Workspace") + && !rewritten.contains("11111111-2222-4333-8444-555555555555"), + "stale workspace identity survived a global rewrite:\n{rewritten}" + ); + assert!(rewritten.contains(&format!("# Workspace: {DEFAULT_WORKSPACE_NAME}"))); + assert!(rewritten.contains(&format!("# Workspace ID: {DEFAULT_WORKSPACE_ID}"))); + assert!(rewritten.starts_with("my own notes above\n"), "{rewritten}"); + assert!(rewritten.contains("my own notes below"), "{rewritten}"); + assert_eq!( + rewritten.matches(CONTEXTSTREAM_START).count(), + 1, + "exactly one managed block expected:\n{rewritten}" + ); + } + #[test] fn test_remove_contextstream_from_directory_style_target() { let temp = tempdir().expect("tempdir"); diff --git a/crates/mcp-tools/src/domains/session.rs b/crates/mcp-tools/src/domains/session.rs index c9cbb60..2c3c8de 100644 --- a/crates/mcp-tools/src/domains/session.rs +++ b/crates/mcp-tools/src/domains/session.rs @@ -4438,6 +4438,73 @@ fn truncate_context_wire_text(value: &str, max_chars: usize) -> String { output } +/// Smallest truncated head worth keeping instead of dropping a block outright: +/// below this the block is a sliver that reads as noise. +const CONTEXT_WIRE_MIN_PARTIAL_BLOCK_TOKENS: usize = 200; + +/// The longest truncated head of `blocks[index]` that, with the other blocks, +/// still fits `target_tokens`. `None` when dropping the block would not fit +/// anyway, when what would fit is only a sliver, or when the block carries +/// `` boundaries (those are only ever kept or dropped whole). +fn truncate_context_wire_block_to_fit( + blocks: &[ContextWireTextBlock], + index: usize, + structured: Option<&Value>, + target_tokens: usize, +) -> Option { + let block = &blocks[index]; + if block.text.contains("") || block.text.contains("") { + return None; + } + + let mut candidate_blocks = blocks.to_vec(); + candidate_blocks[index].text = String::new(); + let base_tokens = estimated_context_tool_wire_tokens_with_optional( + &render_context_wire_text_blocks(&candidate_blocks), + structured, + ); + if base_tokens.saturating_add(CONTEXT_WIRE_MIN_PARTIAL_BLOCK_TOKENS) > target_tokens { + return None; + } + + let ends_with_newline = block.text.ends_with('\n'); + let render_head = |chars: usize| { + let mut head = truncate_context_wire_text(&block.text, chars); + // Keep the next block's tag on its own line. + if ends_with_newline && !head.ends_with('\n') { + head.push('\n'); + } + head + }; + let fits = |head: String, candidate_blocks: &mut Vec| { + candidate_blocks[index].text = head; + estimated_context_tool_wire_tokens_with_optional( + &render_context_wire_text_blocks(candidate_blocks), + structured, + ) <= target_tokens + }; + + let mut low = 0usize; + let mut high = block.text.chars().count(); + let mut best: Option<(usize, String)> = None; + while low <= high { + let mid = low + (high - low) / 2; + if fits(render_head(mid), &mut candidate_blocks) { + best = Some((mid, render_head(mid))); + low = mid.saturating_add(1); + } else if mid == 0 { + break; + } else { + high = mid - 1; + } + } + + // `estimated_context_tool_wire_tokens` is bytes / 4, so a head is worth + // keeping once it is roughly the minimum token count in characters. + best.filter(|(chars, _)| *chars >= CONTEXT_WIRE_MIN_PARTIAL_BLOCK_TOKENS * 2) + .map(|(_, head)| head) +} + fn reduce_context_wire_text( text: &str, structured: Option<&Value>, @@ -4445,10 +4512,14 @@ fn reduce_context_wire_text( ) -> (String, usize, bool) { let mut blocks = context_wire_text_blocks(text); let mut dropped_blocks = 0usize; + let mut truncated_block = false; // Drop complete low/normal/high blocks in priority order. Within a tier, // remove the largest block first; ties remove the later block so retained - // output stays stable and front-loaded. + // output stays stable and front-loaded. When the largest block is the only + // thing standing between the response and its budget, keep a truncated head + // of it instead: removing a 9k-token context pack whole used to leave a 4k + // budget less than a third used. for priority in 0u8..=2 { while estimated_context_tool_wire_tokens_with_optional( &render_context_wire_text_blocks(&blocks), @@ -4469,6 +4540,13 @@ fn reduce_context_wire_text( let Some(remove_at) = remove_at else { break; }; + if let Some(head) = + truncate_context_wire_block_to_fit(&blocks, remove_at, structured, target_tokens) + { + blocks[remove_at].text = head; + truncated_block = true; + break; + } blocks.remove(remove_at); dropped_blocks += 1; } @@ -4476,7 +4554,7 @@ fn reduce_context_wire_text( let rendered = render_context_wire_text_blocks(&blocks); if estimated_context_tool_wire_tokens_with_optional(&rendered, structured) <= target_tokens { - return (rendered, dropped_blocks, false); + return (rendered, dropped_blocks, truncated_block); } // A response can carry more than one critical semantic unit, most often a @@ -4558,14 +4636,28 @@ const CONTEXT_STRUCTURED_DROP_ORDER: &[&str] = &[ "recent_decisions", "memory_nodes", "flash_suggestions", + // `items`, `summary` and `context` repeat what the text already carries + // (the `[CTX]` block and friends) and are by far the largest structured + // fields, so they are expendable before anything an agent acts on. "items", + "summary", + "context", +]; + +/// Small structured fields an agent acts on: skills, lessons, instructions, +/// coordination notices and grounding evidence. They are kept while the text can +/// still be shrunk to make room, and only dropped (in this order) when even the +/// critical text plus these fields cannot fit the budget. +const CONTEXT_STRUCTURED_PROTECTED_ORDER: &[&str] = &[ "matched_skills_typed", "matched_skills", "lessons", "remember_items", "instructions", - "summary", - "context", + // Listed explicitly: unknown fields fall into the lexical pass, where + // `coordination_inbox` sorts early and would be removed before anything + // else. + "coordination_inbox", // This is the structured form of the tool's primary job. Preserve it // after duplicated context/summary fields and every diagnostic surface. "grounding_hits", @@ -4577,6 +4669,71 @@ fn record_context_wire_field(fields: &mut Vec, field: &str) { } } +/// Remove `order`'s structured fields one at a time until the whole wire fits. +fn drop_context_wire_fields_while_over( + order: &[&str], + text: &str, + structured: &mut Value, + target_tokens: usize, + dropped_fields: &mut Vec, +) { + for field in order { + if estimated_context_tool_wire_tokens(text, structured) <= target_tokens { + break; + } + if structured + .as_object_mut() + .and_then(|object| object.remove(*field)) + .is_some() + { + record_context_wire_field(dropped_fields, field); + } + } +} + +/// Whether the structured fields still present fit the budget beside the text's +/// high-priority and critical blocks (instructions, skills, grounding, notices), +/// i.e. whether shrinking only the *low-priority* text can bring the whole wire +/// within budget. +/// +/// The text repeats most of what the structured fields carry, so when it is the +/// oversized part, dropping small structured fields cannot fix the budget; it +/// only loses them (and the ids a follow-up call needs) before the text is ever +/// touched. The check deliberately counts the high-priority text blocks as +/// untouchable too: a client may show the model only the text, so under a budget +/// too tight to keep both copies this keeps the old behaviour (structured fields +/// go first) rather than trading the text's instructions for their duplicates. +fn context_wire_fits_by_shrinking_text( + text: &str, + structured: &Value, + target_tokens: usize, +) -> bool { + let kept_text: String = context_wire_text_blocks(text) + .iter() + .filter(|block| block.priority >= 2) + .map(|block| block.text.as_str()) + .collect(); + estimated_context_tool_wire_tokens(&kept_text, structured) <= target_tokens +} + +/// Notice, then shrink the text to the budget. Returns the dropped block count +/// and whether any block was truncated. +fn compact_context_wire_text( + text: &mut String, + structured: &Value, + target_tokens: usize, + estimated_tokens_before: usize, +) -> (usize, bool) { + if estimated_tokens_before > target_tokens && !text.contains("[WIRE_BUDGET]") { + text.push_str( + "\n[WIRE_BUDGET] Whole-wire context compacted to the requested token envelope.", + ); + } + let reduced = reduce_context_wire_text(text, Some(structured), target_tokens); + *text = reduced.0; + (reduced.1, reduced.2) +} + fn budget_context_wire_payload( mut text: String, mut structured: Value, @@ -4615,19 +4772,38 @@ fn budget_context_wire_payload( ); } - for field in CONTEXT_STRUCTURED_DROP_ORDER { - if estimated_context_tool_wire_tokens(&text, &structured) <= target_tokens { - break; - } - if structured - .as_object_mut() - .and_then(|object| object.remove(*field)) - .is_some() - { - record_context_wire_field(&mut dropped_fields, field); - } + drop_context_wire_fields_while_over( + CONTEXT_STRUCTURED_DROP_ORDER, + &text, + &mut structured, + target_tokens, + &mut dropped_fields, + ); + + // What is left is small, actionable and identifying. If shrinking the text + // can fit the budget beside it, do that first instead of stripping every + // structured field to chase a budget the text is what's blowing. + if estimated_context_tool_wire_tokens(&text, &structured) > target_tokens + && context_wire_fits_by_shrinking_text(&text, &structured, target_tokens) + { + let (dropped, truncated) = compact_context_wire_text( + &mut text, + &structured, + target_tokens, + estimated_tokens_before, + ); + dropped_text_blocks += dropped; + truncated_text |= truncated; } + drop_context_wire_fields_while_over( + CONTEXT_STRUCTURED_PROTECTED_ORDER, + &text, + &mut structured, + target_tokens, + &mut dropped_fields, + ); + // Unknown/forward-compatible fields are still subject to the same wire // budget. Remove them in lexical order, retaining only the report until // text has had a chance to preserve the agent-facing critical guidance. @@ -4658,15 +4834,14 @@ fn budget_context_wire_payload( } if estimated_context_tool_wire_tokens(&text, &structured) > target_tokens { - if estimated_tokens_before > target_tokens { - text.push_str( - "\n[WIRE_BUDGET] Whole-wire context compacted to the requested token envelope.", - ); - } - let reduced = reduce_context_wire_text(&text, Some(&structured), target_tokens); - text = reduced.0; - dropped_text_blocks += reduced.1; - truncated_text |= reduced.2; + let (dropped, truncated) = compact_context_wire_text( + &mut text, + &structured, + target_tokens, + estimated_tokens_before, + ); + dropped_text_blocks += dropped; + truncated_text |= truncated; } if let Some(report) = structured diff --git a/crates/mcp-tools/src/domains/session_tests.rs b/crates/mcp-tools/src/domains/session_tests.rs index f8be02f..306e48d 100644 --- a/crates/mcp-tools/src/domains/session_tests.rs +++ b/crates/mcp-tools/src/domains/session_tests.rs @@ -250,6 +250,33 @@ mod init_index_status_tests { assert!(!init_index_status_reports_ready(&status, true)); assert!(init_index_status_reports_ready(&status, false)); } + + #[test] + fn init_does_not_report_an_empty_project_as_ready() { + // What the hosted status returns for a project that was never indexed: + // a "ready"/"completed" label, zero files, and no `indexed` field. + let status = json!({ + "project_index_state": "ready", + "status": "completed", + "status_detail": "no_files_indexed", + "indexed_files": 0, + "indexed_file_count": 0, + "total_files": 0, + "committed_generation": 0 + }); + + assert!(!init_index_status_reports_ready(&status, true)); + assert!(!init_index_status_reports_ready(&status, false)); + + // The same project once files are committed is ready. + let indexed = json!({ + "project_index_state": "ready", + "status": "completed", + "indexed_files": 41, + "indexed_file_count": 41 + }); + assert!(init_index_status_reports_ready(&indexed, true)); + } } #[test] @@ -841,6 +868,125 @@ M:request metrics 403 handling details ); } + /// A production-sized `context` response: the text carries every block, and + /// the structured copy repeats the large context pack beside a few small, + /// actionable fields (about 15k+ tokens before budgeting). + fn production_sized_context_wire_payload() -> (String, serde_json::Value) { + let pack = "repository context and implementation detail\n".repeat(800); + let instructions = "run search for code discovery ".repeat(60); + let evidence = "prior session evidence ".repeat(100); + let text = format!( + "[SEARCH] {}\n[ACCOUNT_CONTEXT] {}\n[TEAM_CONTEXT] {}\n[CTX]\n{}[/CTX]\n[MATCHED_SKILLS] {}\n[INSTRUCTIONS] {}\n[COORDINATION] 2 pending handoffs: review the index fix, confirm the worktree binding.\n[GROUNDING] {}", + "search reminder ".repeat(200), + "account detail ".repeat(200), + "team detail ".repeat(200), + pack, + "skill guidance ".repeat(60), + instructions, + evidence, + ); + let structured = json!({ + "context": pack, + "instructions": instructions, + "matched_skills": [{"name": "index-repair", "description": "skill guidance ".repeat(20)}], + "coordination_inbox": { + "pending": 2, + "items": [{"title": "review the index fix"}, {"title": "confirm the worktree binding"}] + }, + "grounding_hits": [{"content": evidence}], + "why_this_context": {"trace": "diagnostic detail ".repeat(100)}, + }); + (text, structured) + } + + #[test] + fn oversized_context_pack_does_not_cost_the_small_actionable_fields() { + let requested = 4_000usize; + let (text, structured) = production_sized_context_wire_payload(); + let before = estimated_context_tool_wire_tokens(&text, &structured); + assert!( + before > 4 * requested, + "fixture must be far over budget: {before}" + ); + + let (text, structured) = budget_context_wire_payload(text, structured, requested); + let report = &structured["wire_budget"]; + + assert!( + estimated_context_tool_wire_tokens(&text, &structured) + <= requested + CONTEXT_WIRE_ENVELOPE_TOKENS + ); + // The structured context pack only duplicates the `[CTX]` text block, so + // it is what goes. The few hundred tokens of guidance, the coordination + // inbox and the grounding evidence must not be sacrificed for it. + assert!(structured.get("context").is_none(), "report={report:?}"); + for field in [ + "instructions", + "matched_skills", + "coordination_inbox", + "grounding_hits", + ] { + assert!( + structured.get(field).is_some(), + "{field} was dropped to make room for the duplicate context pack; report={report:?}" + ); + } + for tag in [ + "[INSTRUCTIONS]", + "[COORDINATION]", + "[MATCHED_SKILLS]", + "[GROUNDING]", + ] { + assert!(text.contains(tag), "{tag} missing from {text:.200}"); + } + } + + #[test] + fn budget_is_filled_by_truncating_a_block_rather_than_dropping_it_whole() { + let requested = 4_000usize; + let (text, structured) = production_sized_context_wire_payload(); + let (text, structured) = budget_context_wire_payload(text, structured, requested); + let target = requested + CONTEXT_WIRE_ENVELOPE_TOKENS; + let estimated = estimated_context_tool_wire_tokens(&text, &structured); + + assert!(estimated <= target); + // Dropping the 9k-token `[CTX]` block whole used to leave about 1.3k of a + // 4.1k budget. A truncated head of it is more useful than nothing. + assert!( + estimated * 100 >= target * 70, + "budget left unused: {estimated} of {target} tokens" + ); + assert!(text.contains("[CTX]"), "context pack lost entirely"); + assert!(text.contains("repository context and implementation detail")); + assert_eq!(structured["wire_budget"]["truncated_text"], true); + } + + #[test] + fn budget_too_tight_for_both_copies_keeps_the_text_instructions() { + // A client may show the model only the text. When the budget cannot hold + // a block's text and its structured duplicate, the text must win. + let requested = 1_500usize; + let (text, structured) = production_sized_context_wire_payload(); + let (text, structured) = budget_context_wire_payload(text, structured, requested); + + assert!( + estimated_context_tool_wire_tokens(&text, &structured) + <= requested + CONTEXT_WIRE_ENVELOPE_TOKENS + ); + for tag in [ + "[INSTRUCTIONS]", + "[COORDINATION]", + "[MATCHED_SKILLS]", + "[GROUNDING]", + ] { + assert!( + text.contains(tag), + "{tag} lost from the text while its structured copy was kept: {:?}", + structured["wire_budget"] + ); + } + } + #[test] fn whole_wire_budget_respects_each_requested_size() { for requested in [50usize, 120, 240, 480, 800] { diff --git a/testing/adoption/search_adoption.py b/testing/adoption/search_adoption.py new file mode 100644 index 0000000..99b289f --- /dev/null +++ b/testing/adoption/search_adoption.py @@ -0,0 +1,285 @@ +#!/usr/bin/env python3 +"""ContextStream search adoption from local agent transcripts. + +Answers one question per agent: when an agent searches code, does it use +ContextStream search or shell `rg`/`grep`/`find`? Run it before and after a +change that is meant to move that number and compare the two JSON files. + +Reads Claude Code (`~/.claude/projects/**/*.jsonl`) and Codex +(`~/.codex/sessions/**/*.jsonl`) transcripts; nothing leaves the machine and no +transcript text is printed or stored, only counts. + +The shell classifier is a port of `detect_bash_code_search` in +`crates/mcp-server/src/hook_handlers/pre_tool_use.rs`, so "shell code search" +here means exactly what the PreToolUse hook would nudge. Keep the two in step. + +The Claude Code reader is validated against real transcripts. The Codex reader +accepts the `function_call` / `local_shell_call` shapes seen in recent rollouts +but has not been checked against every Codex release; treat its counts as a +floor and rerun if a Codex update changes the transcript format. +""" +import argparse +import json +import sys +import time +from collections import Counter +from pathlib import Path + +CONTEXTSTREAM_PREFIX = "mcp__contextstream__" +NATIVE_SEARCH_TOOLS = {"Grep", "Glob"} +FIND_NAME_FLAGS = (" -name ", " -iname ", " -path ", " -ipath ", " -regex ", " -iregex ") +CONCRETE_FILE_BLOCKERS = set("*?[]()|\\$^") + + +def _recursive_flag(head): + for token in head.split(): + if token == "--recursive": + return True + if ( + token.startswith("-") + and not token.startswith("--") + and len(token) > 1 + and token[1:].isalpha() + and any(c in "rR" for c in token[1:]) + ): + return True + return False + + +def _log_or_text_target(token): + lower = token.strip("'\"").lower() + return lower.endswith((".log", ".txt")) or lower.startswith("/var/log") + + +def _concrete_file_target(token): + t = token.strip("'\"") + if not t or t.endswith("/") or any(c in CONCRETE_FILE_BLOCKERS for c in t): + return False + last = t.rsplit("/", 1)[-1] + dot = last.rfind(".") + return 0 < dot < len(last) - 1 + + +def shell_code_search(command): + """Return the search tool name when `command` is shell code discovery.""" + head = (command or "").strip() + if not head: + return None + while True: + stripped = head.lstrip() + if stripped.startswith("cd "): + rest = stripped[3:] + for separator in ("&&", ";"): + before, found, after = rest.partition(separator) + if found: + head = after.lstrip() + break + else: + break + continue + break + # Anything piped is filtering, not discovery. + if "|" in head: + return None + tokens = head.split() + if not tokens: + return None + tool, needs_find_flag = { + "grep": ("grep", False), + "egrep": ("grep", False), + "fgrep": ("grep", False), + "rg": ("rg", False), + "ripgrep": ("rg", False), + "ag": ("ag", False), + "find": ("find", True), + "fd": ("fd", False), + "fdfind": ("fd", False), + }.get(tokens[0], (None, False)) + if tool is None: + return None + if needs_find_flag and not any(flag in head for flag in FIND_NAME_FLAGS): + return None + recursive = _recursive_flag(head) + for token in tokens[1:]: + if token.startswith("-"): + continue + if _log_or_text_target(token): + return None + if not recursive and tool in ("grep", "rg", "ag") and _concrete_file_target(token): + return None + return tool + + +def _command_text(arguments): + """The shell command in a tool-call argument object, if any.""" + if isinstance(arguments, str): + try: + arguments = json.loads(arguments) + except ValueError: + return None + if not isinstance(arguments, dict): + return None + command = arguments.get("command", arguments.get("cmd")) + if isinstance(command, list): + # ["bash", "-lc", "rg foo"]: the script is the last element. + command = command[-1] if command else None + return command if isinstance(command, str) else None + + +def _contextstream_tool(name): + """`search` for `mcp__contextstream__search`, `contextstream__search`, ...""" + lowered = (name or "").lower() + if "contextstream" not in lowered: + return None + return lowered.rsplit("__", 1)[-1] or None + + +def records(path): + """JSON objects of a JSONL transcript, skipping lines that are not JSON.""" + with path.open(errors="replace") as handle: + for line in handle: + try: + record = json.loads(line) + except ValueError: + continue + if isinstance(record, dict): + yield record + + +def claude_calls(path): + """Yield (kind, detail) for each tool call in a Claude Code transcript.""" + for record in records(path): + if record.get("type") != "assistant": + continue + content = (record.get("message") or {}).get("content") + if not isinstance(content, list): + continue + for block in content: + if not isinstance(block, dict) or block.get("type") != "tool_use": + continue + name = block.get("name") or "" + arguments = block.get("input") + if name.startswith(CONTEXTSTREAM_PREFIX): + yield "contextstream", name[len(CONTEXTSTREAM_PREFIX):] + elif name in NATIVE_SEARCH_TOOLS: + yield "native_search", name + elif name == "Bash": + tool = shell_code_search(_command_text(arguments)) + if tool: + yield "shell_search", tool + + +def codex_calls(path): + """Yield (kind, detail) for each tool call in a Codex rollout.""" + for record in records(path): + payload = record.get("payload") if isinstance(record.get("payload"), dict) else record + if payload.get("type") not in ("function_call", "local_shell_call", "custom_tool_call"): + continue + name = payload.get("name") or "" + arguments = payload.get("arguments", payload.get("action", payload.get("input"))) + tool = _contextstream_tool(name) + if tool: + yield "contextstream", tool + continue + command = _command_text(arguments) + search = shell_code_search(command) + if search: + yield "shell_search", search + + +READERS = {"claude-code": claude_calls, "codex": codex_calls} + + +def transcripts(root, since): + root = Path(root).expanduser() + if not root.is_dir(): + return + for path in root.rglob("*.jsonl"): + try: + if path.stat().st_mtime >= since: + yield path + except OSError: + continue + + +def measure(agent, root, since): + reader = READERS[agent] + sessions = with_cs_search = 0 + totals = Counter() + cs_tools = Counter() + shell_tools = Counter() + for path in transcripts(root, since): + sessions += 1 + used_search = False + for kind, detail in reader(path): + totals[kind] += 1 + if kind == "contextstream": + cs_tools[detail] += 1 + used_search |= detail == "search" + elif kind == "shell_search": + shell_tools[detail] += 1 + with_cs_search += used_search + cs_search = cs_tools["search"] + searches = cs_search + totals["shell_search"] + totals["native_search"] + return { + "sessions": sessions, + "sessions_with_contextstream_search": with_cs_search, + "contextstream_search": cs_search, + "contextstream_other": totals["contextstream"] - cs_search, + "shell_search": totals["shell_search"], + "native_search": totals["native_search"], + "contextstream_search_share": round(cs_search / searches, 4) if searches else None, + "contextstream_tools": dict(sorted(cs_tools.items())), + "shell_search_tools": dict(sorted(shell_tools.items())), + } + + +def render(report): + columns = ( + ("agent", None), + ("sessions", "sessions"), + ("w/ cs search", "sessions_with_contextstream_search"), + ("cs search", "contextstream_search"), + ("shell search", "shell_search"), + ("native search", "native_search"), + ("cs share", "contextstream_search_share"), + ) + rows = [[title for title, _ in columns]] + for agent, stats in report["agents"].items(): + row = [agent] + for _, key in columns[1:]: + value = stats[key] + row.append("n/a" if value is None else f"{value:.1%}" if key.endswith("share") else str(value)) + rows.append(row) + widths = [max(len(row[i]) for row in rows) for i in range(len(columns))] + lines = [f"window: last {report['days']} days"] + for row in rows: + lines.append(" ".join(cell.ljust(widths[i]) if i == 0 else cell.rjust(widths[i]) for i, cell in enumerate(row))) + return "\n".join(lines) + + +def main(argv=None): + parser = argparse.ArgumentParser(description=__doc__.split("\n\n")[0]) + parser.add_argument("--days", type=float, default=14, help="look back this many days (default 14)") + parser.add_argument("--claude-dir", default="~/.claude/projects") + parser.add_argument("--codex-dir", default="~/.codex/sessions") + parser.add_argument("--json", metavar="FILE", help="also write the report as JSON for before/after comparison") + args = parser.parse_args(argv) + + since = time.time() - args.days * 86400 + report = { + "days": args.days, + "generated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), + "agents": { + "claude-code": measure("claude-code", args.claude_dir, since), + "codex": measure("codex", args.codex_dir, since), + }, + } + print(render(report)) + if args.json: + Path(args.json).write_text(json.dumps(report, indent=2, sort_keys=True) + "\n") + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/testing/adoption/test_search_adoption.py b/testing/adoption/test_search_adoption.py new file mode 100644 index 0000000..c23820b --- /dev/null +++ b/testing/adoption/test_search_adoption.py @@ -0,0 +1,167 @@ +"""Synthetic transcripts only; no real session content is read or stored.""" +import json +import os +from pathlib import Path +import tempfile +import time +import unittest + +import search_adoption as adoption +from search_adoption import measure, shell_code_search + + +def write_jsonl(path, records): + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text("".join(json.dumps(record) + "\n" for record in records)) + + +def claude_call(name, **arguments): + return { + "type": "assistant", + "message": {"content": [{"type": "tool_use", "id": "t", "name": name, "input": arguments}]}, + } + + +class ShellClassifierParityTests(unittest.TestCase): + """Mirrors the `bash_search_*` tests of the Rust PreToolUse hook.""" + + def test_detects_code_discovery(self): + for command, tool in [ + ("grep -rn handle_oauth crates/", "grep"), + ("grep -nE 'pub async fn handle' src/", "grep"), + ("rg --type rust 'PgPool' .", "rg"), + ("fd '\\.rs$' crates/", "fd"), + ("find . -name '*.rs' -type f", "find"), + ("cd /home/foo && grep -rn foo .", "grep"), + ("cd crates; rg -n foo", "rg"), + ]: + self.assertEqual(shell_code_search(command), tool, command) + + def test_ignores_filtering_and_metadata_work(self): + for command in [ + "ps aux | grep contextstream", + "cat /tmp/log | grep ERROR", + "find . -mtime -7 -type f", + "find /var/log -size +10M", + "grep -n ERROR /tmp/app.log", + "grep -rn ERROR /var/log/syslog", + "grep -i warning deploy.txt", + "grep -n foo src/main.rs", + "grep -n handler ./crates/api/lib.rs", + "git status", + "cargo build", + "npm test", + "", + " ", + ]: + self.assertIsNone(shell_code_search(command), command) + + +class AdoptionTests(unittest.TestCase): + def setUp(self): + self.tmp = tempfile.TemporaryDirectory() + self.addCleanup(self.tmp.cleanup) + self.root = Path(self.tmp.name) + + def test_claude_sessions_count_each_search_path(self): + claude = self.root / "claude" + write_jsonl( + claude / "proj" / "a.jsonl", + [ + claude_call("mcp__contextstream__search", mode="hybrid", query="x"), + claude_call("mcp__contextstream__search", mode="keyword", query="y"), + claude_call("mcp__contextstream__context", user_message="hi"), + claude_call("Bash", command="rg -n foo crates"), + claude_call("Bash", command="ps aux | grep node"), + claude_call("Grep", pattern="foo"), + {"type": "user", "message": {"content": "ignored"}}, + ], + ) + write_jsonl( + claude / "proj" / "b.jsonl", + [claude_call("Bash", command="grep -rn foo ."), claude_call("Glob", pattern="**/*.rs")], + ) + + stats = measure("claude-code", claude, since=0) + + self.assertEqual(stats["sessions"], 2) + self.assertEqual(stats["sessions_with_contextstream_search"], 1) + self.assertEqual(stats["contextstream_search"], 2) + self.assertEqual(stats["contextstream_other"], 1) + self.assertEqual(stats["shell_search"], 2) + self.assertEqual(stats["native_search"], 2) + self.assertEqual(stats["shell_search_tools"], {"grep": 1, "rg": 1}) + self.assertEqual(stats["contextstream_search_share"], round(2 / 6, 4)) + + def test_codex_shell_and_mcp_calls(self): + codex = self.root / "codex" + write_jsonl( + codex / "2026" / "rollout-1.jsonl", + [ + {"type": "response_item", "payload": { + "type": "function_call", "name": "shell", + "arguments": json.dumps({"command": ["bash", "-lc", "rg -n foo ."]})}}, + {"type": "response_item", "payload": { + "type": "function_call", "name": "exec_command", + "arguments": json.dumps({"cmd": "grep -rn bar src/"})}}, + {"type": "response_item", "payload": { + "type": "function_call", "name": "mcp__contextstream__search", + "arguments": json.dumps({"query": "x"})}}, + {"type": "response_item", "payload": { + "type": "function_call", "name": "shell", + "arguments": json.dumps({"command": ["bash", "-lc", "git status"]})}}, + {"type": "event_msg", "payload": {"type": "agent_message", "message": "rg foo"}}, + ], + ) + + stats = measure("codex", codex, since=0) + + self.assertEqual(stats["sessions"], 1) + self.assertEqual(stats["sessions_with_contextstream_search"], 1) + self.assertEqual(stats["contextstream_search"], 1) + self.assertEqual(stats["shell_search"], 2) + + def test_window_excludes_old_transcripts_and_missing_dirs_are_empty(self): + claude = self.root / "claude" + old = claude / "old.jsonl" + write_jsonl(old, [claude_call("Grep", pattern="foo")]) + stale = time.time() - 40 * 86400 + os.utime(old, (stale, stale)) + + self.assertEqual(measure("claude-code", claude, since=time.time() - 14 * 86400)["sessions"], 0) + empty = measure("codex", self.root / "does-not-exist", since=0) + self.assertEqual(empty["sessions"], 0) + self.assertIsNone(empty["contextstream_search_share"]) + + def test_malformed_lines_are_skipped(self): + claude = self.root / "claude" + claude.mkdir() + (claude / "bad.jsonl").write_text( + "not json\n" + json.dumps(claude_call("Grep", pattern="foo")) + "\n{\n" + ) + self.assertEqual(measure("claude-code", claude, since=0)["native_search"], 1) + + def test_cli_writes_comparable_json_and_prints_table(self): + claude = self.root / "claude" + write_jsonl(claude / "a.jsonl", [claude_call("mcp__contextstream__search", query="x")]) + out = self.root / "report.json" + import contextlib + import io + + buffer = io.StringIO() + with contextlib.redirect_stdout(buffer): + code = adoption.main([ + "--claude-dir", str(claude), + "--codex-dir", str(self.root / "none"), + "--json", str(out), + ]) + + self.assertEqual(code, 0) + report = json.loads(out.read_text()) + self.assertEqual(report["agents"]["claude-code"]["contextstream_search"], 1) + self.assertIn("claude-code", buffer.getvalue()) + self.assertIn("cs share", buffer.getvalue()) + + +if __name__ == "__main__": + unittest.main()