1use std::any::Any;
6use std::collections::{HashMap, HashSet};
7use std::sync::{
8 atomic::{AtomicBool, Ordering},
9 Arc,
10};
11
12use khive_storage::{
13 AtomicUnitOp, Note, SqlAccess, SqlRow, SqlStatement, SqlValue, SqlWriter, StorageCapability,
14 StorageError, WriterTaskRequestState,
15};
16use khive_types::{Details, KhiveError};
17use serde::Deserialize;
18use serde_json::{json, Value};
19use uuid::Uuid;
20
21use crate::atomic_message::{
22 prepare_atomic_note_requests, AtomicNoteOptions, AtomicNoteRequest, AtomicNoteSpec,
23};
24use crate::atomic_plan::{PlanStatement, PostCommitEffect};
25use crate::atomic_runner::{
26 apply_plan, run_prepared_atomic_unit, AtomicOpFailure, AtomicOpPlan,
27 CommittedPostCommitEffects, PreparedAtomicError, PreparedAtomicOp, PreparedAtomicOutcome,
28};
29use crate::note_write::{
30 NoteFence, NoteFences, NoteWriteConflict, NoteWriteGuard, NoteWriteOptions,
31};
32use crate::{
33 micros_to_iso, DomainDisposition, KhiveRuntime, NamespaceToken, RuntimeError, RuntimeResult,
34 VerbRegistry,
35};
36
37#[cfg(test)]
38mod append_failure_tests;
39
40fn statement(sql: &str, params: Vec<SqlValue>) -> SqlStatement {
41 SqlStatement {
42 sql: sql.into(),
43 params,
44 label: Some("stream".into()),
45 }
46}
47
48fn validate_stream(stream: &str) -> RuntimeResult<()> {
49 if stream.len() > 512 || stream.contains('\0') {
50 return Err(RuntimeError::InvalidInput(
51 "stream must be at most 512 UTF-8 bytes and contain no U+0000".into(),
52 ));
53 }
54 Ok(())
55}
56
57fn integer(row: &SqlRow, name: &str) -> RuntimeResult<i64> {
58 match row.get(name) {
59 Some(SqlValue::Integer(value)) => Ok(*value),
60 _ => Err(RuntimeError::Internal(format!(
61 "stream query missing integer {name}"
62 ))),
63 }
64}
65
66fn text<'a>(row: &'a SqlRow, name: &str) -> RuntimeResult<&'a str> {
67 match row.get(name) {
68 Some(SqlValue::Text(value)) => Ok(value),
69 _ => Err(RuntimeError::Internal(format!(
70 "stream query missing text {name}"
71 ))),
72 }
73}
74
75fn write_failure(message: &str) -> StorageError {
76 StorageError::Conflict {
77 capability: StorageCapability::Notes,
78 operation: "stream.append".into(),
79 message: message.into(),
80 }
81}
82
83#[derive(Clone, Copy, Debug, PartialEq, Eq)]
86pub enum StreamAppendDisposition {
87 NotCommitted,
89 Unknown,
91}
92
93#[derive(Debug, thiserror::Error)]
95#[error("{source}")]
96pub struct StreamAppendFailure {
97 #[source]
98 source: RuntimeError,
99 disposition: StreamAppendDisposition,
100}
101
102impl StreamAppendFailure {
103 fn not_committed(source: RuntimeError) -> Self {
104 Self {
105 source,
106 disposition: StreamAppendDisposition::NotCommitted,
107 }
108 }
109
110 pub fn unknown(source: RuntimeError) -> Self {
113 Self {
114 source,
115 disposition: StreamAppendDisposition::Unknown,
116 }
117 }
118
119 fn after_submission(source: RuntimeError, refused_before_write: bool) -> Self {
120 let disposition = if let Some(context) = source.writer_task_failure_context() {
121 match context.request_state {
122 WriterTaskRequestState::NotStarted
123 | WriterTaskRequestState::TransactionRolledBack => {
124 StreamAppendDisposition::NotCommitted
125 }
126 WriterTaskRequestState::SideEffectsUnknown => StreamAppendDisposition::Unknown,
127 }
128 } else if refused_before_write
129 || source.admission_failure_context().is_some()
130 || matches!(
131 &source,
132 RuntimeError::Storage(StorageError::WriterTaskBusy { .. })
133 )
134 {
135 StreamAppendDisposition::NotCommitted
136 } else {
137 StreamAppendDisposition::Unknown
138 };
139 Self {
140 source,
141 disposition,
142 }
143 }
144
145 pub fn source(&self) -> &RuntimeError {
147 &self.source
148 }
149
150 pub fn disposition(&self) -> StreamAppendDisposition {
152 self.disposition
153 }
154
155 pub fn into_parts(self) -> (RuntimeError, StreamAppendDisposition) {
157 (self.source, self.disposition)
158 }
159
160 pub fn into_source(self) -> RuntimeError {
162 self.source
163 }
164}
165
166pub struct StreamAppendSpec {
169 pub stream: String,
170 pub record: Value,
171 pub expected_seq: Option<i64>,
172 pub embed: Option<bool>,
173 pub embedding_model: Option<String>,
174 pub note_kind: String,
175 pub tags: Option<Vec<String>>,
176 pub fence: Option<NoteFences>,
177}
178
179impl StreamAppendSpec {
180 fn validate_embedding(&self, member: Option<usize>) -> RuntimeResult<()> {
181 if self.embedding_model.is_some() && self.embed != Some(true) {
182 let mut error = KhiveError::invalid_input("embedding_model requires embed=true");
183 if let Some(member) = member {
184 error = error.with_details(Details::new_owned([("member", member.to_string())]));
185 }
186 return Err(error.into());
187 }
188 Ok(())
189 }
190
191 fn note_options(&self) -> AtomicNoteOptions<'_> {
192 AtomicNoteOptions {
193 embed: Some(self.embed.unwrap_or(false)),
194 embedding_model: self.embedding_model.as_deref(),
195 ..Default::default()
196 }
197 }
198}
199
200#[derive(Clone, Debug)]
203pub struct StreamWriteSpec {
204 pub key: String,
205 pub kind: String,
206 pub doc: Value,
207 pub tags: Option<Vec<String>>,
208 pub embed: Option<bool>,
209 pub expected_version: Option<i64>,
210}
211
212#[derive(Clone, Debug)]
215pub struct StreamObservation {
216 pub key: String,
217 pub kind: String,
218 pub version: Option<i64>,
219 pub id: Option<Uuid>,
220 pub live_until: Option<String>,
222}
223
224pub enum StreamBatchMember {
226 Append(StreamAppendSpec),
227 Write(StreamWriteSpec),
228 Refused(KhiveError),
229}
230
231const MAX_BATCH_FENCE_ENTRIES: usize = 100;
234
235fn stream_batch_fence_count(members: &[StreamBatchMember], fence: Option<&NoteFence>) -> usize {
236 usize::from(fence.is_some())
237 + members
238 .iter()
239 .filter_map(|member| match member {
240 StreamBatchMember::Append(spec) => spec.fence.as_ref(),
241 StreamBatchMember::Write(_) | StreamBatchMember::Refused(_) => None,
242 })
243 .map(|fences| fences.entries().len())
244 .sum::<usize>()
245}
246
247fn validate_stream_batch_fence_count(count: usize) -> RuntimeResult<()> {
248 if count > MAX_BATCH_FENCE_ENTRIES {
249 return Err(RuntimeError::InvalidInput(format!(
250 "stream.batch admits at most {MAX_BATCH_FENCE_ENTRIES} total fence entries; this call sent {count}: each fence is a read taken while holding the writer"
251 )));
252 }
253 Ok(())
254}
255
256#[derive(Debug)]
258pub struct StreamBatchRefusal {
259 pub member: usize,
260 pub error: KhiveError,
261}
262
263pub fn refusal_value(error: &KhiveError) -> RuntimeResult<Value> {
266 let mut value = serde_json::to_value(error)
267 .map_err(|e| RuntimeError::Internal(format!("stream refusal serialization: {e}")))?;
268 value["domain_disposition"] = json!(DomainDisposition::NotCommitted.as_str());
269 Ok(value)
270}
271
272fn seq_conflict(stream: &str, expected: i64, next: i64, member: Option<usize>) -> KhiveError {
273 let mut pairs = vec![
274 ("reason", "seq_conflict".to_string()),
275 ("stream", stream.to_string()),
276 ("expected_seq", expected.to_string()),
277 ("next_seq", next.to_string()),
278 ];
279 if let Some(member) = member {
280 pairs.push(("member", member.to_string()));
281 }
282 KhiveError::conflict("stream sequence precondition failed")
283 .with_details(Details::new_owned(pairs))
284}
285
286struct PreparedAppend {
288 stream: String,
289 expected_seq: Option<i64>,
290 note: Note,
291 statements: Vec<PlanStatement>,
292 guard: NoteWriteGuard,
293}
294
295fn append_result(prepared: &PreparedAppend, seq: i64) -> Value {
296 json!({"seq": seq, "id": prepared.note.id, "created_at": micros_to_iso(prepared.note.created_at)})
297}
298
299enum BatchOutcome {
300 Appended(Vec<i64>),
301 Conflict { next: i64 },
302 FenceConflict { conflict: NoteWriteConflict },
303}
304
305struct PreparedBatchMember {
306 index: usize,
307 action: PreparedBatchAction,
308}
309
310enum PreparedBatchAction {
311 Append {
312 stream: String,
313 expected_seq: Option<i64>,
314 note: Box<Note>,
315 plan: AtomicOpPlan,
316 fence: Option<NoteFences>,
317 },
318 Write {
319 key: String,
320 kind: String,
321 id: Uuid,
322 plan: AtomicOpPlan,
323 after_create: Option<Value>,
324 },
325 Refused(KhiveError),
326}
327
328#[derive(Deserialize)]
329struct StreamCreateFields {
330 content: String,
331 name: Option<String>,
332 properties: Option<Value>,
333 tags: Option<Vec<String>>,
334 salience: Option<f64>,
335 decay_factor: Option<f64>,
336 embed: Option<bool>,
337}
338
339fn normalize_stream_create_fields(fields: &mut StreamCreateFields) -> RuntimeResult<()> {
340 let tags = fields.tags.take();
341 if let crate::EffectiveCreateTags::TopLevel(tags) =
342 crate::effective_create_tags(tags.as_deref(), fields.properties.as_ref())
343 {
344 let tags = json!(tags);
345 let mut properties = match fields.properties.take() {
346 None => serde_json::Map::new(),
347 Some(Value::Object(properties)) => properties,
348 Some(_) => {
349 return Err(RuntimeError::InvalidInput(
350 "note tags require object properties".into(),
351 ))
352 }
353 };
354 properties.insert("tags".into(), tags);
355 fields.properties = Some(Value::Object(properties));
356 }
357 Ok(())
358}
359
360struct PreparedStreamCreate {
361 key: String,
362 kind: String,
363 fields: StreamCreateFields,
364 args: Value,
365}
366
367enum StreamBatchPreparation {
368 Ready(Box<PreparedBatchAction>),
369 Append(Box<(StreamAppendSpec, String)>),
370 Create(Box<PreparedStreamCreate>),
371}
372
373fn batch_create_hooks(members: &[PreparedBatchMember]) -> Vec<(Uuid, String, Value)> {
374 members
375 .iter()
376 .filter_map(|member| match &member.action {
377 PreparedBatchAction::Write {
378 id,
379 kind,
380 after_create: Some(args),
381 ..
382 } => Some((*id, kind.clone(), args.clone())),
383 _ => None,
384 })
385 .collect()
386}
387
388fn place_member(error: KhiveError, member: Option<usize>) -> RuntimeResult<KhiveError> {
389 let pairs = member
390 .into_iter()
391 .map(|member| ("member".to_owned(), member.to_string()))
392 .chain(
393 error
394 .details()
395 .into_iter()
396 .flat_map(Details::iter)
397 .filter(|(key, _)| *key != "member")
398 .map(|(key, value)| (key.to_owned(), value.to_owned())),
399 );
400 let details = Details::deserialize(serde::de::value::MapDeserializer::<
401 _,
402 serde::de::value::Error,
403 >::new(pairs))
404 .map_err(|error| RuntimeError::Internal(format!("stream member details: {error}")))?;
405 Ok(error.with_details(details))
406}
407
408fn missing_write(key: &str) -> KhiveError {
409 KhiveError::not_found("note key", key).with_details(Details::new_owned([
410 ("reason", "stream_write_not_found".into()),
411 ("key", key.into()),
412 ]))
413}
414
415async fn check_observed(
416 writer: &mut dyn SqlWriter,
417 namespace: &str,
418 observed: &[StreamObservation],
419 now: Option<i64>,
420) -> Result<Option<KhiveError>, StorageError> {
421 for (index, entry) in observed.iter().enumerate() {
422 let current = crate::fence_identity::read_holder(
423 writer,
424 namespace,
425 &entry.kind,
426 &entry.key,
427 "stream-batch-observed",
428 )
429 .await?;
430 if let (Some(asserted), Some(holder)) = (entry.id, current.as_ref()) {
433 if asserted != holder.id {
434 let version = entry.version.ok_or_else(|| {
435 StorageError::Internal("identity observation missing version".into())
436 })?;
437 let mut details = vec![
438 ("reason", "identity_conflict".into()),
439 ("key", entry.key.clone()),
440 ("kind", entry.kind.clone()),
441 ("version", version.to_string()),
442 ];
443 details.extend(crate::fence_identity::identity_evidence(
444 asserted, holder.id,
445 ));
446 details.push(("index", index.to_string()));
447 return Ok(Some(
448 KhiveError::conflict("stream observation identity precondition failed")
449 .with_details(Details::new_owned(details)),
450 ));
451 }
452 }
453 let current = current.map(|holder| holder.version);
454 if current != entry.version {
455 let mut details = vec![
456 ("reason", "version_conflict".into()),
457 ("key", entry.key.clone()),
458 ("index", index.to_string()),
459 ];
460 if let Some(expected) = entry.version {
461 details.push(("expected_version", expected.to_string()));
462 }
463 if let Some(current) = current {
464 details.push(("current_version", current.to_string()));
465 }
466 return Ok(Some(
467 KhiveError::conflict("stream observation precondition failed")
468 .with_details(Details::new_owned(details)),
469 ));
470 }
471 if let Some(field) = &entry.live_until {
472 let now = now
473 .ok_or_else(|| StorageError::Internal("missing stream observation clock".into()))?;
474 if let Some(refusal) = crate::live_until::evaluate(
475 writer,
476 namespace,
477 &entry.kind,
478 &entry.key,
479 field,
480 now,
481 "stream-batch-live-until",
482 )
483 .await?
484 {
485 let mut details = vec![
486 ("reason", refusal.reason().into()),
487 ("key", entry.key.clone()),
488 ("kind", entry.kind.clone()),
489 ("version", entry.version.unwrap().to_string()),
490 ("field", field.clone()),
491 ("index", index.to_string()),
492 ];
493 details.extend(refusal.details());
494 return Ok(Some(
495 KhiveError::conflict("stream observation time precondition failed")
496 .with_details(Details::new_owned(details)),
497 ));
498 }
499 }
500 }
501 Ok(None)
502}
503
504async fn observation_clock(writer: &mut dyn SqlWriter) -> Result<i64, StorageError> {
505 crate::live_until::writer_clock(writer, "stream-batch-clock").await
506}
507
508struct BatchFailure {
509 member: Option<usize>,
510 error: RuntimeError,
511}
512
513async fn run_prepared_stream_batch(
514 access: &dyn SqlAccess,
515 namespace: String,
516 members: Vec<PreparedBatchMember>,
517 fence: Option<NoteFence>,
518 observed: Vec<StreamObservation>,
519) -> RuntimeResult<Result<(Vec<Value>, CommittedPostCommitEffects), StreamBatchRefusal>> {
520 let op: PreparedAtomicOp<Vec<Value>, BatchFailure> = Box::new(move |writer| {
521 Box::pin(async move {
522 let now = if observed.iter().any(|entry| entry.live_until.is_some()) {
524 Some(observation_clock(writer).await?)
525 } else {
526 None
527 };
528 let guard = NoteWriteGuard {
529 namespace: namespace.clone(),
530 target_id: Uuid::nil(),
531 expected_version: None,
532 fence: fence.map(Into::into),
533 create_key: None,
534 };
535 let predicate_error = if let Some(conflict) = guard.check_fence(writer).await? {
536 Some(conflict.into_error())
537 } else {
538 check_observed(writer, &namespace, &observed, now).await?
539 };
540 if let Some(error) = predicate_error {
541 return Err(PreparedAtomicError::Refused {
542 failure: BatchFailure {
543 member: None,
544 error: error.into(),
545 },
546 message: "stream batch predicate refused".into(),
547 });
548 }
549 let mut results = Vec::with_capacity(members.len());
550 let mut effects = Vec::new();
551 let mut heads = HashMap::new();
552 for member in &members {
555 if let PreparedBatchAction::Append { fence, .. } = &member.action {
556 let guard = NoteWriteGuard {
557 namespace: namespace.clone(),
558 target_id: Uuid::nil(),
559 expected_version: None,
560 fence: fence.clone(),
561 create_key: None,
562 };
563 if let Some(conflict) = guard.check_fence(writer).await? {
564 return Err(PreparedAtomicError::Refused {
565 failure: BatchFailure {
566 member: Some(member.index),
567 error: conflict.into_error().into(),
568 },
569 message: "stream append fence refused".into(),
570 });
571 }
572 }
573 }
574 for member in members {
575 let result =
576 apply_stream_member(writer, &namespace, &mut heads, member.action).await;
577 match result {
578 Ok((value, effect)) => {
579 results.push(value);
580 if let Some(effect) = effect {
581 effects.push(effect);
582 }
583 }
584 Err(error) => {
585 return Err(PreparedAtomicError::Refused {
586 failure: BatchFailure {
587 member: Some(member.index),
588 error,
589 },
590 message: "stream batch member refused".into(),
591 });
592 }
593 }
594 }
595 Ok((results, effects))
596 })
597 });
598 match run_prepared_atomic_unit(access, op).await? {
599 PreparedAtomicOutcome::Committed { value, post_commit } => Ok(Ok((value, post_commit))),
600 PreparedAtomicOutcome::RolledBack(BatchFailure {
601 member: Some(member),
602 error: RuntimeError::Khive(error),
603 }) => Ok(Err(StreamBatchRefusal {
604 member,
605 error: place_member(error, Some(member))?,
606 })),
607 PreparedAtomicOutcome::RolledBack(BatchFailure { error, .. }) => Err(error),
608 }
609}
610
611enum SequenceRefusal {
612 Exhausted,
613 Conflict { expected: i64, next: i64 },
614}
615
616fn allocate_sequence(head: &mut i64, expected: Option<i64>) -> Result<i64, SequenceRefusal> {
617 let next = head.checked_add(1).ok_or(SequenceRefusal::Exhausted)?;
618 if let Some(expected) = expected.filter(|expected| *expected != next) {
619 return Err(SequenceRefusal::Conflict { expected, next });
620 }
621 *head = next;
622 Ok(next)
623}
624
625async fn insert_stream_entry(
626 writer: &mut dyn SqlWriter,
627 namespace: &str,
628 stream: String,
629 seq: i64,
630 note_id: String,
631) -> Result<(), StorageError> {
632 writer
633 .execute(statement(
634 "INSERT INTO note_streams(namespace,stream,seq,note_id) VALUES (?1,?2,?3,?4)",
635 vec![
636 SqlValue::Text(namespace.into()),
637 SqlValue::Text(stream),
638 SqlValue::Integer(seq),
639 SqlValue::Text(note_id),
640 ],
641 ))
642 .await?;
643 Ok(())
644}
645
646async fn apply_stream_member(
647 writer: &mut dyn SqlWriter,
648 namespace: &str,
649 heads: &mut HashMap<String, i64>,
650 action: PreparedBatchAction,
651) -> RuntimeResult<(Value, Option<PostCommitEffect>)> {
652 match action {
653 PreparedBatchAction::Refused(error) => Err(error.into()),
654 PreparedBatchAction::Append {
655 stream,
656 expected_seq,
657 note,
658 plan,
659 ..
660 } => {
661 let head = match heads.entry(stream.clone()) {
662 std::collections::hash_map::Entry::Occupied(entry) => entry.into_mut(),
663 std::collections::hash_map::Entry::Vacant(entry) => {
664 let head = writer.query_scalar(statement(
665 "SELECT COALESCE(MAX(seq), 0) FROM note_streams WHERE namespace=?1 AND stream=?2",
666 vec![SqlValue::Text(namespace.into()), SqlValue::Text(stream.clone())],
667 )).await?;
668 let Some(SqlValue::Integer(head)) = head else {
669 return Err(RuntimeError::Internal("invalid stream head".into()));
670 };
671 entry.insert(head)
672 }
673 };
674 let next = allocate_sequence(head, expected_seq).map_err(|refusal| match refusal {
675 SequenceRefusal::Exhausted => {
676 RuntimeError::InvalidInput("stream sequence exhausted".into())
677 }
678 SequenceRefusal::Conflict { expected, next } => {
679 seq_conflict(&stream, expected, next, None).into()
680 }
681 })?;
682 let applied = apply_plan(writer, &plan, false).await.map_err(|error| {
683 RuntimeError::Internal(format!("stream append plan failed: {error:?}"))
684 })?;
685 insert_stream_entry(writer, namespace, stream, next, note.id.to_string()).await?;
686 Ok((
687 json!({"seq": next, "id": note.id, "created_at": micros_to_iso(note.created_at)}),
688 applied.effect,
689 ))
690 }
691 PreparedBatchAction::Write {
692 key,
693 kind,
694 id,
695 plan,
696 ..
697 } => match apply_plan(writer, &plan, false).await {
698 Ok(applied) => {
699 let stored = writer
702 .query_row(SqlStatement {
703 sql: "SELECT version, updated_at FROM notes WHERE namespace=?1 AND id=?2"
704 .into(),
705 params: vec![
706 SqlValue::Text(namespace.into()),
707 SqlValue::Text(id.to_string()),
708 ],
709 label: Some("stream-batch-write-time".into()),
710 })
711 .await?;
712 let Some(stored) = stored else {
713 return Err(RuntimeError::Internal(
714 "missing stream write timestamp".into(),
715 ));
716 };
717 Ok((
718 json!({"id": id, "version": integer(&stored, "version")?, "updated_at": micros_to_iso(integer(&stored, "updated_at")?)}),
719 applied.effect,
720 ))
721 }
722 Err(AtomicOpFailure::NoteConflict(conflict)) => Err(conflict.into_error().into()),
723 Err(AtomicOpFailure::GuardFailed { .. }) => {
724 let holder = writer.query_row(statement(
725 "SELECT id, version FROM notes WHERE namespace=?1 AND kind=?2 AND key=?3 AND deleted_at IS NULL",
726 vec![SqlValue::Text(namespace.into()), SqlValue::Text(kind), SqlValue::Text(key.clone())],
727 )).await?;
728 let Some(holder) = holder else {
729 return Err(missing_write(&key).into());
730 };
731 if text(&holder, "id")? != id.to_string() {
734 if let AtomicOpPlan::Update(update) = &plan {
735 if let Some(expected) = update
736 .note_guard
737 .as_ref()
738 .and_then(|guard| guard.expected_version)
739 {
740 return Err(NoteWriteConflict::Version {
741 expected,
742 current: integer(&holder, "version")?,
743 }
744 .into_error()
745 .into());
746 }
747 }
748 }
749 Err(crate::curation::stale_note_snapshot_error(id))
750 }
751 Err(error) => Err(RuntimeError::Internal(format!(
752 "stream write plan failed: {error:?}"
753 ))),
754 },
755 }
756}
757
758impl KhiveRuntime {
759 async fn prepare_stream_appends(
763 &self,
764 token: &NamespaceToken,
765 specs: &[&StreamAppendSpec],
766 registry: &VerbRegistry,
767 ) -> RuntimeResult<Vec<PreparedAppend>> {
768 let mut fields = Vec::with_capacity(specs.len());
769 for spec in specs {
770 spec.validate_embedding(None)?;
771 validate_stream(&spec.stream)?;
772 if let Some(fences) = &spec.fence {
773 fences.validate()?;
774 for fence in fences.entries() {
775 self.validate_note_kind(&fence.kind)?;
776 }
777 }
778 let content = serde_json::to_string(&spec.record)
779 .map_err(|e| RuntimeError::InvalidInput(e.to_string()))?;
780 let hook = registry.find_kind_hook(&spec.note_kind);
781 let mut args = if hook.is_some() {
785 match &spec.record {
786 Value::Object(record) => Value::Object(record.clone()),
787 _ => json!({}),
788 }
789 } else {
790 json!({})
791 };
792 args["kind"] = json!("note");
793 args["note_kind"] = json!(&spec.note_kind);
794 args["content"] = json!(content);
795 args["namespace"] = json!(token.namespace().as_str());
796 if let Some(tags) = &spec.tags {
797 args["tags"] = json!(tags);
798 }
799 if let Some(embed) = spec.embed {
800 args["embed"] = json!(embed);
801 }
802 if let Some(hook) = hook {
803 hook.prepare_create(self, &mut args).await?;
804 }
805 let mut note_fields: StreamCreateFields =
806 serde_json::from_value(args).map_err(|error| {
807 RuntimeError::InvalidInput(format!("stream append fields: {error}"))
808 })?;
809 normalize_stream_create_fields(&mut note_fields)?;
810 fields.push(note_fields);
811 }
812 let atomic_specs = specs
813 .iter()
814 .zip(&fields)
815 .map(|(spec, fields)| AtomicNoteRequest {
816 spec: AtomicNoteSpec {
817 token,
818 id: None,
819 kind: &spec.note_kind,
820 name: fields.name.as_deref(),
821 content: &fields.content,
822 properties: fields.properties.clone(),
823 },
824 options: AtomicNoteOptions {
825 salience: fields.salience,
826 decay_factor: fields.decay_factor,
827 embed: Some(fields.embed.or(spec.embed).unwrap_or(false)),
828 embedding_model: spec.embedding_model.as_deref(),
829 ..Default::default()
830 },
831 })
832 .collect();
833 let prepared = prepare_atomic_note_requests(self, atomic_specs).await?;
834 let mut out = Vec::with_capacity(specs.len());
835 for ((spec, note), plan) in specs.iter().zip(prepared.notes).zip(prepared.plans) {
836 let AtomicOpPlan::AddNote(plan) = plan else {
837 return Err(RuntimeError::Internal(
838 "stream preparation did not produce a note".into(),
839 ));
840 };
841 let plan = *plan;
842 let mut guard = plan.note_guard.ok_or_else(|| {
843 RuntimeError::Internal("stream preparation did not produce a note guard".into())
844 })?;
845 guard.fence = spec.fence.clone();
846 out.push(PreparedAppend {
847 stream: spec.stream.clone(),
848 expected_seq: spec.expected_seq,
849 note,
850 statements: plan.statements,
851 guard,
852 });
853 }
854 Ok(out)
855 }
856
857 async fn run_stream_appends(
862 access: &dyn SqlAccess,
863 token: &NamespaceToken,
864 appends: &[PreparedAppend],
865 ) -> Result<BatchOutcome, StreamAppendFailure> {
866 let ns = token.namespace().as_str().to_string();
867 let entries: Vec<_> = appends
868 .iter()
869 .map(|a| {
870 (
871 a.stream.clone(),
872 a.expected_seq,
873 a.note.id.to_string(),
874 a.statements.clone(),
875 a.guard.clone(),
876 )
877 })
878 .collect();
879 let refused_before_write = Arc::new(AtomicBool::new(false));
882 let refusal_proof = Arc::clone(&refused_before_write);
883 let op: AtomicUnitOp = Box::new(move |writer| {
884 Box::pin(async move {
885 let mut writes_started = false;
886 let result = async {
887 let mut heads: Vec<(String, i64)> = Vec::new();
888 for (stream, _, _, _, _) in &entries {
889 if heads.iter().any(|(known, _)| known == stream) {
890 continue;
891 }
892 let head = writer
893 .query_scalar(statement(
894 "SELECT COALESCE(MAX(seq), 0) FROM note_streams WHERE namespace=?1 AND stream=?2",
895 vec![SqlValue::Text(ns.clone()), SqlValue::Text(stream.clone())],
896 ))
897 .await?;
898 let Some(SqlValue::Integer(head)) = head else {
899 return Err(write_failure("invalid stream head"));
900 };
901 heads.push((stream.clone(), head));
902 }
903 let mut assigned = Vec::with_capacity(entries.len());
904 for (stream, expected_seq, _, _, guard) in &entries {
905 if let Some(conflict) = guard.check_fence(writer).await? {
906 return Ok(Box::new(BatchOutcome::FenceConflict { conflict })
907 as Box<dyn Any + Send>);
908 }
909 let head = heads
910 .iter_mut()
911 .find(|(known, _)| known == stream)
912 .map(|(_, head)| head)
913 .ok_or_else(|| write_failure("stream head missing"))?;
914 let next = match allocate_sequence(head, *expected_seq) {
915 Ok(next) => next,
916 Err(SequenceRefusal::Exhausted) => {
917 return Err(write_failure("stream sequence exhausted"));
918 }
919 Err(SequenceRefusal::Conflict { next, .. }) => {
920 return Ok(
921 Box::new(BatchOutcome::Conflict { next }) as Box<dyn Any + Send>
922 );
923 }
924 };
925 assigned.push(next);
926 }
927 for ((stream, _, note_id, statements, _), seq) in
928 entries.into_iter().zip(assigned.iter().copied())
929 {
930 writes_started = true;
933 for planned in statements {
934 let affected = writer.execute(planned.statement).await?;
935 if planned
936 .guard
937 .is_some_and(|guard| !guard.holds_for(affected))
938 {
939 return Err(write_failure("prepared note write guard failed"));
940 }
941 }
942 insert_stream_entry(writer, &ns, stream, seq, note_id).await?;
943 }
944 Ok(Box::new(BatchOutcome::Appended(assigned)) as Box<dyn Any + Send>)
945 }.await;
946 if result.is_err() && !writes_started {
947 refusal_proof.store(true, Ordering::Release);
948 }
949 result
950 })
951 });
952 let outcome = access
953 .atomic_unit(op)
954 .await
955 .map_err(|source| {
956 StreamAppendFailure::after_submission(
957 source.into(),
958 refused_before_write.load(Ordering::Acquire),
959 )
960 })?
961 .downcast::<BatchOutcome>()
962 .map_err(|_| {
963 StreamAppendFailure::unknown(RuntimeError::Internal(
964 "invalid stream append outcome".into(),
965 ))
966 })?;
967 Ok(*outcome)
968 }
969
970 #[allow(clippy::too_many_arguments)]
973 pub async fn stream_append(
974 &self,
975 token: &NamespaceToken,
976 stream: &str,
977 record: &Value,
978 expected_seq: Option<i64>,
979 note_kind: &str,
980 tags: Option<Vec<String>>,
981 fence: Option<NoteFences>,
982 embed: Option<bool>,
983 embedding_model: Option<String>,
984 registry: &VerbRegistry,
985 ) -> RuntimeResult<Value> {
986 self.stream_append_with_outcome(
987 token,
988 stream,
989 record,
990 expected_seq,
991 note_kind,
992 tags,
993 fence,
994 embed,
995 embedding_model,
996 registry,
997 )
998 .await
999 .map_err(StreamAppendFailure::into_source)
1000 }
1001
1002 #[allow(clippy::too_many_arguments)]
1006 pub async fn stream_append_with_outcome(
1007 &self,
1008 token: &NamespaceToken,
1009 stream: &str,
1010 record: &Value,
1011 expected_seq: Option<i64>,
1012 note_kind: &str,
1013 tags: Option<Vec<String>>,
1014 fence: Option<NoteFences>,
1015 embed: Option<bool>,
1016 embedding_model: Option<String>,
1017 registry: &VerbRegistry,
1018 ) -> Result<Value, StreamAppendFailure> {
1019 let spec = StreamAppendSpec {
1020 stream: stream.to_string(),
1021 record: record.clone(),
1022 expected_seq,
1023 note_kind: note_kind.to_string(),
1024 tags,
1025 fence,
1026 embed,
1027 embedding_model,
1028 };
1029 let prepared = self
1030 .prepare_stream_appends(token, &[&spec], registry)
1031 .await
1032 .map_err(StreamAppendFailure::not_committed)?;
1033 let access = self.sql();
1034 match Self::run_stream_appends(access.as_ref(), token, &prepared).await? {
1035 BatchOutcome::Appended(seqs) => Ok(append_result(&prepared[0], seqs[0])),
1036 BatchOutcome::FenceConflict { conflict } => Err(StreamAppendFailure::not_committed(
1037 conflict.into_error().into(),
1038 )),
1039 BatchOutcome::Conflict { next } => Err(StreamAppendFailure::not_committed(
1040 seq_conflict(
1041 stream,
1042 expected_seq.expect("only conditional appends conflict"),
1043 next,
1044 None,
1045 )
1046 .into(),
1047 )),
1048 }
1049 }
1050
1051 fn validate_stream_batch(
1052 &self,
1053 members: &[StreamBatchMember],
1054 atomic: bool,
1055 ) -> RuntimeResult<()> {
1056 if members.is_empty() {
1057 return Err(RuntimeError::InvalidInput(
1058 "stream.batch requires at least one member".into(),
1059 ));
1060 }
1061 let mut keys = HashSet::new();
1062 for (index, member) in members.iter().enumerate() {
1063 match member {
1064 StreamBatchMember::Append(spec) => {
1065 spec.validate_embedding(atomic.then_some(index))?;
1066 validate_stream(&spec.stream)?;
1067 self.validate_note_kind(&spec.note_kind)?;
1068 if let Some(fences) = &spec.fence {
1069 fences.validate()?;
1070 for fence in fences.entries() {
1071 self.validate_note_kind(&fence.kind)?;
1072 }
1073 }
1074 crate::secret_gate::check_at(
1075 &serde_json::to_string(&spec.record)
1076 .map_err(|error| RuntimeError::InvalidInput(error.to_string()))?,
1077 &format!("member[{index}]"),
1078 "record",
1079 )?;
1080 }
1081 StreamBatchMember::Write(spec) => {
1082 self.validate_note_kind(&spec.kind)?;
1083 NoteWriteOptions {
1084 key: Some(spec.key.clone()),
1085 expected_version: spec.expected_version,
1086 embed: spec.embed,
1087 fence: None,
1088 }
1089 .validate()?;
1090 if spec.kind == "scheduled_event" {
1091 return Err(RuntimeError::InvalidInput(
1092 "scheduled_event notes are not writable through stream.batch; use schedule verbs".into(),
1093 ));
1094 }
1095 if !keys.insert((&spec.kind, &spec.key)) {
1096 return Err(RuntimeError::InvalidInput(format!(
1097 "stream.batch repeats write key {:?} for kind {:?}",
1098 spec.key, spec.kind,
1099 )));
1100 }
1101 crate::secret_gate::check_at(
1102 &serde_json::to_string(&spec.doc)
1103 .map_err(|error| RuntimeError::InvalidInput(error.to_string()))?,
1104 &format!("member[{index}]"),
1105 "doc",
1106 )?;
1107 if let Some(tags) = &spec.tags {
1108 crate::secret_gate::check_json_at(
1109 &json!({"tags": tags}),
1110 &format!("member[{index}]"),
1111 "tags",
1112 )?;
1113 }
1114 }
1115 StreamBatchMember::Refused(_) => {}
1116 }
1117 }
1118 Ok(())
1119 }
1120
1121 async fn prepare_stream_write(
1122 &self,
1123 token: &NamespaceToken,
1124 spec: StreamWriteSpec,
1125 registry: &VerbRegistry,
1126 ) -> RuntimeResult<StreamBatchPreparation> {
1127 let content = serde_json::to_string(&spec.doc)
1128 .map_err(|error| RuntimeError::InvalidInput(error.to_string()))?;
1129 if let Some(expected) = spec.expected_version {
1130 let snapshot = match self
1131 .get_note_by_key(token, &spec.key, Some(&spec.kind), false)
1132 .await
1133 {
1134 Ok(note) => note,
1135 Err(RuntimeError::Khive(error))
1136 if error.kind() == khive_types::ErrorKind::NotFound =>
1137 {
1138 return Ok(StreamBatchPreparation::Ready(Box::new(
1139 PreparedBatchAction::Refused(missing_write(&spec.key)),
1140 )));
1141 }
1142 Err(error) => return Err(error),
1143 };
1144 let id = snapshot.id;
1145 snapshot
1146 .version
1147 .checked_add(1)
1148 .ok_or_else(|| RuntimeError::InvalidInput("note version exhausted".into()))?;
1149 let mut args = json!({
1150 "id": id, "kind": "note", "note_kind": spec.kind,
1151 "content": content, "expected_version": expected,
1152 "namespace": token.namespace().as_str(),
1153 });
1154 if let Some(tags) = spec.tags {
1155 args["properties"] = json!({"tags": tags});
1156 }
1157 if let Some(embed) = spec.embed {
1158 args["embed"] = json!(embed);
1159 }
1160 let update_policy = registry
1161 .prepare_note_update_policy(self, token, &snapshot, &mut args)
1162 .await?;
1163 let (_, plan) = crate::atomic_prepare::prepare_update_from_note_snapshot(
1164 self,
1165 token,
1166 &args,
1167 None,
1168 snapshot,
1169 update_policy,
1170 registry,
1171 )
1172 .await?;
1173 return Ok(StreamBatchPreparation::Ready(Box::new(
1174 PreparedBatchAction::Write {
1175 key: spec.key,
1176 kind: spec.kind,
1177 id,
1178 plan,
1179 after_create: None,
1180 },
1181 )));
1182 }
1183 let mut args = json!({
1184 "kind": "note", "note_kind": spec.kind, "key": spec.key,
1185 "content": content, "namespace": token.namespace().as_str(),
1186 });
1187 if let Some(tags) = spec.tags {
1188 args["tags"] = json!(tags);
1189 }
1190 if let Some(embed) = spec.embed {
1191 args["embed"] = json!(embed);
1192 }
1193 if let Some(hook) = registry.find_kind_hook(&spec.kind) {
1194 hook.prepare_create(self, &mut args).await?;
1195 }
1196 let mut fields: StreamCreateFields = serde_json::from_value(args.clone())
1197 .map_err(|error| RuntimeError::InvalidInput(format!("stream write fields: {error}")))?;
1198 normalize_stream_create_fields(&mut fields)?;
1199 let mut candidate = Note::new(token.namespace().as_str(), &spec.kind, &fields.content);
1200 candidate.name = fields.name.clone();
1201 candidate.properties = fields.properties.clone();
1202 crate::note_write::validate_head(&candidate)?;
1203 Ok(StreamBatchPreparation::Create(Box::new(
1204 PreparedStreamCreate {
1205 key: spec.key,
1206 kind: spec.kind,
1207 fields,
1208 args,
1209 },
1210 )))
1211 }
1212
1213 async fn prepare_stream_batch(
1214 &self,
1215 token: &NamespaceToken,
1216 members: Vec<StreamBatchMember>,
1217 registry: &VerbRegistry,
1218 ) -> RuntimeResult<Vec<PreparedBatchMember>> {
1219 let mut pending = Vec::with_capacity(members.len());
1220 for member in members {
1221 let action = match member {
1222 StreamBatchMember::Refused(error) => {
1223 StreamBatchPreparation::Ready(Box::new(PreparedBatchAction::Refused(error)))
1224 }
1225 StreamBatchMember::Write(spec) => {
1226 match self.prepare_stream_write(token, spec, registry).await {
1227 Ok(action) => action,
1228 Err(RuntimeError::Khive(error))
1229 if error.kind() == khive_types::ErrorKind::Conflict =>
1230 {
1231 StreamBatchPreparation::Ready(Box::new(PreparedBatchAction::Refused(
1232 error,
1233 )))
1234 }
1235 Err(error) => return Err(error),
1236 }
1237 }
1238 StreamBatchMember::Append(spec) => {
1239 let content = serde_json::to_string(&spec.record)
1240 .map_err(|error| RuntimeError::InvalidInput(error.to_string()))?;
1241 StreamBatchPreparation::Append(Box::new((spec, content)))
1242 }
1243 };
1244 pending.push(action);
1245 }
1246 let requests = pending
1247 .iter()
1248 .filter_map(|action| match action {
1249 StreamBatchPreparation::Ready(_) => None,
1250 StreamBatchPreparation::Append(append) => {
1251 let (spec, content) = append.as_ref();
1252 Some(AtomicNoteRequest {
1253 spec: AtomicNoteSpec {
1254 token,
1255 id: None,
1256 kind: &spec.note_kind,
1257 name: None,
1258 content,
1259 properties: spec.tags.clone().map(|tags| json!({"tags": tags})),
1260 },
1261 options: spec.note_options(),
1262 })
1263 }
1264 StreamBatchPreparation::Create(create) => Some(AtomicNoteRequest {
1265 spec: AtomicNoteSpec {
1266 token,
1267 id: None,
1268 kind: &create.kind,
1269 name: create.fields.name.as_deref(),
1270 content: &create.fields.content,
1271 properties: create.fields.properties.clone(),
1272 },
1273 options: AtomicNoteOptions {
1274 salience: create.fields.salience,
1275 decay_factor: create.fields.decay_factor,
1276 key: Some(&create.key),
1277 embed: Some(create.fields.embed.unwrap_or(create.kind != "head")),
1278 ..Default::default()
1279 },
1280 }),
1281 })
1282 .collect();
1283 let notes = prepare_atomic_note_requests(self, requests).await?;
1284 let mut note_plans = notes.notes.into_iter().zip(notes.plans);
1285 let mut prepared = Vec::with_capacity(pending.len());
1286 for (index, pending) in pending.into_iter().enumerate() {
1287 let action = match pending {
1288 StreamBatchPreparation::Ready(action) => *action,
1289 StreamBatchPreparation::Append(append) => {
1290 let (spec, _) = *append;
1291 let (note, plan) = note_plans.next().expect("one plan per new note");
1292 PreparedBatchAction::Append {
1293 stream: spec.stream,
1294 expected_seq: spec.expected_seq,
1295 note: Box::new(note),
1296 plan,
1297 fence: spec.fence,
1298 }
1299 }
1300 StreamBatchPreparation::Create(create) => {
1301 let (note, mut plan) = note_plans.next().expect("one plan per new note");
1302 let AtomicOpPlan::AddNote(add) = &mut plan else {
1303 return Err(RuntimeError::Internal(
1304 "stream keyed create did not prepare a note".into(),
1305 ));
1306 };
1307 add.post_commit = PostCommitEffect::NoteChanged {
1308 note_id: note.id,
1309 kind: note.kind.clone(),
1310 };
1311 PreparedBatchAction::Write {
1312 key: create.key,
1313 kind: create.kind,
1314 id: note.id,
1315 plan,
1316 after_create: Some(create.args),
1317 }
1318 }
1319 };
1320 prepared.push(PreparedBatchMember { index, action });
1321 }
1322 Ok(prepared)
1323 }
1324
1325 async fn stream_batch_after_create(
1326 &self,
1327 registry: &VerbRegistry,
1328 hooks: Vec<(Uuid, String, Value)>,
1329 ) {
1330 for (id, kind, args) in hooks {
1331 if let Some(hook) = registry.find_kind_hook(&kind) {
1332 if let Err(error) = hook.after_create(self, id, &args).await {
1333 tracing::warn!(%id, %kind, %error, "stream batch after_create failed after commit");
1334 }
1335 }
1336 }
1337 }
1338
1339 pub async fn stream_batch_atomic(
1342 &self,
1343 token: &NamespaceToken,
1344 members: Vec<StreamBatchMember>,
1345 fence: Option<NoteFence>,
1346 observed: Vec<StreamObservation>,
1347 registry: &VerbRegistry,
1348 ) -> RuntimeResult<Result<Vec<Value>, StreamBatchRefusal>> {
1349 self.validate_stream_batch(&members, true)?;
1350 if let Some(fence) = &fence {
1351 fence.validate()?;
1352 self.validate_note_kind(&fence.kind)?;
1353 }
1354 for entry in &observed {
1355 if entry.id.is_some() && entry.version.is_none() {
1356 return Err(RuntimeError::InvalidInput(
1357 "observed id requires a positive version".into(),
1358 ));
1359 }
1360 if entry.live_until.is_some() && entry.version.is_none() {
1361 return Err(RuntimeError::InvalidInput(
1362 "observed live_until requires a positive version".into(),
1363 ));
1364 }
1365 crate::keyed_memory::validate_memory_key(&entry.key)?;
1366 self.validate_note_kind(&entry.kind)?;
1367 if entry.version.is_some_and(|version| version < 1) {
1368 return Err(RuntimeError::InvalidInput(
1369 "observed version must be positive or null".into(),
1370 ));
1371 }
1372 }
1373 validate_stream_batch_fence_count(stream_batch_fence_count(&members, fence.as_ref()))?;
1374 for (member, item) in members.iter().enumerate() {
1375 if let StreamBatchMember::Refused(error) = item {
1376 return Ok(Err(StreamBatchRefusal {
1377 member,
1378 error: place_member(error.clone(), Some(member))?,
1379 }));
1380 }
1381 }
1382 let prepared = self.prepare_stream_batch(token, members, registry).await?;
1383 let hooks = batch_create_hooks(&prepared);
1384 match run_prepared_stream_batch(
1385 self.sql().as_ref(),
1386 token.namespace().as_str().into(),
1387 prepared,
1388 fence,
1389 observed,
1390 )
1391 .await?
1392 {
1393 Ok((results, effects)) => {
1394 crate::atomic_prepare::apply_post_commit_effects_with_report(self, token, effects)
1395 .await?;
1396 self.stream_batch_after_create(registry, hooks).await;
1397 Ok(Ok(results))
1398 }
1399 Err(refusal) => Ok(Err(refusal)),
1400 }
1401 }
1402
1403 pub async fn stream_batch_per_member(
1408 &self,
1409 token: &NamespaceToken,
1410 members: Vec<StreamBatchMember>,
1411 registry: &VerbRegistry,
1412 ) -> RuntimeResult<Vec<Value>> {
1413 self.validate_stream_batch(&members, false)?;
1414 validate_stream_batch_fence_count(stream_batch_fence_count(&members, None))?;
1415 let mut results = Vec::with_capacity(members.len());
1416 let prepared = self.prepare_stream_batch(token, members, registry).await?;
1417 for member in prepared {
1418 if let PreparedBatchAction::Refused(error) = &member.action {
1419 results.push(refusal_value(&place_member(error.clone(), None)?)?);
1420 continue;
1421 }
1422 let hooks = batch_create_hooks(std::slice::from_ref(&member));
1423 match run_prepared_stream_batch(
1424 self.sql().as_ref(),
1425 token.namespace().as_str().into(),
1426 vec![member],
1427 None,
1428 vec![],
1429 )
1430 .await?
1431 {
1432 Ok((mut values, effects)) => {
1433 crate::atomic_prepare::apply_post_commit_effects_with_report(
1434 self, token, effects,
1435 )
1436 .await?;
1437 self.stream_batch_after_create(registry, hooks).await;
1438 results.append(&mut values);
1439 }
1440 Err(refusal) => results.push(refusal_value(&place_member(refusal.error, None)?)?),
1441 }
1442 }
1443 Ok(results)
1444 }
1445
1446 pub async fn stream_read(
1448 &self,
1449 token: &NamespaceToken,
1450 stream: &str,
1451 after: i64,
1452 limit: i64,
1453 ) -> RuntimeResult<Value> {
1454 validate_stream(stream)?;
1455 if after < 0 || limit < 1 {
1456 return Err(RuntimeError::InvalidInput(
1457 "stream.read requires after >= 0 and limit >= 1".into(),
1458 ));
1459 }
1460 let mut reader = self.sql().reader().await?;
1461 let rows = reader.query_all(statement(
1462 "WITH head AS (SELECT COALESCE(MAX(seq),0) AS head_seq FROM note_streams WHERE namespace=?1 AND stream=?2), \
1463 page AS (SELECT s.seq,n.id,n.content,n.created_at FROM note_streams s JOIN notes n ON n.id=s.note_id \
1464 WHERE s.namespace=?1 AND s.stream=?2 AND s.seq>?3 ORDER BY s.seq LIMIT ?4) \
1465 SELECT head.head_seq,page.seq,page.id,page.content,page.created_at FROM head LEFT JOIN page ON 1=1 ORDER BY page.seq",
1466 vec![SqlValue::Text(token.namespace().as_str().into()), SqlValue::Text(stream.into()), SqlValue::Integer(after), SqlValue::Integer(limit)],
1467 )).await?;
1468 let head = rows
1469 .first()
1470 .map(|row| integer(row, "head_seq"))
1471 .transpose()?
1472 .unwrap_or(0);
1473 let mut entries = Vec::new();
1474 let mut last = None;
1475 for row in &rows {
1476 if matches!(row.get("seq"), Some(SqlValue::Null)) {
1477 continue;
1478 }
1479 let seq = integer(row, "seq")?;
1480 let record: Value = serde_json::from_str(text(row, "content")?)
1481 .map_err(|e| RuntimeError::Internal(format!("invalid stream record JSON: {e}")))?;
1482 entries.push(json!({"seq": seq, "id": text(row, "id")?, "record": record, "created_at": micros_to_iso(integer(row, "created_at")?)}));
1483 last = Some(seq);
1484 }
1485 Ok(
1486 json!({"entries": entries, "head_seq": head, "next_after": last.filter(|last| *last < head)}),
1487 )
1488 }
1489
1490 pub async fn stream_stat(&self, token: &NamespaceToken, stream: &str) -> RuntimeResult<Value> {
1492 validate_stream(stream)?;
1493 let row = self.sql().reader().await?.query_row(statement(
1494 "SELECT COUNT(*) AS count, COALESCE(MAX(seq),0) AS head_seq FROM note_streams WHERE namespace=?1 AND stream=?2",
1495 vec![SqlValue::Text(token.namespace().as_str().into()), SqlValue::Text(stream.into())],
1496 )).await?.ok_or_else(|| RuntimeError::Internal("stream.stat returned no aggregate row".into()))?;
1497 Ok(json!({"head_seq": integer(&row, "head_seq")?, "count": integer(&row, "count")?}))
1498 }
1499
1500 pub(crate) async fn stream_member_error(
1503 &self,
1504 note: &Note,
1505 ) -> RuntimeResult<Option<RuntimeError>> {
1506 let row = self
1507 .sql()
1508 .reader()
1509 .await?
1510 .query_row(statement(
1511 "SELECT stream,seq FROM note_streams WHERE namespace=?1 AND note_id=?2",
1512 vec![
1513 SqlValue::Text(note.namespace.clone()),
1514 SqlValue::Text(note.id.to_string()),
1515 ],
1516 ))
1517 .await?;
1518 row.map(|row| {
1519 Ok(KhiveError::conflict("stream entries are immutable")
1520 .with_details(Details::new_owned([
1521 ("reason", "stream_member".into()),
1522 ("id", note.id.to_string()),
1523 ("stream", text(&row, "stream")?.into()),
1524 ("seq", integer(&row, "seq")?.to_string()),
1525 ]))
1526 .into())
1527 })
1528 .transpose()
1529 }
1530}
1531
1532#[cfg(test)]
1533#[path = "streams_batch_tests.rs"]
1534mod batch_tests;
1535
1536#[cfg(test)]
1537mod tests {
1538 use super::{allocate_sequence, SequenceRefusal};
1539 use crate::atomic_prepare::{prepare_delete, prepare_update};
1540 use crate::atomic_runner::{run_atomic_unit, AtomicRunOutcome};
1541 use crate::{KhiveRuntime, Namespace, NotePatch, RuntimeError, VerbRegistryBuilder};
1542 use serde_json::{json, Value};
1543
1544 #[test]
1545 fn allocate_sequence_exhaustion_preserves_head_at_checked_add_boundary() {
1546 let mut head = i64::MAX - 1;
1551 assert!(matches!(
1552 allocate_sequence(&mut head, None),
1553 Ok(seq) if seq == i64::MAX
1554 ));
1555 assert_eq!(head, i64::MAX);
1556 assert!(matches!(
1557 allocate_sequence(&mut head, None),
1558 Err(SequenceRefusal::Exhausted)
1559 ));
1560 assert_eq!(
1561 head,
1562 i64::MAX,
1563 "refused allocation must not advance the head"
1564 );
1565 }
1566
1567 #[tokio::test]
1568 async fn stream_atomic_metadata_cas_preserves_record_and_refuses_stale_plan() {
1569 let rt = KhiveRuntime::memory().unwrap();
1570 let token = rt.authorize(Namespace::local()).unwrap();
1571 let registry = VerbRegistryBuilder::new().build().unwrap();
1572 let appended = rt
1573 .stream_append(
1574 &token,
1575 "cas",
1576 &json!({"n": 1}),
1577 None,
1578 "observation",
1579 None,
1580 None,
1581 None,
1582 None,
1583 ®istry,
1584 )
1585 .await
1586 .unwrap();
1587 let id = uuid::Uuid::parse_str(appended["id"].as_str().unwrap()).unwrap();
1588 for args in [
1589 json!({"id": id, "content": "changed"}),
1590 json!({"id": id, "properties": {"x": 1}}),
1591 ] {
1592 let err = prepare_update(&rt, &token, &args, None).await.unwrap_err();
1593 assert!(matches!(err, RuntimeError::Khive(ref e) if e.details().is_some()));
1594 assert!(err.to_string().contains("immutable"));
1595 }
1596 for hard in [false, true] {
1597 let err = prepare_delete(&rt, &token, &json!({"id": id, "hard": hard}), None)
1598 .await
1599 .unwrap_err();
1600 assert!(err.to_string().contains("immutable"));
1601 }
1602 let plan = prepare_update(&rt, &token, &json!({"id": id, "salience": 0.7}), None)
1603 .await
1604 .unwrap();
1605 assert!(matches!(
1606 run_atomic_unit(rt.sql().as_ref(), vec![plan])
1607 .await
1608 .unwrap(),
1609 AtomicRunOutcome::Committed { .. }
1610 ));
1611 let stale = prepare_update(&rt, &token, &json!({"id": id, "salience": 0.2}), None)
1612 .await
1613 .unwrap();
1614 rt.update_note(
1615 &token,
1616 id,
1617 NotePatch::new(None, None, Some(Some(0.9)), None, None),
1618 )
1619 .await
1620 .unwrap();
1621 assert!(matches!(
1622 run_atomic_unit(rt.sql().as_ref(), vec![stale])
1623 .await
1624 .unwrap(),
1625 AtomicRunOutcome::RolledBack { .. }
1626 ));
1627 let note = rt
1628 .notes(&token)
1629 .unwrap()
1630 .get_note(id)
1631 .await
1632 .unwrap()
1633 .unwrap();
1634 assert_eq!(note.salience, Some(0.9));
1635 assert_eq!(
1636 serde_json::from_str::<Value>(¬e.content).unwrap(),
1637 json!({"n": 1})
1638 );
1639 }
1640
1641 #[tokio::test]
1642 async fn stream_namespace_isolation_and_real_stat_count() {
1643 let rt = KhiveRuntime::memory().unwrap();
1644 let a = rt.authorize(Namespace::parse("a").unwrap()).unwrap();
1645 let b = rt.authorize(Namespace::parse("b").unwrap()).unwrap();
1646 let registry = VerbRegistryBuilder::new().build().unwrap();
1647 for token in [&a, &b] {
1648 assert_eq!(
1649 rt.stream_append(
1650 token,
1651 "same-name",
1652 &json!(token.namespace().as_str()),
1653 Some(1),
1654 "observation",
1655 None,
1656 None,
1657 None,
1658 None,
1659 ®istry
1660 )
1661 .await
1662 .unwrap()["seq"],
1663 1
1664 );
1665 }
1666 for token in [&a, &b] {
1667 let page = rt.stream_read(token, "same-name", 0, 10).await.unwrap();
1668 assert_eq!(page["entries"][0]["record"], token.namespace().as_str());
1669 }
1670 rt.sql().writer().await.unwrap().execute_script("DROP TRIGGER refuse_stream_ledger_update; UPDATE note_streams SET seq=5 WHERE namespace='a';".into()).await.unwrap();
1671 assert_eq!(
1672 rt.stream_stat(&a, "same-name").await.unwrap(),
1673 json!({"count": 1, "head_seq": 5})
1674 );
1675 assert_eq!(
1676 rt.stream_stat(&b, "same-name").await.unwrap(),
1677 json!({"count": 1, "head_seq": 1})
1678 );
1679 }
1680}