use std::collections::HashMap;
use anyhow::{Context, Result};
use serde_json::{json, Value};
use uuid::Uuid;
use khive_pack_gtd::handlers::{ensure_audit_schema, write_audit_record_with_status};
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, Namespace, NamespaceToken, Resolved,
RuntimeConfig, RuntimeError,
};
use khive_storage::EdgeRelation;
#[cfg(test)]
use khive_storage::{types::SqlValue, SqlStatement};
use khive_types::RefusalReason;
use crate::exec::OpsFileEntry;
#[derive(Clone, Debug)]
struct AtomicFailureDetail {
op_index: usize,
tool: String,
error: String,
reason: Option<RefusalReason>,
}
#[derive(Debug)]
pub(crate) struct AtomicExecFailure {
message: String,
envelope: Value,
}
impl AtomicExecFailure {
fn new(ops: &[OpsFileEntry], message: String, failures: Vec<AtomicFailureDetail>) -> Self {
let by_index: std::collections::BTreeMap<usize, AtomicFailureDetail> = failures
.into_iter()
.map(|failure| (failure.op_index, failure))
.collect();
let first_failed = by_index.keys().next().copied();
let results = ops
.iter()
.enumerate()
.map(|(op_index, op)| {
if let Some(failure) = by_index.get(&op_index) {
let mut entry = json!({
"ok": false,
"tool": failure.tool.as_str(),
"op_index": op_index,
"error": failure.error.as_str(),
});
if let Some(reason) = failure.reason {
entry["reason"] = json!(reason.as_str());
}
entry
} else {
json!({
"ok": false,
"tool": op.tool.as_str(),
"op_index": op_index,
"error": "not applied: atomic pre-commit validation rejected another operation",
})
}
})
.collect::<Vec<_>>();
let envelope = json!({
"results": results,
"summary": {"total": ops.len(), "succeeded": 0, "failed": ops.len()},
"atomic": {
"committed": false,
"rolled_back": false,
"failed_op_index": first_failed,
"error": message.as_str(),
},
});
Self { message, envelope }
}
pub(crate) fn envelope(&self) -> Value {
self.envelope.clone()
}
}
impl std::fmt::Display for AtomicExecFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message)
}
}
impl std::error::Error for AtomicExecFailure {}
#[derive(Clone, Debug)]
struct AtomicDegradation {
stage: &'static str,
op_index: Option<usize>,
tool: Option<String>,
error: String,
}
impl AtomicDegradation {
fn post_commit_reindex(error: anyhow::Error) -> Self {
Self {
stage: "post_commit_reindex",
op_index: None,
tool: None,
error: format!("post-commit reindex after atomic unit commit: {error:#}"),
}
}
fn result_rendering(op_index: usize, tool: &str, error: anyhow::Error) -> Self {
Self {
stage: "result_rendering",
op_index: Some(op_index),
tool: Some(tool.to_string()),
error: format!(
"op {op_index} (`{tool}`) committed but result rendering failed: {error:#}"
),
}
}
fn save_file_publish(error: &anyhow::Error) -> Self {
Self {
stage: "save_file_publish",
op_index: None,
tool: None,
error: format!("atomic unit committed but save-file publication failed: {error:#}"),
}
}
fn as_json(&self) -> Value {
json!({
"stage": self.stage,
"op_index": self.op_index,
"tool": self.tool.as_deref(),
"error": self.error.as_str(),
})
}
}
pub(crate) fn record_save_file_publish_failure(
envelope: &mut Value,
error: &anyhow::Error,
) -> bool {
let Some(atomic) = envelope.get_mut("atomic").and_then(Value::as_object_mut) else {
return false;
};
if atomic.get("committed").and_then(Value::as_bool) != Some(true) {
return false;
}
let degradation = AtomicDegradation::save_file_publish(error).as_json();
atomic.insert(
"status".to_string(),
Value::String("committed_degraded".to_string()),
);
atomic.insert("retryable".to_string(), Value::Bool(false));
match atomic.get_mut("degradations") {
Some(Value::Array(degradations)) => degradations.push(degradation),
_ => {
atomic.insert("degradations".to_string(), Value::Array(vec![degradation]));
}
}
true
}
fn committed_atomic_block(degradations: &[AtomicDegradation]) -> Value {
let mut atomic = json!({
"committed": true,
"rolled_back": false,
"failed_op_index": Value::Null,
"error": Value::Null,
});
if !degradations.is_empty() {
atomic["status"] = json!("committed_degraded");
atomic["retryable"] = json!(false);
atomic["degradations"] = Value::Array(
degradations
.iter()
.map(AtomicDegradation::as_json)
.collect(),
);
}
atomic
}
fn committed_degraded_result_entry(
op_index: usize,
tool: &str,
degradation: &AtomicDegradation,
) -> Value {
json!({
"ok": true,
"tool": tool,
"op_index": op_index,
"result": Value::Null,
"status": "committed_degraded",
"retryable": false,
"degradation": degradation.as_json(),
})
}
fn atomic_failure_error(
ops: &[OpsFileEntry],
message: String,
failures: Vec<AtomicFailureDetail>,
) -> anyhow::Error {
anyhow::Error::new(AtomicExecFailure::new(ops, message, failures))
}
fn build_atomic_preflight_registry(cfg: &RuntimeConfig) -> Result<(VerbRegistry, KhiveRuntime)> {
let mut metadata_cfg = cfg.clone();
metadata_cfg.db_path = None;
metadata_cfg.embedding_model = None;
metadata_cfg.additional_embedding_models.clear();
metadata_cfg.default_namespace =
Namespace::parse("kkernel-atomic-preflight").unwrap_or_else(|_| Namespace::local());
let pack_names = metadata_cfg.packs.clone();
let runtime =
KhiveRuntime::new(metadata_cfg).context("building --atomic preflight registry runtime")?;
let mut builder = VerbRegistryBuilder::new();
PackRegistry::register_packs(&pack_names, runtime.clone(), &mut builder)
.map_err(|error| anyhow::anyhow!("build --atomic preflight registry: {error}"))?;
let registry = builder
.build()
.context("building --atomic preflight VerbRegistry")?;
Ok((registry, runtime))
}
fn classify_atomic_preflight(
ops: &[OpsFileEntry],
cfg: &RuntimeConfig,
) -> Result<Vec<AtomicFailureDetail>> {
let parsed: Vec<khive_request::ParsedOp> = ops
.iter()
.map(|op| khive_request::ParsedOp {
tool: op.tool.clone(),
args: std::collections::BTreeMap::new(),
})
.collect();
let rejections: std::collections::BTreeMap<usize, khive_request::AtomicRejection> =
khive_request::atomic::check_atomic_admissible(&parsed)
.into_iter()
.map(|rejection| (rejection.op_index, rejection))
.collect();
let (registry, _runtime) = build_atomic_preflight_registry(cfg)?;
let mut failures = Vec::new();
for (op_index, op) in ops.iter().enumerate() {
let Some(rejection) = rejections.get(&op_index) else {
continue;
};
if !registry.has_verb(&op.tool) {
let static_detail = rejections
.get(&op_index)
.map(|rejection| format!("; static policy also reports: {rejection}"))
.unwrap_or_default();
let error = format!(
"op {op_index} (`{}`) cannot run under --atomic: verb is unknown or not loaded{static_detail}",
op.tool
);
failures.push(AtomicFailureDetail {
op_index,
tool: op.tool.clone(),
error,
reason: Some(RefusalReason::VerbRefused),
});
} else {
failures.push(AtomicFailureDetail {
op_index,
tool: op.tool.clone(),
error: rejection.to_string(),
reason: None,
});
}
}
Ok(failures)
}
fn refusal_reason_for_prepare_error(error: &anyhow::Error) -> Option<RefusalReason> {
match error.downcast_ref::<RuntimeError>() {
Some(RuntimeError::SecretDetected(_)) => Some(RefusalReason::GateRefusal),
_ => None,
}
}
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> {
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 preflight_failures = classify_atomic_preflight(&ops, &cfg)?;
if !preflight_failures.is_empty() {
let messages: Vec<&str> = preflight_failures
.iter()
.map(|failure| failure.error.as_str())
.collect();
let message = format!(
"--atomic rejected {} op(s) before any write:\n{}",
messages.len(),
messages.join("\n")
);
return Err(atomic_failure_error(&ops, message, preflight_failures));
}
let pack_names: Vec<String> = PackRegistry::discovered_names()
.into_iter()
.map(str::to_string)
.collect();
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();
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() {
match prepare_one(&runtime, &token, &verb_registry, &op.tool, &op.args).await {
Ok((plan, resolved_args)) => {
plans.push(plan);
resolved_args_list.push(resolved_args);
}
Err(error) => {
let reason = refusal_reason_for_prepare_error(&error);
let error_message =
format!("op {op_index} (`{}`) failed to prepare: {error:#}", op.tool);
let failure = AtomicFailureDetail {
op_index,
tool: op.tool.clone(),
error: error_message.clone(),
reason,
};
return Err(atomic_failure_error(&ops, error_message, vec![failure]));
}
}
}
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 } => {
let gtd_audit_outcomes =
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
let mut degradations = Vec::new();
let embedding_outcomes =
match khive_runtime::atomic_prepare::apply_post_commit_effects_with_report(
&runtime,
&token,
post_commit,
)
.await
{
Ok(outcomes) => outcomes,
Err(error) => {
let degradation =
AtomicDegradation::post_commit_reindex(anyhow::Error::new(error));
tracing::warn!(
error = %degradation.error,
"atomic unit committed but post-commit reindex failed"
);
degradations.push(degradation);
Vec::new()
}
};
let mut results: Vec<Value> = Vec::with_capacity(ops.len());
for (idx, op) in ops.iter().enumerate() {
let rendered = build_op_result(
&runtime,
&token,
&op.tool,
&op.args,
&resolved_args_list[idx],
&plans[idx],
>d_audit_outcomes,
)
.await;
match rendered {
Ok(mut result) => {
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}),
);
}
Err(error) => {
let degradation = AtomicDegradation::result_rendering(idx, &op.tool, error);
tracing::warn!(
error = %degradation.error,
op_index = idx,
tool = %op.tool,
"atomic unit committed but operation result rendering failed"
);
results.push(committed_degraded_result_entry(idx, &op.tool, °radation));
degradations.push(degradation);
}
}
}
json!({
"results": results,
"summary": {"total": total, "succeeded": total, "failed": 0},
"atomic": committed_atomic_block(°radations),
})
}
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],
) -> HashMap<Uuid, bool> {
let mut outcomes = HashMap::new();
for effect in effects {
if let PostCommitEffect::GtdAudit {
task_id,
from_status,
to_status,
note,
namespace,
} = effect
{
ensure_audit_schema(runtime).await;
let persisted = write_audit_record_with_status(
runtime,
*task_id,
from_status,
to_status,
note.as_deref(),
namespace,
)
.await;
outcomes.insert(*task_id, persisted);
}
}
outcomes
}
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?;
Ok((plan, args.clone()))
}
"gtd.complete" => {
let plan = prepare_gtd_complete(runtime, token, args).await?;
Ok((plan, args.clone()))
}
"update" => {
let mut resolved = resolve_kg_ids_in_args(runtime, token, tool, args).await?;
let expected_kind = update_expected_kind(&resolved, registry)?;
if resolved
.get("entity_kind")
.is_some_and(|value| !value.is_null())
{
anyhow::bail!(
"entity_kind is immutable; to change kind, delete then re-create the entity, \
or use merge() if this is a deduplication correction"
);
}
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"))?;
let resolved_target = runtime
.resolve_by_id(token, id)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let plan = if let Some(Resolved::Note(note)) = resolved_target {
khive_runtime::atomic_prepare::validate_note_update_expected_kind(
¬e,
&expected_kind,
)
.map_err(anyhow::Error::new)?;
registry
.prepare_note_update_hook(runtime, token, ¬e, &mut resolved)
.await
.map_err(anyhow::Error::new)?;
khive_runtime::atomic_prepare::prepare_update_from_note_snapshot(
runtime,
token,
&resolved,
expected_kind,
note,
)
.await
.map_err(anyhow::Error::new)?
} else {
khive_runtime::atomic_prepare::prepare_update(
runtime,
token,
&resolved,
expected_kind,
)
.await
.map_err(anyhow::Error::new)?
};
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(anyhow::Error::new)?;
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(anyhow::Error::new)?;
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(anyhow::Error::new)?;
Ok((plan, resolved))
}
_ => {
let plan = khive_runtime::atomic_prepare::prepare_op(runtime, token, tool, args)
.await
.map_err(anyhow::Error::new)?;
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,
gtd_audit_outcomes: &HashMap<Uuid, bool>,
) -> 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.is_idempotent_noop() {
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")
})?;
let audit_persisted =
gtd_audit_outcomes
.get(&p.task_id())
.copied()
.ok_or_else(|| {
anyhow::anyhow!("atomic gtd.transition result: missing audit outcome")
})?;
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"],
"audit_persisted": audit_persisted,
}))
}
}
("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 audit_persisted =
gtd_audit_outcomes
.get(&p.task_id())
.copied()
.ok_or_else(|| {
anyhow::anyhow!("atomic gtd.complete result: missing audit outcome")
})?;
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),
"audit_persisted": audit_persisted,
}))
}
(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(anyhow::Error::new)?;
match decision {
khive_pack_gtd::handlers::TransitionDecision::NoOp { note, current, .. } => {
let statement = khive_pack_gtd::handlers::gtd_noop_assertion_statement(¬e, ¤t)
.map_err(|e| anyhow::anyhow!("{e}"))?;
Ok(AtomicOpPlan::GtdTransition(GtdTransitionPlan::new(
note.id,
vec![PlanStatement {
statement,
guard: Some(AffectedRowGuard::exactly(1)),
}],
true,
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(
¬e, ¤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)),
}],
false,
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(anyhow::Error::new)?;
let statement = khive_pack_gtd::handlers::gtd_transition_statement(
&decision.note,
&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")
}
fn full_registry(runtime: &KhiveRuntime) -> VerbRegistry {
let mut 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 builder)
.expect("register packs");
builder.build().expect("registry")
}
async fn assert_task_update_then_lifecycle_rolls_back(lifecycle_tool: &str) {
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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
let (update, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({"id": task_id.to_string(), "content": "new mirrored body"}),
)
.await
.expect("prepare generic task update");
let lifecycle_args = match lifecycle_tool {
"gtd.transition" => json!({"id": task_id.to_string(), "status": "next"}),
"gtd.complete" => json!({"id": task_id.to_string(), "result": "shipped"}),
other => panic!("unexpected lifecycle tool {other}"),
};
let (lifecycle, _) =
prepare_one(&runtime, &token, ®istry, lifecycle_tool, &lifecycle_args)
.await
.unwrap_or_else(|error| panic!("prepare {lifecycle_tool}: {error}"));
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![update, lifecycle],
)
.await
.expect("atomic runner");
assert!(
matches!(
&outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
failure: AtomicOpFailure::GuardFailed { observed: 0, .. },
}
),
"update -> {lifecycle_tool} must fail closed and roll back; got: {outcome:?}"
);
let persisted = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
assert_eq!(persisted.updated_at, before.updated_at);
assert_eq!(persisted.content, "atomic-gtd-test-task");
assert_eq!(
task_properties(&persisted)
.get("status")
.and_then(Value::as_str),
Some("inbox")
);
assert!(task_properties(&persisted).get("description").is_none());
}
#[tokio::test]
async fn atomic_generic_task_update_synchronizes_content_and_description() {
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 registry = full_registry(&runtime);
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({"id": task_id.to_string(), "content": "atomic body"}),
)
.await
.expect("prepare atomic task content update");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit content update");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let after_content = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get task")
.expect("task exists");
assert_eq!(after_content.content, "atomic body");
assert_eq!(
task_properties(&after_content)
.get("description")
.and_then(Value::as_str),
Some("atomic body")
);
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"properties": {"description": "atomic property body"},
}),
)
.await
.expect("prepare atomic task description update");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit description update");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let after_description = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get task")
.expect("task exists");
assert_eq!(after_description.content, "atomic property body");
assert_eq!(
task_properties(&after_description)
.get("description")
.and_then(Value::as_str),
Some("atomic property body")
);
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({"id": task_id.to_string(), "content": null, "properties": null}),
)
.await
.expect("prepare atomic null/no-op patch");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit null/no-op patch");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let after_null = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get task")
.expect("task exists");
assert_eq!(after_null.content, "atomic property body");
assert_eq!(
task_properties(&after_null)
.get("description")
.and_then(Value::as_str),
Some("atomic property body")
);
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"properties": {"description": null},
}),
)
.await
.expect("prepare atomic description clear");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit description clear");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let after_clear = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("get task")
.expect("task exists");
assert_eq!(after_clear.content, "atomic-gtd-test-task");
assert!(task_properties(&after_clear)
.get("description")
.is_some_and(Value::is_null));
}
#[tokio::test]
async fn atomic_task_update_rejects_title_clear_before_description_clear_can_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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task before rejected update")
.expect("task exists");
let err = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"name": null,
"properties": {"description": null},
}),
)
.await
.expect_err("a task title cannot be cleared under --atomic");
assert!(
err.to_string().contains("task title cannot be cleared"),
"title-clear error must identify the task invariant; got: {err}"
);
let after = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task after rejected update")
.expect("task exists");
assert_eq!(after, before, "rejected preparation must not write");
}
#[tokio::test]
async fn atomic_task_update_rejects_lifecycle_properties_during_prepare() {
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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task before rejected update")
.expect("task exists");
let err = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"properties": {"status": "done"},
}),
)
.await
.expect_err("atomic prepare must share lifecycle-owned property rejection");
assert_eq!(
err.to_string(),
"invalid input: properties.status is lifecycle-owned and cannot be patched on a task; use gtd.transition for lifecycle changes or gtd.complete for terminal completion"
);
let after = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task after rejected update")
.expect("task exists");
assert_eq!(after, before, "rejected atomic prepare must not write");
}
#[tokio::test]
async fn atomic_task_update_checks_explicit_kind_before_running_task_hook() {
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 registry = full_registry(&runtime);
let err = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"kind": "observation",
"content": "one body",
"properties": {"description": "another body"},
}),
)
.await
.expect_err("wrong explicit note kind must fail before task normalization");
let message = err.to_string();
assert!(message.contains("not found: note"), "got: {message}");
assert!(
!message.contains("must match"),
"task hook must not run before kind mismatch rejection: {message}"
);
}
#[tokio::test]
async fn atomic_task_update_guard_refuses_snapshot_changed_after_prepare() {
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 registry = full_registry(&runtime);
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({"id": task_id.to_string(), "content": "prepared body"}),
)
.await
.expect("prepare task update");
let mut concurrent = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
concurrent.content = "concurrent body".to_string();
concurrent.properties =
Some(json!({"status": "inbox", "priority": "p2", "description": "concurrent body"}));
concurrent.updated_at = concurrent.updated_at.saturating_add(10);
runtime
.notes(&token)
.expect("note store")
.upsert_note(concurrent)
.await
.expect("concurrent write");
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("atomic runner");
assert!(
matches!(
&outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 0,
failure: AtomicOpFailure::GuardFailed { observed: 0, .. },
}
),
"stale prepared update must roll back; got: {outcome:?}"
);
let persisted = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
assert_eq!(persisted.content, "concurrent body");
assert_eq!(
task_properties(&persisted)
.get("description")
.and_then(Value::as_str),
Some("concurrent body")
);
}
#[tokio::test]
async fn repeated_atomic_task_updates_fail_closed_without_projected_state() {
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 registry = full_registry(&runtime);
let (first, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({"id": task_id.to_string(), "content": "first body"}),
)
.await
.expect("prepare first update");
let (second, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"properties": {"description": "second body"},
}),
)
.await
.expect("prepare second update");
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![first, second],
)
.await
.expect("atomic runner");
assert!(
matches!(
&outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
failure: AtomicOpFailure::GuardFailed { observed: 0, .. },
}
),
"repeated target must fail closed; got: {outcome:?}"
);
let persisted = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
assert_eq!(persisted.content, "atomic-gtd-test-task");
assert!(task_properties(&persisted).get("description").is_none());
}
#[tokio::test]
async fn atomic_generic_update_then_transition_rolls_back_on_shared_snapshot() {
assert_task_update_then_lifecycle_rolls_back("gtd.transition").await;
}
#[tokio::test]
async fn atomic_generic_update_then_complete_rolls_back_on_shared_snapshot() {
assert_task_update_then_lifecycle_rolls_back("gtd.complete").await;
}
#[tokio::test]
async fn atomic_transition_then_stale_noop_rolls_back_whole_unit() {
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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
let (transition, _) = prepare_one(
&runtime,
&token,
®istry,
"gtd.transition",
&json!({"id": task_id.to_string(), "status": "next"}),
)
.await
.expect("prepare real transition");
let (stale_noop, _) = prepare_one(
&runtime,
&token,
®istry,
"gtd.transition",
&json!({"id": task_id.to_string(), "status": "inbox"}),
)
.await
.expect("prepare same-status transition from the original snapshot");
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![transition, stale_noop],
)
.await
.expect("atomic runner");
assert!(
matches!(
&outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
failure: AtomicOpFailure::GuardFailed {
statement_label: Some(label),
observed: 0,
..
},
} if label == "gtd_atomic_noop_assertion"
),
"the stale no-op assertion must fail at op 1; got: {outcome:?}"
);
let persisted = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists after rollback");
assert_eq!(persisted.updated_at, before.updated_at);
assert_eq!(task_properties(&persisted)["status"], "inbox");
}
#[tokio::test]
async fn atomic_delete_then_stale_noop_rolls_back_whole_unit() {
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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
let (delete, _) = prepare_one(
&runtime,
&token,
®istry,
"delete",
&json!({"id": task_id.to_string(), "hard": true}),
)
.await
.expect("prepare hard delete");
let (stale_noop, _) = prepare_one(
&runtime,
&token,
®istry,
"gtd.transition",
&json!({"id": task_id.to_string(), "status": "inbox"}),
)
.await
.expect("prepare same-status transition before delete applies");
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![delete, stale_noop],
)
.await
.expect("atomic runner");
assert!(
matches!(
&outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
failure: AtomicOpFailure::GuardFailed {
statement_label: Some(label),
observed: 0,
..
},
} if label == "gtd_atomic_noop_assertion"
),
"the delete must invalidate the no-op assertion at op 1; got: {outcome:?}"
);
let persisted = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task")
.expect("hard delete must roll back with the unit");
assert_eq!(persisted, before);
}
#[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, "inbox").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("inbox"),
"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_mutation() {
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 AtomicOpPlan::GtdTransition(noop_plan) = &plan else {
panic!("expected gtd transition plan")
};
assert!(noop_plan.is_idempotent_noop());
assert_eq!(noop_plan.statements().len(), 1);
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 no-note no-ops never reach their own audit helper): {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 change updated_at"
);
}
#[tokio::test]
async fn atomic_same_status_transition_with_note_remains_mutation_and_audit_free() {
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("read task")
.expect("task exists");
let args = json!({
"id": task_id.to_string(),
"status": "next",
"note": "canonical-only note event",
});
let plan = prepare_gtd_transition(&runtime, &token, &args)
.await
.expect("prepare atomic same-status no-op");
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![plan.clone()],
)
.await
.expect("commit guarded no-op assertion");
let post_commit = match outcome {
AtomicRunOutcome::Committed { post_commit } => post_commit,
other => panic!("expected committed no-op, got {other:?}"),
};
assert!(post_commit.as_slice().is_empty());
let audit_outcomes = HashMap::new();
let result = build_op_result(
&runtime,
&token,
"gtd.transition",
&args,
&args,
&plan,
&audit_outcomes,
)
.await
.expect("render no-op result");
assert_eq!(result["transitioned"], false);
assert!(result.get("note_recorded").is_none());
assert!(result.get("audit_persisted").is_none());
let after = runtime
.notes(&token)
.expect("notes store")
.get_note(task_id)
.await
.expect("read task")
.expect("task exists");
assert_eq!(after.updated_at, before.updated_at);
assert!(task_properties(&after).get("transition_note").is_none());
}
#[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:?}"),
};
let transition_audit =
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
assert_eq!(transition_audit.get(&transition_task), Some(&true));
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:?}"),
};
let complete_audit =
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
assert_eq!(complete_audit.get(&complete_task), Some(&true));
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:?}"
);
}
#[tokio::test]
async fn atomic_complete_result_reports_failed_audit_append() {
let runtime = scratch_runtime();
let token = runtime
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
ensure_audit_schema(&runtime).await;
{
let mut writer = runtime.sql().writer().await.expect("writer");
writer
.execute_script(
"CREATE TRIGGER reject_atomic_gtd_audit_insert \
BEFORE INSERT ON gtd_lifecycle_audit \
BEGIN SELECT RAISE(FAIL, 'forced atomic audit failure'); END;"
.to_string(),
)
.await
.expect("failure-injection trigger");
}
let task_id = seed_task(&runtime, &token, "next").await;
let args = json!({"id": task_id.to_string(), "result": "shipped"});
let plan = prepare_gtd_complete(&runtime, &token, &args)
.await
.expect("prepare complete");
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![plan.clone()],
)
.await
.expect("commit domain write");
let post_commit = match outcome {
AtomicRunOutcome::Committed { post_commit } => post_commit,
other => panic!("expected committed completion, got {other:?}"),
};
let audit_outcomes =
apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
assert_eq!(audit_outcomes.get(&task_id), Some(&false));
let result = build_op_result(
&runtime,
&token,
"gtd.complete",
&args,
&args,
&plan,
&audit_outcomes,
)
.await
.expect("render committed completion");
assert_eq!(result["completed"], true);
assert_eq!(result["audit_persisted"], false);
}
}