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 }
365 Ok(())
366 }
367}
368
369#[derive(Clone, Debug)]
370pub(crate) struct NoteEmbeddingInheritance {
371 pub vectors: NoteVectors,
372 pub kind: String,
373}
374
375#[derive(Clone, Debug, PartialEq, Eq)]
376pub enum NoteWriteConflict {
377 Version {
378 expected: i64,
379 current: i64,
380 },
381 Fence {
382 key: String,
383 expected: Option<i64>,
384 current: Option<i64>,
385 index: Option<usize>,
386 },
387 Key {
388 key: String,
389 existing_id: String,
390 equal: Option<bool>,
401 },
402 FenceDeadline(Box<FenceDeadline>),
406 FenceIdentity(Box<FenceIdentity>),
410}
411
412#[derive(Clone, Debug, PartialEq, Eq)]
415pub struct FenceIdentity {
416 pub key: String,
417 pub kind: String,
418 pub version: i64,
419 pub index: Option<usize>,
420 pub evidence: Vec<(&'static str, String)>,
421}
422
423#[derive(Clone, Debug, PartialEq, Eq)]
426pub struct FenceDeadline {
427 pub key: String,
428 pub kind: String,
429 pub version: i64,
430 pub field: String,
431 pub index: Option<usize>,
432 pub reason: &'static str,
433 pub evidence: Vec<(&'static str, String)>,
434}
435
436impl NoteWriteConflict {
437 pub fn into_error(self) -> KhiveError {
438 self.into_error_at_member(None)
439 }
440
441 pub(crate) fn into_error_at_member(self, member: Option<usize>) -> KhiveError {
442 let (message, mut details) = match self {
443 Self::Version { expected, current } => (
444 "note version precondition failed",
445 vec![
446 ("reason", "version_conflict".into()),
447 ("expected_version", expected.to_string()),
448 ("current_version", current.to_string()),
449 ],
450 ),
451 Self::Fence {
452 key,
453 expected,
454 current,
455 index,
456 } => {
457 let mut fields = vec![
458 ("reason", "fence_conflict".into()),
459 ("key", key),
460 (
461 "expected_version",
462 expected.map_or_else(|| "absent".into(), |version| version.to_string()),
463 ),
464 ];
465 if let Some(current) = current {
466 fields.push(("current_version", current.to_string()));
467 }
468 if let Some(index) = index {
469 fields.push(("index", index.to_string()));
470 }
471 ("note fence precondition failed", fields)
472 }
473 Self::Key {
474 key,
475 existing_id,
476 equal,
477 } => {
478 let mut fields = vec![
479 ("reason", "key_conflict".into()),
480 ("key", key),
481 ("existing_id", existing_id),
482 ];
483 if let Some(equal) = equal {
484 fields.push(("equal", equal.to_string()));
485 }
486 ("a live note already holds this key", fields)
487 }
488 Self::FenceDeadline(deadline) => {
489 let FenceDeadline {
490 key,
491 kind,
492 version,
493 field,
494 index,
495 reason,
496 evidence,
497 } = *deadline;
498 let mut fields = vec![
499 ("reason", reason.into()),
500 ("key", key),
501 ("kind", kind),
502 ("version", version.to_string()),
503 ("field", field),
504 ];
505 if let Some(index) = index {
506 fields.push(("index", index.to_string()));
507 }
508 fields.extend(evidence);
509 ("note fence time precondition failed", fields)
510 }
511 Self::FenceIdentity(identity) => {
512 let FenceIdentity {
513 key,
514 kind,
515 version,
516 index,
517 evidence,
518 } = *identity;
519 let mut fields = vec![
520 ("reason", "identity_conflict".into()),
521 ("key", key),
522 ("kind", kind),
523 ("version", version.to_string()),
524 ];
525 fields.extend(evidence);
526 if let Some(index) = index {
527 fields.push(("index", index.to_string()));
528 }
529 ("note fence identity precondition failed", fields)
530 }
531 };
532 if let Some(member) = member {
533 details.push(("member", member.to_string()));
534 }
535 KhiveError::conflict(message).with_details(Details::new_owned(details))
536 }
537}
538
539pub(crate) fn statement(sql: impl Into<String>, params: Vec<SqlValue>) -> SqlStatement {
540 SqlStatement {
541 sql: sql.into(),
542 params,
543 label: Some("note-write-guard".into()),
544 }
545}
546
547impl NoteWriteGuard {
548 pub(crate) async fn check_fence(
549 &self,
550 writer: &mut dyn SqlWriter,
551 ) -> Result<Option<NoteWriteConflict>, StorageError> {
552 let Some(fences) = &self.fence else {
553 return Ok(None);
554 };
555 let mut now: Option<i64> = None;
559 for (index, fence) in fences.entries().iter().enumerate() {
560 let holder = crate::fence_identity::read_holder(
561 writer,
562 &self.namespace,
563 &fence.kind,
564 &fence.key,
565 "note-write-guard",
566 )
567 .await?;
568 if let (Some(asserted), Some(holder)) = (fence.id, holder.as_ref()) {
574 if asserted != holder.id {
575 return Ok(Some(NoteWriteConflict::FenceIdentity(Box::new(
576 FenceIdentity {
577 key: fence.key.clone(),
578 kind: fence.kind.clone(),
579 version: fence.expected_version.unwrap_or_default(),
582 index: matches!(fences, NoteFences::Many(_)).then_some(index),
583 evidence: crate::fence_identity::identity_evidence(asserted, holder.id),
584 },
585 ))));
586 }
587 }
588 let current = holder.map(|holder| holder.version);
589 if current != fence.expected_version {
590 return Ok(Some(NoteWriteConflict::Fence {
591 key: fence.key.clone(),
592 expected: fence.expected_version,
593 current,
594 index: matches!(fences, NoteFences::Many(_)).then_some(index),
595 }));
596 }
597 if let Some(field) = &fence.live_until {
598 let clock = match now {
599 Some(clock) => clock,
600 None => {
601 let clock =
602 crate::live_until::writer_clock(writer, "note-write-guard-clock")
603 .await?;
604 now = Some(clock);
605 clock
606 }
607 };
608 if let Some(refusal) = crate::live_until::evaluate(
609 writer,
610 &self.namespace,
611 &fence.kind,
612 &fence.key,
613 field,
614 clock,
615 "note-write-guard-live-until",
616 )
617 .await?
618 {
619 return Ok(Some(NoteWriteConflict::FenceDeadline(Box::new(
620 FenceDeadline {
621 key: fence.key.clone(),
622 kind: fence.kind.clone(),
623 version: fence.expected_version.unwrap_or_default(),
626 field: field.clone(),
627 index: matches!(fences, NoteFences::Many(_)).then_some(index),
628 reason: refusal.reason(),
629 evidence: refusal.details(),
630 },
631 ))));
632 }
633 }
634 }
635 Ok(None)
636 }
637
638 pub(crate) async fn classify_refusal(
639 &self,
640 writer: &mut dyn SqlWriter,
641 ) -> Result<Option<NoteWriteConflict>, StorageError> {
642 if let Some(claim) = &self.create_key {
643 let holder = writer
648 .query_row(statement(
649 "SELECT id, content, properties FROM notes \
650 WHERE namespace=?1 AND kind=?2 AND key=?3 AND deleted_at IS NULL",
651 vec![
652 SqlValue::Text(self.namespace.clone()),
653 SqlValue::Text(claim.kind.clone()),
654 SqlValue::Text(claim.key.clone()),
655 ],
656 ))
657 .await?;
658 if let Some(row) = holder {
659 let existing_id = match row.get("id") {
660 Some(SqlValue::Text(id)) => id.clone(),
661 _ => {
662 return Err(StorageError::Internal(
663 "invalid keyed note holder identity".into(),
664 ))
665 }
666 };
667 let existing_content = match row.get("content") {
668 Some(SqlValue::Text(content)) => content.clone(),
669 _ => {
670 return Err(StorageError::Internal(
671 "invalid keyed note holder content".into(),
672 ))
673 }
674 };
675 let existing_properties: Option<serde_json::Value> = match row.get("properties") {
676 Some(SqlValue::Text(text)) => {
677 Some(serde_json::from_str(text).map_err(|_| {
678 StorageError::Internal("invalid keyed note holder properties".into())
679 })?)
680 }
681 Some(SqlValue::Null) | None => None,
682 _ => {
683 return Err(StorageError::Internal(
684 "invalid keyed note holder properties".into(),
685 ))
686 }
687 };
688 let equal = claim.replay_signal.then(|| {
702 existing_content == claim.content && existing_properties == claim.properties
703 });
704 return Ok(Some(NoteWriteConflict::Key {
705 key: claim.key.clone(),
706 existing_id,
707 equal,
708 }));
709 }
710 }
711 if let Some(expected) = self.expected_version {
712 let current = writer
713 .query_scalar(statement(
714 "SELECT version FROM notes WHERE id=?1 AND deleted_at IS NULL",
715 vec![SqlValue::Text(self.target_id.to_string())],
716 ))
717 .await?;
718 if let Some(SqlValue::Integer(current)) = current {
719 if current != expected {
720 return Ok(Some(NoteWriteConflict::Version { expected, current }));
721 }
722 }
723 }
724 Ok(None)
725 }
726}
727
728async fn check_keyed_create_holder(
742 access: &dyn khive_storage::SqlAccess,
743 guard: NoteWriteGuard,
744) -> RuntimeResult<Option<NoteWriteConflict>> {
745 use crate::atomic_runner::{
746 run_prepared_atomic_unit, PreparedAtomicError, PreparedAtomicOp, PreparedAtomicOutcome,
747 };
748 let op: PreparedAtomicOp<(), NoteWriteConflict> = Box::new(move |writer| {
749 Box::pin(async move {
750 if let Some(conflict) = guard.check_fence(writer).await? {
751 return Err(PreparedAtomicError::Refused {
752 failure: conflict,
753 message: "keyed create: fence stale at holder check".into(),
754 });
755 }
756 if let Some(conflict) = guard.classify_refusal(writer).await? {
757 return Err(PreparedAtomicError::Refused {
758 failure: conflict,
759 message: "keyed create: live holder found at holder check".into(),
760 });
761 }
762 Ok(((), Vec::new()))
763 })
764 });
765 match run_prepared_atomic_unit(access, op).await? {
766 PreparedAtomicOutcome::Committed { .. } => Ok(None),
767 PreparedAtomicOutcome::RolledBack(conflict) => Ok(Some(conflict)),
768 }
769}
770
771pub(crate) fn validate_head(note: &khive_storage::note::Note) -> RuntimeResult<()> {
772 if note.kind != "head" {
773 return Ok(());
774 }
775 if note.name.is_some() {
776 return Err(RuntimeError::InvalidInput("head notes have no name".into()));
777 }
778 serde_json::from_str::<serde_json::Value>(¬e.content).map_err(|error| {
779 RuntimeError::InvalidInput(format!("head content must be JSON text: {error}"))
780 })?;
781 if let Some(tags) = note
782 .properties
783 .as_ref()
784 .and_then(|p| p.get("tags"))
785 .and_then(|t| t.as_array())
786 {
787 for tag in tags.iter().filter_map(|tag| tag.as_str()) {
788 if let Some(kind) = tag.strip_prefix("kind:") {
789 if kind.len() > 64 || kind.contains('\0') {
790 return Err(RuntimeError::InvalidInput(
791 "head document kind must be at most 64 bytes without U+0000".into(),
792 ));
793 }
794 }
795 }
796 }
797 Ok(())
798}
799
800impl KhiveRuntime {
801 pub async fn get_note_by_key(
802 &self,
803 token: &NamespaceToken,
804 key: &str,
805 kind: Option<&str>,
806 after_key: bool,
807 ) -> RuntimeResult<khive_storage::note::Note> {
808 crate::keyed_memory::validate_memory_key(key)?;
809 let mut matches = self
810 .notes(token)?
811 .get_live_notes_by_key(token.namespace().as_str(), key, kind)
812 .await?;
813 match matches.len() {
814 0 => {
815 let error = KhiveError::not_found("note key", key);
816 Err(if after_key {
817 error.with_details(Details::new_owned([
818 ("reason", "after_key_missing".into()),
819 ("key", key.into()),
820 ]))
821 } else {
822 error
823 }
824 .into())
825 }
826 1 => Ok(matches.remove(0)),
827 _ => {
828 let kinds = matches
829 .iter()
830 .map(|note| note.kind.as_str())
831 .collect::<Vec<_>>()
832 .join(",");
833 Err(KhiveError::conflict("note key is ambiguous")
834 .with_details(Details::new_owned([
835 ("reason", "key_ambiguous".into()),
836 ("key", key.into()),
837 ("kinds", kinds),
838 ]))
839 .into())
840 }
841 }
842 }
843
844 pub(crate) async fn prepare_versioned_note_update(
845 &self,
846 token: &NamespaceToken,
847 snapshot: khive_storage::note::Note,
848 patch: crate::curation::NotePatch,
849 ) -> RuntimeResult<(khive_storage::note::Note, crate::atomic_plan::UpdatePlan)> {
850 use crate::atomic_plan::{AffectedRowGuard, PlanStatement, PostCommitEffect, UpdatePlan};
851 let options = patch.write_options.clone();
852 options.validate()?;
853 if options.key.is_some() {
854 return Err(RuntimeError::InvalidInput("key is immutable".into()));
855 }
856 if let Some(fences) = &options.fence {
857 for fence in fences.entries() {
858 self.validate_note_kind(&fence.kind)?;
859 }
860 }
861 let expected_updated_at = snapshot.updated_at;
862 let expected_deleted_at = snapshot.deleted_at;
863 let (mut note, text_changed, changed) = self
864 .prepare_update_note_from_snapshot(token, snapshot, patch)
865 .await?;
866 validate_head(¬e)?;
867 if !changed && options.embed.is_none() && options.expected_version.is_none() {
874 let mut assertion = SqlStatement {
875 sql: "SELECT 1 FROM notes WHERE id=?1 AND updated_at=?2 AND deleted_at IS ?3"
876 .into(),
877 params: vec![
878 SqlValue::Text(note.id.to_string()),
879 SqlValue::Integer(expected_updated_at),
880 expected_deleted_at
881 .map(SqlValue::Integer)
882 .unwrap_or(SqlValue::Null),
883 ],
884 label: Some("note-noop-assertion".into()),
885 };
886 if let Some(version) = options.expected_version {
887 assertion.params.push(SqlValue::Integer(version));
888 assertion
889 .sql
890 .push_str(&format!(" AND version = ?{}", assertion.params.len()));
891 }
892 let plan = UpdatePlan {
893 graph_effects: Vec::new(),
894 target_id: note.id,
895 statements: vec![PlanStatement {
896 statement: assertion,
897 guard: Some(AffectedRowGuard::exactly(1)),
898 }],
899 post_commit: PostCommitEffect::None,
900 edge_natural_key: None,
901 idempotent_noop: true,
902 entity_guard: None,
903 note_guard: Some(NoteWriteGuard {
904 namespace: token.namespace().as_str().into(),
905 target_id: note.id,
906 expected_version: options.expected_version,
907 fence: options.fence,
908 create_key: None,
909 }),
910 note_vector_purge: None,
911 note_embedding_inheritance: None,
912 };
913 return Ok((note, plan));
914 }
915 if !changed {
916 let minimum_updated_at = note.updated_at.checked_add(1).ok_or_else(|| {
921 RuntimeError::Internal(format!(
922 "note {} updated_at is already at i64::MAX and cannot advance",
923 note.id
924 ))
925 })?;
926 note.updated_at = chrono::Utc::now()
927 .timestamp_micros()
928 .max(minimum_updated_at);
929 }
930 let next_version = note
931 .version
932 .checked_add(1)
933 .ok_or_else(|| RuntimeError::InvalidInput("note version exhausted".into()))?;
934 note.version = next_version;
935 let mut update = if self.stream_member_error(¬e).await?.is_some() {
936 khive_db::stores::note::note_metadata_replace_if_unchanged_statement(
937 ¬e,
938 expected_updated_at,
939 expected_deleted_at,
940 )
941 } else {
942 khive_db::stores::note::note_replace_if_unchanged_statement(
943 ¬e,
944 expected_updated_at,
945 expected_deleted_at,
946 )
947 };
948 if let Some(version) = options.expected_version {
949 update.params.push(SqlValue::Integer(version));
950 update
951 .sql
952 .push_str(&format!(" AND version = ?{}", update.params.len()));
953 }
954 let mut statements = vec![PlanStatement {
955 statement: update,
956 guard: Some(AffectedRowGuard::exactly(1)),
957 }];
958 if text_changed {
959 for sql in khive_db::stores::text::delete_document_statements(
960 "fts_notes",
961 ¬e.namespace,
962 note.id,
963 )
964 .into_iter()
965 .chain(khive_db::stores::text::insert_document_statements(
966 "fts_notes",
967 &crate::curation::note_fts_document(¬e),
968 )) {
969 statements.push(PlanStatement {
970 statement: sql,
971 guard: None,
972 });
973 }
974 }
975 statements.extend(crate::atomic_prepare::event_append_statements(
979 token,
980 ¬e.namespace,
981 "update",
982 khive_types::EventKind::NoteUpdated,
983 khive_types::SubstrateKind::Note,
984 note.id,
985 serde_json::json!({
986 "id": note.id,
987 "namespace": note.namespace,
988 "version": note.version,
989 "text_changed": text_changed,
990 }),
991 )?);
992 let post_commit =
995 if options.embed != Some(false) && (text_changed || options.embed == Some(true)) {
996 PostCommitEffect::ReindexNote {
997 note_id: note.id,
998 version: note.version,
999 }
1000 } else if text_changed || options.embed == Some(false) {
1001 PostCommitEffect::NoteChanged {
1002 note_id: note.id,
1003 kind: note.kind.clone(),
1004 }
1005 } else {
1006 PostCommitEffect::None
1007 };
1008 let plan = UpdatePlan {
1009 graph_effects: Vec::new(),
1010 target_id: note.id,
1011 statements,
1012 post_commit,
1013 edge_natural_key: None,
1014 idempotent_noop: false,
1015 entity_guard: None,
1016 note_guard: Some(NoteWriteGuard {
1017 namespace: token.namespace().as_str().into(),
1018 target_id: note.id,
1019 expected_version: options.expected_version,
1020 fence: options.fence,
1021 create_key: None,
1022 }),
1023 note_vector_purge: (options.embed == Some(false))
1024 .then(|| NoteVectors::new(note.namespace.clone(), note.id)),
1025 note_embedding_inheritance: (text_changed && options.embed.is_none()).then(|| {
1026 NoteEmbeddingInheritance {
1027 vectors: NoteVectors::new(note.namespace.clone(), note.id),
1028 kind: note.kind.clone(),
1029 }
1030 }),
1031 };
1032 Ok((note, plan))
1033 }
1034
1035 #[allow(clippy::too_many_arguments)]
1036 pub async fn create_note_with_options(
1037 &self,
1038 token: &NamespaceToken,
1039 kind: &str,
1040 name: Option<&str>,
1041 content: &str,
1042 embedding_content: Option<&str>,
1043 salience: Option<f64>,
1044 decay_factor: Option<f64>,
1045 properties: Option<serde_json::Value>,
1046 annotates: Vec<Uuid>,
1047 embedding_model: Option<&str>,
1048 options: NoteWriteOptions,
1049 ) -> RuntimeResult<(
1050 khive_storage::note::Note,
1051 crate::retrieval::EmbeddingTruncationReport,
1052 )> {
1053 use crate::atomic_message::{AtomicNoteOptions, AtomicNoteSpec};
1054 use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicRunOutcome};
1055 use crate::note_create::{prepare_note_create, KeyPublication};
1056 options.validate()?;
1057 if options.expected_version.is_some() {
1058 return Err(RuntimeError::InvalidInput(
1059 "expected_version applies only to update".into(),
1060 ));
1061 }
1062 if let Some(fences) = &options.fence {
1063 for fence in fences.entries() {
1064 self.validate_note_kind(&fence.kind)?;
1065 }
1066 }
1067 if let Some(prefix) = embedding_content {
1068 if prefix.is_empty() || prefix.len() >= content.len() || !content.starts_with(prefix) {
1069 return Err(RuntimeError::InvalidInput(
1070 "embedding_content must be a non-empty proper prefix of content".into(),
1071 ));
1072 }
1073 crate::secret_gate::check_at(prefix, "note", "embedding_content")?;
1074 }
1075 let mut candidate =
1076 khive_storage::note::Note::new(token.namespace().as_str(), kind, content);
1077 candidate.name = name.map(str::to_owned);
1078 candidate.properties = properties.clone();
1079 validate_head(&candidate)?;
1080
1081 let Some(key) = options.key.as_deref() else {
1086 let (mut prepared, _) = prepare_note_create(
1087 self,
1088 AtomicNoteSpec {
1089 token,
1090 id: None,
1091 kind,
1092 name,
1093 content,
1094 properties,
1095 },
1096 AtomicNoteOptions {
1097 salience,
1098 decay_factor,
1099 embedding_model,
1100 embedding_content,
1101 embed: Some(options.embed.unwrap_or(kind != "head")),
1102 key: None,
1103 replay_receipt: true,
1104 fence: options.fence.as_ref(),
1105 properties_already_derived: false,
1106 },
1107 &annotates,
1108 KeyPublication::AtInsert,
1109 )
1110 .await?;
1111 let note = prepared.notes.remove(0);
1112 return match run_atomic_unit(self.sql().as_ref(), prepared.plans).await {
1113 Ok(AtomicRunOutcome::Committed { .. }) => {
1114 self.fire_note_mutation_hook(¬e.kind, note.id).await;
1115 Ok((note, prepared.embedding_truncation))
1116 }
1117 Ok(AtomicRunOutcome::RolledBack {
1118 failure: AtomicOpFailure::NoteConflict(conflict),
1119 ..
1120 }) => Err(conflict.into_error().into()),
1121 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
1122 format!("note creation rolled back: {failure:?}"),
1123 )),
1124 Err(error) => Err(RuntimeError::Storage(error.0)),
1125 };
1126 };
1127
1128 let derived_properties =
1138 self.derive_note_write_properties(kind, token, properties.clone())?;
1139 let namespace: String = token.namespace().as_str().into();
1140
1141 if let Some(conflict) = check_keyed_create_holder(
1147 self.sql().as_ref(),
1148 NoteWriteGuard {
1149 namespace: namespace.clone(),
1150 target_id: candidate.id,
1151 expected_version: None,
1152 fence: options.fence.clone(),
1153 create_key: Some(CreateKeyClaim {
1154 kind: kind.to_owned(),
1155 key: key.to_owned(),
1156 content: content.to_owned(),
1157 properties: derived_properties.clone(),
1158 replay_signal: true,
1159 }),
1160 },
1161 )
1162 .await?
1163 {
1164 return Err(conflict.into_error().into());
1165 }
1166
1167 let prep_result = prepare_note_create(
1174 self,
1175 AtomicNoteSpec {
1176 token,
1177 id: None,
1178 kind,
1179 name,
1180 content,
1181 properties: derived_properties.clone(),
1182 },
1183 AtomicNoteOptions {
1184 salience,
1185 decay_factor,
1186 embedding_model,
1187 embedding_content,
1188 embed: Some(options.embed.unwrap_or(kind != "head")),
1189 key: Some(key),
1190 replay_receipt: true,
1191 fence: options.fence.as_ref(),
1192 properties_already_derived: true,
1193 },
1194 &annotates,
1195 KeyPublication::AtInsert,
1196 )
1197 .await;
1198
1199 #[cfg(test)]
1205 race_seam::pause_after_prepare().await;
1206
1207 match prep_result {
1208 Ok((mut prepared, _)) => {
1209 let note = prepared.notes.remove(0);
1210 match run_atomic_unit(self.sql().as_ref(), prepared.plans).await {
1211 Ok(AtomicRunOutcome::Committed { .. }) => {
1212 self.fire_note_mutation_hook(¬e.kind, note.id).await;
1213 Ok((note, prepared.embedding_truncation))
1214 }
1215 Ok(AtomicRunOutcome::RolledBack {
1216 failure: AtomicOpFailure::NoteConflict(conflict),
1217 ..
1218 }) => Err(conflict.into_error().into()),
1219 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(
1220 RuntimeError::Internal(format!("note creation rolled back: {failure:?}")),
1221 ),
1222 Err(error) => Err(RuntimeError::Storage(error.0)),
1223 }
1224 }
1225 Err(prep_error) => {
1226 match check_keyed_create_holder(
1232 self.sql().as_ref(),
1233 NoteWriteGuard {
1234 namespace,
1235 target_id: candidate.id,
1236 expected_version: None,
1237 fence: options.fence.clone(),
1238 create_key: Some(CreateKeyClaim {
1239 kind: kind.to_owned(),
1240 key: key.to_owned(),
1241 content: content.to_owned(),
1242 properties: derived_properties,
1243 replay_signal: true,
1244 }),
1245 },
1246 )
1247 .await?
1248 {
1249 Some(conflict) => Err(conflict.into_error().into()),
1250 None => Err(prep_error),
1251 }
1252 }
1253 }
1254 }
1255}