#[allow(clippy::wildcard_imports)]
use super::super::*;
use super::finalize_call_result;
#[allow(clippy::too_many_arguments)]
pub(in crate::server) async fn dispatch_and_post_process(
server: &LeanCtxServer,
name: &str,
args: Option<&serde_json::Map<String, serde_json::Value>>,
minimal: bool,
config: std::sync::Arc<crate::core::config::Config>,
machine_readable: bool,
auto_context: Option<String>,
throttle_warning: Option<String>,
args_fp: String,
mut decision_context: Option<crate::core::decision_loop_runtime::TaskContext>,
) -> Result<CallToolResult, ErrorData> {
let tool_start = std::time::Instant::now();
let shadow_auto_record = config.shadow.enabled && config.shadow.auto_record;
let (mut result_text, tool_saved_tokens, shell_outcome, content_blocks) =
match server.dispatch_tool(name, args, minimal).await {
Ok(tuple) => tuple,
Err(e) => {
if let Ok(mut detector) = tokio::time::timeout(
std::time::Duration::from_secs(1),
server.loop_detector.write(),
)
.await
{
detector.record_error_outcome(name, &args_fp);
}
crate::core::debug_log::log_mcp_error(name, args, &format!("{e:?}"));
if e.code == rmcp::model::ErrorCode::INVALID_PARAMS {
tracing::debug!(
"converting INVALID_PARAMS to soft tool error for '{name}': {}",
e.message
);
let result =
CallToolResult::error(vec![ContentBlock::text(e.message.to_string())]);
record_decision_loop_end(
decision_context.as_ref(),
args,
&result,
false,
shadow_auto_record,
None,
);
return Ok(result);
}
record_decision_loop_end_error(decision_context.as_ref(), args, shadow_auto_record);
return Err(e);
}
};
let task_profile = {
let session = server.session.read().await;
crate::core::decision_loop_runtime::DecisionLoopRuntime::get_or_init()
.profile_for_session(&session.id)
};
result_text =
apply_task_triage_filter(result_text, task_profile.as_ref(), &mut decision_context);
if let Some(blocks) = content_blocks {
let mut result = CallToolResult::success(blocks);
if let Some(outcome) = shell_outcome
&& outcome.is_error()
{
result.is_error = Some(true);
}
record_decision_loop_end(
decision_context.as_ref(),
args,
&result,
result.is_error != Some(true),
shadow_auto_record,
Some(shadow_tokens_for_result(&result)),
);
return Ok(result);
}
let inline_shell = name == "ctx_shell"
&& crate::core::firewall::should_inline_shell(
helpers::get_bool(args, "inline").unwrap_or(false),
result_text.len(),
&config,
);
let is_raw_shell = name == "ctx_shell" && {
let arg_raw = helpers::get_bool(args, "raw").unwrap_or(false);
let arg_bypass = helpers::get_bool(args, "bypass").unwrap_or(false);
arg_raw
|| arg_bypass
|| crate::core::runtime_flags::raw_enabled()
|| inline_shell
|| helpers::get_str(args, "command")
.is_some_and(|c| crate::core::firewall::is_raw_command(&c, &config))
};
let pre_terse_len = result_text.len();
let output_tokens = {
let tokens = crate::core::tokens::count_tokens(&result_text) as u64;
crate::core::budget_tracker::BudgetTracker::global().record_tokens(tokens);
tokens
};
crate::core::anomaly::record_metric("tokens_per_call", output_tokens as f64);
if let Some(ref ir) = server.context_ir {
let tool_duration = tool_start.elapsed();
let source_kind = post_process::context_ir_source_kind(name);
let ir_path = helpers::get_str(args, "path");
let ir_command = helpers::get_str(args, "command");
let ir_mode = helpers::get_str(args, "mode");
let excerpt = if result_text.len() > 200 {
let mut end = 200;
while !result_text.is_char_boundary(end) && end > 0 {
end -= 1;
}
&result_text[..end]
} else {
&result_text
};
let input = crate::core::context_ir::RecordIrInput {
kind: source_kind,
tool: name,
client_name: None,
agent_id: None,
path: ir_path.as_deref(),
command: ir_command.as_deref(),
pattern: ir_mode.as_deref(),
input_tokens: pre_terse_len / 4,
output_tokens: output_tokens as usize,
duration: tool_duration,
content_excerpt: excerpt,
};
ir.write().await.record(input);
}
{
let mut detector = server.loop_detector.write().await;
if name == "ctx_read" {
let path = helpers::get_str(args, "path").unwrap_or_default();
let mode = helpers::get_str(args, "mode").unwrap_or_else(|| "auto".into());
let fresh = helpers::get_bool(args, "fresh").unwrap_or(false);
detector.record_read_for_correction(&path, &mode, fresh);
} else if name == "ctx_shell" {
let cmd = helpers::get_str(args, "command").unwrap_or_default();
detector.record_shell_for_correction(&cmd);
} else if name == "ctx_expand" || name == "ctx_retrieve" {
detector.record_retrieve();
}
let correction_count = detector.correction_count();
let retrieve_count = detector.retrieve_count();
if correction_count > 0 {
crate::core::anomaly::record_metric(
"correction_loop_rate",
f64::from(correction_count),
);
}
if retrieve_count > 0 {
crate::core::anomaly::record_metric("ccr_retrieve_rate", f64::from(retrieve_count));
}
use crate::core::config::CompressionLevel;
CompressionLevel::apply_degrade_action(CompressionLevel::degrade_action(
correction_count,
retrieve_count,
));
detector.prune_corrections();
}
crate::core::anomaly::save_debounced();
let budget_warning = post_process::budget_warning_message();
{
let path_hint = helpers::get_str(args, "path");
let enforced = crate::core::sensitivity::enforce_text(
std::mem::take(&mut result_text),
path_hint.as_deref().map(std::path::Path::new),
&config.sensitivity_effective(),
);
result_text = enforced.into_text();
}
if crate::core::policy::runtime::is_active() {
let (redacted, hits) = policy_guard::redact_result(&result_text);
if hits > 0 {
tracing::debug!(redactions = hits, "context policy redaction applied");
result_text = redacted;
}
}
if let Some(active) = crate::core::policy::runtime::active()
&& active.filters.is_active()
{
let outcome = crate::core::input_filters::apply(&result_text, &active.filters);
if outcome.blocked {
let reason = outcome.block_reason.as_deref().unwrap_or("policy");
tracing::warn!(tool = name, reason, "content blocked by input filter");
policy_guard::audit_filter(name, &outcome.audit, true);
result_text = format!(
"[POLICY BLOCKED] Content withheld by the active context policy pack \
(input filter: {reason}). Adjust .lean-ctx/policy.toml to proceed."
);
} else {
if !outcome.audit.is_empty() {
tracing::debug!(tool = name, "input filters applied");
policy_guard::audit_filter(name, &outcome.audit, false);
}
result_text = outcome.text;
for warning in &outcome.warnings {
result_text = format!("{result_text}\n\n[FILTER] {warning}");
}
}
}
let mut firewalled = false;
let archive_hint = if minimal {
None
} else {
use crate::core::archive;
let archivable = matches!(
name,
"ctx_shell"
| "ctx_read"
| "ctx_multi_read"
| "ctx_smart_read"
| "ctx_execute"
| "ctx_search"
| "ctx_tree"
);
if archivable && archive::should_archive(&result_text) {
let cmd = helpers::get_str(args, "command")
.or_else(|| helpers::get_str(args, "path"))
.unwrap_or_default();
let session_id = server.session.read().await.id.clone();
let to_store = crate::core::redaction::redact_text_if_enabled(&result_text);
let tokens = crate::core::tokens::count_tokens(&to_store);
match archive::store(name, &cmd, &to_store, Some(&session_id)) {
Some(id)
if !is_raw_shell
&& crate::core::firewall::should_firewall(name, tokens, &config) =>
{
result_text =
crate::core::firewall::summarize(&to_store, &id, name, tokens, &cmd);
firewalled = true;
None
}
Some(id) => Some(archive::format_hint(&id, to_store.len(), tokens)),
None => None,
}
} else {
None
}
};
let pre_compression = result_text.clone();
if !firewalled {
result_text = post_process::compress_terse(result_text, name, args, &config, is_raw_shell);
}
let findings_source = result_text.clone();
if !machine_readable
&& !is_raw_shell
&& !firewalled
&& crate::core::cognitive_gate::full_science_enabled()
&& !result_text.is_empty()
{
let task_input = {
let session = server.session.read().await;
session.task.as_ref().map(|t| t.description.clone())
};
if let Some(task) = task_input.filter(|t| !t.is_empty()) {
let report = crate::core::echo_ratio::compute_echo_ratio(&task, &result_text);
if report.ratio > 0.7 {
result_text.push_str(&format!(
"\n[cognitive: high echo ratio ({:.0}%) — consider generating novel content]",
report.ratio * 100.0
));
}
}
}
let active_profile = crate::core::profiles::active_profile();
let profile_hints = active_profile.output_hints.clone();
if !is_raw_shell && !firewalled && profile_hints.verify_footer() {
let verify_cfg = active_profile.verification;
let vr = crate::core::output_verification::verify_output(
&pre_compression,
&result_text,
&verify_cfg,
);
if !vr.warnings.is_empty() {
let msg = format!("[VERIFY] {}", vr.format_compact());
result_text = format!("{result_text}\n\n{msg}");
}
}
if !firewalled
&& profile_hints.archive_hint()
&& let Some(hint) = archive_hint
{
result_text = format!("{result_text}\n{hint}");
}
let had_auto_context = auto_context.is_some();
let had_budget_warning = budget_warning.is_some();
let had_throttle_warning = throttle_warning.is_some();
if !is_raw_shell && let Some(ctx) = auto_context {
let ctx_tokens = crate::core::tokens::count_tokens(&ctx);
if ctx_tokens <= 400 {
result_text = format!("{ctx}\n\n{result_text}");
}
}
if !is_raw_shell
&& name != "ctx_memory"
&& let Some(hint) = crate::core::shared_context::session_start_hint()
{
result_text = format!("{hint}\n\n{result_text}");
}
if let Some(warning) = throttle_warning {
result_text = format!("{result_text}\n\n{warning}");
}
if let Some(bw) = budget_warning {
result_text = format!("{result_text}\n\n{bw}");
}
if matches!(name, "ctx_read" | "ctx_search" | "ctx_compose") {
let query = helpers::get_str(args, "query")
.or_else(|| helpers::get_str(args, "task"))
.unwrap_or_default();
if !query.is_empty() {
let advice = context_gate::knowledge_advice(&query);
if let Some(hint) = advice.additional_context_hint {
result_text = format!("{result_text}\n\n{hint}");
}
}
}
if !machine_readable
&& !server
.rules_stale_checked
.swap(true, std::sync::atomic::Ordering::Relaxed)
{
let client = server.client_name.read().await.clone();
if !client.is_empty() && crate::rules_inject::check_rules_freshness(&client).is_some() {
let _ = tokio::task::spawn_blocking(|| {
if let Some(home) = dirs::home_dir() {
let _ = crate::rules_inject::inject_all_rules(&home);
}
})
.await;
result_text = format!(
"{result_text}\n\n[RULES AUTO-UPDATED] Your lean-ctx rules were written by \
an older version and have been refreshed on disk. Start a new session to \
load them for full compatibility."
);
} else if !server
.rules_tip_shown
.swap(true, std::sync::atomic::Ordering::Relaxed)
{
let cfg = crate::core::config::Config::load();
if !cfg.setup.should_inject_rules() {
result_text = format!(
"{result_text}\n\n\
--- tip: run 'lean-ctx setup' to configure agent rules for optimal AI integration ---"
);
}
}
}
{
let _ = crate::core::slo::evaluate();
}
if name == "ctx_read" {
if let Some(read_path) = args
.as_ref()
.and_then(|args| args.get("path"))
.and_then(serde_json::Value::as_str)
{
let task_class = task_profile
.as_ref()
.map_or("unknown", |profile| profile.task_class.as_str());
let task_class = task_class.to_owned();
let read_path = read_path.to_owned();
let project_root = {
let session = server.session.read().await;
session.project_root.clone()
};
tokio::spawn(async move {
let predictions = crate::server::predictive_preload::record_read_and_predict(
&task_class,
&read_path,
);
crate::server::predictive_preload::warm_paths(
project_root.as_deref().map(std::path::Path::new),
&predictions,
);
});
}
if minimal {
let cache_clone = server.cache.clone();
let autonomy_clone = server.autonomy.clone();
let name_owned = name.to_string();
tokio::spawn(async move {
let result = std::panic::AssertUnwindSafe(async {
let cache_timeout = tokio::time::timeout(
std::time::Duration::from_secs(5),
cache_clone.write(),
)
.await;
if let Ok(mut cache) = cache_timeout {
crate::tools::autonomy::maybe_auto_dedup(
&autonomy_clone,
&mut cache,
&name_owned,
);
} else {
tracing::debug!("background auto_dedup: cache lock timeout (5s), skipping");
}
})
.catch_unwind()
.await;
if let Err(e) = result {
let msg = e
.downcast_ref::<String>()
.map(String::as_str)
.or_else(|| e.downcast_ref::<&str>().copied())
.unwrap_or("unknown");
tracing::error!("background auto_dedup panicked: {msg}");
}
});
} else {
let read_path = server
.resolve_path_or_passthrough(&helpers::get_str(args, "path").unwrap_or_default())
.await;
let project_root = {
let session = server.session.read().await;
session.project_root.clone()
};
let enrich_timeout =
tokio::time::timeout(std::time::Duration::from_secs(3), server.cache.write()).await;
if let Ok(mut cache) = enrich_timeout {
let enrich = crate::tools::autonomy::enrich_after_read(
&server.autonomy,
&mut cache,
&read_path,
project_root.as_deref(),
None,
crate::tools::CrpMode::effective(),
false,
);
if profile_hints.related_hint()
&& let Some(hint) = enrich.related_hint
{
result_text = format!("{result_text}\n{hint}");
}
crate::tools::autonomy::maybe_auto_dedup(&server.autonomy, &mut cache, name);
} else {
tracing::warn!(
"post-dispatch cache lock timeout (3s) for {read_path}, skipping enrichment"
);
}
if std::path::Path::new(&read_path).is_file() {
let ledger_clone = server.ledger.clone();
let session_clone = server.session.clone();
let peer_clone = server.peer.clone();
let read_path_owned = read_path.clone();
let project_root_owned = project_root.clone();
let mode_used =
helpers::get_str(args, "mode").unwrap_or_else(|| "auto".to_string());
let out_tok = output_tokens as usize;
let sent_tok = crate::core::tokens::count_tokens(&result_text);
let wants_eviction = true;
let wants_elicitation = profile_hints.elicitation_hint();
tokio::spawn(async move {
let result = std::panic::AssertUnwindSafe(async {
let active_task = {
let session = session_clone.read().await;
session.task.as_ref().map(|t| t.description.clone())
};
let mut ledger = ledger_clone.write().await;
let overlay = crate::core::context_overlay::OverlayStore::load_project(
&std::path::PathBuf::from(project_root_owned.as_deref().unwrap_or(".")),
);
let gate_result = context_gate::post_dispatch_record_with_task(
&read_path_owned,
&mode_used,
out_tok,
sent_tok,
&mut ledger,
&overlay,
active_task.as_deref(),
project_root_owned.as_deref(),
);
drop(ledger);
if wants_eviction && let Some(hint) = &gate_result.eviction_hint {
tracing::debug!("deferred eviction hint: {hint}");
}
if wants_elicitation && let Some(hint) = &gate_result.elicitation_hint {
tracing::debug!("deferred elicitation hint: {hint}");
}
if let Some(hint) = &gate_result.prefetch_hint {
tracing::debug!("deferred FEP prefetch hint: {hint}");
}
if gate_result.resource_changed
&& let Some(peer) = peer_clone.read().await.as_ref()
{
notifications::send_resource_updated(
peer,
notifications::RESOURCE_URI_SUMMARY,
)
.await;
}
})
.catch_unwind()
.await;
if let Err(e) = result {
let msg = e
.downcast_ref::<String>()
.map(String::as_str)
.or_else(|| e.downcast_ref::<&str>().copied())
.unwrap_or("unknown");
tracing::error!("background post_dispatch panicked: {msg}");
}
});
}
}
}
if !minimal && !is_raw_shell && name == "ctx_shell" {
let cmd = helpers::get_str(args, "command").unwrap_or_default();
if let Some(file_path) = extract_file_read_from_shell(&cmd)
&& let Ok(mut bt) = crate::core::bounce_tracker::global().lock()
{
bt.next_seq();
bt.record_shell_file_access(&file_path);
}
if profile_hints.efficiency_hint() {
let calls = server.tool_calls.read().await;
let last_original = calls.last().map_or(0, |c| c.original_tokens);
drop(calls);
let pre_hint_tokens = crate::core::tokens::count_tokens(&result_text);
if let Some(hint) = crate::tools::autonomy::shell_efficiency_hint(
&server.autonomy,
&cmd,
last_original,
pre_hint_tokens,
) {
result_text = format!("{result_text}\n{hint}");
}
}
}
if !is_raw_shell && bypass_hint::is_enabled() {
if let Ok(data_dir) = crate::core::data_dir::lean_ctx_data_dir() {
let session = server.session.read().await;
bypass_hint::set_session_id(&session.id);
drop(session);
if let Some(hint) = bypass_hint::check(&data_dir) {
result_text = format!("{result_text}\n{hint}");
}
}
bypass_hint::record_lctx_call();
}
let finding_path_hint = helpers::get_str(args, "path");
if let Some(finding) =
crate::core::auto_findings::extract(name, &findings_source, finding_path_hint.as_deref())
{
let mut session = server.session.write().await;
session.add_finding(finding.file.as_deref(), None, &finding.summary);
let project_root = session.project_root.clone();
drop(session);
if let Some(ref root) = project_root {
let f = finding.clone();
let r = root.clone();
std::thread::spawn(move || {
crate::core::auto_capture::capture_finding(&r, &f);
});
}
}
if let Some(extra) = crate::core::auto_capture::extract_extra(name, &findings_source) {
let session = server.session.read().await;
let project_root = session.project_root.clone();
drop(session);
if let Some(ref root) = project_root {
let e = extra.clone();
let r = root.clone();
std::thread::spawn(move || {
crate::core::auto_capture::capture_finding(&r, &e);
});
}
}
{
let tool_name = name.to_string();
let summary = result_text.lines().next().unwrap_or("").to_string();
let dbg_args = args.cloned();
let dbg_bytes = result_text.len();
let dbg_saved = tool_saved_tokens;
let dbg_elapsed = tool_start.elapsed();
std::thread::spawn(move || {
crate::core::journal::maybe_day_separator();
crate::core::journal::log_tool_call(&tool_name, &summary);
crate::core::debug_log::log_mcp_call(
&tool_name,
dbg_args.as_ref(),
&summary,
dbg_bytes,
dbg_saved,
dbg_elapsed,
);
});
}
let output_token_count = post_process::finalize_token_count_and_adjust(
name,
&result_text,
pre_terse_len,
output_tokens,
tool_saved_tokens,
);
record_compression_savings(name, tool_saved_tokens, output_token_count);
let action = helpers::get_str(args, "action");
const K_STALENESS_BOUND: i64 = 10;
if server.session_mode == crate::tools::SessionMode::Shared
&& let Some(ref rt) = server.context_os
{
let latest = rt.bus.latest_id(&server.workspace_id, &server.channel_id);
let cursor = server
.last_seen_event_id
.load(std::sync::atomic::Ordering::Relaxed);
if cursor > 0 && latest - cursor > K_STALENESS_BOUND {
let gap = latest - cursor;
result_text = format!(
"[CONTEXT STALE] {gap} events happened since your last read. \
Use ctx_session(action=\"status\") to sync.\n\n{result_text}"
);
}
server
.last_seen_event_id
.store(latest, std::sync::atomic::Ordering::Relaxed);
}
server
.record_receipt_and_cost(
name,
args,
action.as_deref(),
&result_text,
output_token_count,
)
.await;
if server.session_mode == crate::tools::SessionMode::Shared
&& name == "ctx_knowledge"
&& action.as_deref() == Some("remember")
&& let Some(ref rt) = server.context_os
{
let my_agent = server.agent_id.read().await.clone();
let category = helpers::get_str(args, "category");
let key = helpers::get_str(args, "key");
if let (Some(cat), Some(k)) = (&category, &key) {
let recent = rt.bus.recent_by_kind(
&server.workspace_id,
&server.channel_id,
"knowledge_remembered",
20,
);
for ev in &recent {
let p = &ev.payload;
let ev_cat = p.get("category").and_then(|v| v.as_str());
let ev_key = p.get("key").and_then(|v| v.as_str());
let ev_actor = ev.actor.as_deref();
if ev_cat == Some(cat.as_str())
&& ev_key == Some(k.as_str())
&& ev_actor != my_agent.as_deref()
{
let other = ev_actor.unwrap_or("unknown");
result_text = format!(
"[CONFLICT] Agent '{other}' recently wrote to the same knowledge key \
'{cat}/{k}'. Review before proceeding.\n\n{result_text}"
);
break;
}
}
}
}
server
.persist_shared_context_os(name, action.as_deref(), args)
.await;
let skip_checkpoint = minimal
|| matches!(
name,
"ctx_compress"
| "ctx_metrics"
| "ctx_benchmark"
| "ctx_analyze"
| "ctx_cache"
| "ctx_discover"
| "ctx_dedup"
| "ctx_session"
| "ctx_knowledge"
| "ctx_agent"
| "ctx_share"
| "ctx_gain"
| "ctx_overview"
| "ctx_preload"
| "ctx_cost"
| "ctx_heatmap"
| "ctx_task"
| "ctx_impact"
| "ctx_architecture"
| "ctx_smells"
| "ctx_quality"
| "ctx_workflow"
);
if !skip_checkpoint
&& crate::core::protocol::meta_visible()
&& let Some(nudge) = crate::core::output_echo::take_pending_nudge()
{
result_text.push_str(&nudge);
}
if !skip_checkpoint
&& crate::core::protocol::meta_visible()
&& let Some(hint) = crate::core::version_check::session_update_hint()
{
result_text.push_str("\n\n");
result_text.push_str(&hint);
}
if !skip_checkpoint
&& server.increment_and_check()
&& let Some(checkpoint) = server.auto_checkpoint().await
&& profile_hints.checkpoint_in_output()
&& crate::core::protocol::meta_visible()
{
let combined = format!("{result_text}\n\n--- AUTO CHECKPOINT ---\n{checkpoint}");
let result = finalize_call_result(&combined, shell_outcome);
record_decision_loop_end(
decision_context.as_ref(),
args,
&result,
result.is_error != Some(true),
shadow_auto_record,
Some(shadow_tokens_for_result(&result)),
);
return Ok(result);
}
let current_count = server.call_count.load(std::sync::atomic::Ordering::Relaxed);
if current_count > 0 && current_count.is_multiple_of(100) {
std::thread::spawn(crate::cloud_sync::cloud_background_tasks);
std::thread::spawn(|| {
let _ = crate::core::archive::cleanup();
});
if let Some(root) = server.session.read().await.project_root.clone() {
crate::core::cognition_scheduler::maybe_run(&root);
}
}
if let Some(notice) = crate::server::dynamic_tools::deprecation_notice(name) {
result_text = format!("{notice}\n{result_text}");
}
if machine_readable {
result_text = pre_compression;
}
let budget_limit = crate::core::config::Config::load().turn_fresh_limit_effective();
if budget_limit > 0 {
let (budgeted, action) = crate::core::budget::apply_turn_budget(&result_text, budget_limit);
if let crate::core::budget::BudgetAction::Truncated {
original_tokens,
delivered_tokens,
} = action
{
tracing::debug!(
"budget: truncated {original_tokens} → {delivered_tokens} tokens (limit {budget_limit})"
);
}
result_text = budgeted;
}
let compressed_input_tokens = crate::core::tokens::count_tokens(&result_text) as u64;
let raw_input_tokens = compressed_input_tokens
.saturating_add(u64::try_from(tool_saved_tokens).unwrap_or(u64::MAX));
let mut result = finalize_call_result(&result_text, shell_outcome);
let has_dynamic = had_auto_context || had_budget_warning || had_throttle_warning;
let mut meta = rmcp::model::Meta::new();
meta.0.insert(
"cache_hint".to_owned(),
serde_json::Value::String(if has_dynamic { "ephemeral" } else { "stable" }.to_owned()),
);
result.meta = Some(meta);
let agent_id = if let Some(agent_id) = server.agent_id.read().await.clone() {
agent_id
} else {
server
.presence_agent_id
.read()
.await
.clone()
.unwrap_or_else(|| "mcp-agent".to_owned())
};
crate::core::agent_budget::record_consumption(
&agent_id,
usize::try_from(compressed_input_tokens).unwrap_or(usize::MAX),
);
crate::core::agent_budget::record_turn_delivery(&agent_id, compressed_input_tokens);
record_decision_loop_end(
decision_context.as_ref(),
args,
&result,
result.is_error != Some(true),
shadow_auto_record,
Some((raw_input_tokens, compressed_input_tokens)),
);
Ok(result)
}
fn apply_task_triage_filter(
result_text: String,
profile: Option<&crate::core::triage::profile::TaskProfileLocal>,
decision_context: &mut Option<crate::core::decision_loop_runtime::TaskContext>,
) -> String {
let Some(profile) = profile else {
return result_text;
};
let filtered = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let level = context_gate::triage_filter_level(profile);
(level > 0).then(|| context_gate::apply_triage_filter(&result_text, profile, level))
}));
let Ok(Some((filtered_text, filtered_lines))) = filtered else {
return result_text;
};
if let Some(context) = decision_context {
context.filtered_lines = filtered_lines;
}
filtered_text
}
fn record_decision_loop_end(
context: Option<&crate::core::decision_loop_runtime::TaskContext>,
args: Option<&serde_json::Map<String, serde_json::Value>>,
result: &CallToolResult,
success: bool,
shadow_auto_record: bool,
shadow_tokens: Option<(u64, u64)>,
) {
let input_tokens = args
.and_then(|args| serde_json::to_string(args).ok())
.map_or(0, |input| (input.len() / 4) as u64);
let output_tokens = format!("{result:?}").len() as u64 / 4;
record_decision_loop(
context,
input_tokens,
output_tokens,
success,
shadow_auto_record,
shadow_tokens,
);
}
fn shadow_tokens_for_result(result: &CallToolResult) -> (u64, u64) {
let tokens = format!("{result:?}").len() as u64 / 4;
(tokens, tokens)
}
fn record_compression_savings(name: &str, tool_saved_tokens: usize, output_tokens: usize) {
if let Some((raw_tokens, compressed_tokens)) =
compression_tracker_tokens(name, tool_saved_tokens, output_tokens)
{
crate::core::savings_tracker::record_compression(raw_tokens, compressed_tokens, name);
}
}
fn compression_tracker_tokens(
name: &str,
tool_saved_tokens: usize,
output_tokens: usize,
) -> Option<(u64, u64)> {
matches!(
name,
"ctx_read" | "ctx_shell" | "ctx_search" | "ctx_compose"
)
.then(|| {
let raw_tokens = tool_saved_tokens.saturating_add(output_tokens) as u64;
(raw_tokens, output_tokens as u64)
})
}
fn record_decision_loop_end_error(
context: Option<&crate::core::decision_loop_runtime::TaskContext>,
args: Option<&serde_json::Map<String, serde_json::Value>>,
shadow_auto_record: bool,
) {
let input_tokens = args
.and_then(|args| serde_json::to_string(args).ok())
.map_or(0, |input| (input.len() / 4) as u64);
record_decision_loop(context, input_tokens, 0, false, shadow_auto_record, None);
}
fn record_decision_loop(
context: Option<&crate::core::decision_loop_runtime::TaskContext>,
input_tokens: u64,
output_tokens: u64,
success: bool,
shadow_auto_record: bool,
shadow_tokens: Option<(u64, u64)>,
) {
let Some(context) = context else {
return;
};
if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
crate::core::decision_loop_runtime::DecisionLoopRuntime::get_or_init()
.on_tool_end_with_shadow(
context,
input_tokens,
output_tokens,
"mcp-tool",
success,
shadow_auto_record,
shadow_tokens,
)
}))
.is_err()
{
tracing::warn!("decision loop end panicked");
}
}
#[cfg(test)]
mod savings_tests {
use super::{apply_task_triage_filter, compression_tracker_tokens};
#[test]
fn test_tracker_in_pipeline() {
let mut tracker = crate::core::savings_tracker::SessionSavingsTracker::default();
let (raw, compressed) = compression_tracker_tokens("ctx_read", 75, 25).expect("tracked");
tracker.record_compression(raw, compressed, "ctx_read");
let after = tracker.session_summary();
assert_eq!(
(
after.total_raw,
after.total_compressed,
after.savings_tokens
),
(100, 25, 75)
);
}
#[test]
fn triage_filter_rewrites_raw_output_and_tracks_removed_lines() {
let profile = crate::core::triage::profile::TaskProfileLocal {
confidence_milli: 500,
context_need_milli: 400,
..Default::default()
};
let mut context = Some(crate::core::decision_loop_runtime::TaskContext {
task_id: String::new(),
session_id: String::new(),
triage_class: String::new(),
profile_intent: String::new(),
profile_complexity: String::new(),
filtered_lines: 0,
start_time: std::time::Instant::now(),
});
let raw = format!("// boilerplate\n{}", "content\n".repeat(100));
let filtered = apply_task_triage_filter(raw, Some(&profile), &mut context);
assert!(!filtered.starts_with("// boilerplate"));
assert_eq!(context.as_ref().unwrap().filtered_lines, 1);
}
#[test]
fn triage_filter_fails_open_without_a_profile() {
let raw = "content\n".repeat(100);
let mut context = None;
assert_eq!(
apply_task_triage_filter(raw.clone(), None, &mut context),
raw
);
}
}