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