1use khive_storage::{SqlStatement, SqlValue, SqlWriter, StorageError};
4use khive_types::{Details, KhiveError};
5use serde::{Deserialize, Serialize};
6use uuid::Uuid;
7
8use crate::{KhiveRuntime, NamespaceToken, RuntimeError, RuntimeResult};
9
10#[cfg(test)]
26pub(crate) mod race_seam {
27 use std::sync::Arc;
28 use tokio::sync::Barrier;
29
30 tokio::task_local! {
31 pub(crate) static AFTER_PREPARE_BARRIER: Arc<Barrier>;
32 }
33
34 pub(crate) async fn pause_after_prepare() {
39 if let Ok(barrier) = AFTER_PREPARE_BARRIER.try_with(Arc::clone) {
40 barrier.wait().await;
41 barrier.wait().await;
42 }
43 }
44}
45
46pub const MAX_NOTE_FENCES: usize = 100;
48
49#[derive(Clone, Debug, Serialize)]
50pub struct NoteFence {
51 pub key: String,
52 pub kind: String,
53 pub expected_version: Option<i64>,
55 #[serde(skip_serializing_if = "Option::is_none")]
61 pub live_until: Option<String>,
62 #[serde(skip_serializing_if = "Option::is_none")]
67 pub id: Option<Uuid>,
68}
69
70impl<'de> Deserialize<'de> for NoteFence {
71 fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
72 #[derive(Deserialize)]
73 #[serde(deny_unknown_fields)]
74 struct Fields {
75 key: String,
76 kind: String,
77 #[serde(alias = "version")]
78 expected_version: Option<i64>,
79 live_until: Option<String>,
80 id: Option<Uuid>,
81 }
82
83 let value = serde_json::Value::deserialize(deserializer)?;
84 if value.get("expected_version").is_none() && value.get("version").is_none() {
86 return Err(serde::de::Error::custom(
87 "fence requires expected_version (positive integer or null)",
88 ));
89 }
90 let fields = Fields::deserialize(value).map_err(|error| {
91 serde::de::Error::custom(format!(
92 "fence requires expected_version (positive integer or null): {error}"
93 ))
94 })?;
95 Ok(Self {
96 key: fields.key,
97 kind: fields.kind,
98 expected_version: fields.expected_version,
99 live_until: fields.live_until,
100 id: fields.id,
101 })
102 }
103}
104
105impl NoteFence {
106 pub fn validate(&self) -> RuntimeResult<()> {
107 crate::keyed_memory::validate_memory_key(&self.key)?;
108 if self.kind.is_empty() {
109 return Err(RuntimeError::InvalidInput(
110 "fence requires a note kind".into(),
111 ));
112 }
113 if self.expected_version.is_some_and(|version| version < 1) {
114 return Err(RuntimeError::InvalidInput(
115 "fence requires expected_version (positive integer or null)".into(),
116 ));
117 }
118 if self.live_until.is_some() && self.expected_version.is_none() {
119 return Err(RuntimeError::InvalidInput(
122 "fence live_until requires a positive expected_version".into(),
123 ));
124 }
125 if self.live_until.as_ref().is_some_and(|path| path.is_empty()) {
126 return Err(RuntimeError::InvalidInput(
127 "fence live_until requires a document path".into(),
128 ));
129 }
130 if self.id.is_some() && self.expected_version.is_none() {
131 return Err(RuntimeError::InvalidInput(
134 "fence id requires a positive expected_version".into(),
135 ));
136 }
137 Ok(())
138 }
139}
140
141#[derive(Clone, Debug, Serialize)]
143#[serde(untagged)]
144pub enum NoteFences {
145 One(NoteFence),
146 Many(Vec<NoteFence>),
147}
148
149impl From<NoteFence> for NoteFences {
150 fn from(fence: NoteFence) -> Self {
151 Self::One(fence)
152 }
153}
154
155impl NoteFences {
156 fn validate_count(count: usize) -> RuntimeResult<()> {
157 if count > MAX_NOTE_FENCES {
158 return Err(RuntimeError::InvalidInput(format!(
159 "fence list admits at most {MAX_NOTE_FENCES} entries; this call sent {count}: each fence is a read taken while holding the writer"
160 )));
161 }
162 Ok(())
163 }
164
165 pub fn entries(&self) -> &[NoteFence] {
166 match self {
167 Self::One(fence) => std::slice::from_ref(fence),
168 Self::Many(fences) => fences,
169 }
170 }
171
172 pub fn validate(&self) -> RuntimeResult<()> {
173 Self::validate_count(self.entries().len())?;
174 if self.entries().is_empty() {
175 return Err(RuntimeError::InvalidInput(
176 "fence list must not be empty".into(),
177 ));
178 }
179 let mut seen = std::collections::HashMap::new();
180 for (index, fence) in self.entries().iter().enumerate() {
181 fence.validate()?;
182 if let Some(first) = seen.insert((&fence.kind, &fence.key), index) {
183 return Err(RuntimeError::InvalidInput(format!(
184 "duplicate fence (kind, key) at indices {first} and {index}"
185 )));
186 }
187 }
188 Ok(())
189 }
190}
191
192impl<'de> Deserialize<'de> for NoteFences {
193 fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
194 let value = serde_json::Value::deserialize(deserializer)?;
195 let fences = match value {
196 serde_json::Value::Object(_) => {
197 Self::One(NoteFence::deserialize(value).map_err(serde::de::Error::custom)?)
198 }
199 serde_json::Value::Array(entries) => {
200 Self::validate_count(entries.len()).map_err(serde::de::Error::custom)?;
201 let mut fences = Vec::with_capacity(entries.len());
202 for (index, entry) in entries.into_iter().enumerate() {
203 fences.push(NoteFence::deserialize(entry).map_err(|error| {
204 serde::de::Error::custom(format!("fence entry {index}: {error}"))
205 })?);
206 }
207 Self::Many(fences)
208 }
209 _ => {
210 return Err(serde::de::Error::custom(
211 "fence requires an object or a non-empty list of objects with key, kind and expected_version (positive integer or null)",
212 ));
213 }
214 };
215 fences.validate().map_err(serde::de::Error::custom)?;
216 Ok(fences)
217 }
218}
219
220pub fn deserialize_optional_fences<'de, D: serde::Deserializer<'de>>(
222 deserializer: D,
223) -> Result<Option<NoteFences>, D::Error> {
224 NoteFences::deserialize(deserializer).map(Some)
225}
226
227#[derive(Clone, Debug, Default)]
228pub struct NoteWriteOptions {
229 pub key: Option<String>,
230 pub expected_version: Option<i64>,
231 pub fence: Option<NoteFences>,
232 pub embed: Option<bool>,
233}
234
235impl NoteWriteOptions {
236 pub fn validate(&self) -> RuntimeResult<()> {
237 if let Some(key) = &self.key {
238 crate::keyed_memory::validate_memory_key(key)?;
239 }
240 if self.expected_version.is_some_and(|version| version < 1) {
241 return Err(RuntimeError::InvalidInput(
242 "expected_version must be positive".into(),
243 ));
244 }
245 if let Some(fence) = &self.fence {
246 fence.validate()?;
247 }
248 Ok(())
249 }
250}
251
252#[derive(Clone, Debug)]
258pub(crate) struct CreateKeyClaim {
259 pub kind: String,
260 pub key: String,
261 pub content: String,
262 pub properties: Option<serde_json::Value>,
263 pub replay_signal: bool,
268}
269
270#[derive(Clone, Debug)]
271pub(crate) struct NoteWriteGuard {
272 pub namespace: String,
273 pub target_id: Uuid,
274 pub expected_version: Option<i64>,
275 pub fence: Option<NoteFences>,
276 pub create_key: Option<CreateKeyClaim>,
277}
278
279#[derive(Clone, Debug)]
280pub(crate) struct NoteVectors {
281 namespace: String,
282 subject_id: Uuid,
283}
284
285impl NoteVectors {
286 pub(crate) fn new(namespace: String, subject_id: Uuid) -> Self {
287 Self {
288 namespace,
289 subject_id,
290 }
291 }
292
293 async fn tables(writer: &mut dyn SqlWriter) -> Result<Vec<String>, StorageError> {
294 let tables = writer
297 .query_all(statement(
298 "SELECT name FROM pragma_table_list \
299 WHERE schema='main' AND type='virtual' AND name GLOB 'vec_*' ORDER BY name",
300 vec![],
301 ))
302 .await?;
303 let mut names = Vec::with_capacity(tables.len());
304 for row in tables {
305 let Some(SqlValue::Text(table)) = row.get("name") else {
306 return Err(StorageError::Internal(
307 "invalid vector table catalog row".into(),
308 ));
309 };
310 if !table.strip_prefix("vec_").is_some_and(|key| {
311 !key.is_empty() && key.bytes().all(|c| c.is_ascii_alphanumeric() || c == b'_')
312 }) {
313 return Err(StorageError::Internal(
314 "invalid persisted vector table name".into(),
315 ));
316 }
317 names.push(table.clone());
318 }
319 Ok(names)
320 }
321
322 pub(crate) async fn has_rows(&self, writer: &mut dyn SqlWriter) -> Result<bool, StorageError> {
323 for table in Self::tables(writer).await? {
324 if writer
325 .query_scalar(statement(
326 format!(
327 "SELECT 1 FROM main.{table} WHERE namespace=?1 AND subject_id=?2 LIMIT 1"
328 ),
329 vec![
330 SqlValue::Text(self.namespace.clone()),
331 SqlValue::Text(self.subject_id.to_string()),
332 ],
333 ))
334 .await?
335 .is_some()
336 {
337 return Ok(true);
338 }
339 }
340 Ok(false)
341 }
342
343 pub(crate) async fn apply(&self, writer: &mut dyn SqlWriter) -> Result<(), StorageError> {
344 for table in Self::tables(writer).await? {
345 let scope = vec![
346 SqlValue::Text(self.namespace.clone()),
347 SqlValue::Text(self.subject_id.to_string()),
348 ];
349 writer
350 .execute(statement(
351 format!(
352 "INSERT INTO ann_write_log (namespace,embedding_model,kind,field,subject_id,op) \
353 SELECT namespace,embedding_model,kind,field,subject_id,'delete' \
354 FROM main.{table} WHERE namespace=?1 AND subject_id=?2"),
355 scope.clone(),
356 ))
357 .await?;
358 writer
359 .execute(statement(
360 format!("DELETE FROM main.{table} WHERE namespace=?1 AND subject_id=?2"),
361 scope,
362 ))
363 .await?;
364 let model_key = table.strip_prefix("vec_").ok_or_else(|| {
365 StorageError::Internal("invalid persisted vector table name".into())
366 })?;
367 writer
368 .execute(statement(
369 "DELETE FROM vector_provenance \
370 WHERE model_key=?1 AND namespace=?2 AND subject_id=?3",
371 vec![
372 SqlValue::Text(model_key.to_string()),
373 SqlValue::Text(self.namespace.clone()),
374 SqlValue::Text(self.subject_id.to_string()),
375 ],
376 ))
377 .await?;
378 }
379 Ok(())
380 }
381}
382
383#[derive(Clone, Debug)]
384pub(crate) struct NoteEmbeddingInheritance {
385 pub vectors: NoteVectors,
386 pub kind: String,
387}
388
389#[derive(Clone, Debug, PartialEq, Eq)]
390pub enum NoteWriteConflict {
391 Version {
392 expected: i64,
393 current: i64,
394 },
395 Fence {
396 key: String,
397 expected: Option<i64>,
398 current: Option<i64>,
399 index: Option<usize>,
400 },
401 Key {
402 key: String,
403 existing_id: String,
404 equal: Option<bool>,
415 },
416 FenceDeadline(Box<FenceDeadline>),
420 FenceIdentity(Box<FenceIdentity>),
424}
425
426#[derive(Clone, Debug, PartialEq, Eq)]
429pub struct FenceIdentity {
430 pub key: String,
431 pub kind: String,
432 pub version: i64,
433 pub index: Option<usize>,
434 pub evidence: Vec<(&'static str, String)>,
435}
436
437#[derive(Clone, Debug, PartialEq, Eq)]
440pub struct FenceDeadline {
441 pub key: String,
442 pub kind: String,
443 pub version: i64,
444 pub field: String,
445 pub index: Option<usize>,
446 pub reason: &'static str,
447 pub evidence: Vec<(&'static str, String)>,
448}
449
450impl NoteWriteConflict {
451 pub fn into_error(self) -> KhiveError {
452 self.into_error_at_member(None)
453 }
454
455 pub(crate) fn into_error_at_member(self, member: Option<usize>) -> KhiveError {
456 let (message, mut details) = match self {
457 Self::Version { expected, current } => (
458 "note version precondition failed",
459 vec![
460 ("reason", "version_conflict".into()),
461 ("expected_version", expected.to_string()),
462 ("current_version", current.to_string()),
463 ],
464 ),
465 Self::Fence {
466 key,
467 expected,
468 current,
469 index,
470 } => {
471 let mut fields = vec![
472 ("reason", "fence_conflict".into()),
473 ("key", key),
474 (
475 "expected_version",
476 expected.map_or_else(|| "absent".into(), |version| version.to_string()),
477 ),
478 ];
479 if let Some(current) = current {
480 fields.push(("current_version", current.to_string()));
481 }
482 if let Some(index) = index {
483 fields.push(("index", index.to_string()));
484 }
485 ("note fence precondition failed", fields)
486 }
487 Self::Key {
488 key,
489 existing_id,
490 equal,
491 } => {
492 let mut fields = vec![
493 ("reason", "key_conflict".into()),
494 ("key", key),
495 ("existing_id", existing_id),
496 ];
497 if let Some(equal) = equal {
498 fields.push(("equal", equal.to_string()));
499 }
500 ("a live note already holds this key", fields)
501 }
502 Self::FenceDeadline(deadline) => {
503 let FenceDeadline {
504 key,
505 kind,
506 version,
507 field,
508 index,
509 reason,
510 evidence,
511 } = *deadline;
512 let mut fields = vec![
513 ("reason", reason.into()),
514 ("key", key),
515 ("kind", kind),
516 ("version", version.to_string()),
517 ("field", field),
518 ];
519 if let Some(index) = index {
520 fields.push(("index", index.to_string()));
521 }
522 fields.extend(evidence);
523 ("note fence time precondition failed", fields)
524 }
525 Self::FenceIdentity(identity) => {
526 let FenceIdentity {
527 key,
528 kind,
529 version,
530 index,
531 evidence,
532 } = *identity;
533 let mut fields = vec![
534 ("reason", "identity_conflict".into()),
535 ("key", key),
536 ("kind", kind),
537 ("version", version.to_string()),
538 ];
539 fields.extend(evidence);
540 if let Some(index) = index {
541 fields.push(("index", index.to_string()));
542 }
543 ("note fence identity precondition failed", fields)
544 }
545 };
546 if let Some(member) = member {
547 details.push(("member", member.to_string()));
548 }
549 KhiveError::conflict(message).with_details(Details::new_owned(details))
550 }
551}
552
553pub(crate) fn statement(sql: impl Into<String>, params: Vec<SqlValue>) -> SqlStatement {
554 SqlStatement {
555 sql: sql.into(),
556 params,
557 label: Some("note-write-guard".into()),
558 }
559}
560
561impl NoteWriteGuard {
562 pub(crate) async fn check_fence(
563 &self,
564 writer: &mut dyn SqlWriter,
565 ) -> Result<Option<NoteWriteConflict>, StorageError> {
566 let Some(fences) = &self.fence else {
567 return Ok(None);
568 };
569 let mut now: Option<i64> = None;
573 for (index, fence) in fences.entries().iter().enumerate() {
574 let holder = crate::fence_identity::read_holder(
575 writer,
576 &self.namespace,
577 &fence.kind,
578 &fence.key,
579 "note-write-guard",
580 )
581 .await?;
582 if let (Some(asserted), Some(holder)) = (fence.id, holder.as_ref()) {
588 if asserted != holder.id {
589 return Ok(Some(NoteWriteConflict::FenceIdentity(Box::new(
590 FenceIdentity {
591 key: fence.key.clone(),
592 kind: fence.kind.clone(),
593 version: fence.expected_version.unwrap_or_default(),
596 index: matches!(fences, NoteFences::Many(_)).then_some(index),
597 evidence: crate::fence_identity::identity_evidence(asserted, holder.id),
598 },
599 ))));
600 }
601 }
602 let current = holder.map(|holder| holder.version);
603 if current != fence.expected_version {
604 return Ok(Some(NoteWriteConflict::Fence {
605 key: fence.key.clone(),
606 expected: fence.expected_version,
607 current,
608 index: matches!(fences, NoteFences::Many(_)).then_some(index),
609 }));
610 }
611 if let Some(field) = &fence.live_until {
612 let clock = match now {
613 Some(clock) => clock,
614 None => {
615 let clock =
616 crate::live_until::writer_clock(writer, "note-write-guard-clock")
617 .await?;
618 now = Some(clock);
619 clock
620 }
621 };
622 if let Some(refusal) = crate::live_until::evaluate(
623 writer,
624 &self.namespace,
625 &fence.kind,
626 &fence.key,
627 field,
628 clock,
629 "note-write-guard-live-until",
630 )
631 .await?
632 {
633 return Ok(Some(NoteWriteConflict::FenceDeadline(Box::new(
634 FenceDeadline {
635 key: fence.key.clone(),
636 kind: fence.kind.clone(),
637 version: fence.expected_version.unwrap_or_default(),
640 field: field.clone(),
641 index: matches!(fences, NoteFences::Many(_)).then_some(index),
642 reason: refusal.reason(),
643 evidence: refusal.details(),
644 },
645 ))));
646 }
647 }
648 }
649 Ok(None)
650 }
651
652 pub(crate) async fn classify_refusal(
653 &self,
654 writer: &mut dyn SqlWriter,
655 ) -> Result<Option<NoteWriteConflict>, StorageError> {
656 if let Some(claim) = &self.create_key {
657 let holder = writer
662 .query_row(statement(
663 "SELECT id, content, properties FROM notes \
664 WHERE namespace=?1 AND kind=?2 AND key=?3 AND deleted_at IS NULL",
665 vec![
666 SqlValue::Text(self.namespace.clone()),
667 SqlValue::Text(claim.kind.clone()),
668 SqlValue::Text(claim.key.clone()),
669 ],
670 ))
671 .await?;
672 if let Some(row) = holder {
673 let existing_id = match row.get("id") {
674 Some(SqlValue::Text(id)) => id.clone(),
675 _ => {
676 return Err(StorageError::Internal(
677 "invalid keyed note holder identity".into(),
678 ))
679 }
680 };
681 let existing_content = match row.get("content") {
682 Some(SqlValue::Text(content)) => content.clone(),
683 _ => {
684 return Err(StorageError::Internal(
685 "invalid keyed note holder content".into(),
686 ))
687 }
688 };
689 let existing_properties: Option<serde_json::Value> = match row.get("properties") {
690 Some(SqlValue::Text(text)) => {
691 Some(serde_json::from_str(text).map_err(|_| {
692 StorageError::Internal("invalid keyed note holder properties".into())
693 })?)
694 }
695 Some(SqlValue::Null) | None => None,
696 _ => {
697 return Err(StorageError::Internal(
698 "invalid keyed note holder properties".into(),
699 ))
700 }
701 };
702 let equal = claim.replay_signal.then(|| {
716 existing_content == claim.content && existing_properties == claim.properties
717 });
718 return Ok(Some(NoteWriteConflict::Key {
719 key: claim.key.clone(),
720 existing_id,
721 equal,
722 }));
723 }
724 }
725 if let Some(expected) = self.expected_version {
726 let current = writer
727 .query_scalar(statement(
728 "SELECT version FROM notes WHERE id=?1 AND deleted_at IS NULL",
729 vec![SqlValue::Text(self.target_id.to_string())],
730 ))
731 .await?;
732 if let Some(SqlValue::Integer(current)) = current {
733 if current != expected {
734 return Ok(Some(NoteWriteConflict::Version { expected, current }));
735 }
736 }
737 }
738 Ok(None)
739 }
740}
741
742async fn check_keyed_create_holder(
756 access: &dyn khive_storage::SqlAccess,
757 guard: NoteWriteGuard,
758) -> RuntimeResult<Option<NoteWriteConflict>> {
759 use crate::atomic_runner::{
760 run_prepared_atomic_unit, PreparedAtomicError, PreparedAtomicOp, PreparedAtomicOutcome,
761 };
762 let op: PreparedAtomicOp<(), NoteWriteConflict> = Box::new(move |writer| {
763 Box::pin(async move {
764 if let Some(conflict) = guard.check_fence(writer).await? {
765 return Err(PreparedAtomicError::Refused {
766 failure: conflict,
767 message: "keyed create: fence stale at holder check".into(),
768 });
769 }
770 if let Some(conflict) = guard.classify_refusal(writer).await? {
771 return Err(PreparedAtomicError::Refused {
772 failure: conflict,
773 message: "keyed create: live holder found at holder check".into(),
774 });
775 }
776 Ok(((), Vec::new()))
777 })
778 });
779 match run_prepared_atomic_unit(access, op).await? {
780 PreparedAtomicOutcome::Committed { .. } => Ok(None),
781 PreparedAtomicOutcome::RolledBack(conflict) => Ok(Some(conflict)),
782 }
783}
784
785pub(crate) fn validate_head(note: &khive_storage::note::Note) -> RuntimeResult<()> {
786 if note.kind != "head" {
787 return Ok(());
788 }
789 if note.name.is_some() {
790 return Err(RuntimeError::InvalidInput("head notes have no name".into()));
791 }
792 serde_json::from_str::<serde_json::Value>(¬e.content).map_err(|error| {
793 RuntimeError::InvalidInput(format!("head content must be JSON text: {error}"))
794 })?;
795 if let Some(tags) = note
796 .properties
797 .as_ref()
798 .and_then(|p| p.get("tags"))
799 .and_then(|t| t.as_array())
800 {
801 for tag in tags.iter().filter_map(|tag| tag.as_str()) {
802 if let Some(kind) = tag.strip_prefix("kind:") {
803 if kind.len() > 64 || kind.contains('\0') {
804 return Err(RuntimeError::InvalidInput(
805 "head document kind must be at most 64 bytes without U+0000".into(),
806 ));
807 }
808 }
809 }
810 }
811 Ok(())
812}
813
814impl KhiveRuntime {
815 pub async fn get_note_by_key(
816 &self,
817 token: &NamespaceToken,
818 key: &str,
819 kind: Option<&str>,
820 after_key: bool,
821 ) -> RuntimeResult<khive_storage::note::Note> {
822 self.get_note_by_key_in_scope(token, key, kind, after_key, None)
823 .await
824 }
825
826 pub async fn get_note_by_key_in_scope(
830 &self,
831 token: &NamespaceToken,
832 key: &str,
833 kind: Option<&str>,
834 after_key: bool,
835 mailbox: Option<&khive_storage::note::NoteMailboxScope>,
836 ) -> RuntimeResult<khive_storage::note::Note> {
837 crate::keyed_memory::validate_memory_key(key)?;
838 let mut matches = self
839 .notes(token)?
840 .get_live_notes_by_key(token.namespace().as_str(), key, kind)
841 .await?;
842 if let Some(scope) = mailbox {
843 matches.retain(|note| crate::MailboxView::scope_permits_message_note(scope, note));
844 }
845 match matches.len() {
846 0 => {
847 let error = KhiveError::not_found("note key", key);
848 Err(if after_key {
849 error.with_details(Details::new_owned([
850 ("reason", "after_key_missing".into()),
851 ("key", key.into()),
852 ]))
853 } else {
854 error
855 }
856 .into())
857 }
858 1 => Ok(matches.remove(0)),
859 _ => {
860 let kinds = matches
861 .iter()
862 .map(|note| note.kind.as_str())
863 .collect::<Vec<_>>()
864 .join(",");
865 Err(KhiveError::conflict("note key is ambiguous")
866 .with_details(Details::new_owned([
867 ("reason", "key_ambiguous".into()),
868 ("key", key.into()),
869 ("kinds", kinds),
870 ]))
871 .into())
872 }
873 }
874 }
875
876 pub(crate) async fn prepare_versioned_note_update(
877 &self,
878 token: &NamespaceToken,
879 snapshot: khive_storage::note::Note,
880 patch: crate::curation::NotePatch,
881 ) -> RuntimeResult<(khive_storage::note::Note, crate::atomic_plan::UpdatePlan)> {
882 use crate::atomic_plan::{AffectedRowGuard, PlanStatement, PostCommitEffect, UpdatePlan};
883 let options = patch.write_options.clone();
884 options.validate()?;
885 if options.key.is_some() {
886 return Err(RuntimeError::InvalidInput("key is immutable".into()));
887 }
888 if let Some(fences) = &options.fence {
889 for fence in fences.entries() {
890 self.validate_note_kind(&fence.kind)?;
891 }
892 }
893 let expected_updated_at = snapshot.updated_at;
894 let expected_deleted_at = snapshot.deleted_at;
895 let expected_snapshot_version = snapshot.version;
896 let (mut note, text_changed, changed) = self
897 .prepare_update_note_from_snapshot(token, snapshot, patch)
898 .await?;
899 validate_head(¬e)?;
900 if !changed && options.embed.is_none() && options.expected_version.is_none() {
907 let assertion = SqlStatement {
908 sql: "SELECT 1 FROM notes WHERE id=?1 AND updated_at=?2 AND deleted_at IS ?3 AND version=?4"
909 .into(),
910 params: vec![
911 SqlValue::Text(note.id.to_string()),
912 SqlValue::Integer(expected_updated_at),
913 expected_deleted_at
914 .map(SqlValue::Integer)
915 .unwrap_or(SqlValue::Null),
916 SqlValue::Integer(expected_snapshot_version),
917 ],
918 label: Some("note-noop-assertion".into()),
919 };
920 let plan = UpdatePlan {
921 graph_effects: Vec::new(),
922 target_id: note.id,
923 statements: vec![PlanStatement {
924 statement: assertion,
925 guard: Some(AffectedRowGuard::exactly(1)),
926 }],
927 post_commit: PostCommitEffect::None,
928 edge_natural_key: None,
929 idempotent_noop: true,
930 entity_guard: None,
931 note_guard: Some(NoteWriteGuard {
932 namespace: token.namespace().as_str().into(),
933 target_id: note.id,
934 expected_version: options.expected_version,
935 fence: options.fence,
936 create_key: None,
937 }),
938 note_vector_purge: None,
939 note_embedding_inheritance: None,
940 };
941 return Ok((note, plan));
942 }
943 if !changed {
944 let minimum_updated_at = note.updated_at.checked_add(1).ok_or_else(|| {
949 RuntimeError::Internal(format!(
950 "note {} updated_at is already at i64::MAX and cannot advance",
951 note.id
952 ))
953 })?;
954 note.updated_at = chrono::Utc::now()
955 .timestamp_micros()
956 .max(minimum_updated_at);
957 }
958 let next_version = note
959 .version
960 .checked_add(1)
961 .ok_or_else(|| RuntimeError::InvalidInput("note version exhausted".into()))?;
962 note.version = next_version;
963 let mut update = if self.stream_member_error(¬e).await?.is_some() {
964 khive_db::stores::note::note_metadata_replace_if_unchanged_statement(
965 ¬e,
966 expected_updated_at,
967 expected_deleted_at,
968 )
969 } else {
970 khive_db::stores::note::note_replace_if_unchanged_statement(
971 ¬e,
972 expected_updated_at,
973 expected_deleted_at,
974 )
975 };
976 update
980 .params
981 .push(SqlValue::Integer(expected_snapshot_version));
982 update
983 .sql
984 .push_str(&format!(" AND version = ?{}", update.params.len()));
985 if let Some(version) = options.expected_version {
986 update.params.push(SqlValue::Integer(version));
987 update
988 .sql
989 .push_str(&format!(" AND version = ?{}", update.params.len()));
990 }
991 let mut statements = vec![PlanStatement {
992 statement: update,
993 guard: Some(AffectedRowGuard::exactly(1)),
994 }];
995 if text_changed {
996 for sql in khive_db::stores::text::delete_document_statements(
997 "fts_notes",
998 ¬e.namespace,
999 note.id,
1000 )
1001 .into_iter()
1002 .chain(khive_db::stores::text::insert_document_statements(
1003 "fts_notes",
1004 &crate::curation::note_fts_document(¬e),
1005 )) {
1006 statements.push(PlanStatement {
1007 statement: sql,
1008 guard: None,
1009 });
1010 }
1011 }
1012 statements.extend(crate::atomic_prepare::event_append_statements(
1016 token,
1017 ¬e.namespace,
1018 "update",
1019 khive_types::EventKind::NoteUpdated,
1020 khive_types::SubstrateKind::Note,
1021 note.id,
1022 serde_json::json!({
1023 "id": note.id,
1024 "namespace": note.namespace,
1025 "version": note.version,
1026 "text_changed": text_changed,
1027 }),
1028 )?);
1029 let post_commit =
1032 if options.embed != Some(false) && (text_changed || options.embed == Some(true)) {
1033 PostCommitEffect::ReindexNote {
1034 note_id: note.id,
1035 version: note.version,
1036 }
1037 } else if text_changed || options.embed == Some(false) {
1038 PostCommitEffect::NoteChanged {
1039 note_id: note.id,
1040 kind: note.kind.clone(),
1041 }
1042 } else {
1043 PostCommitEffect::None
1044 };
1045 let plan = UpdatePlan {
1046 graph_effects: Vec::new(),
1047 target_id: note.id,
1048 statements,
1049 post_commit,
1050 edge_natural_key: None,
1051 idempotent_noop: false,
1052 entity_guard: None,
1053 note_guard: Some(NoteWriteGuard {
1054 namespace: token.namespace().as_str().into(),
1055 target_id: note.id,
1056 expected_version: options.expected_version,
1057 fence: options.fence,
1058 create_key: None,
1059 }),
1060 note_vector_purge: (options.embed == Some(false))
1061 .then(|| NoteVectors::new(note.namespace.clone(), note.id)),
1062 note_embedding_inheritance: (text_changed && options.embed.is_none()).then(|| {
1063 NoteEmbeddingInheritance {
1064 vectors: NoteVectors::new(note.namespace.clone(), note.id),
1065 kind: note.kind.clone(),
1066 }
1067 }),
1068 };
1069 Ok((note, plan))
1070 }
1071
1072 #[allow(clippy::too_many_arguments)]
1073 pub async fn create_note_with_options(
1074 &self,
1075 token: &NamespaceToken,
1076 kind: &str,
1077 name: Option<&str>,
1078 content: &str,
1079 embedding_content: Option<&str>,
1080 salience: Option<f64>,
1081 decay_factor: Option<f64>,
1082 properties: Option<serde_json::Value>,
1083 annotates: Vec<Uuid>,
1084 embedding_model: Option<&str>,
1085 options: NoteWriteOptions,
1086 ) -> RuntimeResult<(
1087 khive_storage::note::Note,
1088 crate::retrieval::EmbeddingTruncationReport,
1089 )> {
1090 self.create_note_with_options_resolving_annotations(
1091 token,
1092 kind,
1093 name,
1094 content,
1095 embedding_content,
1096 salience,
1097 decay_factor,
1098 properties,
1099 std::future::ready(Ok(annotates)),
1100 embedding_model,
1101 options,
1102 )
1103 .await
1104 }
1105
1106 #[allow(clippy::too_many_arguments)]
1110 pub async fn create_note_with_options_resolving_annotations<F>(
1111 &self,
1112 token: &NamespaceToken,
1113 kind: &str,
1114 name: Option<&str>,
1115 content: &str,
1116 embedding_content: Option<&str>,
1117 salience: Option<f64>,
1118 decay_factor: Option<f64>,
1119 properties: Option<serde_json::Value>,
1120 annotations: F,
1121 embedding_model: Option<&str>,
1122 options: NoteWriteOptions,
1123 ) -> RuntimeResult<(
1124 khive_storage::note::Note,
1125 crate::retrieval::EmbeddingTruncationReport,
1126 )>
1127 where
1128 F: std::future::Future<Output = RuntimeResult<Vec<Uuid>>> + Send,
1129 {
1130 use crate::atomic_message::{AtomicNoteOptions, AtomicNoteSpec};
1131 use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicRunOutcome};
1132 use crate::note_create::{prepare_note_create, KeyPublication};
1133 options.validate()?;
1134 if options.expected_version.is_some() {
1135 return Err(RuntimeError::InvalidInput(
1136 "expected_version applies only to update".into(),
1137 ));
1138 }
1139 if let Some(fences) = &options.fence {
1140 for fence in fences.entries() {
1141 self.validate_note_kind(&fence.kind)?;
1142 }
1143 }
1144 if let Some(prefix) = embedding_content {
1145 if prefix.is_empty() || prefix.len() >= content.len() || !content.starts_with(prefix) {
1146 return Err(RuntimeError::InvalidInput(
1147 "embedding_content must be a non-empty proper prefix of content".into(),
1148 ));
1149 }
1150 crate::secret_gate::check_at(prefix, "note", "embedding_content")?;
1151 }
1152 let mut candidate =
1153 khive_storage::note::Note::new(token.namespace().as_str(), kind, content);
1154 candidate.name = name.map(str::to_owned);
1155 candidate.properties = properties.clone();
1156 validate_head(&candidate)?;
1157
1158 let Some(key) = options.key.as_deref() else {
1163 let annotates = annotations.await?;
1164 let (mut prepared, _) = prepare_note_create(
1165 self,
1166 AtomicNoteSpec {
1167 token,
1168 id: None,
1169 kind,
1170 name,
1171 content,
1172 properties,
1173 },
1174 AtomicNoteOptions {
1175 salience,
1176 decay_factor,
1177 embedding_model,
1178 embedding_content,
1179 embed: Some(options.embed.unwrap_or(kind != "head")),
1180 key: None,
1181 memory_visibility_receipt: false,
1182 replay_receipt: true,
1183 fence: options.fence.as_ref(),
1184 properties_already_derived: false,
1185 },
1186 &annotates,
1187 KeyPublication::AtInsert,
1188 )
1189 .await?;
1190 let note = prepared.notes.remove(0);
1191 return match run_atomic_unit(self.sql().as_ref(), prepared.plans).await {
1192 Ok(AtomicRunOutcome::Committed { .. }) => {
1193 self.fire_note_mutation_hook(¬e.kind, note.id).await;
1194 Ok((note, prepared.embedding_truncation))
1195 }
1196 Ok(AtomicRunOutcome::RolledBack {
1197 failure: AtomicOpFailure::NoteConflict(conflict),
1198 ..
1199 }) => Err(conflict.into_error().into()),
1200 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
1201 format!("note creation rolled back: {failure:?}"),
1202 )),
1203 Err(error) => Err(RuntimeError::Storage(error.0)),
1204 };
1205 };
1206
1207 let derived_properties =
1217 self.derive_note_write_properties(kind, token, properties.clone())?;
1218 let namespace: String = token.namespace().as_str().into();
1219
1220 if let Some(conflict) = check_keyed_create_holder(
1226 self.sql().as_ref(),
1227 NoteWriteGuard {
1228 namespace: namespace.clone(),
1229 target_id: candidate.id,
1230 expected_version: None,
1231 fence: options.fence.clone(),
1232 create_key: Some(CreateKeyClaim {
1233 kind: kind.to_owned(),
1234 key: key.to_owned(),
1235 content: content.to_owned(),
1236 properties: derived_properties.clone(),
1237 replay_signal: true,
1238 }),
1239 },
1240 )
1241 .await?
1242 {
1243 return Err(conflict.into_error().into());
1244 }
1245
1246 let prep_result = async {
1253 let annotates = annotations.await?;
1254 prepare_note_create(
1255 self,
1256 AtomicNoteSpec {
1257 token,
1258 id: None,
1259 kind,
1260 name,
1261 content,
1262 properties: derived_properties.clone(),
1263 },
1264 AtomicNoteOptions {
1265 salience,
1266 decay_factor,
1267 embedding_model,
1268 embedding_content,
1269 embed: Some(options.embed.unwrap_or(kind != "head")),
1270 key: Some(key),
1271 memory_visibility_receipt: false,
1272 replay_receipt: true,
1273 fence: options.fence.as_ref(),
1274 properties_already_derived: true,
1275 },
1276 &annotates,
1277 KeyPublication::AtInsert,
1278 )
1279 .await
1280 }
1281 .await;
1282
1283 #[cfg(test)]
1289 race_seam::pause_after_prepare().await;
1290
1291 match prep_result {
1292 Ok((mut prepared, _)) => {
1293 let note = prepared.notes.remove(0);
1294 match run_atomic_unit(self.sql().as_ref(), prepared.plans).await {
1295 Ok(AtomicRunOutcome::Committed { .. }) => {
1296 self.fire_note_mutation_hook(¬e.kind, note.id).await;
1297 Ok((note, prepared.embedding_truncation))
1298 }
1299 Ok(AtomicRunOutcome::RolledBack {
1300 failure: AtomicOpFailure::NoteConflict(conflict),
1301 ..
1302 }) => Err(conflict.into_error().into()),
1303 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(
1304 RuntimeError::Internal(format!("note creation rolled back: {failure:?}")),
1305 ),
1306 Err(error) => Err(RuntimeError::Storage(error.0)),
1307 }
1308 }
1309 Err(prep_error) => {
1310 match check_keyed_create_holder(
1316 self.sql().as_ref(),
1317 NoteWriteGuard {
1318 namespace,
1319 target_id: candidate.id,
1320 expected_version: None,
1321 fence: options.fence.clone(),
1322 create_key: Some(CreateKeyClaim {
1323 kind: kind.to_owned(),
1324 key: key.to_owned(),
1325 content: content.to_owned(),
1326 properties: derived_properties,
1327 replay_signal: true,
1328 }),
1329 },
1330 )
1331 .await?
1332 {
1333 Some(conflict) => Err(conflict.into_error().into()),
1334 None => Err(prep_error),
1335 }
1336 }
1337 }
1338 }
1339}