From 8c620889918a2a490e0c9bf2ef91204b9805a347 Mon Sep 17 00:00:00 2001 From: Divyam Talwar Date: Mon, 24 Aug 2026 00:36:01 +0530 Subject: [PATCH 1/6] perf: reuse fresh specialist research artifacts --- docs/research-cache.md | 9 ++ src-tauri/src/jobs/queue.rs | 42 ++++++- src-tauri/src/lib.rs | 26 ++-- src-tauri/src/orchestration.rs | 71 ++++++++++- src-tauri/src/research_cache.rs | 208 ++++++++++++++++++++++++++++++++ 5 files changed, 333 insertions(+), 23 deletions(-) create mode 100644 docs/research-cache.md create mode 100644 src-tauri/src/research_cache.rs diff --git a/docs/research-cache.md b/docs/research-cache.md new file mode 100644 index 0000000..783f593 --- /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, the complete job prompt, and a digest of every agent template. A prompt, model, depth, company, or template change therefore produces a cold miss. + +Only specialist JSON is cached. 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 manifests, symlinks, non-files, logs, 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..e2a47c2 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, + "claude-opus-4-7", + &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..a0887a6 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; @@ -152,24 +153,19 @@ pub fn run() { .build(tauri::generate_context!()) .expect("error while building tauri application") .run(|app_handle, event| { - #[cfg(target_os = "macos")] + // Handle macOS dock icon click + if let tauri::RunEvent::Reopen { + has_visible_windows, + .. + } = event { - // Handle macOS dock icon click. - if let tauri::RunEvent::Reopen { - has_visible_windows, - .. - } = event - { - if !has_visible_windows { - if let Some(window) = app_handle.get_webview_window("main") { - let _ = window.show(); - let _ = window.unminimize(); - let _ = window.set_focus(); - } + if !has_visible_windows { + if let Some(window) = app_handle.get_webview_window("main") { + let _ = window.show(); + let _ = window.unminimize(); + let _ = window.set_focus(); } } } - #[cfg(not(target_os = "macos"))] - let _ = (app_handle, event); }); } 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..6833dc2 --- /dev/null +++ b/src-tauri/src/research_cache.rs @@ -0,0 +1,208 @@ +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 || 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; + } + fs::copy(&path, destination.join(entry.file_name())).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; + } + fs::copy(&path, staging.join(entry.file_name())).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") +} + +#[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("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_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); + } +} From b9ab21bdf29106284be01e426f3b572aa204b90f Mon Sep 17 00:00:00 2001 From: Divyam Talwar Date: Mon, 24 Aug 2026 00:39:06 +0530 Subject: [PATCH 2/6] fix: include browser tool profile in cache identity --- docs/research-cache.md | 2 +- src-tauri/src/jobs/queue.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/research-cache.md b/docs/research-cache.md index 783f593..72e32ec 100644 --- a/docs/research-cache.md +++ b/docs/research-cache.md @@ -1,6 +1,6 @@ # 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, the complete job prompt, and a digest of every agent template. A prompt, model, depth, company, or template change therefore produces a cold miss. +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, depth, company, or template change therefore produces a cold miss. Only specialist JSON is cached. Stream logs, credentials, Apollo contact results, verifier output, and synthesizer output are excluded. The verifier and synthesizer always rerun over the selected specialist set. diff --git a/src-tauri/src/jobs/queue.rs b/src-tauri/src/jobs/queue.rs index e2a47c2..873b189 100644 --- a/src-tauri/src/jobs/queue.rs +++ b/src-tauri/src/jobs/queue.rs @@ -409,7 +409,7 @@ impl JobQueue { Some(crate::orchestration::research_cache_key( &entity_label, execution_depth, - "claude-opus-4-7", + &format!("claude-opus-4-7;chrome={}", settings.use_chrome), &prompt, )) } else { From 93b67da43d980c4097d352b6d4d7ca651d2e0ed6 Mon Sep 17 00:00:00 2001 From: Divyam Talwar Date: Mon, 24 Aug 2026 00:40:40 +0530 Subject: [PATCH 3/6] fix: preserve the cross-platform application gate --- src-tauri/src/lib.rs | 25 +++++++++++++++---------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index a0887a6..07006cf 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -153,19 +153,24 @@ pub fn run() { .build(tauri::generate_context!()) .expect("error while building tauri application") .run(|app_handle, event| { - // Handle macOS dock icon click - if let tauri::RunEvent::Reopen { - has_visible_windows, - .. - } = event + #[cfg(target_os = "macos")] { - if !has_visible_windows { - if let Some(window) = app_handle.get_webview_window("main") { - let _ = window.show(); - let _ = window.unminimize(); - let _ = window.set_focus(); + // Handle macOS dock icon click. + if let tauri::RunEvent::Reopen { + has_visible_windows, + .. + } = event + { + if !has_visible_windows { + if let Some(window) = app_handle.get_webview_window("main") { + let _ = window.show(); + let _ = window.unminimize(); + let _ = window.set_focus(); + } } } } + #[cfg(not(target_os = "macos"))] + let _ = (app_handle, event); }); } From aff56b16278769a8484c318a23127eda4242d2e9 Mon Sep 17 00:00:00 2001 From: Divyam Talwar Date: Mon, 24 Aug 2026 00:41:47 +0530 Subject: [PATCH 4/6] test: reject unsafe cache artifacts and timestamps --- docs/research-cache.md | 6 +++--- src-tauri/src/research_cache.rs | 24 +++++++++++++++++++++--- 2 files changed, 24 insertions(+), 6 deletions(-) diff --git a/docs/research-cache.md b/docs/research-cache.md index 72e32ec..2e6e467 100644 --- a/docs/research-cache.md +++ b/docs/research-cache.md @@ -1,9 +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, depth, company, or template change therefore produces a cold miss. +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 is cached. Stream logs, credentials, Apollo contact results, verifier output, and synthesizer output are excluded. The verifier and synthesizer always rerun over the selected specialist set. +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 manifests, symlinks, non-files, logs, and derived-agent files. Cache failure never converts a failed research job to success. +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/research_cache.rs b/src-tauri/src/research_cache.rs index 6833dc2..b00fb75 100644 --- a/src-tauri/src/research_cache.rs +++ b/src-tauri/src/research_cache.rs @@ -69,7 +69,10 @@ impl ResearchCache { .get("createdAtMs") .and_then(|value| value.as_i64()) .unwrap_or(0); - if created_at <= 0 || now_ms.saturating_sub(created_at) > self.ttl_ms { + if created_at <= 0 + || created_at > now_ms + || now_ms.saturating_sub(created_at) > self.ttl_ms + { return Ok(CacheRestore { agents: Vec::new() }); } @@ -94,7 +97,11 @@ impl ResearchCache { if is_derived_agent(agent) { continue; } - fs::copy(&path, destination.join(entry.file_name())).map_err(|e| e.to_string())?; + 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(); @@ -126,7 +133,11 @@ impl ResearchCache { if is_derived_agent(agent) { continue; } - fs::copy(&path, staging.join(entry.file_name())).map_err(|e| e.to_string())?; + 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(); @@ -167,6 +178,11 @@ 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::*; @@ -192,11 +208,13 @@ mod tests { 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()); From db8680abe727151c1ee195f3e3197ce881b7f7f0 Mon Sep 17 00:00:00 2001 From: Divyam Talwar Date: Mon, 24 Aug 2026 00:42:53 +0530 Subject: [PATCH 5/6] feat: enforce a citation integrity gateway --- README.md | 2 + docs/citation-contract.md | 16 ++ src-tauri/src/citation_integrity.rs | 278 +++++++++++++++++++++++ src-tauri/src/db/mod.rs | 20 ++ src-tauri/src/jobs/queue.rs | 60 ++++- src-tauri/src/lib.rs | 1 + src-tauri/src/prompts/agents/verifier.md | 3 +- 7 files changed, 375 insertions(+), 5 deletions(-) create mode 100644 docs/citation-contract.md create mode 100644 src-tauri/src/citation_integrity.rs diff --git a/README.md b/README.md index 0d4cdaa..7eae505 100644 --- a/README.md +++ b/README.md @@ -76,6 +76,8 @@ Augur OS attacks all three problems at once: it does the research fast, it groun - **Reviewable evidence.** A dedicated verifier stage reviews every finding the research agents produce, checks its supplied URL and evidence metadata, weighs confidence and freshness, and rejects or weakens unsupported, stale, or contradicted findings. The verifier is bounded and does not independently fetch every URL; open important citations before acting on them. +- **Enforced citation contract.** Standard and Deep runs must produce a strict verifier ledger. The Rust completion gate refuses malformed ledgers and will not accept a `verified` claim without a credential-free HTTPS URL, a substantive evidence quote, a bounded confidence score, and a valid date. Accepted receipts are persisted separately from the generated profile for audit. This validates the evidence contract; it still does not independently prove the source content. + - **Signal-based lead scoring.** Define your ideal customer once as a rubric. Augur grades every company from 0 to 100 against that rubric and shows the full reasoning behind the number, so the score is something you can defend in a pipeline review rather than a black box. - **Buying-committee discovery.** Augur finds the people who actually matter at each target, with their roles and the context for why each one is worth reaching, instead of dumping a flat list of names. diff --git a/docs/citation-contract.md b/docs/citation-contract.md new file mode 100644 index 0000000..e54ceb2 --- /dev/null +++ b/docs/citation-contract.md @@ -0,0 +1,16 @@ +# Citation contract gateway + +Standard and Deep research jobs do not become successful merely because the Claude process exits zero. Before profile parsing, Augur requires `outputs/specialists/verifier.json` to match a strict Rust schema. + +For every `verified` claim the gateway requires: + +- non-empty claim and source-agent identity; +- a credential-free HTTPS URL with a host; +- a substantive evidence quote; +- finite confidence from 0 through 1; +- `YYYY-MM-DD` or null for the evidence date; +- a unique SHA-256 receipt over the normalized claim, source-agent identity, and source URL. + +Weak and conflicting claims may omit evidence, but any URL they do retain must pass the same credential-free HTTPS boundary. Rejected claims must retain a reason. A malformed ledger fails the job, marks the entity failed, and prevents generated output from reaching the database. + +Accepted claim receipts are stored in SQLite's `evidence_receipts` table with the job and entity ids. This makes the prompt contract application-enforced and auditable. It does **not** independently fetch the URL or prove that the quote occurs on the source page; source snapshotting and quote-span verification remain future work. diff --git a/src-tauri/src/citation_integrity.rs b/src-tauri/src/citation_integrity.rs new file mode 100644 index 0000000..b683e92 --- /dev/null +++ b/src-tauri/src/citation_integrity.rs @@ -0,0 +1,278 @@ +use chrono::NaiveDate; +use rusqlite::{params, Connection}; +use serde::Deserialize; +use sha2::{Digest, Sha256}; +use std::collections::HashSet; +use std::fs; +use std::path::Path; + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct Ledger { + agent: String, + summary: String, + verified_claims: Vec, + rejected_claims: Vec, + conflicts: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct Claim { + claim: String, + source_agent: String, + evidence_url: Option, + evidence_quote: Option, + confidence: f64, + as_of_date: Option, + verdict: String, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct RejectedClaim { + claim: String, + source_agent: String, + reason: String, +} + +#[derive(Debug, Clone, PartialEq)] +pub struct EvidenceReceipt { + pub claim_hash: String, + pub claim: String, + pub source_agent: String, + pub evidence_url: Option, + pub evidence_quote: Option, + pub confidence: f64, + pub as_of_date: Option, + pub verdict: String, +} + +#[derive(Debug, Clone, PartialEq)] +pub struct IntegrityReport { + pub receipts: Vec, + pub rejected_count: usize, + pub conflict_count: usize, +} + +pub fn validate_ledger_file(path: &Path) -> Result { + let bytes = fs::read(path).map_err(|e| format!("cannot read verifier ledger: {e}"))?; + validate_ledger(&bytes) +} + +pub fn validate_ledger(bytes: &[u8]) -> Result { + let ledger: Ledger = serde_json::from_slice(bytes) + .map_err(|e| format!("verifier ledger does not match the strict schema: {e}"))?; + if ledger.agent != "verifier" || ledger.summary.trim().is_empty() { + return Err("verifier ledger must identify the verifier and include a summary".into()); + } + let mut receipts = Vec::with_capacity(ledger.verified_claims.len()); + let mut unique = HashSet::new(); + for (index, claim) in ledger.verified_claims.into_iter().enumerate() { + let label = format!("verified_claims[{index}]"); + if claim.claim.trim().is_empty() || claim.source_agent.trim().is_empty() { + return Err(format!("{label} must include a claim and source_agent")); + } + if !claim.confidence.is_finite() || !(0.0..=1.0).contains(&claim.confidence) { + return Err(format!("{label}.confidence must be between 0 and 1")); + } + if !matches!(claim.verdict.as_str(), "verified" | "weak" | "conflicting") { + return Err(format!("{label}.verdict is not supported")); + } + if let Some(ref date) = claim.as_of_date { + NaiveDate::parse_from_str(date, "%Y-%m-%d") + .map_err(|_| format!("{label}.as_of_date must be YYYY-MM-DD or null"))?; + } + if let Some(url) = claim.evidence_url.as_deref() { + let parsed = reqwest::Url::parse(url) + .map_err(|_| format!("{label}.evidence_url is not a valid URL"))?; + if parsed.scheme() != "https" + || parsed.host_str().is_none() + || !parsed.username().is_empty() + || parsed.password().is_some() + { + return Err(format!("{label}.evidence_url must be a credential-free HTTPS URL")); + } + } + if claim.verdict == "verified" && claim.evidence_url.is_none() { + return Err(format!("{label} cannot be verified without an evidence_url")); + } + if claim.verdict == "verified" { + if claim.evidence_quote.as_deref().map(str::trim).unwrap_or("").len() < 8 { + return Err(format!("{label} cannot be verified without a substantive evidence_quote")); + } + } + + let mut digest = Sha256::new(); + digest.update(claim.claim.trim().as_bytes()); + digest.update(b"\0"); + digest.update(claim.source_agent.trim().as_bytes()); + digest.update(b"\0"); + digest.update(claim.evidence_url.as_deref().unwrap_or("").as_bytes()); + let claim_hash = format!("{:x}", digest.finalize()); + if !unique.insert(claim_hash.clone()) { + return Err(format!("{label} duplicates an earlier claim/source receipt")); + } + receipts.push(EvidenceReceipt { + claim_hash, + claim: claim.claim.trim().to_string(), + source_agent: claim.source_agent.trim().to_string(), + evidence_url: claim.evidence_url, + evidence_quote: claim.evidence_quote, + confidence: claim.confidence, + as_of_date: claim.as_of_date, + verdict: claim.verdict, + }); + } + for (index, claim) in ledger.rejected_claims.iter().enumerate() { + if claim.claim.trim().is_empty() || claim.source_agent.trim().is_empty() || claim.reason.trim().is_empty() { + return Err(format!("rejected_claims[{index}] must include claim, source_agent, and reason")); + } + } + Ok(IntegrityReport { + receipts, + rejected_count: ledger.rejected_claims.len(), + conflict_count: ledger.conflicts.len(), + }) +} + +pub fn persist_receipts( + conn: &mut Connection, + job_id: &str, + entity_id: i64, + report: &IntegrityReport, +) -> Result<(), String> { + let transaction = conn.transaction().map_err(|e| e.to_string())?; + transaction + .execute("DELETE FROM evidence_receipts WHERE job_id = ?1", [job_id]) + .map_err(|e| e.to_string())?; + for receipt in &report.receipts { + transaction + .execute( + "INSERT INTO evidence_receipts ( + job_id, entity_id, claim_hash, claim, source_agent, evidence_url, + evidence_quote, confidence, as_of_date, verdict, created_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)", + params![ + job_id, + entity_id, + receipt.claim_hash, + receipt.claim, + receipt.source_agent, + receipt.evidence_url, + receipt.evidence_quote, + receipt.confidence, + receipt.as_of_date, + receipt.verdict, + chrono::Utc::now().timestamp_millis(), + ], + ) + .map_err(|e| e.to_string())?; + } + transaction.commit().map_err(|e| e.to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn ledger(claim: serde_json::Value) -> Vec { + serde_json::to_vec(&serde_json::json!({ + "agent": "verifier", + "summary": "bounded review", + "verified_claims": [claim], + "rejected_claims": [], + "conflicts": [] + })).unwrap() + } + + #[test] + fn verified_claim_requires_https_url_and_quote() { + let base = serde_json::json!({ + "claim": "Acme launched a product", + "source_agent": "trigger-signal-analyst", + "evidence_url": "https://example.com/launch", + "evidence_quote": "Acme today announced its new product.", + "confidence": 0.9, + "as_of_date": "2026-08-24", + "verdict": "verified" + }); + assert_eq!(validate_ledger(&ledger(base.clone())).unwrap().receipts.len(), 1); + let mut missing_quote = base; + missing_quote["evidence_quote"] = serde_json::Value::Null; + assert!(validate_ledger(&ledger(missing_quote)).unwrap_err().contains("evidence_quote")); + } + + #[test] + fn weak_claim_can_preserve_uncertain_evidence() { + let report = validate_ledger(&ledger(serde_json::json!({ + "claim": "Acme may be hiring", + "source_agent": "trigger-signal-analyst", + "evidence_url": null, + "evidence_quote": null, + "confidence": 0.4, + "as_of_date": null, + "verdict": "weak" + }))).unwrap(); + assert_eq!(report.receipts[0].verdict, "weak"); + } + + #[test] + fn optional_evidence_must_still_be_safe_and_sources_remain_distinct() { + let unsafe_report = ledger(serde_json::json!({ + "claim": "Acme may be hiring", + "source_agent": "trigger-signal-analyst", + "evidence_url": "https://user:secret@example.com/jobs", + "evidence_quote": null, + "confidence": 0.4, + "as_of_date": null, + "verdict": "weak" + })); + assert!(validate_ledger(&unsafe_report) + .unwrap_err() + .contains("credential-free HTTPS URL")); + + let mut value: serde_json::Value = serde_json::from_slice(&ledger(serde_json::json!({ + "claim": "Acme may be hiring", + "source_agent": "trigger-signal-analyst", + "evidence_url": "https://example.com/jobs", + "evidence_quote": null, + "confidence": 0.4, + "as_of_date": null, + "verdict": "weak" + }))).unwrap(); + let mut second = value["verified_claims"][0].clone(); + second["source_agent"] = serde_json::json!("people-finder"); + value["verified_claims"].as_array_mut().unwrap().push(second); + assert_eq!(validate_ledger(&serde_json::to_vec(&value).unwrap()).unwrap().receipts.len(), 2); + } + + #[test] + fn persists_a_deduplicated_job_receipt() { + let report = validate_ledger(&ledger(serde_json::json!({ + "claim": "Acme launched a product", + "source_agent": "trigger-signal-analyst", + "evidence_url": "https://example.com/launch", + "evidence_quote": "Acme today announced its new product.", + "confidence": 0.9, + "as_of_date": "2026-08-24", + "verdict": "verified" + }))).unwrap(); + let mut conn = Connection::open_in_memory().unwrap(); + conn.execute_batch(r#" + CREATE TABLE evidence_receipts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, job_id TEXT NOT NULL, + entity_id INTEGER NOT NULL, claim_hash TEXT NOT NULL, claim TEXT NOT NULL, + source_agent TEXT NOT NULL, evidence_url TEXT, evidence_quote TEXT, + confidence REAL NOT NULL, as_of_date TEXT, verdict TEXT NOT NULL, + created_at INTEGER NOT NULL, UNIQUE(job_id, claim_hash) + ); + "#).unwrap(); + persist_receipts(&mut conn, "job-1", 42, &report).unwrap(); + persist_receipts(&mut conn, "job-1", 42, &report).unwrap(); + let count: i64 = conn + .query_row("SELECT COUNT(*) FROM evidence_receipts WHERE job_id = 'job-1'", [], |row| row.get(0)) + .unwrap(); + assert_eq!(count, 1); + } +} diff --git a/src-tauri/src/db/mod.rs b/src-tauri/src/db/mod.rs index df746d6..6a1c70a 100644 --- a/src-tauri/src/db/mod.rs +++ b/src-tauri/src/db/mod.rs @@ -171,6 +171,26 @@ fn init_schema(conn: &Connection) -> SqliteResult<()> { CREATE INDEX IF NOT EXISTS idx_job_logs_job_id ON job_logs(job_id); CREATE INDEX IF NOT EXISTS idx_job_logs_sequence ON job_logs(job_id, sequence); + -- Machine-validated receipts from orchestrated verifier ledgers. + CREATE TABLE IF NOT EXISTS evidence_receipts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + job_id TEXT NOT NULL REFERENCES jobs(id) ON DELETE CASCADE, + entity_id INTEGER NOT NULL, + claim_hash TEXT NOT NULL, + claim TEXT NOT NULL, + source_agent TEXT NOT NULL, + evidence_url TEXT, + evidence_quote TEXT, + confidence REAL NOT NULL, + as_of_date TEXT, + verdict TEXT NOT NULL CHECK (verdict IN ('verified', 'weak', 'conflicting')), + created_at INTEGER NOT NULL, + UNIQUE(job_id, claim_hash) + ); + + CREATE INDEX IF NOT EXISTS idx_evidence_receipts_entity ON evidence_receipts(entity_id); + CREATE INDEX IF NOT EXISTS idx_evidence_receipts_job ON evidence_receipts(job_id); + CREATE TABLE IF NOT EXISTS apollo_usage ( id INTEGER PRIMARY KEY AUTOINCREMENT, job_id TEXT, diff --git a/src-tauri/src/jobs/queue.rs b/src-tauri/src/jobs/queue.rs index 873b189..d253f27 100644 --- a/src-tauri/src/jobs/queue.rs +++ b/src-tauri/src/jobs/queue.rs @@ -824,13 +824,65 @@ impl JobQueue { } // Finalize stream processor and get completion context - let completion_ctx = stream_processor.finalize(result.2, result.1).await; + let mut completion_ctx = stream_processor.finalize(result.2, result.1).await; + + let integrity_error = if result.2 + && execution_depth != crate::orchestration::ResearchDepth::Light + { + let ledger_path = PathBuf::from(&effective_working_dir) + .join("outputs") + .join("specialists") + .join("verifier.json"); + match crate::citation_integrity::validate_ledger_file(&ledger_path) { + Ok(report) => { + let persisted = db_conn + .lock() + .map_err(|e| e.to_string()) + .and_then(|mut conn| { + crate::citation_integrity::persist_receipts( + &mut conn, + &job_id_clone, + metadata.entity_id, + &report, + ) + }); + match persisted { + Ok(()) => { + eprintln!( + "[job_queue] job_id={} Citation contract accepted {} claims, {} rejected, {} conflicts", + job_id_clone, + report.receipts.len(), + report.rejected_count, + report.conflict_count + ); + None + } + Err(error) => Some(format!( + "could not persist citation receipts: {}", + error + )), + } + } + Err(error) => Some(format!("citation contract failed: {}", error)), + } + } else { + None + }; + if integrity_error.is_some() { + completion_ctx.success = false; + } // Process completion atomically using CompletionHandler let completion_handler = CompletionHandler::new(db_conn.clone(), app_clone.clone()); - let mut final_status = result.0.clone(); - let mut final_success = result.2; - let mut final_error_msg = if !result.2 { + let mut final_status = if integrity_error.is_some() { + "error".to_string() + } else { + result.0.clone() + }; + let mut final_success = result.2 && integrity_error.is_none(); + let mut final_error_msg = if let Some(error) = integrity_error { + Some(error) + } else if !result.2 { Some(format!("Job {} with code {:?}", result.0, result.1)) } else { None diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index 07006cf..8d7b6f2 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -1,5 +1,6 @@ mod apollo; mod commands; +mod citation_integrity; mod db; mod events; mod jobs; diff --git a/src-tauri/src/prompts/agents/verifier.md b/src-tauri/src/prompts/agents/verifier.md index 5033ada..865a679 100644 --- a/src-tauri/src/prompts/agents/verifier.md +++ b/src-tauri/src/prompts/agents/verifier.md @@ -31,6 +31,7 @@ The file must contain JSON: "claim": "verified claim", "source_agent": "agent name", "evidence_url": "https://...", + "evidence_quote": "substantive source excerpt supporting the claim", "confidence": 0.0, "as_of_date": "YYYY-MM-DD or null", "verdict": "verified | weak | conflicting" @@ -47,4 +48,4 @@ The file must contain JSON: } ``` -When in doubt, mark weak or reject. The final synthesizer must treat your verdict as authoritative. +The application enforces this envelope after the run. A `verified` claim without a valid credential-free HTTPS URL and a substantive evidence quote fails the job instead of reaching the profile. Weak/conflicting claims may retain null evidence fields. When in doubt, mark weak or reject. The final synthesizer must treat your verdict as authoritative. From fcd508d41b400dd618167a715666b0fc923d9eb8 Mon Sep 17 00:00:00 2001 From: Divyam Talwar Date: Mon, 24 Aug 2026 00:45:54 +0530 Subject: [PATCH 6/6] fix: satisfy strict citation lint gate --- src-tauri/src/citation_integrity.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/src-tauri/src/citation_integrity.rs b/src-tauri/src/citation_integrity.rs index b683e92..e99d1e6 100644 --- a/src-tauri/src/citation_integrity.rs +++ b/src-tauri/src/citation_integrity.rs @@ -97,10 +97,10 @@ pub fn validate_ledger(bytes: &[u8]) -> Result { if claim.verdict == "verified" && claim.evidence_url.is_none() { return Err(format!("{label} cannot be verified without an evidence_url")); } - if claim.verdict == "verified" { - if claim.evidence_quote.as_deref().map(str::trim).unwrap_or("").len() < 8 { - return Err(format!("{label} cannot be verified without a substantive evidence_quote")); - } + if claim.verdict == "verified" + && claim.evidence_quote.as_deref().map(str::trim).unwrap_or("").len() < 8 + { + return Err(format!("{label} cannot be verified without a substantive evidence_quote")); } let mut digest = Sha256::new();