Skip to main content

khive_runtime/
streams.rs

1//! Ordered note streams. Embedding preparation precedes one SQL-only writer
2//! transaction containing the sequence predicates, note/index writes and ledger
3//! rows. A batch is that transaction over several members (atomic mode) or one
4//! such transaction per member, in list order (per-member mode).
5use 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/// What the append path can prove about its own failed append, independently
58/// of an enclosing dispatch or a nested error's wire disposition.
59#[derive(Clone, Copy, Debug, PartialEq, Eq)]
60pub enum StreamAppendDisposition {
61    /// This append either never reached a write or its transaction rolled back.
62    NotCommitted,
63    /// The append may have committed; no absence of a record is asserted.
64    Unknown,
65}
66
67/// An unchanged source error together with evidence local to one append.
68#[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    /// Retain an error when no append-local committedness proof is available.
85    /// This constructor cannot assert that an append was dropped.
86    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    /// Borrow the original error without rewriting its fields or provenance.
120    pub fn source(&self) -> &RuntimeError {
121        &self.source
122    }
123
124    /// Evidence about this append, not the enclosing handler's disposition.
125    pub fn disposition(&self) -> StreamAppendDisposition {
126        self.disposition
127    }
128
129    /// Separate the unchanged source from append-local committedness evidence.
130    pub fn into_parts(self) -> (RuntimeError, StreamAppendDisposition) {
131        (self.source, self.disposition)
132    }
133
134    /// Preserve the compatibility error path without an additional wrapper.
135    pub fn into_source(self) -> RuntimeError {
136        self.source
137    }
138}
139
140/// One append, shape-validated by the verb layer; the runtime validates the
141/// stream name and the record's serialization before any write.
142pub 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/// A keyed document write. Kinds are canonical note-kind names. A missing
175/// expected version creates only; a positive version updates only.
176#[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/// An exact live-key observation; a null version asserts that the key is unheld.
187/// An optional identity pins which live note must hold the key at that version.
188#[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    /// Dotted document path whose RFC 3339 value must exceed the writer clock.
195    pub live_until: Option<String>,
196}
197
198/// A batch member or its already established refusal. The mode places it.
199pub enum StreamBatchMember {
200    Append(StreamAppendSpec),
201    Write(StreamWriteSpec),
202    Refused(KhiveError),
203}
204
205// Like the per-list fence and observation bounds, this limits precondition
206// reads while holding the writer; member count alone does not bound their sum.
207const 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/// The member refusal that stopped an atomic batch; nothing was written.
231#[derive(Debug)]
232pub struct StreamBatchRefusal {
233    pub member: usize,
234    pub error: KhiveError,
235}
236
237/// A member refusal as a value: the error object a refused op carries, plus
238/// the disposition the consumer rule reads without a special case.
239pub 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
260/// A prepared append: the note, its planned statements and the ledger key.
261struct 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        // No live holder is the existing version conflict. Identity conflict
405        // specifically names a replacement, so current_id is always present.
406        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            // One SQL clock read after writer admission, shared by the list.
497            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            // Member fences observe the transaction's initial state, even when
527            // an earlier keyed write changes a fenced note in this batch.
528            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                // Read the row under this writer, before commit or any later
677                // writer. Preserve its existing timestamp/CAS semantics.
678                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                // Versions are local to a note identity. A replacement can be
709                // at the expected version without being the prepared target.
710                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    /// Validate every spec and prepare every note before any write. All
744    /// allocations for embedding, note normalization and SQL plans complete
745    /// here, so only bounded statement driving occurs while the writer is held.
746    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            // A hooked kind may use its ordinary create fields from an object record (for
766            // example, `task` reads `title`), while an unhooked kind keeps the append record as
767            // opaque JSON content exactly as before.
768            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    /// Run prepared appends as one writer transaction. Every stream head is
842    /// read, every number assigned in list order and every `expected_seq`
843    /// checked before the first write; then notes and ledger rows land in
844    /// list order, so appends to one stream take consecutive numbers.
845    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        // A false flag is never proof: the backend may have accepted a closure
864        // that has not run yet. Only a completed pre-write error sets it true.
865        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                    // Set before invoking the first write, including a write
915                    // whose driver returns an error with ambiguous effects.
916                    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    /// Append a JSON value as an immutable note. The sequence precondition and
955    /// every note/index/ledger statement share the same writer transaction.
956    #[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    /// Append once while retaining evidence about the failed append's own
987    /// committedness. The original error remains available for stop semantics
988    /// or the shared structured error projection; this method never retries.
989    #[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    /// All members share one transaction. Batch predicates precede member DML;
1324    /// effects become executable only after the complete unit commits.
1325    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    /// One writer transaction per member, in list order. Every member is
1388    /// prepared before the first write; a member's refusal is returned as its
1389    /// value and its siblings stand, so numbers on one stream increase with
1390    /// list position but another writer's append may fall between them.
1391    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    /// Read one ordered page and its head from one SQL snapshot.
1431    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    /// Independently count entries and read the head in the same statement.
1483    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    /// Membership lookup only after the caller has resolved an accessible note.
1495    /// Match the stored namespace as well as its globally unique id.
1496    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        // Overflow is logically possible, but no bounded real-store fixture
1562        // reaches it: stream_gap requires consecutive inserts from 1, and
1563        // ledger UPDATE/DELETE are forbidden. Supply the head at this unit
1564        // boundary instead of fabricating an exhausted persisted stream.
1565        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                &registry,
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>(&note.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                    &registry
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}