use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::process::Command;
use anyhow::{bail, Context, Result};
use khive_changeset::{ChangeSet, CreateTarget, Op};
use super::types::{CommitArgs, CommitReport, OutputFormat, RuleResult, Violation};
use super::validate;
pub(super) fn cmd_commit(args: CommitArgs) -> Result<()> {
let changeset_text = std::fs::read_to_string(&args.changeset)
.with_context(|| format!("reading change-set {}", args.changeset.display()))?;
let changeset = khive_changeset::from_ndjson(&changeset_text).with_context(|| {
format!(
"parsing change-set {} as NDJSON-delta",
args.changeset.display()
)
})?;
if !args.rules.exists() {
bail!(
"rules file not found: {} — `kg commit` requires an explicit rules \
file (ADR-102 D2: the tier predicate and its rule pass live in this \
file, not a default)",
args.rules.display()
);
}
let rule_results =
run_commit_time_rules(&changeset, &args.rules).context("running commit-time rule pass")?;
let errors: usize = rule_results
.iter()
.filter(|r| r.severity == "error" && !r.passed)
.count();
let warnings: usize = rule_results
.iter()
.filter(|r| r.severity == "warning" && !r.passed)
.count();
let info: usize = rule_results
.iter()
.filter(|r| r.severity == "info" && !r.passed)
.count();
let passed = errors == 0;
print_report(&args.format, &rule_results, errors, warnings, info, passed);
if !passed {
std::process::exit(1);
}
let repo = args
.repo
.canonicalize()
.with_context(|| format!("resolving repo path {}", args.repo.display()))?;
if !repo.join(".git").exists() {
bail!(
"{} is not a git repository (no .git); run `git init` in the \
staged change-set/snapshot repo first",
repo.display()
);
}
ensure_no_remote(&repo)?;
let rel_path = stage_changeset_file(&repo, &args.changeset)?;
run_git_ok(&repo, &["add", "--", &rel_path])?;
let producer_batch = match &changeset.envelope.batch_id {
Some(batch_id) => batch_id.clone(),
None => format!(
"{}@{}us",
changeset.envelope.producer,
changeset.envelope.staged_at.as_micros()
),
};
let full_message = format!(
"{}\n\nChange-Set-Producer: {}\nChange-Set-Producer-Batch: {}\n",
args.message.trim_end(),
sanitize_trailer_value(&changeset.envelope.producer),
sanitize_trailer_value(&producer_batch),
);
let message_file =
tempfile::NamedTempFile::new().context("creating commit-message temp file")?;
std::fs::write(message_file.path(), &full_message).context("writing commit message")?;
let message_path = message_file
.path()
.to_str()
.context("commit-message temp path is not valid UTF-8")?;
run_git_ok(&repo, &["commit", "-F", message_path])?;
let commit_sha = run_git_ok(&repo, &["rev-parse", "HEAD"])?;
let report = CommitReport {
commit_sha,
changeset_path: rel_path,
ops: changeset.ops.len(),
producer: changeset.envelope.producer.clone(),
producer_batch,
};
println!(
"{}",
serde_json::to_string_pretty(&report).expect("serialize CommitReport")
);
Ok(())
}
struct ProjectedNdjson {
entities: String,
notes: String,
edges: String,
duplicate_ids: Vec<String>,
}
fn project_changeset(changeset: &ChangeSet) -> ProjectedNdjson {
let mut entities = String::new();
let mut notes = String::new();
let mut edges = String::new();
let mut seen_ids: HashMap<String, ()> = HashMap::new();
let mut duplicate_ids = Vec::new();
for op in &changeset.ops {
match op {
Op::Create(create) => {
let id_str = create.id.to_string();
if seen_ids.insert(id_str.clone(), ()).is_some() {
duplicate_ids.push(id_str.clone());
}
match &create.target {
CreateTarget::Entity(fields) => {
let rec = serde_json::json!({
"id": id_str,
"kind": serde_json::to_value(fields.entity_kind)
.expect("EntityKind serializes"),
"entity_type": fields.entity_type,
"name": fields.name,
"description": fields.description,
"properties": serde_json::to_value(&fields.properties)
.expect("properties serialize"),
"tags": fields.tags,
});
entities.push_str(&serde_json::to_string(&rec).expect("json line"));
entities.push('\n');
}
CreateTarget::Note(fields) => {
let rec = serde_json::json!({
"id": id_str,
"kind": fields.note_kind,
"content": fields.content,
"properties": serde_json::to_value(&fields.properties)
.expect("properties serialize"),
"tags": fields.tags,
"salience": fields.salience,
"decay_factor": fields.decay_factor,
});
notes.push_str(&serde_json::to_string(&rec).expect("json line"));
notes.push('\n');
}
}
}
Op::Link(link) => {
let rec = serde_json::json!({
"edge_id": link.id.to_string(),
"source_id": link.source.to_string(),
"target_id": link.target.to_string(),
"relation": serde_json::to_value(link.relation)
.expect("EdgeRelation serializes"),
"weight": link.weight,
"properties": serde_json::to_value(&link.properties)
.expect("properties serialize"),
});
edges.push_str(&serde_json::to_string(&rec).expect("json line"));
edges.push('\n');
}
Op::Update(_) | Op::Delete(_) | Op::Merge(_) => {}
}
}
ProjectedNdjson {
entities,
notes,
edges,
duplicate_ids,
}
}
fn check_no_duplicate_stage_ids(duplicate_ids: &[String]) -> RuleResult {
let violations: Vec<Violation> = duplicate_ids
.iter()
.map(|id| Violation {
entity_id: Some(id.clone()),
entity_name: None,
entity_kind: None,
rule_id: "no-duplicate-uuids".into(),
severity: "error",
message: format!("Duplicate stage-time id within change-set: {id}"),
fixable: false,
})
.collect();
RuleResult {
id: "no-duplicate-uuids".into(),
severity: "error",
passed: violations.is_empty(),
violations,
}
}
fn run_commit_time_rules(changeset: &ChangeSet, rules_path: &Path) -> Result<Vec<RuleResult>> {
let projected = project_changeset(changeset);
let tmp = tempfile::TempDir::new().context("creating projection temp dir")?;
let entities_path = tmp.path().join("entities.ndjson");
let notes_path = tmp.path().join("notes.ndjson");
let edges_path = tmp.path().join("edges.ndjson");
std::fs::write(&entities_path, &projected.entities).context("writing projected entities")?;
std::fs::write(¬es_path, &projected.notes).context("writing projected notes")?;
std::fs::write(&edges_path, &projected.edges).context("writing projected edges")?;
let taxonomy = validate::build_taxonomy().context("building KG taxonomy")?;
let mut results = vec![
check_no_duplicate_stage_ids(&projected.duplicate_ids),
validate::check_valid_entity_kinds(&entities_path, &taxonomy.entity_kinds),
validate::check_valid_note_kinds(¬es_path, &taxonomy.note_kinds),
];
if rules_path.exists() {
let configurable = validate::configurable_rule_checks_partial_view(
&entities_path,
&edges_path,
¬es_path,
rules_path,
)
.context("evaluating configurable rules")?;
results.extend(configurable);
}
Ok(results)
}
fn print_report(
format: &OutputFormat,
rule_results: &[RuleResult],
errors: usize,
warnings: usize,
info: usize,
passed: bool,
) {
match format {
OutputFormat::Json => {
let payload = serde_json::json!({
"rules": rule_results,
"summary": {
"errors": errors,
"warnings": warnings,
"info": info,
"passed": passed,
},
});
println!(
"{}",
serde_json::to_string_pretty(&payload).expect("serialize commit report")
);
}
OutputFormat::Github => {
for r in rule_results {
for v in &r.violations {
let level = if r.severity == "error" {
"error"
} else {
"warning"
};
println!("::{level} ::{}", v.message);
}
}
}
OutputFormat::Text => {
for r in rule_results {
let symbol = if r.passed {
"\u{2713}"
} else if r.severity == "error" {
"\u{2717}"
} else {
"\u{26a0}"
};
println!(" {symbol} {}: {} violation(s)", r.id, r.violations.len());
for v in &r.violations {
println!(" - {}", v.message);
}
}
println!("\nSummary: {errors} error(s), {warnings} warning(s), {info} info");
if passed {
println!("clean pass — proceeding to commit");
} else {
println!("refusing to commit: {errors} error-severity finding(s)");
}
}
}
}
fn ensure_no_remote(repo: &Path) -> Result<()> {
let stdout = run_git_ok(repo, &["remote"])?;
let remotes: Vec<&str> = stdout
.lines()
.map(str::trim)
.filter(|l| !l.is_empty())
.collect();
if !remotes.is_empty() {
bail!(
"refusing to commit: {} has configured git remote(s) [{}] — ADR-102 D6 \
requires the staged change-set/snapshot repository to be local-only; \
remove the remote(s) before running `kg commit`",
repo.display(),
remotes.join(", ")
);
}
Ok(())
}
fn stage_changeset_file(repo: &Path, changeset_src: &Path) -> Result<String> {
let repo_abs = repo
.canonicalize()
.with_context(|| format!("resolving repo path {}", repo.display()))?;
let src_abs = changeset_src
.canonicalize()
.with_context(|| format!("resolving change-set path {}", changeset_src.display()))?;
if let Ok(rel) = src_abs.strip_prefix(&repo_abs) {
return Ok(rel.to_string_lossy().into_owned());
}
let file_name = changeset_src
.file_name()
.context("change-set path has no file name")?;
let dest_dir = repo.join(".khive/kg/changesets");
std::fs::create_dir_all(&dest_dir)
.with_context(|| format!("creating {}", dest_dir.display()))?;
let dest = dest_dir.join(file_name);
std::fs::copy(&src_abs, &dest).with_context(|| {
format!(
"copying change-set {} into {}",
src_abs.display(),
dest.display()
)
})?;
let rel: PathBuf = [".khive", "kg", "changesets"]
.iter()
.collect::<PathBuf>()
.join(file_name);
Ok(rel.to_string_lossy().into_owned())
}
fn run_git_ok(repo: &Path, args: &[&str]) -> Result<String> {
let out = Command::new("git")
.args(args)
.current_dir(repo)
.env("KHIVE_ALLOW_DATA", "1")
.output()
.with_context(|| format!("running git {} in {}", args.join(" "), repo.display()))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
bail!(
"git {} failed in {}: {}",
args.join(" "),
repo.display(),
stderr.trim()
);
}
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
}
fn sanitize_trailer_value(value: &str) -> String {
value.replace(['\n', '\r'], " ").trim().to_string()
}
#[cfg(test)]
mod tests {
use khive_changeset::{
CreateOp, CreateTarget, EntityCreateFields, Envelope, LinkOp, NoteCreateFields,
};
use khive_types::{EdgeRelation, EntityKind, Id128, Namespace, Timestamp};
use tempfile::TempDir;
use super::*;
fn run_git(dir: &Path, args: &[&str]) {
let status = std::process::Command::new("git")
.args(["-c", "core.hooksPath=/dev/null"])
.args(args)
.current_dir(dir)
.status()
.unwrap_or_else(|e| panic!("git {} failed to spawn: {e}", args.join(" ")));
assert!(
status.success(),
"git {} exited with {}",
args.join(" "),
status
);
}
fn init_repo(dir: &Path) {
run_git(dir, &["init", "-b", "main"]);
run_git(dir, &["config", "user.email", "test@example.com"]);
run_git(dir, &["config", "user.name", "Test"]);
run_git(dir, &["config", "core.hooksPath", "/dev/null"]);
run_git(dir, &["commit", "--allow-empty", "-m", "init"]);
}
fn sample_changeset() -> ChangeSet {
let envelope = Envelope::new("agent:test", "family:sonnet", Timestamp::from_secs(1));
let create = Op::Create(CreateOp {
id: Id128::from_u128(1),
namespace: Namespace::local(),
target: CreateTarget::Entity(EntityCreateFields {
entity_kind: EntityKind::Concept,
entity_type: None,
name: "X".into(),
description: None,
properties: Default::default(),
tags: vec![],
}),
});
ChangeSet::new(envelope, vec![create])
}
fn write_changeset(dir: &Path, cs: &ChangeSet) -> PathBuf {
let path = dir.join("changeset.ndjson");
std::fs::write(&path, khive_changeset::to_ndjson(cs).unwrap()).unwrap();
path
}
#[test]
fn project_changeset_emits_entity_and_edge_records() {
let a = Id128::from_u128(1);
let b = Id128::from_u128(2);
let envelope = Envelope::new("agent:test", "family:sonnet", Timestamp::from_secs(1));
let ops = vec![
Op::Create(CreateOp {
id: a,
namespace: Namespace::local(),
target: CreateTarget::Entity(EntityCreateFields {
entity_kind: EntityKind::Concept,
entity_type: None,
name: "A".into(),
description: None,
properties: Default::default(),
tags: vec![],
}),
}),
Op::Create(CreateOp {
id: b,
namespace: Namespace::local(),
target: CreateTarget::Entity(EntityCreateFields {
entity_kind: EntityKind::Concept,
entity_type: None,
name: "B".into(),
description: None,
properties: Default::default(),
tags: vec![],
}),
}),
Op::Link(LinkOp {
id: Id128::from_u128(3),
namespace: Namespace::local(),
source: a,
target: b,
relation: EdgeRelation::Extends,
weight: 1.0,
properties: Default::default(),
}),
];
let cs = ChangeSet::new(envelope, ops);
let projected = project_changeset(&cs);
assert_eq!(projected.entities.lines().count(), 2);
assert_eq!(projected.edges.lines().count(), 1);
assert!(projected.duplicate_ids.is_empty());
assert!(projected.edges.contains("extends"));
}
#[test]
fn project_changeset_flags_duplicate_stage_ids() {
let dup = Id128::from_u128(1);
let envelope = Envelope::new("agent:test", "family:sonnet", Timestamp::from_secs(1));
let mk = || {
Op::Create(CreateOp {
id: dup,
namespace: Namespace::local(),
target: CreateTarget::Entity(EntityCreateFields {
entity_kind: EntityKind::Concept,
entity_type: None,
name: "dup".into(),
description: None,
properties: Default::default(),
tags: vec![],
}),
})
};
let cs = ChangeSet::new(envelope, vec![mk(), mk()]);
let projected = project_changeset(&cs);
assert_eq!(projected.duplicate_ids.len(), 1);
}
#[test]
fn commit_time_rules_reject_invalid_note_kind() {
let envelope = Envelope::new("agent:test", "family:sonnet", Timestamp::from_secs(1));
let ops = vec![Op::Create(CreateOp {
id: Id128::from_u128(1),
namespace: Namespace::local(),
target: CreateTarget::Note(NoteCreateFields {
note_kind: "not_a_real_kind".into(),
content: "hello".into(),
properties: Default::default(),
tags: vec![],
salience: None,
decay_factor: None,
}),
})];
let cs = ChangeSet::new(envelope, ops);
let tmp = TempDir::new().unwrap();
let rules_path = tmp.path().join("rules.toml");
std::fs::write(&rules_path, "").unwrap();
let results = run_commit_time_rules(&cs, &rules_path).unwrap();
let note_rule = results
.iter()
.find(|r| r.id == "valid-note-kinds")
.expect("valid-note-kinds must run");
assert!(!note_rule.passed);
}
#[test]
fn commit_time_rules_exclude_dangling_refs_from_results() {
let envelope = Envelope::new("agent:test", "family:sonnet", Timestamp::from_secs(1));
let ops = vec![Op::Link(LinkOp {
id: Id128::from_u128(3),
namespace: Namespace::local(),
source: Id128::from_u128(100),
target: Id128::from_u128(200),
relation: EdgeRelation::Extends,
weight: 1.0,
properties: Default::default(),
})];
let cs = ChangeSet::new(envelope, ops);
let tmp = TempDir::new().unwrap();
let rules_path = tmp.path().join("rules.toml");
std::fs::write(
&rules_path,
"[dangling_refs]\nenabled = true\nseverity = \"error\"\n",
)
.unwrap();
let results = run_commit_time_rules(&cs, &rules_path).unwrap();
assert!(
!results.iter().any(|r| r.id == "dangling-refs"),
"dangling-refs must be excluded from commit-time results: {results:?}"
);
}
#[test]
fn ensure_no_remote_passes_for_remote_free_repo() {
let tmp = TempDir::new().unwrap();
init_repo(tmp.path());
ensure_no_remote(tmp.path()).unwrap();
}
#[test]
fn ensure_no_remote_refuses_when_remote_configured() {
let tmp = TempDir::new().unwrap();
init_repo(tmp.path());
run_git(
tmp.path(),
&[
"remote",
"add",
"origin",
"https://example.invalid/repo.git",
],
);
let err = ensure_no_remote(tmp.path()).unwrap_err();
assert!(err.to_string().contains("local-only"), "{err}");
}
#[test]
fn stage_changeset_file_copies_external_file_into_repo() {
let repo_tmp = TempDir::new().unwrap();
init_repo(repo_tmp.path());
let outside_tmp = TempDir::new().unwrap();
let cs = sample_changeset();
let src = write_changeset(outside_tmp.path(), &cs);
let rel = stage_changeset_file(repo_tmp.path(), &src).unwrap();
assert_eq!(rel, ".khive/kg/changesets/changeset.ndjson");
assert!(repo_tmp.path().join(&rel).exists());
}
#[test]
fn stage_changeset_file_uses_in_place_path_when_already_inside_repo() {
let repo_tmp = TempDir::new().unwrap();
init_repo(repo_tmp.path());
let cs = sample_changeset();
let src = write_changeset(repo_tmp.path(), &cs);
let rel = stage_changeset_file(repo_tmp.path(), &src).unwrap();
assert_eq!(rel, "changeset.ndjson");
}
#[test]
fn sanitize_trailer_value_collapses_newlines() {
assert_eq!(sanitize_trailer_value("a\nb\r\nc"), "a b c");
}
#[test]
fn cmd_commit_lands_clean_changeset_and_carries_trailers() {
let repo_tmp = TempDir::new().unwrap();
init_repo(repo_tmp.path());
let stage_tmp = TempDir::new().unwrap();
let cs = sample_changeset();
let changeset_path = write_changeset(stage_tmp.path(), &cs);
let rules_path = stage_tmp.path().join("rules.toml");
std::fs::write(&rules_path, "").unwrap();
let args = CommitArgs {
changeset: changeset_path,
rules: rules_path,
repo: repo_tmp.path().to_path_buf(),
message: "stage batch 1".into(),
format: OutputFormat::Json,
};
cmd_commit(args).expect("clean change-set must commit");
let log = run_git_ok(repo_tmp.path(), &["log", "-1", "--pretty=%B"]).unwrap();
assert!(log.contains("stage batch 1"));
assert!(log.contains("Change-Set-Producer: agent:test"));
assert!(log.contains("Change-Set-Producer-Batch: agent:test@1000000us"));
assert!(repo_tmp
.path()
.join(".khive/kg/changesets/changeset.ndjson")
.exists());
}
#[test]
fn cmd_commit_prefers_envelope_batch_id_over_derived_form() {
let repo_tmp = TempDir::new().unwrap();
init_repo(repo_tmp.path());
let stage_tmp = TempDir::new().unwrap();
let mut cs = sample_changeset();
cs.envelope = cs.envelope.with_batch_id("batch-explicit-42");
let changeset_path = write_changeset(stage_tmp.path(), &cs);
let rules_path = stage_tmp.path().join("rules.toml");
std::fs::write(&rules_path, "").unwrap();
let args = CommitArgs {
changeset: changeset_path,
rules: rules_path,
repo: repo_tmp.path().to_path_buf(),
message: "stage batch 2".into(),
format: OutputFormat::Json,
};
cmd_commit(args).expect("clean change-set must commit");
let log = run_git_ok(repo_tmp.path(), &["log", "-1", "--pretty=%B"]).unwrap();
assert!(log.contains("Change-Set-Producer-Batch: batch-explicit-42"));
assert!(!log.contains("agent:test@1000000us"));
}
#[test]
fn cmd_commit_refuses_repo_with_remote() {
let repo_tmp = TempDir::new().unwrap();
init_repo(repo_tmp.path());
run_git(
repo_tmp.path(),
&[
"remote",
"add",
"origin",
"https://example.invalid/repo.git",
],
);
let stage_tmp = TempDir::new().unwrap();
let cs = sample_changeset();
let changeset_path = write_changeset(stage_tmp.path(), &cs);
let rules_path = stage_tmp.path().join("rules.toml");
std::fs::write(&rules_path, "").unwrap();
let args = CommitArgs {
changeset: changeset_path,
rules: rules_path,
repo: repo_tmp.path().to_path_buf(),
message: "should not land".into(),
format: OutputFormat::Json,
};
let err = cmd_commit(args).unwrap_err();
assert!(err.to_string().contains("local-only"), "{err}");
let log = run_git_ok(repo_tmp.path(), &["log", "--oneline"]).unwrap();
assert_eq!(log.lines().count(), 1, "only the init commit must exist");
}
}