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