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