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(
298 SqlStatement::new(
299 "SELECT name FROM pragma_table_list \
300 WHERE schema='main' AND type='virtual' AND name GLOB 'vec_*' ORDER BY name",
301 vec![],
302 )
303 .labelled("note-write-guard"),
304 )
305 .await?;
306 let mut names = Vec::with_capacity(tables.len());
307 for row in tables {
308 let Some(SqlValue::Text(table)) = row.get("name") else {
309 return Err(StorageError::Internal(
310 "invalid vector table catalog row".into(),
311 ));
312 };
313 if !table.strip_prefix("vec_").is_some_and(|key| {
314 !key.is_empty() && key.bytes().all(|c| c.is_ascii_alphanumeric() || c == b'_')
315 }) {
316 return Err(StorageError::Internal(
317 "invalid persisted vector table name".into(),
318 ));
319 }
320 names.push(table.clone());
321 }
322 Ok(names)
323 }
324
325 pub(crate) async fn has_rows(&self, writer: &mut dyn SqlWriter) -> Result<bool, StorageError> {
326 for table in Self::tables(writer).await? {
327 if writer
328 .query_scalar(
329 SqlStatement::new(
330 format!(
331 "SELECT 1 FROM main.{table} WHERE namespace=?1 AND subject_id=?2 LIMIT 1"
332 ),
333 vec![
334 SqlValue::Text(self.namespace.clone()),
335 SqlValue::Text(self.subject_id.to_string()),
336 ],
337 )
338 .labelled("note-write-guard"),
339 )
340 .await?
341 .is_some()
342 {
343 return Ok(true);
344 }
345 }
346 Ok(false)
347 }
348
349 pub(crate) async fn apply(&self, writer: &mut dyn SqlWriter) -> Result<(), StorageError> {
350 for table in Self::tables(writer).await? {
351 let scope = vec![
352 SqlValue::Text(self.namespace.clone()),
353 SqlValue::Text(self.subject_id.to_string()),
354 ];
355 writer
356 .execute(
357 SqlStatement::new(
358 format!(
359 "INSERT INTO ann_write_log (namespace,embedding_model,kind,field,subject_id,op) \
360 SELECT namespace,embedding_model,kind,field,subject_id,'delete' \
361 FROM main.{table} WHERE namespace=?1 AND subject_id=?2"),
362 scope.clone(),
363 )
364 .labelled("note-write-guard"),
365 )
366 .await?;
367 writer
368 .execute(
369 SqlStatement::new(
370 format!("DELETE FROM main.{table} WHERE namespace=?1 AND subject_id=?2"),
371 scope,
372 )
373 .labelled("note-write-guard"),
374 )
375 .await?;
376 let model_key = table.strip_prefix("vec_").ok_or_else(|| {
377 StorageError::Internal("invalid persisted vector table name".into())
378 })?;
379 writer
380 .execute(
381 SqlStatement::new(
382 "DELETE FROM vector_provenance \
383 WHERE model_key=?1 AND namespace=?2 AND subject_id=?3",
384 vec![
385 SqlValue::Text(model_key.to_string()),
386 SqlValue::Text(self.namespace.clone()),
387 SqlValue::Text(self.subject_id.to_string()),
388 ],
389 )
390 .labelled("note-write-guard"),
391 )
392 .await?;
393 }
394 Ok(())
395 }
396}
397
398#[derive(Clone, Debug)]
399pub(crate) struct NoteEmbeddingInheritance {
400 pub vectors: NoteVectors,
401 pub kind: String,
402}
403
404#[derive(Clone, Debug, PartialEq, Eq)]
405pub enum NoteWriteConflict {
406 Version {
407 expected: i64,
408 current: i64,
409 },
410 Fence {
411 key: String,
412 expected: Option<i64>,
413 current: Option<i64>,
414 index: Option<usize>,
415 },
416 Key {
417 key: String,
418 existing_id: String,
419 equal: Option<bool>,
430 },
431 FenceDeadline(Box<FenceDeadline>),
435 FenceIdentity(Box<FenceIdentity>),
439}
440
441#[derive(Clone, Debug, PartialEq, Eq)]
444pub struct FenceIdentity {
445 pub key: String,
446 pub kind: String,
447 pub version: i64,
448 pub index: Option<usize>,
449 pub evidence: Vec<(&'static str, String)>,
450}
451
452#[derive(Clone, Debug, PartialEq, Eq)]
455pub struct FenceDeadline {
456 pub key: String,
457 pub kind: String,
458 pub version: i64,
459 pub field: String,
460 pub index: Option<usize>,
461 pub reason: &'static str,
462 pub evidence: Vec<(&'static str, String)>,
463}
464
465impl NoteWriteConflict {
466 pub fn into_error(self) -> KhiveError {
467 self.into_error_at_member(None)
468 }
469
470 pub(crate) fn into_error_at_member(self, member: Option<usize>) -> KhiveError {
471 let (message, mut details) = match self {
472 Self::Version { expected, current } => (
473 "note version precondition failed",
474 vec![
475 ("reason", "version_conflict".into()),
476 ("expected_version", expected.to_string()),
477 ("current_version", current.to_string()),
478 ],
479 ),
480 Self::Fence {
481 key,
482 expected,
483 current,
484 index,
485 } => {
486 let mut fields = vec![
487 ("reason", "fence_conflict".into()),
488 ("key", key),
489 (
490 "expected_version",
491 expected.map_or_else(|| "absent".into(), |version| version.to_string()),
492 ),
493 ];
494 if let Some(current) = current {
495 fields.push(("current_version", current.to_string()));
496 }
497 if let Some(index) = index {
498 fields.push(("index", index.to_string()));
499 }
500 ("note fence precondition failed", fields)
501 }
502 Self::Key {
503 key,
504 existing_id,
505 equal,
506 } => {
507 let mut fields = vec![
508 ("reason", "key_conflict".into()),
509 ("key", key),
510 ("existing_id", existing_id),
511 ];
512 if let Some(equal) = equal {
513 fields.push(("equal", equal.to_string()));
514 }
515 ("a live note already holds this key", fields)
516 }
517 Self::FenceDeadline(deadline) => {
518 let FenceDeadline {
519 key,
520 kind,
521 version,
522 field,
523 index,
524 reason,
525 evidence,
526 } = *deadline;
527 let mut fields = vec![
528 ("reason", reason.into()),
529 ("key", key),
530 ("kind", kind),
531 ("version", version.to_string()),
532 ("field", field),
533 ];
534 if let Some(index) = index {
535 fields.push(("index", index.to_string()));
536 }
537 fields.extend(evidence);
538 ("note fence time precondition failed", fields)
539 }
540 Self::FenceIdentity(identity) => {
541 let FenceIdentity {
542 key,
543 kind,
544 version,
545 index,
546 evidence,
547 } = *identity;
548 let mut fields = vec![
549 ("reason", "identity_conflict".into()),
550 ("key", key),
551 ("kind", kind),
552 ("version", version.to_string()),
553 ];
554 fields.extend(evidence);
555 if let Some(index) = index {
556 fields.push(("index", index.to_string()));
557 }
558 ("note fence identity precondition failed", fields)
559 }
560 };
561 if let Some(member) = member {
562 details.push(("member", member.to_string()));
563 }
564 KhiveError::conflict(message).with_details(Details::new_owned(details))
565 }
566}
567
568impl NoteWriteGuard {
569 pub(crate) async fn check_fence(
570 &self,
571 writer: &mut dyn SqlWriter,
572 ) -> Result<Option<NoteWriteConflict>, StorageError> {
573 let Some(fences) = &self.fence else {
574 return Ok(None);
575 };
576 let mut now: Option<i64> = None;
580 for (index, fence) in fences.entries().iter().enumerate() {
581 let holder = crate::fence_identity::read_holder(
582 writer,
583 &self.namespace,
584 &fence.kind,
585 &fence.key,
586 "note-write-guard",
587 )
588 .await?;
589 if let (Some(asserted), Some(holder)) = (fence.id, holder.as_ref()) {
595 if asserted != holder.id {
596 return Ok(Some(NoteWriteConflict::FenceIdentity(Box::new(
597 FenceIdentity {
598 key: fence.key.clone(),
599 kind: fence.kind.clone(),
600 version: fence.expected_version.unwrap_or_default(),
603 index: matches!(fences, NoteFences::Many(_)).then_some(index),
604 evidence: crate::fence_identity::identity_evidence(asserted, holder.id),
605 },
606 ))));
607 }
608 }
609 let current = holder.map(|holder| holder.version);
610 if current != fence.expected_version {
611 return Ok(Some(NoteWriteConflict::Fence {
612 key: fence.key.clone(),
613 expected: fence.expected_version,
614 current,
615 index: matches!(fences, NoteFences::Many(_)).then_some(index),
616 }));
617 }
618 if let Some(field) = &fence.live_until {
619 let clock = match now {
620 Some(clock) => clock,
621 None => {
622 let clock =
623 crate::live_until::writer_clock(writer, "note-write-guard-clock")
624 .await?;
625 now = Some(clock);
626 clock
627 }
628 };
629 if let Some(refusal) = crate::live_until::evaluate(
630 writer,
631 &self.namespace,
632 &fence.kind,
633 &fence.key,
634 field,
635 clock,
636 "note-write-guard-live-until",
637 )
638 .await?
639 {
640 return Ok(Some(NoteWriteConflict::FenceDeadline(Box::new(
641 FenceDeadline {
642 key: fence.key.clone(),
643 kind: fence.kind.clone(),
644 version: fence.expected_version.unwrap_or_default(),
647 field: field.clone(),
648 index: matches!(fences, NoteFences::Many(_)).then_some(index),
649 reason: refusal.reason(),
650 evidence: refusal.details(),
651 },
652 ))));
653 }
654 }
655 }
656 Ok(None)
657 }
658
659 pub(crate) async fn classify_refusal(
660 &self,
661 writer: &mut dyn SqlWriter,
662 ) -> Result<Option<NoteWriteConflict>, StorageError> {
663 if let Some(claim) = &self.create_key {
664 let holder = writer
669 .query_row(
670 SqlStatement::new(
671 "SELECT id, content, properties FROM notes \
672 WHERE namespace=?1 AND kind=?2 AND key=?3 AND deleted_at IS NULL",
673 vec![
674 SqlValue::Text(self.namespace.clone()),
675 SqlValue::Text(claim.kind.clone()),
676 SqlValue::Text(claim.key.clone()),
677 ],
678 )
679 .labelled("note-write-guard"),
680 )
681 .await?;
682 if let Some(row) = holder {
683 let existing_id = match row.get("id") {
684 Some(SqlValue::Text(id)) => id.clone(),
685 _ => {
686 return Err(StorageError::Internal(
687 "invalid keyed note holder identity".into(),
688 ))
689 }
690 };
691 let existing_content = match row.get("content") {
692 Some(SqlValue::Text(content)) => content.clone(),
693 _ => {
694 return Err(StorageError::Internal(
695 "invalid keyed note holder content".into(),
696 ))
697 }
698 };
699 let existing_properties: Option<serde_json::Value> = match row.get("properties") {
700 Some(SqlValue::Text(text)) => {
701 Some(serde_json::from_str(text).map_err(|_| {
702 StorageError::Internal("invalid keyed note holder properties".into())
703 })?)
704 }
705 Some(SqlValue::Null) | None => None,
706 _ => {
707 return Err(StorageError::Internal(
708 "invalid keyed note holder properties".into(),
709 ))
710 }
711 };
712 let equal = claim.replay_signal.then(|| {
726 existing_content == claim.content && existing_properties == claim.properties
727 });
728 return Ok(Some(NoteWriteConflict::Key {
729 key: claim.key.clone(),
730 existing_id,
731 equal,
732 }));
733 }
734 }
735 if let Some(expected) = self.expected_version {
736 let current = writer
737 .query_scalar(
738 SqlStatement::new(
739 "SELECT version FROM notes WHERE id=?1 AND deleted_at IS NULL",
740 vec![SqlValue::Text(self.target_id.to_string())],
741 )
742 .labelled("note-write-guard"),
743 )
744 .await?;
745 if let Some(SqlValue::Integer(current)) = current {
746 if current != expected {
747 return Ok(Some(NoteWriteConflict::Version { expected, current }));
748 }
749 }
750 }
751 Ok(None)
752 }
753}
754
755async fn check_keyed_create_holder(
769 access: &dyn khive_storage::SqlAccess,
770 guard: NoteWriteGuard,
771) -> RuntimeResult<Option<NoteWriteConflict>> {
772 use crate::atomic_runner::{
773 run_prepared_atomic_unit, PreparedAtomicError, PreparedAtomicOp, PreparedAtomicOutcome,
774 };
775 let op: PreparedAtomicOp<(), NoteWriteConflict> = Box::new(move |writer| {
776 Box::pin(async move {
777 if let Some(conflict) = guard.check_fence(writer).await? {
778 return Err(PreparedAtomicError::Refused {
779 failure: conflict,
780 message: "keyed create: fence stale at holder check".into(),
781 });
782 }
783 if let Some(conflict) = guard.classify_refusal(writer).await? {
784 return Err(PreparedAtomicError::Refused {
785 failure: conflict,
786 message: "keyed create: live holder found at holder check".into(),
787 });
788 }
789 Ok(((), Vec::new()))
790 })
791 });
792 match run_prepared_atomic_unit(access, op).await? {
793 PreparedAtomicOutcome::Committed { .. } => Ok(None),
794 PreparedAtomicOutcome::RolledBack(conflict) => Ok(Some(conflict)),
795 }
796}
797
798pub(crate) fn validate_head(note: &khive_storage::note::Note) -> RuntimeResult<()> {
799 if note.kind != "head" {
800 return Ok(());
801 }
802 if note.name.is_some() {
803 return Err(RuntimeError::InvalidInput("head notes have no name".into()));
804 }
805 serde_json::from_str::<serde_json::Value>(¬e.content).map_err(|error| {
806 RuntimeError::InvalidInput(format!("head content must be JSON text: {error}"))
807 })?;
808 if let Some(tags) = note
809 .properties
810 .as_ref()
811 .and_then(|p| p.get("tags"))
812 .and_then(|t| t.as_array())
813 {
814 for tag in tags.iter().filter_map(|tag| tag.as_str()) {
815 if let Some(kind) = tag.strip_prefix("kind:") {
816 if kind.len() > 64 || kind.contains('\0') {
817 return Err(RuntimeError::InvalidInput(
818 "head document kind must be at most 64 bytes without U+0000".into(),
819 ));
820 }
821 }
822 }
823 }
824 Ok(())
825}
826
827impl KhiveRuntime {
828 pub async fn get_note_by_key(
829 &self,
830 token: &NamespaceToken,
831 key: &str,
832 kind: Option<&str>,
833 after_key: bool,
834 ) -> RuntimeResult<khive_storage::note::Note> {
835 self.get_note_by_key_in_scope(token, key, kind, after_key, None)
836 .await
837 }
838
839 pub async fn get_note_by_key_in_scope(
843 &self,
844 token: &NamespaceToken,
845 key: &str,
846 kind: Option<&str>,
847 after_key: bool,
848 mailbox: Option<&khive_storage::note::NoteMailboxScope>,
849 ) -> RuntimeResult<khive_storage::note::Note> {
850 crate::keyed_memory::validate_memory_key(key)?;
851 let mut matches = self
852 .notes(token)?
853 .get_live_notes_by_key(token.namespace().as_str(), key, kind)
854 .await?;
855 if let Some(scope) = mailbox {
856 matches.retain(|note| crate::MailboxView::scope_permits_message_note(scope, note));
857 }
858 match matches.len() {
859 0 => {
860 let error = KhiveError::not_found("note key", key);
861 Err(if after_key {
862 error.with_details(Details::new_owned([
863 ("reason", "after_key_missing".into()),
864 ("key", key.into()),
865 ]))
866 } else {
867 error
868 }
869 .into())
870 }
871 1 => Ok(matches.remove(0)),
872 _ => {
873 let kinds = matches
874 .iter()
875 .map(|note| note.kind.as_str())
876 .collect::<Vec<_>>()
877 .join(",");
878 Err(KhiveError::conflict("note key is ambiguous")
879 .with_details(Details::new_owned([
880 ("reason", "key_ambiguous".into()),
881 ("key", key.into()),
882 ("kinds", kinds),
883 ]))
884 .into())
885 }
886 }
887 }
888
889 pub(crate) async fn prepare_versioned_note_update(
890 &self,
891 token: &NamespaceToken,
892 snapshot: khive_storage::note::Note,
893 patch: crate::curation::NotePatch,
894 ) -> RuntimeResult<(khive_storage::note::Note, crate::atomic_plan::UpdatePlan)> {
895 use crate::atomic_plan::{AffectedRowGuard, PlanStatement, PostCommitEffect, UpdatePlan};
896 let options = patch.write_options.clone();
897 options.validate()?;
898 if options.key.is_some() {
899 return Err(RuntimeError::InvalidInput("key is immutable".into()));
900 }
901 if let Some(fences) = &options.fence {
902 for fence in fences.entries() {
903 self.validate_note_kind(&fence.kind)?;
904 }
905 }
906 let expected_updated_at = snapshot.updated_at;
907 let expected_deleted_at = snapshot.deleted_at;
908 let expected_snapshot_version = snapshot.version;
909 let (mut note, text_changed, changed) = self
910 .prepare_update_note_from_snapshot(token, snapshot, patch)
911 .await?;
912 validate_head(¬e)?;
913 if !changed && options.embed.is_none() && options.expected_version.is_none() {
920 let assertion = SqlStatement {
921 sql: "SELECT 1 FROM notes WHERE id=?1 AND updated_at=?2 AND deleted_at IS ?3 AND version=?4"
922 .into(),
923 params: vec![
924 SqlValue::Text(note.id.to_string()),
925 SqlValue::Integer(expected_updated_at),
926 expected_deleted_at
927 .map(SqlValue::Integer)
928 .unwrap_or(SqlValue::Null),
929 SqlValue::Integer(expected_snapshot_version),
930 ],
931 label: Some("note-noop-assertion".into()),
932 };
933 let plan = UpdatePlan {
934 graph_effects: Vec::new(),
935 target_id: note.id,
936 statements: vec![PlanStatement {
937 statement: assertion,
938 guard: Some(AffectedRowGuard::exactly(1)),
939 }],
940 post_commit: PostCommitEffect::None,
941 edge_natural_key: None,
942 idempotent_noop: true,
943 entity_guard: None,
944 note_guard: Some(NoteWriteGuard {
945 namespace: token.namespace().as_str().into(),
946 target_id: note.id,
947 expected_version: options.expected_version,
948 fence: options.fence,
949 create_key: None,
950 }),
951 note_vector_purge: None,
952 note_embedding_inheritance: None,
953 };
954 return Ok((note, plan));
955 }
956 if !changed {
957 let minimum_updated_at = note.updated_at.checked_add(1).ok_or_else(|| {
962 RuntimeError::Internal(format!(
963 "note {} updated_at is already at i64::MAX and cannot advance",
964 note.id
965 ))
966 })?;
967 note.updated_at = chrono::Utc::now()
968 .timestamp_micros()
969 .max(minimum_updated_at);
970 }
971 let next_version = note
972 .version
973 .checked_add(1)
974 .ok_or_else(|| RuntimeError::InvalidInput("note version exhausted".into()))?;
975 note.version = next_version;
976 let mut update = if self.stream_member_error(¬e).await?.is_some() {
977 khive_db::stores::note::note_metadata_replace_if_unchanged_statement(
978 ¬e,
979 expected_updated_at,
980 expected_deleted_at,
981 )
982 } else {
983 khive_db::stores::note::note_replace_if_unchanged_statement(
984 ¬e,
985 expected_updated_at,
986 expected_deleted_at,
987 )
988 };
989 update
993 .params
994 .push(SqlValue::Integer(expected_snapshot_version));
995 update
996 .sql
997 .push_str(&format!(" AND version = ?{}", update.params.len()));
998 if let Some(version) = options.expected_version {
999 update.params.push(SqlValue::Integer(version));
1000 update
1001 .sql
1002 .push_str(&format!(" AND version = ?{}", update.params.len()));
1003 }
1004 let mut statements = vec![PlanStatement {
1005 statement: update,
1006 guard: Some(AffectedRowGuard::exactly(1)),
1007 }];
1008 if text_changed {
1009 for sql in khive_db::stores::text::delete_document_statements(
1010 "fts_notes",
1011 ¬e.namespace,
1012 note.id,
1013 )
1014 .into_iter()
1015 .chain(khive_db::stores::text::insert_document_statements(
1016 "fts_notes",
1017 &crate::curation::note_fts_document(¬e),
1018 )) {
1019 statements.push(PlanStatement {
1020 statement: sql,
1021 guard: None,
1022 });
1023 }
1024 }
1025 statements.extend(crate::atomic_prepare::event_append_statements(
1029 token,
1030 ¬e.namespace,
1031 "update",
1032 khive_types::EventKind::NoteUpdated,
1033 khive_types::SubstrateKind::Note,
1034 note.id,
1035 serde_json::json!({
1036 "id": note.id,
1037 "namespace": note.namespace,
1038 "version": note.version,
1039 "text_changed": text_changed,
1040 }),
1041 )?);
1042 let post_commit =
1045 if options.embed != Some(false) && (text_changed || options.embed == Some(true)) {
1046 PostCommitEffect::ReindexNote {
1047 note_id: note.id,
1048 version: note.version,
1049 }
1050 } else if text_changed || options.embed == Some(false) {
1051 PostCommitEffect::NoteChanged {
1052 note_id: note.id,
1053 kind: note.kind.clone(),
1054 }
1055 } else {
1056 PostCommitEffect::None
1057 };
1058 let plan = UpdatePlan {
1059 graph_effects: Vec::new(),
1060 target_id: note.id,
1061 statements,
1062 post_commit,
1063 edge_natural_key: None,
1064 idempotent_noop: false,
1065 entity_guard: None,
1066 note_guard: Some(NoteWriteGuard {
1067 namespace: token.namespace().as_str().into(),
1068 target_id: note.id,
1069 expected_version: options.expected_version,
1070 fence: options.fence,
1071 create_key: None,
1072 }),
1073 note_vector_purge: (options.embed == Some(false))
1074 .then(|| NoteVectors::new(note.namespace.clone(), note.id)),
1075 note_embedding_inheritance: (text_changed && options.embed.is_none()).then(|| {
1076 NoteEmbeddingInheritance {
1077 vectors: NoteVectors::new(note.namespace.clone(), note.id),
1078 kind: note.kind.clone(),
1079 }
1080 }),
1081 };
1082 Ok((note, plan))
1083 }
1084
1085 #[allow(clippy::too_many_arguments)]
1086 pub async fn create_note_with_options(
1087 &self,
1088 token: &NamespaceToken,
1089 kind: &str,
1090 name: Option<&str>,
1091 content: &str,
1092 embedding_content: Option<&str>,
1093 salience: Option<f64>,
1094 decay_factor: Option<f64>,
1095 properties: Option<serde_json::Value>,
1096 annotates: Vec<Uuid>,
1097 embedding_model: Option<&str>,
1098 options: NoteWriteOptions,
1099 ) -> RuntimeResult<(
1100 khive_storage::note::Note,
1101 crate::retrieval::EmbeddingTruncationReport,
1102 )> {
1103 self.create_note_with_options_resolving_annotations(
1104 token,
1105 kind,
1106 name,
1107 content,
1108 embedding_content,
1109 salience,
1110 decay_factor,
1111 properties,
1112 std::future::ready(Ok(annotates)),
1113 embedding_model,
1114 options,
1115 )
1116 .await
1117 }
1118
1119 #[allow(clippy::too_many_arguments)]
1123 pub async fn create_note_with_options_resolving_annotations<F>(
1124 &self,
1125 token: &NamespaceToken,
1126 kind: &str,
1127 name: Option<&str>,
1128 content: &str,
1129 embedding_content: Option<&str>,
1130 salience: Option<f64>,
1131 decay_factor: Option<f64>,
1132 properties: Option<serde_json::Value>,
1133 annotations: F,
1134 embedding_model: Option<&str>,
1135 options: NoteWriteOptions,
1136 ) -> RuntimeResult<(
1137 khive_storage::note::Note,
1138 crate::retrieval::EmbeddingTruncationReport,
1139 )>
1140 where
1141 F: std::future::Future<Output = RuntimeResult<Vec<Uuid>>> + Send,
1142 {
1143 use crate::atomic_message::{AtomicNoteOptions, AtomicNoteSpec};
1144 use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicRunOutcome};
1145 use crate::note_create::{prepare_note_create, KeyPublication};
1146 options.validate()?;
1147 if options.expected_version.is_some() {
1148 return Err(RuntimeError::InvalidInput(
1149 "expected_version applies only to update".into(),
1150 ));
1151 }
1152 if let Some(fences) = &options.fence {
1153 for fence in fences.entries() {
1154 self.validate_note_kind(&fence.kind)?;
1155 }
1156 }
1157 if let Some(prefix) = embedding_content {
1158 if prefix.is_empty() || prefix.len() >= content.len() || !content.starts_with(prefix) {
1159 return Err(RuntimeError::InvalidInput(
1160 "embedding_content must be a non-empty proper prefix of content".into(),
1161 ));
1162 }
1163 crate::secret_gate::check_at(prefix, "note", "embedding_content")?;
1164 }
1165 let mut candidate =
1166 khive_storage::note::Note::new(token.namespace().as_str(), kind, content);
1167 candidate.name = name.map(str::to_owned);
1168 candidate.properties = properties.clone();
1169 validate_head(&candidate)?;
1170
1171 let Some(key) = options.key.as_deref() else {
1176 let annotates = annotations.await?;
1177 let (mut prepared, _) = prepare_note_create(
1178 self,
1179 AtomicNoteSpec {
1180 token,
1181 id: None,
1182 kind,
1183 name,
1184 content,
1185 properties,
1186 },
1187 AtomicNoteOptions {
1188 salience,
1189 decay_factor,
1190 embedding_model,
1191 embedding_content,
1192 embed: Some(options.embed.unwrap_or(kind != "head")),
1193 key: None,
1194 memory_visibility_receipt: false,
1195 replay_receipt: true,
1196 fence: options.fence.as_ref(),
1197 properties_already_derived: false,
1198 },
1199 &annotates,
1200 KeyPublication::AtInsert,
1201 )
1202 .await?;
1203 let note = prepared.notes.remove(0);
1204 return match run_atomic_unit(self.sql().as_ref(), prepared.plans).await {
1205 Ok(AtomicRunOutcome::Committed { .. }) => {
1206 self.fire_note_mutation_hook(¬e.kind, note.id).await;
1207 Ok((note, prepared.embedding_truncation))
1208 }
1209 Ok(AtomicRunOutcome::RolledBack {
1210 failure: AtomicOpFailure::NoteConflict(conflict),
1211 ..
1212 }) => Err(conflict.into_error().into()),
1213 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
1214 format!("note creation rolled back: {failure:?}"),
1215 )),
1216 Err(error) => Err(RuntimeError::Storage(error.0)),
1217 };
1218 };
1219
1220 let derived_properties =
1230 self.derive_note_write_properties(kind, token, properties.clone())?;
1231 let namespace: String = token.namespace().as_str().into();
1232
1233 if let Some(conflict) = check_keyed_create_holder(
1239 self.sql().as_ref(),
1240 NoteWriteGuard {
1241 namespace: namespace.clone(),
1242 target_id: candidate.id,
1243 expected_version: None,
1244 fence: options.fence.clone(),
1245 create_key: Some(CreateKeyClaim {
1246 kind: kind.to_owned(),
1247 key: key.to_owned(),
1248 content: content.to_owned(),
1249 properties: derived_properties.clone(),
1250 replay_signal: true,
1251 }),
1252 },
1253 )
1254 .await?
1255 {
1256 return Err(conflict.into_error().into());
1257 }
1258
1259 let prep_result = async {
1266 let annotates = annotations.await?;
1267 prepare_note_create(
1268 self,
1269 AtomicNoteSpec {
1270 token,
1271 id: None,
1272 kind,
1273 name,
1274 content,
1275 properties: derived_properties.clone(),
1276 },
1277 AtomicNoteOptions {
1278 salience,
1279 decay_factor,
1280 embedding_model,
1281 embedding_content,
1282 embed: Some(options.embed.unwrap_or(kind != "head")),
1283 key: Some(key),
1284 memory_visibility_receipt: false,
1285 replay_receipt: true,
1286 fence: options.fence.as_ref(),
1287 properties_already_derived: true,
1288 },
1289 &annotates,
1290 KeyPublication::AtInsert,
1291 )
1292 .await
1293 }
1294 .await;
1295
1296 #[cfg(test)]
1302 race_seam::pause_after_prepare().await;
1303
1304 match prep_result {
1305 Ok((mut prepared, _)) => {
1306 let note = prepared.notes.remove(0);
1307 match run_atomic_unit(self.sql().as_ref(), prepared.plans).await {
1308 Ok(AtomicRunOutcome::Committed { .. }) => {
1309 self.fire_note_mutation_hook(¬e.kind, note.id).await;
1310 Ok((note, prepared.embedding_truncation))
1311 }
1312 Ok(AtomicRunOutcome::RolledBack {
1313 failure: AtomicOpFailure::NoteConflict(conflict),
1314 ..
1315 }) => Err(conflict.into_error().into()),
1316 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(
1317 RuntimeError::Internal(format!("note creation rolled back: {failure:?}")),
1318 ),
1319 Err(error) => Err(RuntimeError::Storage(error.0)),
1320 }
1321 }
1322 Err(prep_error) => {
1323 match check_keyed_create_holder(
1329 self.sql().as_ref(),
1330 NoteWriteGuard {
1331 namespace,
1332 target_id: candidate.id,
1333 expected_version: None,
1334 fence: options.fence.clone(),
1335 create_key: Some(CreateKeyClaim {
1336 kind: kind.to_owned(),
1337 key: key.to_owned(),
1338 content: content.to_owned(),
1339 properties: derived_properties,
1340 replay_signal: true,
1341 }),
1342 },
1343 )
1344 .await?
1345 {
1346 Some(conflict) => Err(conflict.into_error().into()),
1347 None => Err(prep_error),
1348 }
1349 }
1350 }
1351 }
1352}