use crate::sql::sql;
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use anyhow::{Context, Result};
use chrono::Utc;
use clap::Parser;
use serde::Serialize;
use serde_json::json;
use khive_db::StorageBackend;
use khive_mcp::serve::{resolve_runtime_config, RuntimeConfigInputs};
use khive_pack_code::{ingest_findings_json, CodeIngestBatch, CodeIngestOptions};
use khive_runtime::{
entity_fts_document, note_fts_document, secret_gate, GateRef, IngestAuditStore,
InterceptedDispatchResult, KhiveRuntime, Namespace, PackRegistry, RuntimeError, VerbRegistry,
};
use khive_storage::{SqlStatement, SqlValue, SubstrateKind};
const WRITER_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
const CODE_FINDINGS_INGEST_GATE_VERB: &str = "code.findings_ingest";
#[derive(Parser, Debug)]
pub struct CodeIngestArgs {
pub findings: PathBuf,
#[arg(long = "source-run")]
pub source_run: Option<String>,
#[arg(long, env = "KHIVE_DB")]
pub db: Option<String>,
#[arg(long, default_value = "local")]
pub namespace: String,
#[arg(long)]
pub dry_run: bool,
#[arg(long)]
pub human: bool,
}
#[derive(Debug, Default, Serialize)]
pub struct CodeIngestReport {
pub dry_run: bool,
pub entities_created: u64,
pub entities_skipped_existing: u64,
pub notes_created: u64,
pub notes_skipped_existing: u64,
pub edges_created: u64,
pub edges_skipped_existing: u64,
pub truncation_by_model: BTreeMap<String, khive_runtime::retrieval::EmbeddingTruncationReport>,
}
pub async fn run_code_ingest(args: CodeIngestArgs) -> Result<()> {
let human = args.human;
let report = code_ingest_batch(args).await?;
if human {
println!(
"entities: {} created, {} skipped\nnotes: {} created, {} skipped\nedges: {} created, {} skipped{}",
report.entities_created,
report.entities_skipped_existing,
report.notes_created,
report.notes_skipped_existing,
report.edges_created,
report.edges_skipped_existing,
if report.dry_run {
"\n(dry run: nothing written)"
} else {
""
},
);
if report
.truncation_by_model
.values()
.any(khive_runtime::retrieval::EmbeddingTruncationReport::any_truncated)
{
println!(
"embedding truncation: {}",
serde_json::to_string(&report.truncation_by_model)?
);
}
} else {
println!("{}", serde_json::to_string_pretty(&report)?);
}
Ok(())
}
async fn code_ingest_batch(args: CodeIngestArgs) -> Result<CodeIngestReport> {
code_ingest_batch_with_runtime_setup(args, |_| Ok(())).await
}
async fn code_ingest_batch_with_runtime_setup<F>(
args: CodeIngestArgs,
runtime_setup: F,
) -> Result<CodeIngestReport>
where
F: FnOnce(&KhiveRuntime) -> Result<()>,
{
code_ingest_batch_with_config_and_runtime_setup(args, None, runtime_setup).await
}
fn build_code_ingest_registry(runtime: &KhiveRuntime) -> Result<VerbRegistry> {
PackRegistry::build_ingest_registry(runtime, IngestAuditStore::Attach)
.map_err(|e| anyhow::anyhow!("{e}"))
}
async fn code_ingest_batch_with_config_and_runtime_setup<F>(
args: CodeIngestArgs,
config: Option<&Path>,
runtime_setup: F,
) -> Result<CodeIngestReport>
where
F: FnOnce(&KhiveRuntime) -> Result<()>,
{
code_ingest_batch_with_config_runtime_setup_and_gate(args, config, None, runtime_setup).await
}
async fn code_ingest_batch_with_config_runtime_setup_and_gate<F>(
args: CodeIngestArgs,
config: Option<&Path>,
gate_override: Option<GateRef>,
runtime_setup: F,
) -> Result<CodeIngestReport>
where
F: FnOnce(&KhiveRuntime) -> Result<()>,
{
let bytes = std::fs::read(&args.findings)
.with_context(|| format!("failed to read {}", args.findings.display()))?;
let ns = Namespace::parse(&args.namespace).map_err(|e| anyhow::anyhow!("{e}"))?;
let mut cfg = resolve_runtime_config(RuntimeConfigInputs {
db: args.db.as_deref(),
config,
namespace: ns,
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})?;
if let Some(gate) = gate_override {
cfg.gate = gate;
}
if !cfg.packs.iter().any(|p| p == "code") {
anyhow::bail!(
"the `code` pack is not in the configured pack set {:?}; `finding` notes require it \
to be loaded (set KHIVE_PACKS to include `code`, or drop --pack overrides)",
cfg.packs
);
}
let batch = ingest_findings_json(
&bytes,
CodeIngestOptions {
namespace: cfg.default_namespace.as_str(),
observed_at: Utc::now(),
source_run: args.source_run.as_deref(),
},
)
.with_context(|| format!("{} failed validation", args.findings.display()))?;
preflight_secret_gate(&batch)?;
if args.dry_run {
return dry_run_report(cfg.db_path.as_deref(), &batch).await;
}
let runtime = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
let registry = build_code_ingest_registry(&runtime)?;
let gate_args = json!({
"input_kind": "findings_json",
"record_counts": {
"entities": batch.entities.len(),
"notes": batch.notes.len(),
"edges": batch.edges.len(),
}
});
let gated_result = registry
.dispatch_intercepted_with_metadata_with_identity(
CODE_FINDINGS_INGEST_GATE_VERB,
&gate_args,
None,
|resolved_ns| async {
let direct_result: Result<CodeIngestReport> = async {
runtime_setup(&runtime)?;
let token = runtime
.authorize(resolved_ns)
.map_err(|e| anyhow::anyhow!("{e}"))
.context("failed to authorize namespace")?;
let embedding_model_names = runtime.registered_embedding_model_names();
let mut report = CodeIngestReport {
dry_run: false,
..CodeIngestReport::default()
};
let entities = runtime
.entities(&token)
.map_err(|e| anyhow::anyhow!("{e}"))?;
for entity in &batch.entities {
let existing = entities
.get_entity_including_deleted(entity.id)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if existing.is_some() {
report.entities_skipped_existing += 1;
continue;
}
report.entities_created += 1;
entities
.upsert_entity(entity.clone())
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let doc = entity_fts_document(entity);
let embed_body = doc.body.clone();
if let Ok(fts) = runtime.text(&token) {
if let Err(e) = fts.upsert_document(doc).await {
tracing::warn!(
entity_id = %entity.id,
error = %e,
"code-ingest: entity FTS indexing failed (non-fatal)"
);
}
}
for model_name in &embedding_model_names {
match runtime
.embed_document_with_model_outcome(model_name, &embed_body)
.await
{
Ok(outcome) => {
report
.truncation_by_model
.entry(model_name.clone())
.or_default()
.observe(&outcome);
if let Ok(vs) = runtime.vectors_for_model(&token, model_name) {
if let Err(e) = vs
.insert(
entity.id,
SubstrateKind::Entity,
token.namespace().as_str(),
"entity.body",
vec![outcome.vector],
)
.await
{
tracing::warn!(
entity_id = %entity.id,
model = %model_name,
error = %e,
"code-ingest: entity vector insert failed (non-fatal)"
);
}
}
}
Err(e) => tracing::warn!(
entity_id = %entity.id,
model = %model_name,
error = %e,
"code-ingest: entity embedding failed (non-fatal)"
),
}
}
}
let notes = runtime.notes(&token).map_err(|e| anyhow::anyhow!("{e}"))?;
for note in &batch.notes {
let existing = notes
.get_note_including_deleted(note.id)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if existing.is_some() {
report.notes_skipped_existing += 1;
continue;
}
report.notes_created += 1;
notes
.upsert_note(note.clone())
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if let Ok(fts) = runtime.text_for_notes(&token) {
if let Err(e) = fts.upsert_document(note_fts_document(note)).await {
tracing::warn!(
note_id = %note.id,
error = %e,
"code-ingest: note FTS indexing failed (non-fatal)"
);
}
}
for model_name in &embedding_model_names {
match runtime
.embed_document_with_model_outcome(model_name, ¬e.content)
.await
{
Ok(outcome) => {
report
.truncation_by_model
.entry(model_name.clone())
.or_default()
.observe(&outcome);
if let Ok(vs) = runtime.vectors_for_model(&token, model_name) {
if let Err(e) = vs
.insert(
note.id,
SubstrateKind::Note,
token.namespace().as_str(),
"note.content",
vec![outcome.vector],
)
.await
{
tracing::warn!(
note_id = %note.id,
model = %model_name,
error = %e,
"code-ingest: note vector insert failed (non-fatal)"
);
}
}
}
Err(e) => tracing::warn!(
note_id = %note.id,
model = %model_name,
error = %e,
"code-ingest: note embedding failed (non-fatal)"
),
}
}
}
let graph = runtime.graph(&token).map_err(|e| anyhow::anyhow!("{e}"))?;
for edge in &batch.edges {
let existing = graph
.get_edge_including_deleted(edge.id)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if existing.is_some() {
report.edges_skipped_existing += 1;
continue;
}
report.edges_created += 1;
graph
.upsert_edge(edge.clone())
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
Ok(report)
}
.await;
direct_result
.map(|report| {
let audit_result = json!({
"entities_created": report.entities_created,
"entities_skipped_existing": report.entities_skipped_existing,
"notes_created": report.notes_created,
"notes_skipped_existing": report.notes_skipped_existing,
"edges_created": report.edges_created,
"edges_skipped_existing": report.edges_skipped_existing,
});
InterceptedDispatchResult::new(audit_result, report)
})
.map_err(|err| RuntimeError::Internal(err.to_string()))
},
)
.await;
let ingest_result = gated_result
.map(|outcome| outcome.metadata)
.map_err(|err| anyhow::anyhow!("{err}"));
drop(registry);
let writer_join = take_writer_task_join_or_warn(runtime.backend().pool());
drop(runtime);
settle_writer_drain(ingest_result, writer_join, WRITER_DRAIN_TIMEOUT).await
}
fn take_writer_task_join_or_warn(
pool: &khive_db::ConnectionPool,
) -> Option<tokio::task::JoinHandle<()>> {
let writer_join = pool.take_writer_task_join();
if writer_join.is_none() && pool.write_queue_active() && pool.writer_task_join_was_stored() {
tracing::warn!(
"writer-task JoinHandle was stored but is absent at drain time; \
another caller may have taken the one-shot drain handle, so \
'return implies settled' is not enforced by this call"
);
}
writer_join
}
async fn settle_writer_drain(
ingest_result: Result<CodeIngestReport>,
writer_join: Option<tokio::task::JoinHandle<()>>,
timeout: std::time::Duration,
) -> Result<CodeIngestReport> {
if let Some(join) = writer_join {
let drained = tokio::time::timeout(timeout, join).await;
match (&ingest_result, drained) {
(_, Ok(Ok(()))) => {}
(Err(_), Ok(Err(join_err))) => {
tracing::warn!(error = %join_err, "writer task terminated abnormally after failed ingest");
}
(Err(_), Err(_elapsed)) => {
tracing::warn!(
timeout = ?timeout,
"writer task did not drain after failed ingest; database file state may still be unsettled"
);
}
(Ok(_), Ok(Err(join_err))) => anyhow::bail!(
"writer task terminated abnormally after ingest completed: {join_err}"
),
(Ok(_), Err(_elapsed)) => anyhow::bail!(
"writer task did not drain within {timeout:?} after ingest; \
database file state may still be unsettled"
),
}
}
ingest_result
}
fn preflight_secret_gate(batch: &CodeIngestBatch) -> Result<()> {
for (index, entity) in batch.entities.iter().enumerate() {
let record = format!("entity[{index}]");
secret_gate::check_at(&entity.name, &record, "name").map_err(|e| anyhow::anyhow!("{e}"))?;
if let Some(description) = &entity.description {
secret_gate::check_at(description, &record, "description")
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
if let Some(properties) = &entity.properties {
secret_gate::check_json_at(properties, &record, "properties")
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())
.map_err(|e| anyhow::anyhow!("{e}"))?;
secret_gate::check_tags_at(&entity.tags, &record, "tags")
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
for (index, note) in batch.notes.iter().enumerate() {
let record = format!("note[{index}]");
secret_gate::check_at(¬e.content, &record, "content")
.map_err(|e| anyhow::anyhow!("{e}"))?;
if let Some(name) = ¬e.name {
secret_gate::check_at(name, &record, "name").map_err(|e| anyhow::anyhow!("{e}"))?;
}
if let Some(properties) = ¬e.properties {
secret_gate::check_json_at(properties, &record, "properties")
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
secret_gate::reject_reserved_secret_gate_property(note.properties.as_ref())
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
Ok(())
}
async fn dry_run_report(
db_path: Option<&Path>,
batch: &CodeIngestBatch,
) -> Result<CodeIngestReport> {
let mut report = CodeIngestReport {
dry_run: true,
..CodeIngestReport::default()
};
let existing_path = db_path.filter(|p| p.exists());
let Some(db_path) = existing_path else {
report.entities_created = batch.entities.len() as u64;
report.notes_created = batch.notes.len() as u64;
report.edges_created = batch.edges.len() as u64;
return Ok(report);
};
let (backend, _snapshot_dir) = open_read_only_snapshot(db_path)?;
let sql = backend.sql();
let mut reader = sql.reader().await.map_err(|e| anyhow::anyhow!("{e}"))?;
for entity in &batch.entities {
let row = reader
.query_scalar(SqlStatement {
sql: sql!("entities_exists").to_string(),
params: vec![SqlValue::Uuid(entity.id)],
label: Some("code-ingest dry-run entity existence".to_string()),
})
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if row.is_some() {
report.entities_skipped_existing += 1;
} else {
report.entities_created += 1;
}
}
for note in &batch.notes {
let row = reader
.query_scalar(SqlStatement {
sql: sql!("notes_exists").to_string(),
params: vec![SqlValue::Uuid(note.id)],
label: Some("code-ingest dry-run note existence".to_string()),
})
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if row.is_some() {
report.notes_skipped_existing += 1;
} else {
report.notes_created += 1;
}
}
for edge in &batch.edges {
let row = reader
.query_scalar(SqlStatement {
sql: sql!("graph_edges_exists").to_string(),
params: vec![SqlValue::Uuid(uuid::Uuid::from(edge.id))],
label: Some("code-ingest dry-run edge existence".to_string()),
})
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
if row.is_some() {
report.edges_skipped_existing += 1;
} else {
report.edges_created += 1;
}
}
Ok(report)
}
pub(crate) fn open_read_only_snapshot(
db_path: &Path,
) -> Result<(StorageBackend, tempfile::TempDir)> {
let snapshot_dir = tempfile::TempDir::new()
.context("failed to create a scratch directory for the dry-run db snapshot")?;
let file_name = db_path
.file_name()
.with_context(|| format!("{} has no file name component", db_path.display()))?;
let snapshot_db = snapshot_dir.path().join(file_name);
std::fs::copy(db_path, &snapshot_db)
.with_context(|| format!("failed to snapshot {} for dry-run", db_path.display()))?;
let wal_path = wal_sidecar_path(db_path);
if wal_path.exists() {
let snapshot_wal = wal_sidecar_path(&snapshot_db);
std::fs::copy(&wal_path, &snapshot_wal)
.with_context(|| format!("failed to snapshot {} for dry-run", wal_path.display()))?;
}
let shm_path = shm_sidecar_path(db_path);
if shm_path.exists() {
let snapshot_shm = shm_sidecar_path(&snapshot_db);
std::fs::copy(&shm_path, &snapshot_shm)
.with_context(|| format!("failed to snapshot {} for dry-run", shm_path.display()))?;
}
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
for sidecar in [
wal_sidecar_path(&snapshot_db),
shm_sidecar_path(&snapshot_db),
] {
if sidecar.exists() {
let mut permissions = std::fs::metadata(&sidecar)
.with_context(|| format!("stat snapshot sidecar {}", sidecar.display()))?
.permissions();
permissions.set_mode(0o444);
std::fs::set_permissions(&sidecar, permissions)
.with_context(|| format!("freeze snapshot sidecar {}", sidecar.display()))?;
}
}
}
let backend =
StorageBackend::sqlite_read_only(&snapshot_db).map_err(|e| anyhow::anyhow!("{e}"))?;
Ok((backend, snapshot_dir))
}
fn wal_sidecar_path(db_path: &Path) -> PathBuf {
let mut name = db_path.as_os_str().to_owned();
name.push("-wal");
PathBuf::from(name)
}
fn shm_sidecar_path(db_path: &Path) -> PathBuf {
let mut name = db_path.as_os_str().to_owned();
name.push("-shm");
PathBuf::from(name)
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use khive_runtime::{
EmbedderProvider, Gate, GateDecision, GateError, GateRef, GateRequest, RuntimeError,
};
use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService, MAX_TEXT_BYTES};
use serial_test::serial;
use super::*;
struct WriteQueueEnvGuard {
previous: Option<std::ffi::OsString>,
}
impl WriteQueueEnvGuard {
fn unset() -> Self {
let previous = std::env::var_os("KHIVE_WRITE_QUEUE");
std::env::remove_var("KHIVE_WRITE_QUEUE");
Self { previous }
}
}
impl Drop for WriteQueueEnvGuard {
fn drop(&mut self) {
match &self.previous {
Some(value) => std::env::set_var("KHIVE_WRITE_QUEUE", value),
None => std::env::remove_var("KHIVE_WRITE_QUEUE"),
}
}
}
struct PacksEnvGuard {
previous: Option<std::ffi::OsString>,
}
impl PacksEnvGuard {
fn pin(packs: &[&str]) -> Self {
let previous = std::env::var_os("KHIVE_PACKS");
std::env::set_var("KHIVE_PACKS", packs.join(","));
Self { previous }
}
}
impl Drop for PacksEnvGuard {
fn drop(&mut self) {
match &self.previous {
Some(value) => std::env::set_var("KHIVE_PACKS", value),
None => std::env::remove_var("KHIVE_PACKS"),
}
}
}
#[derive(Clone, Default)]
struct Capture(Arc<std::sync::Mutex<Vec<u8>>>);
impl std::io::Write for Capture {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
struct MakeCapture(Capture);
impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for MakeCapture {
type Writer = Capture;
fn make_writer(&'a self) -> Self::Writer {
self.0.clone()
}
}
struct FixedEmbeddingService {
dimensions: usize,
}
#[async_trait]
impl EmbeddingService for FixedEmbeddingService {
async fn embed(
&self,
texts: &[String],
_model: EmbeddingModel,
) -> std::result::Result<Vec<Vec<f32>>, EmbedError> {
Ok(texts
.iter()
.map(|_| vec![1.0_f32; self.dimensions])
.collect())
}
fn supports_model(&self, _model: EmbeddingModel) -> bool {
true
}
fn name(&self) -> &'static str {
"code-ingest-test"
}
}
struct FixedEmbeddingProvider {
name: String,
dimensions: usize,
}
#[async_trait]
impl EmbedderProvider for FixedEmbeddingProvider {
fn name(&self) -> &str {
&self.name
}
fn dimensions(&self) -> usize {
self.dimensions
}
async fn build(&self) -> std::result::Result<Arc<dyn EmbeddingService>, RuntimeError> {
Ok(Arc::new(FixedEmbeddingService {
dimensions: self.dimensions,
}))
}
}
fn base_args(findings: PathBuf, db: PathBuf) -> CodeIngestArgs {
CodeIngestArgs {
findings,
source_run: Some("test-run".to_string()),
db: Some(db.display().to_string()),
namespace: "local".to_string(),
dry_run: false,
human: false,
}
}
fn write_empty_test_config(dir: &Path) -> PathBuf {
let path = dir.join("empty-khive-config.toml");
std::fs::write(&path, "").expect("write isolated empty config");
path
}
fn write_test_config_with_packs(dir: &Path, packs: &[&str]) -> PathBuf {
let path = dir.join("packs-khive-config.toml");
let list = packs
.iter()
.map(|p| format!("\"{p}\""))
.collect::<Vec<_>>()
.join(", ");
std::fs::write(&path, format!("[runtime]\npacks = [{list}]\n"))
.expect("write isolated pack-scoped config");
path
}
async fn code_ingest_batch(args: CodeIngestArgs) -> Result<CodeIngestReport> {
code_ingest_batch_with_runtime_setup(args, |_| Ok(())).await
}
fn install_test_embedders(runtime: &KhiveRuntime) -> Result<()> {
for name in runtime.registered_embedding_model_names() {
let dimensions = runtime.resolve_embedding_model(Some(&name))?.dimensions();
runtime.register_embedder(FixedEmbeddingProvider { name, dimensions });
}
Ok(())
}
#[tokio::test]
async fn default_ingest_fixture_writes_without_model_cache() {
if crate::test_process::run_in_child() {
return;
}
let dir = tempfile::tempdir().expect("ingest fixture");
let findings = write_valid_findings(dir.path());
let report = code_ingest_batch(base_args(findings, dir.path().join("ingest.db")))
.await
.expect("offline fixture ingest");
assert_eq!(report.entities_created, 1);
assert_eq!(report.notes_created, 1);
assert_eq!(report.edges_created, 1);
let home = std::env::var_os("HOME").expect("private child HOME");
assert!(
std::fs::read_dir(home).unwrap().next().is_none(),
"default ingest fixture must not create a native model cache"
);
}
async fn code_ingest_batch_with_runtime_setup<F>(
args: CodeIngestArgs,
runtime_setup: F,
) -> Result<CodeIngestReport>
where
F: FnOnce(&KhiveRuntime) -> Result<()>,
{
let config = write_empty_test_config(
args.findings
.parent()
.expect("findings fixture must have a parent directory"),
);
super::code_ingest_batch_with_config_and_runtime_setup(args, Some(&config), |runtime| {
install_test_embedders(runtime)?;
runtime_setup(runtime)
})
.await
}
fn write_valid_findings(dir: &std::path::Path) -> PathBuf {
let path = dir.join("findings.json");
std::fs::write(
&path,
r#"{
"audit": {
"date": "2026-07-11",
"scope": "khive-pack-code",
"repo": "ohdearquant/khive",
"branch": "feat/adr085-code-ingest-admin",
"commit": "abc1234",
"standards_file": "docs/standards.md"
},
"findings": [
{
"id": "F-001",
"title": "Example finding for a CLI integration test",
"severity": "medium",
"confidence": "high",
"failure_scenario": "Reproduced by running kkernel code-ingest twice.",
"evidence": "code_ingest.rs test",
"impact": "none, this is a test fixture"
}
]
}"#,
)
.expect("write findings.json fixture");
path
}
const ALL_PACKS: &[&str] = &[
"kg",
"gtd",
"memory",
"brain",
"comm",
"schedule",
"knowledge",
"session",
"git",
"code",
"workspace",
"blob",
"tool",
"exec",
];
#[test]
fn all_packs_fixture_matches_the_shipping_declaration() {
let mut fixture: Vec<String> = ALL_PACKS.iter().map(|p| p.to_string()).collect();
let mut declared = khive_runtime::RuntimeConfig::built_in_packs();
fixture.sort();
declared.sort();
assert_eq!(
fixture, declared,
"ALL_PACKS must equal RuntimeConfig::built_in_packs(); add the pack to the \
declaration and this fixture together, or derive the fixture from it"
);
}
#[derive(Debug, Default)]
struct DenyFindingsIngestGate {
requests: Mutex<Vec<GateRequest>>,
}
impl Gate for DenyFindingsIngestGate {
fn check(&self, request: &GateRequest) -> std::result::Result<GateDecision, GateError> {
self.requests
.lock()
.expect("gate requests lock")
.push(request.clone());
if request.verb == CODE_FINDINGS_INGEST_GATE_VERB {
Ok(GateDecision::deny("findings ingest denied by test policy"))
} else {
Ok(GateDecision::allow())
}
}
fn impl_name(&self) -> &'static str {
"DenyFindingsIngestGate"
}
}
#[derive(Debug, Default)]
struct AllowFindingsIngestGate;
impl Gate for AllowFindingsIngestGate {
fn check(&self, _request: &GateRequest) -> std::result::Result<GateDecision, GateError> {
Ok(GateDecision::allow())
}
fn impl_name(&self) -> &'static str {
"AllowFindingsIngestGate"
}
}
#[derive(Debug, Default)]
struct ErroringFindingsIngestGate;
impl Gate for ErroringFindingsIngestGate {
fn check(&self, request: &GateRequest) -> std::result::Result<GateDecision, GateError> {
if request.verb == CODE_FINDINGS_INGEST_GATE_VERB {
Err(GateError::Internal(
"test gate backend unreachable".to_string(),
))
} else {
Ok(GateDecision::allow())
}
}
fn impl_name(&self) -> &'static str {
"ErroringFindingsIngestGate"
}
}
async fn only_findings_ingest_event(db: &Path, config: &Path) -> khive_storage::Event {
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(db.to_str().expect("utf8 path")),
config: Some(config),
namespace: Namespace::parse("local").expect("valid namespace"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve runtime config for event read-back");
let runtime = KhiveRuntime::new(cfg).expect("open runtime for event read-back");
let token = runtime
.authorize(Namespace::parse("local").expect("valid namespace"))
.expect("authorize local namespace for event read-back");
let events = runtime.events(&token).expect("events store handle");
let page = events
.query_events(
khive_storage::EventFilter {
verbs: vec![CODE_FINDINGS_INGEST_GATE_VERB.to_string()],
..Default::default()
},
khive_storage::PageRequest::default(),
)
.await
.expect("query code.findings_ingest audit events");
assert_eq!(
page.items.len(),
1,
"expected exactly one code.findings_ingest audit event, got {:?}",
page.items
);
page.items.into_iter().next().expect("checked len == 1")
}
#[serial]
#[test]
fn packs_env_guard_pins_the_fixture_pack_set_over_an_ambient_override() {
if crate::test_process::run_in_child() {
return;
}
let tmp = tempfile::TempDir::new().expect("temp dir");
let config = write_test_config_with_packs(tmp.path(), ALL_PACKS);
let db = tmp.path().join("hermetic.db");
let resolve = || {
resolve_runtime_config(RuntimeConfigInputs {
db: Some(db.to_str().expect("utf8 path")),
config: Some(&config),
namespace: Namespace::parse("local").expect("valid namespace"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve runtime config")
.packs
};
let expected: Vec<String> = ALL_PACKS.iter().map(|p| p.to_string()).collect();
let previous = std::env::var_os("KHIVE_PACKS");
std::env::set_var("KHIVE_PACKS", "kg,formal");
assert_eq!(resolve(), vec!["kg".to_string(), "formal".to_string()]);
{
let _packs = PacksEnvGuard::pin(ALL_PACKS);
assert_eq!(
resolve(),
expected,
"the pinned list must win over the ambient value"
);
}
assert_eq!(
std::env::var("KHIVE_PACKS").as_deref(),
Ok("kg,formal"),
"the guard must restore the ambient value it replaced"
);
match previous {
Some(value) => std::env::set_var("KHIVE_PACKS", value),
None => std::env::remove_var("KHIVE_PACKS"),
}
}
#[serial]
#[tokio::test]
async fn code_ingest_batch_gate_denial_precedes_setup_and_record_writes() {
if crate::test_process::run_in_child() {
return;
}
let tmp = tempfile::TempDir::new().expect("temp dir");
let _packs = PacksEnvGuard::pin(ALL_PACKS);
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("gate-denied.db");
let config = write_test_config_with_packs(tmp.path(), ALL_PACKS);
let gate = Arc::new(DenyFindingsIngestGate::default());
let setup_called = Arc::new(std::sync::atomic::AtomicBool::new(false));
let setup_called_in_closure = Arc::clone(&setup_called);
let err = super::code_ingest_batch_with_config_runtime_setup_and_gate(
base_args(findings, db.clone()),
Some(&config),
Some(Arc::clone(&gate) as GateRef),
move |_| {
setup_called_in_closure.store(true, std::sync::atomic::Ordering::SeqCst);
Ok(())
},
)
.await
.expect_err("the operation-specific gate must refuse the direct-store batch");
assert!(
err.to_string().contains(CODE_FINDINGS_INGEST_GATE_VERB),
"typed refusal must name the governed batch operation: {err}"
);
assert!(
!setup_called.load(std::sync::atomic::Ordering::SeqCst),
"gate denial must occur before the direct-write section begins"
);
{
let requests = gate.requests.lock().expect("gate requests lock");
assert_eq!(
requests.len(),
1,
"one batch must make one gate consultation"
);
let request = &requests[0];
assert_eq!(request.verb, CODE_FINDINGS_INGEST_GATE_VERB);
assert_eq!(request.args["record_counts"]["entities"], 1);
assert_eq!(request.args["record_counts"]["notes"], 1);
assert_eq!(request.args["record_counts"]["edges"], 1);
}
let backend = StorageBackend::sqlite(&db).expect("open denied-ingest database");
let sql = backend.sql();
let mut reader = sql.reader().await.expect("acquire denied-ingest reader");
for table in ["entities", "notes", "graph_edges"] {
let count = reader
.query_scalar(SqlStatement {
sql: format!("SELECT COUNT(*) FROM {table}"),
params: vec![],
label: Some(format!("count denied-ingest {table}")),
})
.await
.expect("count denied-ingest rows");
assert!(
matches!(count, Some(SqlValue::Integer(0))),
"gate denial must leave {table} empty, got {count:?}"
);
}
let event = only_findings_ingest_event(&db, &config).await;
assert_eq!(event.outcome, khive_storage::EventOutcome::Denied);
assert_eq!(event.verb, CODE_FINDINGS_INGEST_GATE_VERB);
assert_eq!(event.payload["decision"], "deny");
}
#[serial]
#[tokio::test]
async fn code_ingest_batch_gate_allow_persists_success_audit_event() {
if crate::test_process::run_in_child() {
return;
}
let tmp = tempfile::TempDir::new().expect("temp dir");
let _packs = PacksEnvGuard::pin(ALL_PACKS);
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("gate-allowed.db");
let config = write_test_config_with_packs(tmp.path(), ALL_PACKS);
let gate = Arc::new(AllowFindingsIngestGate);
let report = super::code_ingest_batch_with_config_runtime_setup_and_gate(
base_args(findings, db.clone()),
Some(&config),
Some(gate as GateRef),
install_test_embedders,
)
.await
.expect("an allowing gate must let the batch proceed");
assert_eq!(report.entities_created, 1);
assert_eq!(report.notes_created, 1);
assert_eq!(report.edges_created, 1);
let event = only_findings_ingest_event(&db, &config).await;
assert_eq!(event.outcome, khive_storage::EventOutcome::Success);
assert_eq!(event.verb, CODE_FINDINGS_INGEST_GATE_VERB);
assert_eq!(event.payload["decision"], "allow");
}
#[serial]
#[tokio::test]
async fn code_ingest_batch_gate_error_precedes_setup_and_record_writes() {
if crate::test_process::run_in_child() {
return;
}
let tmp = tempfile::TempDir::new().expect("temp dir");
let _packs = PacksEnvGuard::pin(ALL_PACKS);
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("gate-errored.db");
let config = write_test_config_with_packs(tmp.path(), ALL_PACKS);
let gate = Arc::new(ErroringFindingsIngestGate);
let setup_called = Arc::new(std::sync::atomic::AtomicBool::new(false));
let setup_called_in_closure = Arc::clone(&setup_called);
let err = super::code_ingest_batch_with_config_runtime_setup_and_gate(
base_args(findings, db.clone()),
Some(&config),
Some(gate as GateRef),
move |_| {
setup_called_in_closure.store(true, std::sync::atomic::Ordering::SeqCst);
Ok(())
},
)
.await
.expect_err("a gate backend failure must refuse the direct-store batch");
assert!(
err.to_string().contains(CODE_FINDINGS_INGEST_GATE_VERB),
"gate-unavailable refusal must name the governed batch operation: {err}"
);
assert!(
!setup_called.load(std::sync::atomic::Ordering::SeqCst),
"gate error must occur before the direct-write section begins"
);
let backend = StorageBackend::sqlite(&db).expect("open errored-ingest database");
let sql = backend.sql();
let mut reader = sql.reader().await.expect("acquire errored-ingest reader");
for table in ["entities", "notes", "graph_edges"] {
let count = reader
.query_scalar(SqlStatement {
sql: format!("SELECT COUNT(*) FROM {table}"),
params: vec![],
label: Some(format!("count errored-ingest {table}")),
})
.await
.expect("count errored-ingest rows");
assert!(
matches!(count, Some(SqlValue::Integer(0))),
"gate unavailability must leave {table} empty, got {count:?}"
);
}
let event = only_findings_ingest_event(&db, &config).await;
assert_eq!(event.outcome, khive_storage::EventOutcome::Error);
assert_eq!(event.verb, CODE_FINDINGS_INGEST_GATE_VERB);
assert_eq!(event.payload["decision"], "gate_unavailable");
}
const TOMBSTONE_WITNESS: i64 = 1_772_812_800_000_000;
fn mapped_batch(findings: &Path) -> CodeIngestBatch {
let bytes = std::fs::read(findings).expect("read findings fixture");
ingest_findings_json(
&bytes,
CodeIngestOptions {
namespace: "local",
observed_at: Utc::now(),
source_run: Some("test-run"),
},
)
.expect("map valid findings fixture")
}
async fn soft_delete_mapped_batch(db: &Path, batch: &CodeIngestBatch) {
let backend = StorageBackend::sqlite(db).expect("open tombstone writer");
let sql = backend.sql();
let mut writer = sql.writer().await.expect("acquire tombstone writer");
let rows = [
(
"entities",
SqlValue::Uuid(batch.entities[0].id),
"soft-delete mapped entity",
),
(
"notes",
SqlValue::Uuid(batch.notes[0].id),
"soft-delete mapped note",
),
(
"graph_edges",
SqlValue::Uuid(uuid::Uuid::from(batch.edges[0].id)),
"soft-delete mapped edge",
),
];
for (table, id, label) in rows {
let version = if table == "entities" {
", version = version + 1"
} else {
""
};
let changed = writer
.execute(SqlStatement {
sql: format!("UPDATE {table} SET deleted_at = ?1{version} WHERE id = ?2"),
params: vec![SqlValue::Integer(TOMBSTONE_WITNESS), id],
label: Some(label.to_string()),
})
.await
.expect("soft-delete mapped row");
assert_eq!(changed, 1, "fixture must tombstone exactly one {table} row");
}
}
async fn assert_mapped_batch_remains_tombstoned(db: &Path, batch: &CodeIngestBatch) {
let backend = StorageBackend::sqlite(db).expect("open tombstone reader");
let sql = backend.sql();
let mut reader = sql.reader().await.expect("acquire tombstone reader");
let rows = [
("entities", SqlValue::Uuid(batch.entities[0].id)),
("notes", SqlValue::Uuid(batch.notes[0].id)),
(
"graph_edges",
SqlValue::Uuid(uuid::Uuid::from(batch.edges[0].id)),
),
];
for (table, id) in rows {
let marker = reader
.query_scalar(SqlStatement {
sql: format!("SELECT deleted_at FROM {table} WHERE id = ?1"),
params: vec![id],
label: Some(format!("read mapped {table} tombstone")),
})
.await
.expect("read mapped tombstone");
assert!(
matches!(&marker, Some(SqlValue::Integer(value)) if *value == TOMBSTONE_WITNESS),
"re-ingest must preserve the exact {table} tombstone marker, got {marker:?}"
);
}
}
#[serial]
#[tokio::test]
async fn code_ingest_creates_once_then_skips_on_rerun() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("scratch.db");
let first = code_ingest_batch(base_args(findings.clone(), db.clone()))
.await
.expect("first ingest must succeed");
assert_eq!(first.entities_created, 1);
assert_eq!(first.notes_created, 1);
assert_eq!(first.edges_created, 1);
assert_eq!(first.entities_skipped_existing, 0);
assert_eq!(first.notes_skipped_existing, 0);
assert_eq!(first.edges_skipped_existing, 0);
let second = code_ingest_batch(base_args(findings, db))
.await
.expect("re-ingesting the same sweep must succeed");
assert_eq!(
second.notes_created, 0,
"content-derived ids must make a re-ingest a no-op, not a duplicate write"
);
assert_eq!(second.notes_skipped_existing, 1);
assert_eq!(second.entities_skipped_existing, 1);
assert_eq!(second.edges_skipped_existing, 1);
}
#[serial]
#[tokio::test]
async fn code_ingest_never_reactivates_consumed_tombstone_ids() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("tombstones.db");
let batch = mapped_batch(&findings);
code_ingest_batch(base_args(findings.clone(), db.clone()))
.await
.expect("initial ingest must succeed");
soft_delete_mapped_batch(&db, &batch).await;
assert_mapped_batch_remains_tombstoned(&db, &batch).await;
let mut dry_args = base_args(findings.clone(), db.clone());
dry_args.dry_run = true;
let dry = code_ingest_batch(dry_args)
.await
.expect("dry-run over tombstones must succeed");
assert!(dry.dry_run);
assert_eq!(dry.entities_created, 0);
assert_eq!(dry.entities_skipped_existing, 1);
assert_eq!(dry.notes_created, 0);
assert_eq!(dry.notes_skipped_existing, 1);
assert_eq!(dry.edges_created, 0);
assert_eq!(dry.edges_skipped_existing, 1);
assert_mapped_batch_remains_tombstoned(&db, &batch).await;
let real = code_ingest_batch(base_args(findings, db.clone()))
.await
.expect("real re-ingest over tombstones must succeed");
assert!(!real.dry_run);
assert_eq!(real.entities_created, 0);
assert_eq!(real.entities_skipped_existing, 1);
assert_eq!(real.notes_created, 0);
assert_eq!(real.notes_skipped_existing, 1);
assert_eq!(real.edges_created, 0);
assert_eq!(real.edges_skipped_existing, 1);
assert_mapped_batch_remains_tombstoned(&db, &batch).await;
}
#[serial]
#[tokio::test]
async fn code_ingest_with_no_queue_writes_does_not_warn_on_drain() {
if crate::test_process::run_in_child() {
return;
}
let _write_queue_env = WriteQueueEnvGuard::unset();
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("no_queue_writes.db");
code_ingest_batch(base_args(findings.clone(), db.clone()))
.await
.expect("initial ingest must succeed");
let capture = Capture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCapture(capture.clone()))
.with_ansi(false)
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let report = code_ingest_batch(base_args(findings, db))
.await
.expect("all-skipped ingest must succeed");
drop(guard);
assert_eq!(report.entities_created, 0);
assert_eq!(report.notes_created, 0);
assert_eq!(report.edges_created, 0);
let log = String::from_utf8_lossy(&capture.0.lock().unwrap()).to_string();
assert!(
!log.contains("writer-task JoinHandle"),
"an ingest with no queue-routed writes must not warn at drain time: {log}"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_dry_run_writes_nothing() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("scratch.db");
let mut args = base_args(findings, db.clone());
args.dry_run = true;
let report = code_ingest_batch(args)
.await
.expect("dry-run must validate successfully");
assert!(report.dry_run);
assert_eq!(
report.notes_created, 1,
"dry-run still reports what would be created"
);
assert_eq!(
finding_note_count(&db).await,
0,
"a dry run must never persist the finding note"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_rejects_invalid_document_before_any_write() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let path = tmp.path().join("bad.json");
std::fs::write(
&path,
r#"{
"audit": {
"date": "2026-07-11",
"scope": "x",
"repo": "r",
"branch": "b",
"commit": "c",
"standards_file": "s"
},
"findings": [
{"id": "F-002", "title": "bad", "severity": "high", "confidence": "low"}
]
}"#,
)
.expect("write invalid fixture");
let db = tmp.path().join("scratch.db");
let err = code_ingest_batch(base_args(path, db.clone()))
.await
.expect_err("missing failure_scenario for a high-severity finding must be rejected");
assert!(
err.to_string().contains("failed validation"),
"error must name the failing document: {err}"
);
assert_eq!(
finding_note_count(&db).await,
0,
"whole-document validation must reject the sweep before any finding note is written"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_dry_run_against_nonexistent_db_creates_no_file() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("does-not-exist.db");
assert!(!db.exists());
let mut args = base_args(findings, db.clone());
args.dry_run = true;
let report = code_ingest_batch(args).await.expect("dry-run must succeed");
assert!(report.dry_run);
assert_eq!(report.entities_created, 1);
assert_eq!(report.notes_created, 1);
assert_eq!(report.edges_created, 1);
assert!(
!db.exists(),
"a dry run against a nonexistent db path must not create it"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_dry_run_against_existing_db_does_not_mutate_it() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("scratch.db");
code_ingest_batch(base_args(findings.clone(), db.clone()))
.await
.expect("initial ingest must succeed");
let bytes_before = std::fs::read(&db).expect("read db bytes before dry run");
let mut args = base_args(findings, db.clone());
args.dry_run = true;
let report = code_ingest_batch(args)
.await
.expect("dry-run against an existing db must succeed");
assert!(report.dry_run);
assert_eq!(
report.entities_skipped_existing, 1,
"the record from the prior real ingest must be reported as already existing"
);
assert_eq!(report.notes_skipped_existing, 1);
assert_eq!(report.edges_skipped_existing, 1);
let bytes_after = std::fs::read(&db).expect("read db bytes after dry run");
assert_eq!(
bytes_before, bytes_after,
"a dry run against an existing db must not change a single byte of it"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_return_implies_settled_file_state() {
if crate::test_process::run_in_child() {
return;
}
let _write_queue_env = WriteQueueEnvGuard::unset();
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("settled.db");
code_ingest_batch(base_args(findings, db.clone()))
.await
.expect("real ingest must succeed");
let wal_path = wal_sidecar_path(&db);
let shm_path = shm_sidecar_path(&db);
assert!(
!wal_path.exists(),
"ingest return must imply the -wal sidecar was checkpointed and removed"
);
assert!(
!shm_path.exists(),
"ingest return must imply the -shm sidecar was removed"
);
let bytes_at_return = std::fs::read(&db).expect("read db at the return boundary");
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let bytes_later = std::fs::read(&db).expect("re-read db after the settle window");
assert_eq!(
bytes_at_return, bytes_later,
"no byte of the database may move after code_ingest_batch has returned"
);
}
fn shm_sidecar_path(db_path: &std::path::Path) -> PathBuf {
let mut name = db_path.as_os_str().to_owned();
name.push("-shm");
PathBuf::from(name)
}
#[serial]
#[tokio::test]
async fn code_ingest_dry_run_against_existing_wal_db_leaves_sidecars_untouched() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("wal_scratch.db");
code_ingest_batch(base_args(findings.clone(), db.clone()))
.await
.expect("initial ingest must succeed");
let pin = StorageBackend::sqlite(&db).expect("open pin backend");
{
let sql = pin.sql();
let mut writer = sql.writer().await.expect("pin writer");
writer
.execute_script(
"CREATE TABLE IF NOT EXISTS wal_pin_probe(x INTEGER); \
INSERT INTO wal_pin_probe VALUES (1);"
.to_string(),
)
.await
.expect("pin write to keep the wal open");
}
let wal_path = wal_sidecar_path(&db);
let shm_path = shm_sidecar_path(&db);
assert!(
wal_path.exists(),
"expected a live -wal sidecar before dry-run"
);
assert!(
shm_path.exists(),
"expected a live -shm sidecar before dry-run"
);
let db_before = std::fs::read(&db).expect("read db before dry run");
let wal_before = std::fs::read(&wal_path).expect("read -wal before dry run");
let shm_before = std::fs::read(&shm_path).expect("read -shm before dry run");
let mut args = base_args(findings, db.clone());
args.dry_run = true;
let report = code_ingest_batch(args)
.await
.expect("dry-run against an existing WAL db must succeed");
assert!(report.dry_run);
assert!(
wal_path.exists(),
"the existing -wal sidecar must not disappear"
);
assert!(
shm_path.exists(),
"the existing -shm sidecar must not disappear"
);
let db_after = std::fs::read(&db).expect("read db after dry run");
let wal_after = std::fs::read(&wal_path).expect("read -wal after dry run");
let shm_after = std::fs::read(&shm_path).expect("read -shm after dry run");
assert_eq!(
db_before, db_after,
"dry-run must not touch the main db file"
);
assert_eq!(
wal_before, wal_after,
"dry-run must not touch the existing -wal sidecar"
);
assert_eq!(
shm_before, shm_after,
"dry-run must not touch the existing -shm sidecar"
);
drop(pin);
}
#[serial]
#[tokio::test]
async fn code_ingest_rejects_secret_bearing_evidence_before_any_write() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let path = tmp.path().join("secret.json");
std::fs::write(
&path,
r#"{
"audit": {
"date": "2026-07-11",
"scope": "khive-pack-code",
"repo": "ohdearquant/khive",
"branch": "feat/adr085-code-ingest-admin",
"commit": "abc1234",
"standards_file": "docs/standards.md"
},
"findings": [
{
"id": "F-003",
"title": "Example finding carrying a leaked credential",
"severity": "high",
"confidence": "high",
"failure_scenario": "A scanner captured a live AWS key in evidence.",
"evidence": "AKIAFAKEKEY1234567890",
"impact": "credential AKIAFAKEKEY1234567890 must never persist verbatim"
}
]
}"#,
)
.expect("write secret-bearing fixture");
let db = tmp.path().join("scratch.db");
let err = code_ingest_batch(base_args(path, db.clone()))
.await
.expect_err("a secret-shaped evidence value must be rejected before any write");
assert!(
err.to_string().to_lowercase().contains("secret"),
"error must name the secret-gate rejection: {err}"
);
assert!(
!db.exists(),
"rejecting a secret-bearing document must leave the db path untouched"
);
}
#[test]
fn preflight_secret_gate_rejects_reserved_key_on_entity_properties() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let mut batch = mapped_batch(&findings);
let mut props = batch.entities[0]
.properties
.clone()
.and_then(|v| v.as_object().cloned())
.unwrap_or_default();
props.insert(
"khive:secret_gate".to_string(),
serde_json::json!("exempted:content-sha256-manifest-v1"),
);
batch.entities[0].properties = Some(serde_json::Value::Object(props));
let err = preflight_secret_gate(&batch)
.expect_err("a caller-supplied reserved key on an entity must be rejected");
assert!(
err.to_string().contains("khive:secret_gate")
&& err.to_string().contains("runtime-owned"),
"error must name the reservation rejection: {err}"
);
}
#[test]
fn preflight_secret_gate_rejects_reserved_key_on_note_properties() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let mut batch = mapped_batch(&findings);
let mut props = batch.notes[0]
.properties
.clone()
.and_then(|v| v.as_object().cloned())
.unwrap_or_default();
props.insert(
"khive:secret_gate".to_string(),
serde_json::json!("exempted:content-sha256-manifest-v1"),
);
batch.notes[0].properties = Some(serde_json::Value::Object(props));
let err = preflight_secret_gate(&batch)
.expect_err("a caller-supplied reserved key on a note must be rejected");
assert!(
err.to_string().contains("khive:secret_gate")
&& err.to_string().contains("runtime-owned"),
"error must name the reservation rejection: {err}"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_fails_loud_when_code_pack_not_configured() {
if crate::test_process::run_in_child() {
return;
}
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("scratch.db");
let prior = std::env::var("KHIVE_PACKS").ok();
unsafe {
std::env::set_var("KHIVE_PACKS", "kg");
}
let result = code_ingest_batch(base_args(findings, db.clone())).await;
unsafe {
match &prior {
Some(v) => std::env::set_var("KHIVE_PACKS", v),
None => std::env::remove_var("KHIVE_PACKS"),
}
}
let err = result.expect_err("a pack set without `code` must be rejected");
assert!(
err.to_string().contains("code"),
"error must name the missing `code` pack: {err}"
);
assert!(
!db.exists(),
"rejecting a misconfigured pack set must leave the db path untouched"
);
}
#[serial]
#[tokio::test]
async fn code_ingest_entity_vector_uses_canonical_body_field_label() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("scratch.db");
code_ingest_batch_with_runtime_setup(base_args(findings, db.clone()), |runtime| {
let model_names = runtime.registered_embedding_model_names();
assert!(
!model_names.is_empty(),
"test requires at least one configured embedding model"
);
for name in model_names {
let dimensions = runtime.resolve_embedding_model(Some(&name))?.dimensions();
runtime.register_embedder(FixedEmbeddingProvider { name, dimensions });
}
Ok(())
})
.await
.expect("ingest must succeed");
let config = write_empty_test_config(tmp.path());
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(db.to_str().expect("utf8 path")),
config: Some(&config),
namespace: Namespace::parse("local").expect("valid namespace"),
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve runtime config");
let runtime = KhiveRuntime::new(cfg).expect("runtime");
let sql = runtime.sql();
let mut reader = sql.reader().await.expect("reader");
let tables = reader
.query_all(SqlStatement {
sql: "SELECT name FROM sqlite_master WHERE type='table' AND name LIKE 'vec_%'"
.to_string(),
params: vec![],
label: None,
})
.await
.expect("list vec tables");
assert!(
!tables.is_empty(),
"expected at least one vector table after ingest"
);
let mut saw_entity_row = false;
for table in &tables {
let table_name = match table.get("name") {
Some(SqlValue::Text(s)) => s.clone(),
other => panic!("unexpected table name column: {other:?}"),
};
let Ok(rows) = reader
.query_all(SqlStatement {
sql: format!("SELECT field FROM {table_name} WHERE kind = 'entity'"),
params: vec![],
label: None,
})
.await
else {
continue;
};
for row in rows {
if let Some(SqlValue::Text(field)) = row.get("field") {
assert_eq!(
field, "entity.body",
"entity vector metadata must use the canonical 'entity.body' field \
label to match khive-runtime/src/operations.rs, got {field:?}"
);
saw_entity_row = true;
}
}
}
assert!(saw_entity_row, "expected at least one entity vector row");
}
#[serial]
#[tokio::test]
async fn code_ingest_reports_actual_embedding_truncation_by_model() {
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("scratch.db");
let mut document: serde_json::Value =
serde_json::from_slice(&std::fs::read(&findings).expect("read findings fixture"))
.expect("parse findings fixture");
document["findings"][0]["impact"] =
serde_json::Value::String("x".repeat(MAX_TEXT_BYTES + 1));
std::fs::write(
&findings,
serde_json::to_vec(&document).expect("serialize long findings fixture"),
)
.expect("write long findings fixture");
let report = code_ingest_batch_with_runtime_setup(base_args(findings, db), |runtime| {
for name in runtime.registered_embedding_model_names() {
let dimensions = runtime.resolve_embedding_model(Some(&name))?.dimensions();
runtime.register_embedder(FixedEmbeddingProvider { name, dimensions });
}
Ok(())
})
.await
.expect("ingest with bounded embedding input must succeed");
assert!(
!report.truncation_by_model.is_empty(),
"configured models that received embeddings must appear in the report"
);
assert!(
report
.truncation_by_model
.values()
.any(|truncation| { truncation.truncated > 0 && truncation.discarded_bytes > 0 }),
"the report must reflect the long finding content actually bounded by an embedder: \
{:?}",
report.truncation_by_model
);
}
fn ok_report() -> CodeIngestReport {
CodeIngestReport {
dry_run: false,
..CodeIngestReport::default()
}
}
#[tokio::test]
async fn settle_writer_drain_join_error_after_success_bails_naming_writer_task() {
let join = tokio::spawn(async {
panic!("synthetic writer task explosion");
});
let err = settle_writer_drain(
Ok(ok_report()),
Some(join),
std::time::Duration::from_secs(5),
)
.await
.expect_err("a panicked writer task must fail the settled-file contract");
let msg = err.to_string();
assert!(
msg.contains("writer task terminated abnormally after ingest completed"),
"the error must name the writer task: {msg}"
);
}
#[tokio::test]
async fn settle_writer_drain_timeout_after_success_bails() {
let join = tokio::spawn(std::future::pending::<()>());
let err = settle_writer_drain(
Ok(ok_report()),
Some(join),
std::time::Duration::from_millis(50),
)
.await
.expect_err("an undrained writer task must fail the settled-file contract");
let msg = err.to_string();
assert!(
msg.contains("did not drain within"),
"the error must name the drain timeout: {msg}"
);
}
#[tokio::test]
async fn settle_writer_drain_failed_ingest_error_survives_join_error() {
let join = tokio::spawn(async {
panic!("synthetic writer task explosion");
});
let primary: Result<CodeIngestReport> = Err(anyhow::anyhow!("primary ingest failure"));
let err = settle_writer_drain(primary, Some(join), std::time::Duration::from_secs(5))
.await
.expect_err("the primary ingest error must surface");
assert_eq!(err.to_string(), "primary ingest failure");
}
#[tokio::test]
async fn settle_writer_drain_failed_ingest_error_survives_timeout() {
let join = tokio::spawn(std::future::pending::<()>());
let primary: Result<CodeIngestReport> = Err(anyhow::anyhow!("primary ingest failure"));
let err = settle_writer_drain(primary, Some(join), std::time::Duration::from_millis(50))
.await
.expect_err("the primary ingest error must surface");
assert_eq!(err.to_string(), "primary ingest failure");
}
#[tokio::test]
async fn settle_writer_drain_clean_passes_report_through() {
let join = tokio::spawn(async {});
let out = settle_writer_drain(
Ok(ok_report()),
Some(join),
std::time::Duration::from_secs(5),
)
.await
.expect("clean drain must pass the report through");
assert!(!out.dry_run);
let out = settle_writer_drain(Ok(ok_report()), None, std::time::Duration::from_secs(5))
.await
.expect("no handle means nothing to await");
assert!(!out.dry_run);
}
#[serial]
#[tokio::test]
async fn take_writer_task_join_or_warn_already_taken_warns_loudly() {
if crate::test_process::run_in_child() {
return;
}
let _write_queue_env = WriteQueueEnvGuard::unset();
let dir = tempfile::TempDir::new().expect("temp dir");
let pool = khive_db::ConnectionPool::new(khive_db::PoolConfig {
path: Some(dir.path().join("already_taken.db")),
..khive_db::PoolConfig::default()
})
.expect("file-backed pool should open");
assert!(
pool.write_queue_active(),
"unset preference on a file-backed pool must resolve the queue ON"
);
pool.writer_task_handle()
.expect("runtime is present")
.expect("queue-ON file-backed pool must spawn");
let taken = pool
.take_writer_task_join()
.expect("the first take must return the handle");
let capture = Capture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCapture(capture.clone()))
.with_ansi(false)
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let second = take_writer_task_join_or_warn(&pool);
drop(guard);
assert!(
second.is_none(),
"the one-shot take must yield None once the handle is gone"
);
let log = String::from_utf8_lossy(&capture.0.lock().unwrap()).to_string();
assert!(
log.contains("JoinHandle was stored but is absent"),
"the already-taken arm must warn loudly about the drain contract gap: {log}"
);
drop(pool);
tokio::time::timeout(std::time::Duration::from_secs(5), taken)
.await
.expect("writer task must exit once every handle clone is dropped")
.expect("writer task must not panic");
}
#[serial]
#[tokio::test]
async fn code_ingest_failed_ingest_still_drains_and_surfaces_primary_error() {
if crate::test_process::run_in_child() {
return;
}
let _write_queue_env = WriteQueueEnvGuard::unset();
let tmp = tempfile::TempDir::new().expect("temp dir");
let findings = write_valid_findings(tmp.path());
let db = tmp.path().join("failed_ingest_drain.db");
let err =
code_ingest_batch_with_runtime_setup(base_args(findings, db.clone()), |runtime| {
let writer = runtime
.backend()
.pool()
.try_writer()
.map_err(|e| anyhow::anyhow!("{e}"))?;
writer
.execute_batch(
"CREATE TRIGGER code_ingest_test_block_notes \
BEFORE INSERT ON notes \
BEGIN SELECT RAISE(ABORT, 'synthetic notes insert failure'); END",
)
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(())
})
.await
.expect_err("ingest into a trigger-blocked notes table must fail");
let msg = err.to_string();
assert!(
msg.contains("synthetic notes insert failure"),
"the PRIMARY ingest error (the notes write) must surface, got: {msg}"
);
assert!(
!msg.contains("writer task did not drain")
&& !msg.contains("writer task terminated abnormally"),
"a drain problem must not mask the primary ingest error: {msg}"
);
assert!(
!wal_sidecar_path(&db).exists(),
"a failed ingest must still drain the writer task and settle the file"
);
assert!(!shm_sidecar_path(&db).exists());
}
async fn finding_note_count(db: &std::path::Path) -> u64 {
let config = write_empty_test_config(
db.parent()
.expect("scratch database must have a parent directory"),
);
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(db.to_str().expect("utf8 path")),
config: Some(&config),
namespace: Namespace::parse("local").expect("valid namespace"),
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve runtime config");
let runtime = KhiveRuntime::new(cfg).expect("runtime");
let token = runtime
.authorize(runtime.config().default_namespace.clone())
.expect("authorize");
runtime
.notes(&token)
.expect("notes store")
.count_notes("local", Some("finding"))
.await
.expect("count notes")
}
}