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