diff --git a/docs/research-cache.md b/docs/research-cache.md new file mode 100644 index 0000000..2e6e467 --- /dev/null +++ b/docs/research-cache.md @@ -0,0 +1,9 @@ +# Content-addressed research cache + +Standard and Deep runs can reuse fresh specialist JSON artifacts from a prior equivalent run. The cache key includes normalized company identity, research depth, fixed model and browser-tool profile, the complete job prompt, and a digest of every agent template. A prompt, model, tool profile, depth, company, or template change therefore produces a cold miss. + +Only specialist JSON objects are cached. Malformed JSON, stream logs, credentials, Apollo contact results, verifier output, and synthesizer output are excluded. The verifier and synthesizer always rerun over the selected specialist set. + +Entries expire after seven days. Restoration rejects malformed keys, schema mismatches, stale or future-dated manifests, symlinks, non-files, logs, malformed JSON, and derived-agent files. Cache failure never converts a failed research job to success. + +The unit fixture covers deterministic normalization, material-input invalidation, TTL expiry, and derived-output exclusion. Operational targets for a representative unchanged Standard rerun are at least 70% fewer specialist model calls and p95 under 30 seconds; measure those in release telemetry before making a public performance claim. diff --git a/src-tauri/src/jobs/queue.rs b/src-tauri/src/jobs/queue.rs index d8584b5..873b189 100644 --- a/src-tauri/src/jobs/queue.rs +++ b/src-tauri/src/jobs/queue.rs @@ -405,16 +405,31 @@ impl JobQueue { crate::orchestration::ResearchDepth::Light }; + let cache_key = if execution_depth != crate::orchestration::ResearchDepth::Light { + Some(crate::orchestration::research_cache_key( + &entity_label, + execution_depth, + &format!("claude-opus-4-7;chrome={}", settings.use_chrome), + &prompt, + )) + } else { + None + }; let effective_working_dir = if execution_depth != crate::orchestration::ResearchDepth::Light { - let prepared = - crate::orchestration::prepare_job_workspace(&app, &job_id, execution_depth)?; + let prepared = crate::orchestration::prepare_job_workspace( + &app, + &job_id, + execution_depth, + cache_key.as_deref(), + )?; eprintln!( - "[job_queue] job_id={} Prepared orchestration workspace at {:?} with agents: {:?}", + "[job_queue] job_id={} Prepared orchestration workspace at {:?} with agents: {:?}; cache hits: {:?}", job_id, prepared.path, - prepared.agent_names + prepared.agent_names, + prepared.cache_hit_agents, ); prepared.path.to_string_lossy().to_string() } else { @@ -438,10 +453,10 @@ impl JobQueue { output_path: Some(metadata.primary_output_path.to_string_lossy().to_string()), }; crate::db::insert_job(&conn, &new_job).map_err(|e| e.to_string())?; - Ok((settings, effective_working_dir, execution_depth)) + Ok((settings, effective_working_dir, execution_depth, cache_key)) }; - let (settings, effective_working_dir, execution_depth) = match setup_result { + let (settings, effective_working_dir, execution_depth, cache_key) = match setup_result { Ok(result) => result, Err(e) => { active_jobs.lock().await.remove(&job_id); @@ -841,6 +856,21 @@ impl JobQueue { update_job_status(&final_status, result.1, final_error_msg.as_deref()); + if final_success { + if let Some(ref key) = cache_key { + if let Err(error) = crate::orchestration::persist_research_cache( + &app_clone, + key, + &PathBuf::from(&effective_working_dir), + ) { + eprintln!( + "[job_queue] job_id={} Could not update research cache: {}", + job_id_clone, error + ); + } + } + } + if let Err(e) = on_event.send(StreamEvent { job_id: job_id_clone.clone(), event_type: final_status.clone(), diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index 40236f4..07006cf 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -5,6 +5,7 @@ mod events; mod jobs; mod orchestration; mod prompts; +mod research_cache; use db::{get_db_path, DbState}; use jobs::JobQueue; diff --git a/src-tauri/src/orchestration.rs b/src-tauri/src/orchestration.rs index edaf7ea..b11a9e3 100644 --- a/src-tauri/src/orchestration.rs +++ b/src-tauri/src/orchestration.rs @@ -5,6 +5,9 @@ use tauri::{AppHandle, Manager}; use crate::db::Settings; use crate::prompts; +use crate::research_cache::ResearchCache; + +const RESEARCH_CACHE_TTL_MS: i64 = 7 * 24 * 60 * 60 * 1_000; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum ResearchDepth { @@ -59,6 +62,7 @@ impl ResearchDepth { pub struct PreparedWorkspace { pub path: PathBuf, pub agent_names: Vec<&'static str>, + pub cache_hit_agents: Vec, } struct AgentTemplate { @@ -75,6 +79,7 @@ pub fn prepare_job_workspace( app: &AppHandle, job_id: &str, depth: ResearchDepth, + cache_key: Option<&str>, ) -> Result { let workspace = app .path() @@ -97,11 +102,31 @@ pub fn prepare_job_workspace( fs::write(path, template.content).map_err(|e| e.to_string())?; } + let cache_hit_agents = if let Some(cache_key) = cache_key { + let cache = ResearchCache::new( + app.path() + .app_data_dir() + .unwrap_or_else(|_| PathBuf::from(".")) + .join("research-cache"), + RESEARCH_CACHE_TTL_MS, + )?; + cache + .restore( + cache_key, + &workspace.join("outputs").join("specialists"), + chrono::Utc::now().timestamp_millis(), + )? + .agents + } else { + Vec::new() + }; + fs::write( workspace.join("README.md"), format!( - "# Augur OS Research Job\n\nDepth: `{}`\n\nSpecialist outputs must be written to `outputs/specialists/`.\n", - depth.as_str() + "# Augur OS Research Job\n\nDepth: `{}`\n\nSpecialist outputs must be written to `outputs/specialists/`.\n\nFresh cached specialists: {}. The verifier and synthesizer are never cached and must always run.\n", + depth.as_str(), + if cache_hit_agents.is_empty() { "none".to_string() } else { cache_hit_agents.join(", ") } ), ) .map_err(|e| e.to_string())?; @@ -109,9 +134,50 @@ pub fn prepare_job_workspace( Ok(PreparedWorkspace { path: workspace, agent_names, + cache_hit_agents, }) } +pub fn research_cache_key( + company: &str, + depth: ResearchDepth, + model: &str, + prompt: &str, +) -> String { + use sha2::{Digest, Sha256}; + let mut version = Sha256::new(); + for template in templates_for_depth(depth) { + version.update(template.name.as_bytes()); + version.update(template.content.as_bytes()); + } + ResearchCache::key( + company, + depth.as_str(), + model, + prompt, + &format!("{:x}", version.finalize()), + ) +} + +pub fn persist_research_cache( + app: &AppHandle, + cache_key: &str, + workspace: &std::path::Path, +) -> Result, String> { + let cache = ResearchCache::new( + app.path() + .app_data_dir() + .unwrap_or_else(|_| PathBuf::from(".")) + .join("research-cache"), + RESEARCH_CACHE_TTL_MS, + )?; + cache.store( + cache_key, + &workspace.join("outputs").join("specialists"), + chrono::Utc::now().timestamp_millis(), + ) +} + pub fn cleanup_job_workspace(app: &AppHandle, job_id: &str) { let workspace = app .path() @@ -163,6 +229,7 @@ Execution plan: Rules: - Run independent Wave 1 specialists in parallel. +- Before launching a specialist, check whether its JSON artifact already exists from the fresh content-addressed cache. Reuse a valid cached specialist artifact; never reuse `verifier.json` or synthesizer output. - Each specialist must write its strict JSON envelope to `outputs/specialists/.json`. - Each specialist must write a short progress log to `outputs/specialists/.stream.log`. - Each specialist must return only this compact pointer: `{{"agent":"","status":"completed","path":"outputs/specialists/.json","streamLog":"outputs/specialists/.stream.log"}}`. diff --git a/src-tauri/src/research_cache.rs b/src-tauri/src/research_cache.rs new file mode 100644 index 0000000..b00fb75 --- /dev/null +++ b/src-tauri/src/research_cache.rs @@ -0,0 +1,226 @@ +use sha2::{Digest, Sha256}; +use std::fs; +use std::path::{Path, PathBuf}; + +const MANIFEST: &str = "cache-manifest.json"; +const CACHE_SCHEMA: u32 = 1; + +#[derive(Debug, Clone)] +pub struct ResearchCache { + root: PathBuf, + ttl_ms: i64, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct CacheRestore { + pub agents: Vec, +} + +impl ResearchCache { + pub fn new(root: PathBuf, ttl_ms: i64) -> Result { + if ttl_ms <= 0 { + return Err("research cache ttl must be positive".into()); + } + Ok(Self { root, ttl_ms }) + } + + pub fn key( + company: &str, + depth: &str, + model: &str, + prompt: &str, + orchestration_version: &str, + ) -> String { + let normalized_company = company.split_whitespace().collect::>().join(" "); + let mut digest = Sha256::new(); + for part in [ + normalized_company.to_lowercase(), + depth.to_string(), + model.to_string(), + prompt.to_string(), + orchestration_version.to_string(), + ] { + digest.update((part.len() as u64).to_be_bytes()); + digest.update(part.as_bytes()); + } + format!("{:x}", digest.finalize()) + } + + pub fn restore( + &self, + key: &str, + destination: &Path, + now_ms: i64, + ) -> Result { + validate_key(key)?; + let source = self.root.join(key); + let manifest_path = source.join(MANIFEST); + let manifest: serde_json::Value = match fs::read_to_string(&manifest_path) { + Ok(value) => serde_json::from_str(&value).map_err(|e| e.to_string())?, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + return Ok(CacheRestore { agents: Vec::new() }); + } + Err(error) => return Err(error.to_string()), + }; + if manifest.get("schema").and_then(|value| value.as_u64()) != Some(CACHE_SCHEMA as u64) { + return Ok(CacheRestore { agents: Vec::new() }); + } + let created_at = manifest + .get("createdAtMs") + .and_then(|value| value.as_i64()) + .unwrap_or(0); + if created_at <= 0 + || created_at > now_ms + || now_ms.saturating_sub(created_at) > self.ttl_ms + { + return Ok(CacheRestore { agents: Vec::new() }); + } + + fs::create_dir_all(destination).map_err(|e| e.to_string())?; + let mut agents = Vec::new(); + for entry in fs::read_dir(&source).map_err(|e| e.to_string())? { + let entry = entry.map_err(|e| e.to_string())?; + let file_type = entry.file_type().map_err(|e| e.to_string())?; + if !file_type.is_file() || file_type.is_symlink() { + continue; + } + let path = entry.path(); + if path.extension().and_then(|value| value.to_str()) != Some("json") + || path.file_name().and_then(|value| value.to_str()) == Some(MANIFEST) + { + continue; + } + let agent = path + .file_stem() + .and_then(|value| value.to_str()) + .ok_or_else(|| "cache artifact has an invalid filename".to_string())?; + if is_derived_agent(agent) { + continue; + } + let artifact = fs::read(&path).map_err(|e| e.to_string())?; + if !is_json_object(&artifact) { + continue; + } + fs::write(destination.join(entry.file_name()), artifact).map_err(|e| e.to_string())?; + agents.push(agent.to_string()); + } + agents.sort(); + Ok(CacheRestore { agents }) + } + + pub fn store(&self, key: &str, source: &Path, now_ms: i64) -> Result, String> { + validate_key(key)?; + fs::create_dir_all(&self.root).map_err(|e| e.to_string())?; + let staging = self.root.join(format!(".{}.{}.tmp", key, uuid::Uuid::new_v4())); + fs::create_dir(&staging).map_err(|e| e.to_string())?; + + let result = (|| { + let mut agents = Vec::new(); + for entry in fs::read_dir(source).map_err(|e| e.to_string())? { + let entry = entry.map_err(|e| e.to_string())?; + let file_type = entry.file_type().map_err(|e| e.to_string())?; + if !file_type.is_file() || file_type.is_symlink() { + continue; + } + let path = entry.path(); + if path.extension().and_then(|value| value.to_str()) != Some("json") { + continue; + } + let agent = path + .file_stem() + .and_then(|value| value.to_str()) + .ok_or_else(|| "specialist artifact has an invalid filename".to_string())?; + if is_derived_agent(agent) { + continue; + } + let artifact = fs::read(&path).map_err(|e| e.to_string())?; + if !is_json_object(&artifact) { + continue; + } + fs::write(staging.join(entry.file_name()), artifact).map_err(|e| e.to_string())?; + agents.push(agent.to_string()); + } + agents.sort(); + fs::write( + staging.join(MANIFEST), + serde_json::to_vec_pretty(&serde_json::json!({ + "schema": CACHE_SCHEMA, + "createdAtMs": now_ms, + "agents": agents, + })) + .map_err(|e| e.to_string())?, + ) + .map_err(|e| e.to_string())?; + + let target = self.root.join(key); + if target.exists() { + fs::remove_dir_all(&target).map_err(|e| e.to_string())?; + } + fs::rename(&staging, target).map_err(|e| e.to_string())?; + Ok(agents) + })(); + + if result.is_err() { + let _ = fs::remove_dir_all(&staging); + } + result + } +} + +fn validate_key(key: &str) -> Result<(), String> { + if key.len() != 64 || !key.bytes().all(|byte| byte.is_ascii_hexdigit()) { + return Err("research cache key must be a sha256 digest".into()); + } + Ok(()) +} + +fn is_derived_agent(agent: &str) -> bool { + matches!(agent, "verifier" | "synthesizer") +} + +fn is_json_object(bytes: &[u8]) -> bool { + serde_json::from_slice::(bytes) + .is_ok_and(|value| value.is_object()) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn temp_root(name: &str) -> PathBuf { + std::env::temp_dir().join(format!("augur-cache-{}-{}", name, uuid::Uuid::new_v4())) + } + + #[test] + fn key_is_normalized_and_invalidates_material_inputs() { + let one = ResearchCache::key(" Acme Corp ", "standard", "opus", "prompt", "agents-v1"); + let two = ResearchCache::key("acme corp", "standard", "opus", "prompt", "agents-v1"); + assert_eq!(one, two); + assert_ne!(one, ResearchCache::key("acme corp", "deep", "opus", "prompt", "agents-v1")); + assert_ne!(one, ResearchCache::key("acme corp", "standard", "opus", "changed", "agents-v1")); + assert_ne!(one, ResearchCache::key("acme corp", "standard", "opus", "prompt", "agents-v2")); + } + + #[test] + fn restores_only_fresh_reusable_specialists() { + let root = temp_root("restore"); + let source = root.join("source"); + let destination = root.join("destination"); + fs::create_dir_all(&source).unwrap(); + fs::write(source.join("pain.json"), "{}").unwrap(); + fs::write(source.join("broken.json"), "{").unwrap(); + fs::write(source.join("verifier.json"), "{}").unwrap(); + fs::write(source.join("notes.log"), "secret stream").unwrap(); + let cache = ResearchCache::new(root.join("cache"), 1_000).unwrap(); + let key = "a".repeat(64); + assert_eq!(cache.store(&key, &source, 10_000).unwrap(), vec!["pain"]); + assert!(cache.restore(&key, &destination, 9_999).unwrap().agents.is_empty()); + assert_eq!(cache.restore(&key, &destination, 10_500).unwrap().agents, vec!["pain"]); + assert!(destination.join("pain.json").exists()); + assert!(!destination.join("verifier.json").exists()); + + let stale = root.join("stale"); + assert!(cache.restore(&key, &stale, 11_001).unwrap().agents.is_empty()); + let _ = fs::remove_dir_all(root); + } +}