use super::{
canonical_edge_endpoints, delete_record_attachments_statement, edge_hard_delete_statement,
edge_replace_if_unchanged_statement, edge_soft_delete_statement,
edge_symmetric_absorb_or_update_inplace_statement, edge_symmetric_delete_if_conflict_statement,
entity_hard_delete_statement, entity_soft_delete_statement, event_append_statements,
hard_delete_lineage_warning_statements, note_hard_delete_statement, note_soft_delete_statement,
obj, optional_f64, optional_properties, optional_str, parse_edge_relation,
purge_incident_edges_statement, push_index_purge_statements, refuse_pack_registry_tags,
reject_inapplicable_update_fields, require_uuid, AffectedRowGuard, AtomicOpPlan,
AttachmentSubstrate, DeletePlan, EdgeNaturalKey, EventKind, KhiveRuntime, NamespaceToken,
PlanStatement, PostCommitEffect, Resolved, RuntimeError, RuntimeResult, SubstrateKind,
UpdatePlan, Uuid, Value,
};
pub(super) async fn prepare_update_edge(
runtime: &KhiveRuntime,
token: &NamespaceToken,
id: Uuid,
mut edge: khive_storage::types::Edge,
args: &Value,
) -> RuntimeResult<AtomicOpPlan> {
reject_inapplicable_update_fields(args, "edge")?;
let expected_updated_at = edge.updated_at;
let expected_deleted_at = edge.deleted_at;
let relation_raw = optional_str(args, "relation");
let weight = optional_f64(args, "weight")?;
let properties = optional_properties(args, "properties")?;
if let Some(ref p) = properties {
crate::secret_gate::check_json_at(p, "edge", "properties")?;
}
crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
let namespace = edge.namespace.clone();
let record_tok = token.with_namespace(
khive_types::Namespace::parse(&namespace)
.map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
);
let mut changed_fields: Vec<&'static str> = Vec::new();
if let Some(raw) = relation_raw {
let relation = parse_edge_relation(raw)?;
runtime
.validate_edge_relation_endpoints(&record_tok, edge.source_id, edge.target_id, relation)
.await?;
edge.relation = relation;
changed_fields.push("relation");
}
if let Some(w) = weight {
if !w.is_finite() || !(0.0..=1.0).contains(&w) {
return Err(RuntimeError::InvalidInput(format!(
"edge weight must be a finite value in [0.0, 1.0]; got {w}"
)));
}
edge.weight = w;
changed_fields.push("weight");
}
if let Some(p) = properties {
edge.metadata = Some(p);
changed_fields.push("properties");
}
let (canon_src, canon_tgt) =
canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
let now = chrono::Utc::now();
let mut statements: Vec<PlanStatement> = Vec::new();
let mut edge_natural_key: Option<EdgeNaturalKey> = None;
if edge.relation.is_symmetric() {
let metadata_str = edge
.metadata
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
let minimum_updated_at_micros = expected_updated_at
.timestamp_micros()
.checked_add(1)
.ok_or_else(|| {
RuntimeError::Internal(format!(
"edge {id} updated_at is already at i64::MAX and cannot advance"
))
})?;
let symmetric_updated_at_micros = now.timestamp_micros().max(minimum_updated_at_micros);
let expected_deleted_at_micros = expected_deleted_at.map(|v| v.timestamp_micros());
statements.push(PlanStatement {
statement: edge_symmetric_delete_if_conflict_statement(
&namespace,
id,
canon_src,
canon_tgt,
edge.relation,
expected_updated_at.timestamp_micros(),
expected_deleted_at_micros,
),
guard: Some(AffectedRowGuard {
expected_min: 0,
expected_max: Some(1),
}),
});
statements.push(PlanStatement {
statement: edge_symmetric_absorb_or_update_inplace_statement(
&namespace,
id,
canon_src,
canon_tgt,
edge.relation,
edge.weight,
symmetric_updated_at_micros,
metadata_str.as_deref(),
edge.target_backend.as_deref(),
expected_updated_at.timestamp_micros(),
expected_deleted_at_micros,
),
guard: Some(AffectedRowGuard::exactly(1)),
});
edge_natural_key = Some(EdgeNaturalKey {
namespace: namespace.clone(),
canon_source_id: canon_src,
canon_target_id: canon_tgt,
relation: edge.relation,
});
} else {
let minimum_updated_at_micros = expected_updated_at
.timestamp_micros()
.checked_add(1)
.ok_or_else(|| {
RuntimeError::Internal(format!(
"edge {id} updated_at is already at i64::MAX and cannot advance"
))
})?;
let now_micros = now.timestamp_micros().max(minimum_updated_at_micros);
edge.updated_at = chrono::DateTime::from_timestamp_micros(now_micros).ok_or_else(|| {
RuntimeError::Internal(format!(
"edge {id}: computed updated_at {now_micros} is not a valid timestamp"
))
})?;
statements.push(PlanStatement {
statement: edge_replace_if_unchanged_statement(
&edge,
expected_updated_at,
expected_deleted_at,
),
guard: Some(AffectedRowGuard::exactly(1)),
});
}
statements.extend(event_append_statements(
token,
&namespace,
"update",
EventKind::EdgeUpdated,
SubstrateKind::Entity,
id,
serde_json::json!({"id": id, "namespace": namespace, "changed_fields": changed_fields}),
)?);
Ok(AtomicOpPlan::Update(Box::new(UpdatePlan {
graph_effects: Vec::new(),
note_vector_purge: None,
note_embedding_inheritance: None,
entity_guard: None,
note_guard: None,
target_id: id,
statements,
post_commit: PostCommitEffect::None,
edge_natural_key,
idempotent_noop: false,
})))
}
pub enum AtomicDeleteKind {
Entity {
specific: Option<String>,
entity_type: Option<String>,
},
Note {
specific: Option<String>,
},
Edge,
}
pub async fn prepare_delete(
runtime: &KhiveRuntime,
token: &NamespaceToken,
args: &Value,
expected_kind: Option<AtomicDeleteKind>,
) -> RuntimeResult<AtomicOpPlan> {
let id = require_uuid(args, "id")?;
let actor = format!("{}:{}", token.actor().kind, token.actor().id);
let hard = obj(args)?
.get("hard")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let resolved = if hard {
runtime.resolve_by_id_including_deleted(token, id).await?
} else {
runtime.resolve_by_id(token, id).await?
};
match resolved {
Some(Resolved::Entity(entity)) => {
match &expected_kind {
None => {}
Some(AtomicDeleteKind::Entity {
specific: Some(expected),
..
}) if &entity.kind != expected => {
return Err(RuntimeError::NotFound(format!("{expected} {id}")));
}
Some(AtomicDeleteKind::Entity { .. }) => {}
Some(AtomicDeleteKind::Note { .. }) => {
return Err(RuntimeError::NotFound(format!("note {id}")));
}
Some(AtomicDeleteKind::Edge) => {
return Err(RuntimeError::NotFound(format!("edge {id}")));
}
}
if let Some(AtomicDeleteKind::Entity {
entity_type: Some(expected),
..
}) = &expected_kind
{
if entity
.entity_type
.as_deref()
.is_some_and(|actual| actual != expected.as_str())
{
return Err(RuntimeError::NotFound(format!("entity {id}")));
}
}
refuse_pack_registry_tags(&entity.tags, "delete")?;
let namespace = entity.namespace.clone();
let mut statements = if hard {
vec![
PlanStatement {
statement: delete_record_attachments_statement(
id,
AttachmentSubstrate::Entity,
),
guard: None,
},
PlanStatement {
statement: entity_hard_delete_statement(id),
guard: Some(AffectedRowGuard::exactly(1)),
},
]
} else {
let deleted_at = chrono::Utc::now().timestamp_micros();
vec![PlanStatement {
statement: entity_soft_delete_statement(id, deleted_at),
guard: Some(AffectedRowGuard::exactly(1)),
}]
};
if hard {
statements.extend(
hard_delete_lineage_warning_statements(
&namespace,
&actor,
id,
SubstrateKind::Entity,
)
.into_iter()
.map(|statement| PlanStatement {
statement,
guard: None,
}),
);
statements.push(PlanStatement {
statement: purge_incident_edges_statement(id),
guard: None,
});
}
push_index_purge_statements(
runtime,
&mut statements,
"fts_entities",
&namespace,
id,
"atomic-delete-entity",
)
.await?;
statements.extend(event_append_statements(
token,
&namespace,
"delete",
EventKind::EntityDeleted,
SubstrateKind::Entity,
id,
serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
)?);
Ok(AtomicOpPlan::Delete(DeletePlan {
target_id: id,
statements,
post_commit: PostCommitEffect::None,
}))
}
Some(Resolved::Note(note)) => {
match &expected_kind {
None => {}
Some(AtomicDeleteKind::Note {
specific: Some(expected),
}) if ¬e.kind != expected => {
return Err(RuntimeError::NotFound(format!("{expected} {id}")));
}
Some(AtomicDeleteKind::Note { .. }) => {}
Some(AtomicDeleteKind::Entity { .. }) => {
return Err(RuntimeError::NotFound(format!("entity {id}")));
}
Some(AtomicDeleteKind::Edge) => {
return Err(RuntimeError::NotFound(format!("edge {id}")));
}
}
if let Some(error) = runtime.stream_member_error(¬e).await? {
return Err(error);
}
let namespace = note.namespace.clone();
let mut statements = if hard {
vec![
PlanStatement {
statement: delete_record_attachments_statement(
id,
AttachmentSubstrate::Note,
),
guard: None,
},
PlanStatement {
statement: note_hard_delete_statement(id),
guard: Some(AffectedRowGuard::exactly(1)),
},
]
} else {
let deleted_at = chrono::Utc::now().timestamp_micros();
vec![PlanStatement {
statement: note_soft_delete_statement(id, deleted_at),
guard: Some(AffectedRowGuard::exactly(1)),
}]
};
if hard {
statements.extend(
hard_delete_lineage_warning_statements(
&namespace,
&actor,
id,
SubstrateKind::Note,
)
.into_iter()
.map(|statement| PlanStatement {
statement,
guard: None,
}),
);
statements.push(PlanStatement {
statement: purge_incident_edges_statement(id),
guard: None,
});
}
push_index_purge_statements(
runtime,
&mut statements,
"fts_notes",
&namespace,
id,
"atomic-delete-note",
)
.await?;
statements.extend(event_append_statements(
token,
&namespace,
"delete",
EventKind::NoteDeleted,
SubstrateKind::Note,
id,
serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
)?);
Ok(AtomicOpPlan::Delete(DeletePlan {
target_id: id,
statements,
post_commit: PostCommitEffect::NoteDeleted {
note_id: id,
kind: note.kind.clone(),
},
}))
}
Some(_) => Err(RuntimeError::InvalidInput(format!(
"delete target {id} must be an entity, note, or edge"
))),
None => match &expected_kind {
Some(AtomicDeleteKind::Entity { .. }) => {
Err(RuntimeError::NotFound(format!("entity/note {id}")))
}
Some(AtomicDeleteKind::Note { .. }) => {
Err(RuntimeError::NotFound(format!("entity/note {id}")))
}
Some(AtomicDeleteKind::Edge) | None => {
let edge = if hard {
runtime.get_edge_including_deleted(token, id).await?
} else {
runtime.get_edge(token, id).await?
};
match edge {
Some(edge) => prepare_delete_edge(token, id, edge, hard, &actor).await,
None => Err(RuntimeError::NotFound(format!("entity/note/edge {id}"))),
}
}
},
}
}
pub(super) async fn prepare_delete_edge(
token: &NamespaceToken,
id: Uuid,
edge: khive_storage::types::Edge,
hard: bool,
actor: &str,
) -> RuntimeResult<AtomicOpPlan> {
let namespace = edge.namespace.clone();
let mut statements: Vec<PlanStatement> = Vec::new();
if hard {
statements.extend(
hard_delete_lineage_warning_statements(&namespace, actor, id, SubstrateKind::Entity)
.into_iter()
.map(|statement| PlanStatement {
statement,
guard: None,
}),
);
statements.push(PlanStatement {
statement: purge_incident_edges_statement(id),
guard: None,
});
statements.push(PlanStatement {
statement: edge_hard_delete_statement(id),
guard: Some(AffectedRowGuard::exactly(1)),
});
} else {
let now = chrono::Utc::now().timestamp_micros();
statements.push(PlanStatement {
statement: edge_soft_delete_statement(id, now),
guard: Some(AffectedRowGuard::exactly(1)),
});
}
statements.extend(event_append_statements(
token,
&namespace,
"delete",
EventKind::EdgeDeleted,
SubstrateKind::Entity,
id,
serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
)?);
Ok(AtomicOpPlan::Delete(DeletePlan {
target_id: id,
statements,
post_commit: PostCommitEffect::None,
}))
}