use anyhow::{Context, Result};
use serde_json::{json, Value};
use uuid::Uuid;
use khive_pack_gtd::handlers::{ensure_audit_schema, write_audit_record};
use khive_pack_gtd::schema::{is_terminal, normalize_status};
use khive_runtime::atomic_plan::{
AffectedRowGuard, GtdCompletePlan, GtdTransitionPlan, PlanStatement, PostCommitEffect,
};
use khive_runtime::atomic_runner::{AtomicOpFailure, AtomicOpPlan, AtomicRunOutcome};
use khive_runtime::pack::{PackRegistry, VerbRegistry, VerbRegistryBuilder};
use khive_runtime::{
EdgeListFilter, KhiveConfig, KhiveRuntime, LinkSpec, NamespaceToken, Resolved, RuntimeConfig,
};
use khive_storage::EdgeRelation;
#[cfg(test)]
use khive_storage::{types::SqlValue, SqlStatement};
use crate::exec::OpsFileEntry;
fn add_post_commit_embedding_warning(
result: &mut Value,
effect: Option<&PostCommitEffect>,
outcomes: &[khive_runtime::atomic_prepare::PostCommitEmbeddingOutcome],
) {
let truncated = effect.is_some_and(|effect| {
outcomes
.iter()
.filter(|outcome| &outcome.effect == effect)
.any(|outcome| outcome.truncation.any_truncated())
});
if !truncated {
return;
}
if let Some(object) = result.as_object_mut() {
object.insert(
"warnings".to_string(),
json!([khive_runtime::retrieval::EMBEDDING_INPUT_TRUNCATED_WARNING]),
);
}
}
pub(crate) async fn execute_atomic_ops_file(
ops: Vec<OpsFileEntry>,
cfg: RuntimeConfig,
khive_cfg: &KhiveConfig,
max_ops: usize,
) -> Result<Value> {
let parsed_for_check: Vec<khive_request::ParsedOp> = ops
.iter()
.map(|op| khive_request::ParsedOp {
tool: op.tool.clone(),
args: std::collections::BTreeMap::new(),
})
.collect();
let rejections = khive_request::atomic::check_atomic_admissible(&parsed_for_check);
if !rejections.is_empty() {
let messages: Vec<String> = rejections.iter().map(|r| r.to_string()).collect();
anyhow::bail!(
"--atomic rejected {} op(s) before any write:\n{}",
messages.len(),
messages.join("\n")
);
}
if ops.len() > max_ops {
anyhow::bail!(
"--atomic op count {} exceeds the configured maximum {max_ops}; \
split the file or raise --atomic-max-ops",
ops.len()
);
}
if !khive_cfg.backends.is_empty() {
anyhow::bail!(
"--atomic does not support a multi-backend [[backends]] topology in v1; \
found {} declared backend(s)",
khive_cfg.backends.len()
);
}
let boot_guard = crate::exec::acquire_local_construction_guard(&cfg)?;
let namespace = cfg.default_namespace.clone();
let runtime = KhiveRuntime::new(cfg).context("build in-process runtime for --atomic")?;
drop(boot_guard);
let token = runtime
.authorize(namespace)
.context("authorize namespace for --atomic")?;
let mut verb_registry_builder = VerbRegistryBuilder::new();
let pack_names: Vec<String> = PackRegistry::discovered_names()
.into_iter()
.map(str::to_string)
.collect();
PackRegistry::register_packs(&pack_names, runtime.clone(), &mut verb_registry_builder)
.map_err(|n| anyhow::anyhow!("pack {n:?} declared in inventory but factory missing"))?;
let verb_registry = verb_registry_builder
.build()
.context("building VerbRegistry for --atomic kind resolution")?;
runtime.install_edge_rules(verb_registry.all_edge_rules());
verb_registry.call_register_note_mutation_hooks(&runtime);
let mut plans: Vec<AtomicOpPlan> = Vec::with_capacity(ops.len());
let mut resolved_args_list: Vec<Value> = Vec::with_capacity(ops.len());
for (op_index, op) in ops.iter().enumerate() {
let (plan, resolved_args) =
prepare_one(&runtime, &token, &verb_registry, &op.tool, &op.args)
.await
.with_context(|| format!("op {op_index} (`{}`) failed to prepare", op.tool))?;
plans.push(plan);
resolved_args_list.push(resolved_args);
}
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), plans.clone())
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let total = ops.len();
let envelope = match outcome {
AtomicRunOutcome::Committed { post_commit } => {
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
let embedding_outcomes =
khive_runtime::atomic_prepare::apply_post_commit_effects_with_report(
&runtime,
&token,
post_commit,
)
.await
.context("post-commit reindex after atomic unit commit")?;
let mut results: Vec<Value> = Vec::with_capacity(ops.len());
for (idx, op) in ops.iter().enumerate() {
let mut result = build_op_result(
&runtime,
&token,
&op.tool,
&op.args,
&resolved_args_list[idx],
&plans[idx],
)
.await
.with_context(|| {
format!(
"op {idx} (`{}`) committed but result rendering failed",
op.tool
)
})?;
let embedding_effect = match &plans[idx] {
AtomicOpPlan::Update(plan) => Some(plan.post_commit()),
_ => None,
};
add_post_commit_embedding_warning(
&mut result,
embedding_effect,
&embedding_outcomes,
);
results
.push(json!({"ok": true, "tool": op.tool, "op_index": idx, "result": result}));
}
json!({
"results": results,
"summary": {"total": total, "succeeded": total, "failed": 0},
"atomic": {
"committed": true,
"rolled_back": false,
"failed_op_index": Value::Null,
"error": Value::Null,
},
})
}
AtomicRunOutcome::RolledBack {
failed_op_index,
failure,
} => {
let error_message = describe_failure(&failure);
let results: Vec<Value> = ops
.iter()
.enumerate()
.map(|(idx, op)| {
if idx == failed_op_index {
json!({"ok": false, "tool": op.tool, "op_index": idx, "error": error_message})
} else {
json!({"ok": false, "tool": op.tool, "op_index": idx, "error": "not applied: whole atomic unit rolled back"})
}
})
.collect();
json!({
"results": results,
"summary": {"total": total, "succeeded": 0, "failed": total},
"atomic": {
"committed": false,
"rolled_back": true,
"failed_op_index": failed_op_index,
"error": error_message,
},
})
}
};
Ok(envelope)
}
async fn apply_gtd_audit_post_commit_effects(runtime: &KhiveRuntime, effects: &[PostCommitEffect]) {
for effect in effects {
if let PostCommitEffect::GtdAudit {
task_id,
from_status,
to_status,
note,
namespace,
} = effect
{
ensure_audit_schema(runtime).await;
write_audit_record(
runtime,
*task_id,
from_status,
to_status,
note.as_deref(),
namespace,
)
.await;
}
}
}
fn describe_failure(failure: &AtomicOpFailure) -> String {
match failure {
AtomicOpFailure::GuardFailed {
statement_label,
expected,
observed,
} => format!(
"guard failed on statement {statement_label:?}: expected {}..{:?} affected rows, observed {observed}",
expected.expected_min, expected.expected_max
),
AtomicOpFailure::SqlError {
statement_label,
message,
} => format!("sql error on statement {statement_label:?}: {message}"),
}
}
fn validate_atomic_args(tool: &str, args: &Value) -> anyhow::Result<()> {
fn reject<T: serde::de::DeserializeOwned>(args: &Value) -> anyhow::Result<()> {
serde_json::from_value::<T>(args.clone())
.map(|_| ())
.map_err(|e| anyhow::anyhow!("bad params: {e}"))
}
match tool {
"update" => reject::<khive_pack_kg::handlers::UpdateParams>(args),
"delete" => reject::<khive_pack_kg::handlers::DeleteParams>(args),
"link" => reject::<khive_pack_kg::handlers::LinkParams>(args),
"gtd.transition" => reject::<khive_pack_gtd::handlers::TransitionParams>(args),
"gtd.complete" => reject::<khive_pack_gtd::handlers::CompleteParams>(args),
_ => Ok(()),
}
}
async fn prepare_one(
runtime: &KhiveRuntime,
token: &NamespaceToken,
registry: &VerbRegistry,
tool: &str,
args: &Value,
) -> anyhow::Result<(AtomicOpPlan, Value)> {
validate_atomic_args(tool, args)?;
match tool {
"gtd.transition" => {
let plan = prepare_gtd_transition(runtime, token, args)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok((plan, args.clone()))
}
"gtd.complete" => {
let plan = prepare_gtd_complete(runtime, token, args)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok((plan, args.clone()))
}
"update" => {
let resolved = resolve_kg_ids_in_args(runtime, token, tool, args).await?;
let expected_kind = update_expected_kind(&resolved, registry)?;
let plan = khive_runtime::atomic_prepare::prepare_update(
runtime,
token,
&resolved,
expected_kind,
)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let id = resolved
.get("id")
.and_then(Value::as_str)
.and_then(|raw| Uuid::parse_str(raw).ok())
.ok_or_else(|| anyhow::anyhow!("resolved update id must be a full UUID"))?;
if let Some(Resolved::Note(note)) = runtime
.resolve_by_id(token, id)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?
{
let properties = resolved.get("properties").filter(|value| !value.is_null());
registry
.validate_note_update_hook(runtime, token, ¬e, properties)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
}
Ok((plan, resolved))
}
"link" => {
let resolved = resolve_kg_ids_in_args(runtime, token, tool, args).await?;
let plan = khive_runtime::atomic_prepare::prepare_op(runtime, token, tool, &resolved)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let source_id = resolved
.get("source_id")
.and_then(Value::as_str)
.and_then(|raw| Uuid::parse_str(raw).ok())
.ok_or_else(|| anyhow::anyhow!("resolved link source_id must be a full UUID"))?;
let target_id = resolved
.get("target_id")
.and_then(Value::as_str)
.and_then(|raw| Uuid::parse_str(raw).ok())
.ok_or_else(|| anyhow::anyhow!("resolved link target_id must be a full UUID"))?;
let relation = resolved
.get("relation")
.and_then(Value::as_str)
.ok_or_else(|| anyhow::anyhow!("resolved link relation must be a string"))?
.parse::<EdgeRelation>()
.map_err(|e| anyhow::anyhow!("{e}"))?;
let spec = LinkSpec {
namespace: Some(token.namespace().as_str().to_owned()),
source_id,
target_id,
relation,
weight: resolved
.get("weight")
.and_then(Value::as_f64)
.unwrap_or(1.0),
metadata: resolved.get("metadata").cloned(),
};
registry
.validate_link_hooks(runtime, token, std::slice::from_ref(&spec))
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok((plan, resolved))
}
"delete" => {
let resolved = resolve_kg_ids_in_args(runtime, token, tool, args).await?;
let expected_kind = delete_expected_kind(&resolved, registry)?;
let plan = khive_runtime::atomic_prepare::prepare_delete(
runtime,
token,
&resolved,
expected_kind,
)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok((plan, resolved))
}
_ => {
let plan = khive_runtime::atomic_prepare::prepare_op(runtime, token, tool, args)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok((plan, args.clone()))
}
}
}
async fn resolve_kg_ids_in_args(
runtime: &KhiveRuntime,
token: &NamespaceToken,
tool: &str,
args: &Value,
) -> anyhow::Result<Value> {
let mut out = args.clone();
let hard = out
.as_object()
.and_then(|o| o.get("hard"))
.and_then(|v| v.as_bool())
.unwrap_or(false);
async fn rewrite(
obj: &mut serde_json::Map<String, Value>,
key: &str,
runtime: &KhiveRuntime,
token: &NamespaceToken,
including_deleted: bool,
) -> anyhow::Result<()> {
let Some(Value::String(raw)) = obj.get(key).cloned() else {
return Ok(());
};
let resolved = if including_deleted {
khive_pack_kg::handlers::resolve_uuid_unfiltered_including_deleted(&raw, runtime, token)
.await
} else {
khive_pack_kg::handlers::resolve_uuid_unfiltered(&raw, runtime, token).await
}
.map_err(|e| anyhow::anyhow!("{e}"))?;
obj.insert(key.to_string(), json!(resolved.to_string()));
Ok(())
}
let obj = out
.as_object_mut()
.ok_or_else(|| anyhow::anyhow!("op args must be a JSON object"))?;
match tool {
"update" => rewrite(obj, "id", runtime, token, false).await?,
"delete" => rewrite(obj, "id", runtime, token, hard).await?,
"link" => {
rewrite(obj, "source_id", runtime, token, false).await?;
rewrite(obj, "target_id", runtime, token, false).await?;
}
_ => {}
}
Ok(out)
}
fn delete_expected_kind(
args: &Value,
registry: &VerbRegistry,
) -> anyhow::Result<Option<khive_runtime::atomic_prepare::AtomicDeleteKind>> {
let raw = match args
.as_object()
.and_then(|o| o.get("kind"))
.and_then(|v| v.as_str())
{
Some(k) => k,
None => return Ok(None),
};
let spec = khive_pack_kg::handlers::resolve_kind_spec(raw, registry)
.map_err(|e| anyhow::anyhow!("{e}"))?;
match spec {
khive_pack_kg::handlers::KindSpec::Entity { specific } => Ok(Some(
khive_runtime::atomic_prepare::AtomicDeleteKind::Entity { specific },
)),
khive_pack_kg::handlers::KindSpec::Note { specific } => Ok(Some(
khive_runtime::atomic_prepare::AtomicDeleteKind::Note { specific },
)),
khive_pack_kg::handlers::KindSpec::Edge => {
Ok(Some(khive_runtime::atomic_prepare::AtomicDeleteKind::Edge))
}
khive_pack_kg::handlers::KindSpec::Event | khive_pack_kg::handlers::KindSpec::Proposal => {
Err(anyhow::anyhow!(
"kind {raw:?} not supported under --atomic delete; only entity/note/edge \
substrates are v1-admissible"
))
}
}
}
fn update_expected_kind(
args: &Value,
registry: &VerbRegistry,
) -> anyhow::Result<Option<khive_runtime::atomic_prepare::AtomicUpdateKind>> {
let raw = match args
.as_object()
.and_then(|o| o.get("kind"))
.and_then(|v| v.as_str())
{
Some(k) => k,
None => return Ok(None),
};
let spec = khive_pack_kg::handlers::resolve_kind_spec(raw, registry)
.map_err(|e| anyhow::anyhow!("{e}"))?;
match spec {
khive_pack_kg::handlers::KindSpec::Entity { specific } => Ok(Some(
khive_runtime::atomic_prepare::AtomicUpdateKind::Entity { specific },
)),
khive_pack_kg::handlers::KindSpec::Note { specific } => Ok(Some(
khive_runtime::atomic_prepare::AtomicUpdateKind::Note { specific },
)),
khive_pack_kg::handlers::KindSpec::Edge => {
Ok(Some(khive_runtime::atomic_prepare::AtomicUpdateKind::Edge))
}
khive_pack_kg::handlers::KindSpec::Event | khive_pack_kg::handlers::KindSpec::Proposal => {
Err(anyhow::anyhow!(
"kind {raw:?} not supported under --atomic update; only entity/note/edge \
substrates are v1-admissible"
))
}
}
}
fn gtd_audit_from_to(effect: &PostCommitEffect) -> Option<(String, String)> {
match effect {
PostCommitEffect::GtdAudit {
from_status,
to_status,
..
} => Some((from_status.clone(), to_status.clone())),
_ => None,
}
}
async fn build_op_result(
runtime: &KhiveRuntime,
token: &NamespaceToken,
tool: &str,
original_args: &Value,
resolved_args: &Value,
plan: &AtomicOpPlan,
) -> anyhow::Result<Value> {
match (tool, plan) {
("update", AtomicOpPlan::Update(p)) if p.edge_natural_key().is_some() => {
let key = p.edge_natural_key().expect("checked by guard above");
let edge = runtime
.get_edge_by_natural_key_including_deleted(
token,
key.namespace(),
key.canon_source_id(),
key.canon_target_id(),
key.relation(),
)
.await?
.ok_or_else(|| {
anyhow::anyhow!(
"atomic update result: committed symmetric edge not found by natural key \
({}, {}, {})",
key.canon_source_id(),
key.canon_target_id(),
key.relation()
)
})?;
Ok(serde_json::to_value(&edge)?)
}
("update", AtomicOpPlan::Update(p)) => match runtime
.resolve_by_id(token, p.target_id())
.await?
{
Some(Resolved::Entity(entity)) => {
Ok(khive_pack_kg::handlers::normalize_entity_timestamps(
serde_json::to_value(&entity)?,
))
}
Some(Resolved::Note(note)) => Ok(khive_pack_kg::handlers::normalize_entity_timestamps(
serde_json::to_value(¬e)?,
)),
None => {
let edge = runtime
.get_edge(token, p.target_id())
.await?
.ok_or_else(|| {
anyhow::anyhow!(
"atomic update result: target {} not found post-commit",
p.target_id()
)
})?;
Ok(serde_json::to_value(&edge)?)
}
_ => anyhow::bail!(
"atomic update result: target {} not found post-commit",
p.target_id()
),
},
("delete", AtomicOpPlan::Delete(_)) => {
let id_val = original_args
.as_object()
.and_then(|o| o.get("id"))
.cloned()
.unwrap_or(Value::Null);
let kind_val = original_args
.as_object()
.and_then(|o| o.get("kind"))
.cloned()
.unwrap_or(Value::Null);
Ok(json!({"deleted": true, "id": id_val, "kind": kind_val}))
}
("link", AtomicOpPlan::Link(p)) => {
let relation_str = resolved_args
.as_object()
.and_then(|o| o.get("relation"))
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("atomic link result: missing relation"))?;
let relation: EdgeRelation = relation_str
.parse()
.map_err(|e| anyhow::anyhow!("atomic link result: unknown relation: {e}"))?;
let edges = runtime
.list_edges(
token,
EdgeListFilter {
source_id: Some(p.source_id()),
target_id: Some(p.target_id()),
relations: vec![relation],
..Default::default()
},
1,
0,
)
.await?;
let edge = edges.into_iter().next().ok_or_else(|| {
anyhow::anyhow!("atomic link result: committed edge not found by natural key")
})?;
let mut raw = serde_json::to_value(&edge)?;
if relation.is_symmetric() {
if let Some(obj) = raw.as_object_mut() {
let orig_source = resolved_args
.as_object()
.and_then(|o| o.get("source_id"))
.cloned()
.unwrap_or(Value::Null);
let orig_target = resolved_args
.as_object()
.and_then(|o| o.get("target_id"))
.cloned()
.unwrap_or(Value::Null);
obj.insert("source_id".to_string(), orig_source);
obj.insert("target_id".to_string(), orig_target);
}
}
Ok(raw)
}
("gtd.transition", AtomicOpPlan::GtdTransition(p)) => {
let note = runtime
.notes(token)?
.get_note(p.task_id())
.await?
.ok_or_else(|| {
anyhow::anyhow!("atomic gtd.transition result: task not found post-commit")
})?;
let task = khive_pack_gtd::handlers::render_task(¬e);
if p.statements().is_empty() {
let raw_status = original_args
.as_object()
.and_then(|o| o.get("status"))
.and_then(|v| v.as_str())
.ok_or_else(|| {
anyhow::anyhow!("atomic gtd.transition result: missing status")
})?;
let target = normalize_status(raw_status);
Ok(json!({
"transitioned": false,
"id": task["id"],
"full_id": task["full_id"],
"from": target,
"to": target,
"note": "already in target status",
}))
} else {
let (from_status, to_status) =
gtd_audit_from_to(p.post_commit()).ok_or_else(|| {
anyhow::anyhow!("atomic gtd.transition result: missing audit effect")
})?;
Ok(json!({
"transitioned": true,
"id": task["id"],
"full_id": task["full_id"],
"from": from_status,
"to": to_status,
"is_terminal": is_terminal(&to_status),
"title": task["title"],
"priority": task["priority"],
"assignee": task["assignee"],
"due": task["due"],
}))
}
}
("gtd.complete", AtomicOpPlan::GtdComplete(p)) => {
let note = runtime
.notes(token)?
.get_note(p.task_id())
.await?
.ok_or_else(|| {
anyhow::anyhow!("atomic gtd.complete result: task not found post-commit")
})?;
let task = khive_pack_gtd::handlers::render_task(¬e);
let (from_status, to_status) = gtd_audit_from_to(p.post_commit()).ok_or_else(|| {
anyhow::anyhow!("atomic gtd.complete result: missing audit effect")
})?;
let completed_at = note
.properties
.as_ref()
.and_then(|props| props.get("completed_at"))
.and_then(|v| v.as_str())
.map(str::to_string)
.ok_or_else(|| {
anyhow::anyhow!("atomic gtd.complete result: missing completed_at")
})?;
Ok(json!({
"completed": true,
"id": task["id"],
"full_id": task["full_id"],
"from": from_status,
"to": to_status,
"completed_at": completed_at,
"is_terminal": is_terminal(&to_status),
}))
}
(other, _) => anyhow::bail!(
"atomic result rendering: no canonical-shape renderer for {other:?} \
(this is a bug — every v1 --atomic-admissible verb must have one)"
),
}
}
fn require_str<'a>(args: &'a Value, key: &str) -> anyhow::Result<&'a str> {
args.as_object()
.and_then(|o| o.get(key))
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("missing required field {key:?}"))
}
async fn prepare_gtd_transition(
runtime: &KhiveRuntime,
token: &NamespaceToken,
args: &Value,
) -> anyhow::Result<AtomicOpPlan> {
let raw_id = require_str(args, "id")?;
let raw_status = require_str(args, "status")?;
let note_arg = args
.as_object()
.and_then(|o| o.get("note"))
.and_then(|v| v.as_str());
let decision =
khive_pack_gtd::handlers::prepare_transition(runtime, token, raw_id, raw_status, note_arg)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
match decision {
khive_pack_gtd::handlers::TransitionDecision::NoOp { note, .. } => {
Ok(AtomicOpPlan::GtdTransition(GtdTransitionPlan::new(
note.id,
vec![],
PostCommitEffect::None,
)))
}
khive_pack_gtd::handlers::TransitionDecision::Write {
note,
current,
target,
props,
updated_at,
transition_note,
} => {
let statement = khive_pack_gtd::handlers::gtd_transition_statement(
note.id, ¤t, &target, &props, updated_at,
)
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(AtomicOpPlan::GtdTransition(GtdTransitionPlan::new(
note.id,
vec![PlanStatement {
statement,
guard: Some(AffectedRowGuard::exactly(1)),
}],
PostCommitEffect::GtdAudit {
task_id: note.id,
from_status: current,
to_status: target,
note: transition_note,
namespace: token.namespace().as_str().to_string(),
},
)))
}
}
}
async fn prepare_gtd_complete(
runtime: &KhiveRuntime,
token: &NamespaceToken,
args: &Value,
) -> anyhow::Result<AtomicOpPlan> {
let raw_id = require_str(args, "id")?;
let status_arg = args
.as_object()
.and_then(|o| o.get("status"))
.and_then(|v| v.as_str());
let result_arg = args
.as_object()
.and_then(|o| o.get("result"))
.and_then(|v| v.as_str());
let decision =
khive_pack_gtd::handlers::prepare_complete(runtime, token, raw_id, status_arg, result_arg)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let statement = khive_pack_gtd::handlers::gtd_transition_statement(
decision.note.id,
&decision.current,
decision.target,
&decision.props,
decision.updated_at,
)
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(AtomicOpPlan::GtdComplete(GtdCompletePlan::new(
decision.note.id,
vec![PlanStatement {
statement,
guard: Some(AffectedRowGuard::exactly(1)),
}],
PostCommitEffect::GtdAudit {
task_id: decision.note.id,
from_status: decision.current,
to_status: decision.target.to_string(),
note: None,
namespace: token.namespace().as_str().to_string(),
},
)))
}
#[cfg(test)]
mod validate_atomic_args_tests {
use super::{add_post_commit_embedding_warning, validate_atomic_args, PostCommitEffect, Uuid};
use serde_json::json;
#[test]
fn atomic_update_result_uses_matching_post_commit_truncation_outcome() {
let note_id = Uuid::new_v4();
let effect = PostCommitEffect::ReindexNote { note_id };
let outcomes = vec![khive_runtime::atomic_prepare::PostCommitEmbeddingOutcome {
effect: effect.clone(),
truncation: khive_runtime::retrieval::EmbeddingTruncationReport {
truncated: 1,
discarded_bytes: 17,
},
}];
let mut result = json!({"id": note_id});
add_post_commit_embedding_warning(&mut result, Some(&effect), &outcomes);
assert_eq!(
result["warnings"],
json!([khive_runtime::retrieval::EMBEDDING_INPUT_TRUNCATED_WARNING])
);
let mut unrelated = json!({"id": Uuid::new_v4()});
let other_effect = PostCommitEffect::ReindexEntity {
entity_id: Uuid::new_v4(),
};
add_post_commit_embedding_warning(&mut unrelated, Some(&other_effect), &outcomes);
assert!(unrelated.get("warnings").is_none());
}
#[test]
fn atomic_duplicate_update_effect_aggregates_later_truncation_outcome() {
let note_id = Uuid::new_v4();
let effect = PostCommitEffect::ReindexNote { note_id };
let outcomes = vec![
khive_runtime::atomic_prepare::PostCommitEmbeddingOutcome {
effect: effect.clone(),
truncation: khive_runtime::retrieval::EmbeddingTruncationReport::default(),
},
khive_runtime::atomic_prepare::PostCommitEmbeddingOutcome {
effect: effect.clone(),
truncation: khive_runtime::retrieval::EmbeddingTruncationReport {
truncated: 1,
discarded_bytes: 23,
},
},
];
let mut first_result = json!({"id": note_id});
let mut second_result = json!({"id": note_id});
add_post_commit_embedding_warning(&mut first_result, Some(&effect), &outcomes);
add_post_commit_embedding_warning(&mut second_result, Some(&effect), &outcomes);
let expected = json!([khive_runtime::retrieval::EMBEDDING_INPUT_TRUNCATED_WARNING]);
assert_eq!(first_result["warnings"], expected);
assert_eq!(second_result["warnings"], expected);
}
#[test]
fn update_rejects_unknown_field() {
let err = validate_atomic_args("update", &json!({"id": "x", "conten": "hello"}))
.expect_err("typo'd `conten` must be rejected");
assert!(err.to_string().contains("unknown field"), "error: {err}");
}
#[test]
fn update_accepts_well_formed_args() {
validate_atomic_args("update", &json!({"id": "x", "content": "hello"}))
.expect("well-formed update args must be accepted");
}
#[test]
fn delete_rejects_unknown_field() {
let err = validate_atomic_args("delete", &json!({"id": "x", "hardd": true}))
.expect_err("typo'd `hardd` must be rejected");
assert!(err.to_string().contains("unknown field"), "error: {err}");
}
#[test]
fn delete_accepts_well_formed_args() {
validate_atomic_args("delete", &json!({"id": "x", "hard": true}))
.expect("well-formed delete args must be accepted");
}
#[test]
fn link_rejects_unknown_field() {
let err = validate_atomic_args(
"link",
&json!({
"source_id": "a",
"target_id": "b",
"relation": "extends",
"targt_backend": "x",
}),
)
.expect_err("typo'd `targt_backend` must be rejected");
assert!(err.to_string().contains("unknown field"), "error: {err}");
}
#[test]
fn link_accepts_well_formed_args() {
validate_atomic_args(
"link",
&json!({"source_id": "a", "target_id": "b", "relation": "extends"}),
)
.expect("well-formed link args must be accepted");
}
#[test]
fn gtd_transition_rejects_unknown_field() {
let err = validate_atomic_args(
"gtd.transition",
&json!({"id": "x", "status": "next", "notee": "typo"}),
)
.expect_err("typo'd `notee` must be rejected");
assert!(err.to_string().contains("unknown field"), "error: {err}");
}
#[test]
fn gtd_transition_accepts_well_formed_args() {
validate_atomic_args(
"gtd.transition",
&json!({"id": "x", "status": "next", "note": "ok"}),
)
.expect("well-formed gtd.transition args must be accepted");
}
#[test]
fn gtd_complete_rejects_unknown_field() {
let err = validate_atomic_args("gtd.complete", &json!({"id": "x", "resutl": "typo"}))
.expect_err("typo'd `resutl` must be rejected");
assert!(err.to_string().contains("unknown field"), "error: {err}");
}
#[test]
fn gtd_complete_accepts_well_formed_args() {
validate_atomic_args("gtd.complete", &json!({"id": "x", "result": "ok"}))
.expect("well-formed gtd.complete args must be accepted");
}
}
#[cfg(test)]
mod tests {
use super::*;
use khive_types::Namespace;
fn scratch_runtime() -> KhiveRuntime {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("atomic_apply_gtd.db");
let rt = KhiveRuntime::new(RuntimeConfig {
db_path: Some(path),
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
})
.expect("runtime");
std::mem::forget(dir);
rt
}
async fn seed_task(runtime: &KhiveRuntime, token: &NamespaceToken, status: &str) -> Uuid {
let mut note = khive_storage::note::Note::new("local", "task", "atomic-gtd-test-task");
note.name = Some("atomic-gtd-test-task".to_string());
note.properties = Some(json!({"status": status, "priority": "p2"}));
let id = note.id;
runtime
.notes(token)
.expect("notes store")
.upsert_note(note)
.await
.expect("seed task");
id
}
fn task_properties(note: &khive_storage::note::Note) -> &Value {
note.properties
.as_ref()
.expect("task must carry properties")
}
#[tokio::test]
async fn atomic_gtd_transition_persists_transition_note() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let task_id = seed_task(&runtime, &token, "inbox").await;
let plan = prepare_gtd_transition(
&runtime,
&token,
&json!({"id": task_id.to_string(), "status": "next", "note": "handed off to reviewer"}),
)
.await
.expect("prepare transition");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit ok");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let note = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must still exist");
let props = task_properties(¬e);
assert_eq!(props.get("status").and_then(|v| v.as_str()), Some("next"));
assert_eq!(
props.get("transition_note").and_then(|v| v.as_str()),
Some("handed off to reviewer"),
"transition_note must be persisted into properties: {props:?}"
);
}
#[tokio::test]
async fn atomic_gtd_transition_rejects_secret_in_note_before_any_write() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let task_id = seed_task(&runtime, &token, "inbox").await;
let err = prepare_gtd_transition(
&runtime,
&token,
&json!({
"id": task_id.to_string(),
"status": "next",
"note": "leaked key AKIAFAKEKEY1234567890",
}),
)
.await
.expect_err("a secret in the transition note must be rejected at prepare");
assert!(
err.to_string().contains("write blocked"),
"expected a secret_gate rejection, got: {err}"
);
let note = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must still exist");
assert_eq!(
task_properties(¬e)
.get("status")
.and_then(|v| v.as_str()),
Some("inbox"),
"rejected prepare must not have mutated the task"
);
}
#[tokio::test]
async fn atomic_gtd_complete_rejects_secret_in_result_and_persists_clean_result() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let task_id = seed_task(&runtime, &token, "next").await;
let err = prepare_gtd_complete(
&runtime,
&token,
&json!({
"id": task_id.to_string(),
"result": "shipped using AKIAFAKEKEY1234567890",
}),
)
.await
.expect_err("a secret in the complete result must be rejected at prepare");
assert!(
err.to_string().contains("write blocked"),
"expected a secret_gate rejection, got: {err}"
);
let note = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must still exist");
assert_eq!(
task_properties(¬e)
.get("status")
.and_then(|v| v.as_str()),
Some("next"),
"rejected prepare must not have mutated the task"
);
let plan = prepare_gtd_complete(
&runtime,
&token,
&json!({"id": task_id.to_string(), "result": "shipped clean"}),
)
.await
.expect("prepare complete");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit ok");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let note = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must still exist");
let props = task_properties(¬e);
assert_eq!(props.get("status").and_then(|v| v.as_str()), Some("done"));
assert_eq!(
props.get("result").and_then(|v| v.as_str()),
Some("shipped clean")
);
}
#[tokio::test]
async fn atomic_gtd_transition_normalizes_status_alias() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let task_id = seed_task(&runtime, &token, "active").await;
let plan = prepare_gtd_transition(
&runtime,
&token,
&json!({"id": task_id.to_string(), "status": "finished"}),
)
.await
.expect("prepare transition with aliased status must succeed");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit ok");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let note = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must still exist");
assert_eq!(
task_properties(¬e)
.get("status")
.and_then(|v| v.as_str()),
Some("done"),
"the \"finished\" alias must normalize to \"done\", parity with canonical"
);
}
#[tokio::test]
async fn atomic_gtd_transition_idempotent_noop_performs_no_write() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let task_id = seed_task(&runtime, &token, "next").await;
let before = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must exist");
let updated_at_before = before.updated_at;
let plan = prepare_gtd_transition(
&runtime,
&token,
&json!({"id": task_id.to_string(), "status": "next"}),
)
.await
.expect("prepare idempotent transition must succeed (no-op, not an error)");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit ok");
let post_commit = match outcome {
AtomicRunOutcome::Committed { post_commit } => post_commit,
other => panic!("idempotent no-op must still succeed as Committed, got {other:?}"),
};
assert!(
post_commit.as_slice().is_empty(),
"an idempotent no-op transition must produce no post-commit effect (no audit row \
either — canonical never reaches its own write_audit_record call): {post_commit:?}"
);
let after = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get_note")
.expect("task must still exist");
assert_eq!(
after.updated_at, updated_at_before,
"an idempotent transition must not touch updated_at — no write happened"
);
}
#[tokio::test]
async fn atomic_gtd_transition_and_complete_write_lifecycle_audit_rows() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let transition_task = seed_task(&runtime, &token, "inbox").await;
let plan = prepare_gtd_transition(
&runtime,
&token,
&json!({"id": transition_task.to_string(), "status": "next", "note": "audit me"}),
)
.await
.expect("prepare transition");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit ok");
let post_commit = match outcome {
AtomicRunOutcome::Committed { post_commit } => post_commit,
other => panic!("expected Committed, got {other:?}"),
};
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
let complete_task = seed_task(&runtime, &token, "next").await;
let plan = prepare_gtd_complete(
&runtime,
&token,
&json!({"id": complete_task.to_string(), "result": "shipped"}),
)
.await
.expect("prepare complete");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit ok");
let post_commit = match outcome {
AtomicRunOutcome::Committed { post_commit } => post_commit,
other => panic!("expected Committed, got {other:?}"),
};
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
let mut reader = runtime.sql().reader().await.expect("reader");
let rows = reader
.query_all(SqlStatement {
sql: "SELECT note_id, from_state, to_state, namespace FROM gtd_lifecycle_audit \
ORDER BY at ASC"
.to_string(),
params: vec![],
label: Some("test-gtd-audit-rows".to_string()),
})
.await
.expect("query gtd_lifecycle_audit");
assert_eq!(
rows.len(),
2,
"both the transition and the complete must each write exactly one audit row: {rows:?}"
);
let transition_task_str = transition_task.to_string();
let complete_task_str = complete_task.to_string();
let transition_row_present = rows.iter().any(|r| {
matches!(r.get("note_id"), Some(SqlValue::Text(id)) if id == &transition_task_str)
&& matches!(r.get("from_state"), Some(SqlValue::Text(s)) if s == "inbox")
&& matches!(r.get("to_state"), Some(SqlValue::Text(s)) if s == "next")
&& matches!(r.get("namespace"), Some(SqlValue::Text(ns)) if ns == "local")
});
assert!(
transition_row_present,
"expected an audit row for the transition op: {rows:?}"
);
let complete_row_present = rows.iter().any(|r| {
matches!(r.get("note_id"), Some(SqlValue::Text(id)) if id == &complete_task_str)
&& matches!(r.get("from_state"), Some(SqlValue::Text(s)) if s == "next")
&& matches!(r.get("to_state"), Some(SqlValue::Text(s)) if s == "done")
&& matches!(r.get("namespace"), Some(SqlValue::Text(ns)) if ns == "local")
});
assert!(
complete_row_present,
"expected an audit row for the complete op: {rows:?}"
);
}
}