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());
if reason == RefusalReason::PolicyRefusal {
entry["domain_disposition"] = json!("not_committed");
}
}
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 atomic_preparation_pack_names(cfg: &RuntimeConfig) -> Vec<String> {
PackRegistry::discovered_names()
.into_iter()
.filter(|name| *name != "telemetry" || cfg.packs.iter().any(|pack| pack == "telemetry"))
.map(str::to_string)
.collect()
}
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)
}
pub(crate) fn preflight_atomic_ops_file(
ops: &[OpsFileEntry],
cfg: &RuntimeConfig,
khive_cfg: &KhiveConfig,
max_ops: usize,
) -> Result<()> {
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));
}
Ok(())
}
fn refusal_reason_for_prepare_error(error: &anyhow::Error) -> Option<RefusalReason> {
match error.downcast_ref::<RuntimeError>() {
Some(RuntimeError::SecretDetected(_)) => Some(RefusalReason::GateRefusal),
Some(error) if error.is_stream_policy_refusal() => Some(RefusalReason::PolicyRefusal),
_ => 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> {
preflight_atomic_ops_file(&ops, &cfg, khive_cfg, max_ops)?;
let pack_names = atomic_preparation_pack_names(&cfg);
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() {
let operation = khive_types::OperationAttribution {
op_index: u32::try_from(op_index).expect("ops-file bounds operation count"),
ref_resolution: khive_types::RefResolution::Literal,
};
match khive_storage::operation_context::scope_operation_attribution(
operation,
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 error_value = match &failure {
AtomicOpFailure::NoteConflict(conflict) => json!(conflict.clone().into_error()),
AtomicOpFailure::EntityConflict(conflict) => json!(conflict.clone().into_error()),
_ => json!(error_message),
};
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_value})
} 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::NoteConflict(conflict) => conflict.clone().into_error().to_string(),
AtomicOpFailure::EntityConflict(conflict) => conflict.clone().into_error().to_string(),
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)?;
let update_policy = registry
.prepare_note_update_policy(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,
update_policy,
registry,
)
.await
.map_err(anyhow::Error::new)?
.1
} 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(),
resurrect: resolved
.get("resurrect")
.and_then(Value::as_bool)
.unwrap_or(false),
};
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)) => {
let mut value = khive_pack_kg::handlers::remap_note_status(
khive_pack_kg::handlers::normalize_entity_timestamps(serde_json::to_value(
¬e,
)?),
);
if p.is_idempotent_noop() {
value["unchanged"] = json!(true);
}
Ok(value)
}
Some(Resolved::Event(_)) | Some(Resolved::PackRecord { .. }) => {
Err(anyhow::anyhow!(
"atomic update result: target {} resolved to an unsupported record",
p.target_id()
))
}
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)?)
}
}
}
("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);
}
}
if let Some(obj) = raw.as_object_mut() {
obj.insert("mutation".to_string(), json!(p.disposition().name()));
}
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,
"note_recorded": false,
"id": task["id"],
"full_id": task["full_id"],
"from": target,
"to": target,
"reason": "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 ignore_dependencies = args
.get("ignore_dependencies")
.and_then(Value::as_bool)
.unwrap_or(false);
let decision = khive_pack_gtd::handlers::prepare_transition(
runtime,
token,
raw_id,
raw_status,
note_arg,
khive_pack_gtd::handlers::DependencyOptions {
ignore_dependencies,
},
)
.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 ignore_dependencies = args
.get("ignore_dependencies")
.and_then(Value::as_bool)
.unwrap_or(false);
let decision = khive_pack_gtd::handlers::prepare_complete(
runtime,
token,
raw_id,
status_arg,
result_arg,
khive_pack_gtd::handlers::DependencyOptions {
ignore_dependencies,
},
)
.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,
version: 2,
};
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,
version: 2,
};
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;
struct TestRuntime {
runtime: KhiveRuntime,
_temp_dir: tempfile::TempDir,
}
impl std::ops::Deref for TestRuntime {
type Target = KhiveRuntime;
fn deref(&self) -> &Self::Target {
&self.runtime
}
}
fn scratch_runtime() -> TestRuntime {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("atomic_apply_gtd.db");
let runtime = KhiveRuntime::new(RuntimeConfig {
db_path: Some(path),
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
})
.expect("runtime");
TestRuntime {
runtime,
_temp_dir: dir,
}
}
#[test]
fn atomic_preparation_keeps_discovered_hooks_without_activating_unrequested_telemetry() {
let runtime = KhiveRuntime::new(RuntimeConfig {
db_path: None,
packs: vec!["kg".into()],
brain_profile: None,
..RuntimeConfig::no_embeddings()
})
.expect("runtime");
assert_eq!(runtime.config().telemetry.default_carrier, None);
let names = atomic_preparation_pack_names(runtime.config());
let discovered = PackRegistry::discovered_names();
assert!(discovered.contains(&"telemetry"));
for name in &discovered {
assert_eq!(
names.iter().any(|selected| selected == name),
*name != "telemetry"
);
}
let registry = full_registry(&runtime);
assert!(registry.all_note_kinds().contains(&"task"));
assert!(registry.all_note_kinds().contains(&"memory"));
assert!(!registry.has_verb("telemetry.emit"));
assert!(registry.has_verb("gtd.transition"));
let (preflight, _runtime) = build_atomic_preflight_registry(runtime.config())
.expect("unrequested telemetry needs no declaration");
assert!(preflight.has_verb("create"));
let mut configured = runtime.config().clone();
configured.packs.push("telemetry".into());
assert!(atomic_preparation_pack_names(&configured)
.iter()
.any(|name| name == "telemetry"));
let error = build_atomic_preflight_registry(&configured)
.err()
.expect("requested telemetry requires its declared default");
assert!(format!("{error:#}").contains("telemetry.default_carrier"));
configured.telemetry.default_carrier = Some(khive_runtime::TelemetryCarrier::Ephemeral);
let (preflight, _runtime) = build_atomic_preflight_registry(&configured)
.expect("requested telemetry has its declared default");
assert!(preflight.has_verb("telemetry.emit"));
}
#[test]
fn telemetry_atomic_exemption_has_no_vocabulary_hooks_or_admissible_verbs() {
use khive_types::Pack;
type Telemetry = khive_pack_telemetry::TelemetryPack;
assert!(Telemetry::ENTITY_KINDS.is_empty());
assert!(Telemetry::NOTE_KINDS.is_empty());
assert!(Telemetry::EDGE_RULES.is_empty());
assert!(Telemetry::ENTITY_TYPES.is_empty());
let implementation = include_str!("../../khive-pack-telemetry/src/pack.rs")
.split_once("impl PackRuntime for TelemetryPack {")
.expect("telemetry runtime implementation")
.1;
let implementation = implementation.split("\n}\n").next().unwrap();
let mut methods = Vec::new();
for line in implementation.lines() {
let method = line
.strip_prefix(" fn ")
.or_else(|| line.strip_prefix(" async fn "));
if let Some(method) = method {
methods.push(method.split('(').next().unwrap());
}
}
methods.sort_unstable();
assert_eq!(
methods,
[
"dispatch",
"entity_kinds",
"handlers",
"name",
"note_kinds",
"requires",
"validate_config"
],
"new runtime hooks require reconsidering the atomic telemetry exemption"
);
let cfg = RuntimeConfig {
db_path: None,
packs: vec!["kg".into(), "telemetry".into()],
telemetry: khive_runtime::TelemetryConfig {
default_carrier: Some(khive_runtime::TelemetryCarrier::Ephemeral),
..khive_runtime::TelemetryConfig::default()
},
..RuntimeConfig::no_embeddings()
};
for handler in Telemetry::HANDLERS {
let ops = vec![OpsFileEntry {
tool: handler.name.into(),
args: json!({}),
}];
assert!(!classify_atomic_preflight(&ops, &cfg).unwrap().is_empty());
}
}
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 = atomic_preparation_pack_names(runtime.config());
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");
}
async fn assert_issue_2675_atomic_derived_property_refused(field: &str, value: Value) {
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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task before refused preparation")
.expect("task exists");
let error = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"content": "body that must not be committed",
"properties": {field: value, "priority": "p0"},
}),
)
.await
.expect_err("atomic preparation must refuse a derived task property");
let message = error.to_string();
assert!(
message.contains(&format!("properties.{field}")),
"refusal must name the supplied diagnostic field; got: {message}"
);
assert!(
message.contains("properties.depends_on"),
"refusal must direct callers to the writable dependency property; got: {message}"
);
let after = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task after refused preparation")
.expect("task exists");
assert_eq!(
after, before,
"refused preparation must preserve the entire task, including body, properties, timestamps and version"
);
}
#[tokio::test]
async fn issue_2675_atomic_task_update_refuses_blocked_by() {
assert_issue_2675_atomic_derived_property_refused(
"blocked_by",
json!(["11111111-1111-4111-8111-111111111111"]),
)
.await;
}
#[tokio::test]
async fn issue_2675_atomic_task_update_refuses_null_blocked_by() {
assert_issue_2675_atomic_derived_property_refused("blocked_by", Value::Null).await;
}
#[tokio::test]
async fn issue_2675_atomic_task_update_refuses_dependency_state() {
assert_issue_2675_atomic_derived_property_refused("dependency_state", json!("blocked"))
.await;
}
#[tokio::test]
async fn issue_2675_atomic_task_update_refuses_null_dependency_state() {
assert_issue_2675_atomic_derived_property_refused("dependency_state", Value::Null).await;
}
#[tokio::test]
async fn issue_2675_atomic_task_update_refuses_actionable() {
assert_issue_2675_atomic_derived_property_refused("actionable", json!(false)).await;
}
#[tokio::test]
async fn issue_2675_atomic_task_update_refuses_null_actionable() {
assert_issue_2675_atomic_derived_property_refused("actionable", Value::Null).await;
}
#[tokio::test]
async fn issue_2675_atomic_task_update_applies_supported_properties() {
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 registry = full_registry(&runtime);
let before = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task before supported preparation")
.expect("task exists");
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
"update",
&json!({
"id": task_id.to_string(),
"content": "supported atomic task body",
"properties": {"planning_label": "reviewed", "priority": "p1"},
}),
)
.await
.expect("prepare legitimate task property update");
let prepared = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read task after supported preparation")
.expect("task exists");
assert_eq!(
prepared, before,
"preparation alone must not persist the patch"
);
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.expect("commit legitimate task property update");
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let after = runtime
.notes(&token)
.expect("note store")
.get_note(task_id)
.await
.expect("read committed task")
.expect("task exists");
assert_eq!(after.content, "supported atomic task body");
assert_eq!(
task_properties(&after)["description"],
"supported atomic task body"
);
assert_eq!(task_properties(&after)["planning_label"], "reviewed");
assert_eq!(task_properties(&after)["priority"], "p1");
assert_eq!(task_properties(&after)["status"], "next");
assert_ne!(
after, before,
"the positive control must actually persist its patch"
);
}
#[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_version_only_change_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)
.unwrap()
.get_note(task_id)
.await
.unwrap()
.unwrap();
let (stale_noop, _) = prepare_one(
&runtime,
&token,
®istry,
"gtd.transition",
&json!({"id": task_id.to_string(), "status": "inbox"}),
)
.await
.unwrap();
let equal_update = AtomicOpPlan::GtdTransition(GtdTransitionPlan::new(
task_id,
vec![PlanStatement {
statement: SqlStatement {
sql: "UPDATE notes SET content=content WHERE id=?1".into(),
params: vec![SqlValue::Text(task_id.to_string())],
label: Some("equal-value-note-update".into()),
},
guard: Some(AffectedRowGuard::exactly(1)),
}],
false,
PostCommitEffect::None,
));
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![equal_update, stale_noop],
)
.await
.unwrap();
assert!(
matches!(
outcome,
AtomicRunOutcome::RolledBack {
failed_op_index: 1,
failure: AtomicOpFailure::GuardFailed { observed: 0, .. },
}
),
"stale version must invalidate the no-op: {outcome:?}"
);
let after = runtime
.notes(&token)
.unwrap()
.get_note(task_id)
.await
.unwrap()
.unwrap();
assert_eq!(after, before);
}
#[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"
);
}
async fn seed_dependency_pair(
runtime: &KhiveRuntime,
token: &NamespaceToken,
blocker_status: &str,
) -> (Uuid, Uuid) {
let store = runtime.notes(token).unwrap();
let blocker_id = seed_task(runtime, token, blocker_status).await;
let task_id = seed_task(runtime, token, "next").await;
let mut task = store.get_note(task_id).await.unwrap().unwrap();
task.properties.as_mut().unwrap()["depends_on"] = json!([blocker_id.to_string()]);
store.upsert_note(task).await.unwrap();
(blocker_id, task_id)
}
async fn dependency_audit_count(runtime: &KhiveRuntime, task_id: Uuid) -> i64 {
ensure_audit_schema(runtime).await;
let mut reader = runtime.sql().reader().await.unwrap();
let rows = reader
.query_all(SqlStatement {
sql: "SELECT COUNT(*) AS count FROM gtd_lifecycle_audit WHERE note_id = ?1".into(),
params: vec![SqlValue::Text(task_id.to_string())],
label: None,
})
.await
.unwrap();
let Some(SqlValue::Integer(count)) = rows[0].get("count") else {
panic!("audit count must be an integer");
};
*count
}
#[tokio::test]
async fn dependency_completion_atomic_prepare_refuses_before_writes() {
for verb in ["gtd.complete", "gtd.transition"] {
let runtime = scratch_runtime();
let token = runtime.authorize(Namespace::local()).unwrap();
let registry = full_registry(&runtime);
let store = runtime.notes(&token).unwrap();
let (blocker_id, task_id) = seed_dependency_pair(&runtime, &token, "next").await;
let before = store.get_note(task_id).await.unwrap().unwrap();
assert_eq!(dependency_audit_count(&runtime, task_id).await, 0);
let error = prepare_one(
&runtime,
&token,
®istry,
verb,
&json!({"id": task_id.to_string(), "status": "done"}),
)
.await
.expect_err("blocked preparation");
let message = error.to_string();
assert!(message.contains(&blocker_id.to_string()));
assert!(message.contains("ignore_dependencies=true"));
let Some(khive_runtime::RuntimeError::Khive(error)) =
error.downcast_ref::<khive_runtime::RuntimeError>()
else {
panic!("structured dependency error must survive atomic prepare: {error:?}");
};
assert_eq!(error.kind(), khive_types::ErrorKind::Conflict);
let details = error.details().unwrap();
assert_eq!(details.get("reason"), Some("dependency_blocked"));
assert_eq!(
details.get("dependency_ids"),
Some(blocker_id.to_string().as_str())
);
assert_eq!(store.get_note(task_id).await.unwrap().unwrap(), before);
assert_eq!(dependency_audit_count(&runtime, task_id).await, 0);
}
}
#[tokio::test]
async fn dependency_completion_atomic_override_is_audited() {
for verb in ["gtd.complete", "gtd.transition"] {
let runtime = scratch_runtime();
let token = runtime.authorize(Namespace::local()).unwrap();
let registry = full_registry(&runtime);
let store = runtime.notes(&token).unwrap();
let (_, task_id) = seed_dependency_pair(&runtime, &token, "next").await;
let mut args =
json!({"id": task_id.to_string(), "status": "done", "ignore_dependencies": "true"});
assert!(prepare_one(&runtime, &token, ®istry, verb, &args)
.await
.is_err());
args["ignore_dependencies"] = json!(true);
let (plan, _) = prepare_one(&runtime, &token, ®istry, verb, &args)
.await
.unwrap();
let outcome =
khive_runtime::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
.await
.unwrap();
let AtomicRunOutcome::Committed { post_commit } = outcome else {
panic!("expected committed override, got {outcome:?}");
};
let audit = apply_gtd_audit_post_commit_effects(&runtime, post_commit.as_slice()).await;
assert_eq!(audit.get(&task_id), Some(&true));
assert_eq!(dependency_audit_count(&runtime, task_id).await, 1);
let mut reader = runtime.sql().reader().await.unwrap();
let rows = reader
.query_all(SqlStatement {
sql: "SELECT from_state, to_state FROM gtd_lifecycle_audit WHERE note_id = ?1"
.into(),
params: vec![SqlValue::Text(task_id.to_string())],
label: None,
})
.await
.unwrap();
assert!(matches!(
rows[0].get("from_state"),
Some(SqlValue::Text(value)) if value == "next"
));
assert!(matches!(
rows[0].get("to_state"),
Some(SqlValue::Text(value)) if value == "done"
));
drop(reader);
let after = store.get_note(task_id).await.unwrap().unwrap();
assert_eq!(task_properties(&after)["status"], "done");
}
}
#[tokio::test]
async fn dependency_completion_atomic_cancellation_and_ready_controls() {
for verb in ["gtd.complete", "gtd.transition"] {
for (blocker_status, target) in [("next", "cancelled"), ("done", "done")] {
let runtime = scratch_runtime();
let token = runtime.authorize(Namespace::local()).unwrap();
let registry = full_registry(&runtime);
let (_, task_id) = seed_dependency_pair(&runtime, &token, blocker_status).await;
let (plan, _) = prepare_one(
&runtime,
&token,
®istry,
verb,
&json!({"id": task_id.to_string(), "status": target}),
)
.await
.unwrap();
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![plan],
)
.await
.unwrap();
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let after = runtime
.notes(&token)
.unwrap()
.get_note(task_id)
.await
.unwrap()
.unwrap();
assert_eq!(task_properties(&after)["status"], target);
}
}
}
#[tokio::test]
async fn dependency_completion_atomic_guard_uses_preparation_time_readiness() {
let runtime = scratch_runtime();
let token = runtime.authorize(Namespace::local()).unwrap();
let registry = full_registry(&runtime);
let (blocker_id, task_id) = seed_dependency_pair(&runtime, &token, "done").await;
let (delete, _) = prepare_one(
&runtime,
&token,
®istry,
"delete",
&json!({"id": blocker_id.to_string(), "hard": true}),
)
.await
.unwrap();
let (complete, _) = prepare_one(
&runtime,
&token,
®istry,
"gtd.complete",
&json!({"id": task_id.to_string()}),
)
.await
.unwrap();
let outcome = khive_runtime::atomic_runner::run_atomic_unit(
runtime.sql().as_ref(),
vec![delete, complete],
)
.await
.unwrap();
assert!(matches!(outcome, AtomicRunOutcome::Committed { .. }));
let store = runtime.notes(&token).unwrap();
assert!(store.get_note(blocker_id).await.unwrap().is_none());
let after = store.get_note(task_id).await.unwrap().unwrap();
assert_eq!(task_properties(&after)["status"], "done");
}
#[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_is_refused_at_prepare() {
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 err = prepare_gtd_transition(&runtime, &token, &args)
.await
.expect_err("a note on a same-status transition must be refused at prepare");
assert!(
err.to_string().contains("already in status"),
"the refusal must name the condition: {err}"
);
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_eq!(
after.version, before.version,
"a refused prepare must not mutate the note revision"
);
assert_eq!(
serde_json::to_value(&after).unwrap(),
serde_json::to_value(&before).unwrap()
);
assert!(task_properties(&after).get("transition_note").is_none());
}
#[tokio::test]
async fn atomic_same_status_transition_without_note_renders_reason_and_note_recorded() {
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"});
let plan = prepare_gtd_transition(&runtime, &token, &args)
.await
.expect("a bare same-status prepare still succeeds");
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_eq!(result["note_recorded"], false);
assert_eq!(result["reason"], "already in target status");
assert!(
result.get("note").is_none(),
"the explanation belongs under `reason`; `note` is the caller's key: {result:?}"
);
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!(
serde_json::to_value(&after).unwrap(),
serde_json::to_value(&before).unwrap(),
"the assertion must leave the row untouched"
);
}
#[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);
}
include!("atomic_project_origin_tests.rs");
}