Skip to main content

notedthat_write/
commit.rs

1//! Shared commit operations for object writes and deletes.
2
3use notedthat_core::{
4    ConditionalHeaders, CopyObjectOptions, KbSlug, ObjectEvent, ObjectPath, PutOutcome, StagedBody,
5    Storage, StorageError,
6};
7use notedthat_indexer::IndexEvent;
8use tokio::sync::mpsc::Sender;
9use tokio::sync::mpsc::error::TrySendError;
10
11use crate::WriteError;
12use crate::error::WriteEffect;
13use crate::mime::sniff_content_type;
14use crate::sinks::WriteSinks;
15
16/// Maximum upload size accepted by shared write paths: 5 GiB.
17pub const MAX_UPLOAD_BYTES: u64 = 5 * 1024 * 1024 * 1024;
18
19/// Validate an upload size against a byte limit.
20pub fn check_size(size: u64, limit: u64) -> Result<(), WriteError> {
21    if size > limit {
22        Err(WriteError::TooLarge { size, limit })
23    } else {
24        Ok(())
25    }
26}
27
28/// Store an object, enqueue a best-effort index upsert event and publish the
29/// change.
30pub async fn commit<B: Into<StagedBody>>(
31    storage: &dyn Storage,
32    sinks: &WriteSinks<'_>,
33    kb: &KbSlug,
34    path: &ObjectPath,
35    body: B,
36    caller_content_type: Option<&str>,
37    conditionals: ConditionalHeaders,
38) -> Result<PutOutcome, WriteError> {
39    let body = body.into();
40    let size = body.len();
41    check_size(size, MAX_UPLOAD_BYTES)?;
42    crate::manifest::check_staged_manifest(kb, path, &body).await?;
43    let _prefix = body.prefix(512).await.map_err(|source| {
44        WriteError::Storage(StorageError::Other {
45            source: Box::new(source),
46        })
47    })?;
48    let mime = sniff_content_type(caller_content_type, path);
49    let outcome = storage
50        .put_staged_object(kb, path, body, Some(&mime), conditionals)
51        .await?;
52
53    after_write(sinks, kb, path, &outcome, size, &mime).await?;
54    Ok(outcome)
55}
56
57/// Copy an object natively, enqueue its destination for indexing and publish
58/// the change.
59pub async fn commit_copy(
60    storage: &dyn Storage,
61    sinks: &WriteSinks<'_>,
62    kb: &KbSlug,
63    source: &ObjectPath,
64    destination: &ObjectPath,
65    options: CopyObjectOptions,
66) -> Result<PutOutcome, WriteError> {
67    let content_type = options.content_type.clone();
68    if crate::manifest::is_manifest(destination) {
69        // A copy never has the bytes in hand, so the source is read to check
70        // them — bounded by its HEAD first, so an oversized source is refused
71        // unread. The source can change between this read and the copy; the
72        // `.notedthat` namespace is the credential holder's alone (D51), so
73        // that is a race with oneself, not a bypass.
74        let meta = storage
75            .head_object(kb, source, ConditionalHeaders::default())
76            .await?;
77        if meta.size > crate::manifest::MANIFEST_MAX_BYTES {
78            return Err(WriteError::InvalidManifest {
79                message: format!(
80                    "manifest is {} bytes; a manifest is at most {}",
81                    meta.size,
82                    crate::manifest::MANIFEST_MAX_BYTES
83                ),
84            });
85        }
86        let read = storage
87            .get_object(kb, source, None, ConditionalHeaders::default())
88            .await?;
89        crate::manifest::check_manifest_bytes(kb, destination, &read.bytes)?;
90    }
91    let outcome = storage
92        .copy_object(kb, source, destination, options)
93        .await?;
94    // A copy returns only the new ETag. The event wants the size too, and only
95    // a subscriber cares, so the extra HEAD is paid only when one can exist.
96    // The copy is already durable by now, so a HEAD that fails — a racing
97    // delete, a transient error — degrades the stamp rather than the copy.
98    let (size, mime) = if sinks.events.is_some() {
99        match storage
100            .head_object(kb, destination, ConditionalHeaders::default())
101            .await
102        {
103            Ok(meta) => (meta.size, meta.content_type.or(content_type)),
104            Err(error) => {
105                tracing::warn!(
106                    target: "notedthat::events",
107                    kb = %kb, path = %destination, %error,
108                    "could not HEAD the copy destination for its event; publishing without a size"
109                );
110                (0, content_type)
111            }
112        }
113    } else {
114        (0, content_type)
115    };
116    after_write(
117        sinks,
118        kb,
119        destination,
120        &outcome,
121        size,
122        mime.as_deref().unwrap_or_default(),
123    )
124    .await?;
125    Ok(outcome)
126}
127
128/// Holds the `Upsert` for a write whose bytes are already stored, and enqueues
129/// it from `Drop` if the normal path never got there.
130///
131/// Only ever fires on cancellation: [`after_write`] disarms it the moment it
132/// takes the event to send properly, so the ordinary path — including every
133/// queue-full and queue-closed case, which need the health view and the error
134/// this cannot produce — is untouched. What it covers is the future being
135/// dropped between the store and the enqueue.
136struct EnqueueOnDrop {
137    indexer_tx: Sender<IndexEvent>,
138    event: Option<IndexEvent>,
139}
140
141impl EnqueueOnDrop {
142    /// Take the event for the normal enqueue; nothing happens on drop after.
143    fn disarm(&mut self) -> IndexEvent {
144        self.event.take().expect("disarmed once")
145    }
146}
147
148impl Drop for EnqueueOnDrop {
149    fn drop(&mut self) {
150        let Some(event) = self.event.take() else {
151            return;
152        };
153        // Best effort by construction: there is no caller left to return an
154        // error to, and a full queue here means the same as it does anywhere
155        // else — the reconciliation pass is the backstop.
156        if self.indexer_tx.try_send(event).is_err() {
157            tracing::warn!(
158                target: "notedthat::indexing",
159                "INDEX_ENQUEUE_LOST_ON_CANCEL: a write was cancelled after its bytes were stored \
160                 and its Upsert could not be queued; reconciliation will repair it"
161            );
162        }
163    }
164}
165
166/// Report a durable write: publish, then index.
167///
168/// Publishing first is what lets the worker's `object.indexed` for this
169/// `ETag` always follow `object.written` in the log (D65): the write is on
170/// the log before its `Upsert` is even in the queue. The enqueue happens
171/// whether or not the publish succeeded, so the index never depends on the
172/// broker; a refused publish is still the write's answer, and the retry
173/// re-announces and re-indexes. D38's promise holds either way — either
174/// failure returns an error after the bytes are stored, and a retried write
175/// re-runs both — so an event may be published twice, never zero times.
176pub(crate) async fn after_write(
177    sinks: &WriteSinks<'_>,
178    kb: &KbSlug,
179    path: &ObjectPath,
180    outcome: &PutOutcome,
181    size: u64,
182    mime: &str,
183) -> Result<(), WriteError> {
184    // Publish first, so `object.indexed` for this `ETag` always follows the
185    // `object.written` it answers (D65) — but enqueue regardless of the
186    // outcome: the index is what search depends on and the broker is the
187    // optional component (§5 principle 9), so a refused publish must not
188    // leave stored bytes unsearchable until a client happens to retry. The
189    // refusal is still the write's answer, after the enqueue; the retry then
190    // announces the version and, an `Upsert` never being skipped, indexes it
191    // again — an `object.indexed` with no `object.written` before it is the
192    // same class of thing as a duplicate, and subscribers tolerate both.
193    // Armed before the publish, disarmed by the enqueue below.
194    //
195    // `publish` can await for seconds (the NATS client's own timeout), and a
196    // caller that gives up in that window — a client disconnect, or the
197    // request timeout D71 added — drops this future where it stands. The bytes
198    // are already stored by then, so without this the object would be durable
199    // and unsearchable, with nothing to repair it on `s3` until the next
200    // reconciliation pass. `try_send` is synchronous, so `Drop` can still do
201    // it; this is the same bargain D69 strikes for metered calls, for the same
202    // reason.
203    let mut enqueue = EnqueueOnDrop {
204        indexer_tx: sinks.indexer_tx.clone(),
205        event: Some(IndexEvent::Upsert {
206            kb: kb.clone(),
207            object_key: path.clone(),
208            etag: outcome.etag.clone().unwrap_or_default(),
209            mtime: current_unix_seconds(),
210        }),
211    };
212
213    let published = match sinks.events {
214        Some(events) => {
215            let event = ObjectEvent::written(
216                kb.clone(),
217                path.clone(),
218                outcome.etag.clone().unwrap_or_default(),
219                size,
220                mime.to_string(),
221                current_unix_seconds(),
222                sinks.source,
223            );
224            publish(events, event, WriteEffect::Stored).await
225        }
226        None => Ok(()),
227    };
228
229    let event = enqueue.disarm();
230    match sinks.indexer_tx.try_send(event) {
231        Ok(()) => {
232            if let Some(health) = sinks.index_health {
233                health.enqueued(kb.as_str());
234            }
235        }
236        Err(TrySendError::Full(ev)) => {
237            tracing::warn!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_FULL");
238            let _ = ev;
239            if let Some(health) = sinks.index_health {
240                health.backpressured(kb.as_str());
241            }
242            return Err(WriteError::IndexerBackpressureUpsert);
243        }
244        Err(TrySendError::Closed(ev)) => {
245            tracing::error!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_CLOSED");
246            let _ = ev;
247            if let Some(health) = sinks.index_health {
248                health.worker_stopped();
249            }
250            // Closed = indexer worker ended (shutdown OR panic). v1 preserves success-with-error-log
251            // until post-v1 worker liveness detection is added.
252        }
253    }
254    published
255}
256
257/// Delete an object idempotently, publish the change and enqueue a
258/// best-effort tombstone event — in that order, the one rule `after_write`
259/// follows (D65).
260pub async fn commit_delete(
261    storage: &dyn Storage,
262    sinks: &WriteSinks<'_>,
263    kb: &KbSlug,
264    path: &ObjectPath,
265    conditionals: ConditionalHeaders,
266) -> Result<(), WriteError> {
267    match storage.delete_object(kb, path, conditionals).await {
268        Ok(()) | Err(StorageError::NotFound { .. }) => {}
269        Err(e) => return Err(WriteError::Storage(e)),
270    }
271
272    // As in `after_write`: publish first, enqueue regardless, answer with the
273    // publish's outcome — a refused publish must not leave a deleted object's
274    // points searchable until a client retries.
275    let published = match sinks.events {
276        Some(events) => {
277            let event = ObjectEvent::deleted(kb.clone(), path.clone(), sinks.source);
278            publish(events, event, WriteEffect::Deleted).await
279        }
280        None => Ok(()),
281    };
282
283    let event = IndexEvent::Tombstone {
284        kb: kb.clone(),
285        object_key: path.clone(),
286    };
287    match sinks.indexer_tx.try_send(event) {
288        Ok(()) => {
289            if let Some(health) = sinks.index_health {
290                health.enqueued(kb.as_str());
291            }
292        }
293        Err(TrySendError::Full(ev)) => {
294            tracing::warn!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_FULL");
295            let _ = ev;
296            if let Some(health) = sinks.index_health {
297                health.backpressured(kb.as_str());
298            }
299            return Err(WriteError::IndexerBackpressureTombstone);
300        }
301        // Closed means the indexer worker task ended via shutdown OR panic; v1 deliberately preserves success-with-error-log so writes are not blocked by a crashed worker. The health view reports the stopped worker instead (#97).
302        Err(TrySendError::Closed(ev)) => {
303            tracing::error!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_CLOSED");
304            let _ = ev;
305            if let Some(health) = sinks.index_health {
306                health.worker_stopped();
307            }
308        }
309    }
310    published
311}
312
313async fn publish(
314    events: &dyn notedthat_core::EventPublisher,
315    event: ObjectEvent,
316    after: WriteEffect,
317) -> Result<(), WriteError> {
318    let (kb, path) = (event.kb.clone(), event.object_key.clone());
319    match events.publish(event).await {
320        Ok(_) => Ok(()),
321        Err(error) => {
322            tracing::error!(target: "notedthat::events", kb = %kb, path = %path, %error, "EVENT_PUBLISH_FAILED");
323            Err(WriteError::EventPublishFailed { after })
324        }
325    }
326}
327
328pub(crate) fn current_unix_seconds() -> i64 {
329    std::time::SystemTime::now()
330        .duration_since(std::time::UNIX_EPOCH)
331        .map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
332}
333
334#[cfg(test)]
335mod tests {
336    use super::*;
337    use crate::WriteEffect;
338    use async_trait::async_trait;
339    use bytes::Bytes;
340    use notedthat_core::{
341        EventId, EventPublisher, EventSource, EventStream, KbManifest, ListResponse, ObjectEvent,
342        ObjectMeta, ObjectRead, PublishError, SubscribeError,
343    };
344    use std::collections::HashMap;
345    use std::sync::{Arc, Mutex};
346    use tokio::sync::mpsc;
347
348    #[derive(Default)]
349    struct TestStorage {
350        objects: Mutex<HashMap<String, String>>,
351    }
352
353    #[async_trait]
354    impl Storage for TestStorage {
355        async fn probe(&self, _kb: &KbSlug) -> Result<(), StorageError> {
356            Ok(())
357        }
358
359        async fn ensure_bucket(&self, _kb: &KbSlug) -> Result<(), StorageError> {
360            unimplemented!()
361        }
362
363        async fn read_manifest(&self, _kb: &KbSlug) -> Result<KbManifest, StorageError> {
364            unimplemented!()
365        }
366
367        async fn write_manifest(
368            &self,
369            _kb: &KbSlug,
370            _manifest: &KbManifest,
371        ) -> Result<(), StorageError> {
372            unimplemented!()
373        }
374
375        async fn head_object(
376            &self,
377            kb: &KbSlug,
378            path: &ObjectPath,
379            _conditionals: ConditionalHeaders,
380        ) -> Result<ObjectMeta, StorageError> {
381            let key = format!("{}/{}", kb.as_str(), path.as_str());
382            let objects = self.objects.lock().expect("mutex not poisoned");
383            let etag = objects
384                .get(&key)
385                .ok_or_else(|| StorageError::NotFound { key: key.clone() })?;
386            Ok(ObjectMeta {
387                key: path.as_str().to_string(),
388                size: 7,
389                last_modified: Some(1_700_000_000),
390                content_type: Some("text/markdown".into()),
391                etag: Some(etag.clone()),
392            })
393        }
394
395        async fn get_object(
396            &self,
397            _kb: &KbSlug,
398            _path: &ObjectPath,
399            _range: Option<notedthat_core::ByteRange>,
400            _conditionals: ConditionalHeaders,
401        ) -> Result<ObjectRead, StorageError> {
402            unimplemented!()
403        }
404
405        async fn get_object_stream(
406            &self,
407            _kb: &KbSlug,
408            _path: &ObjectPath,
409            _range: Option<notedthat_core::ByteRange>,
410            _conditionals: ConditionalHeaders,
411        ) -> Result<notedthat_core::ObjectStream, StorageError> {
412            unimplemented!()
413        }
414
415        async fn put_object(
416            &self,
417            kb: &KbSlug,
418            path: &ObjectPath,
419            _bytes: Bytes,
420            _content_type: Option<&str>,
421            conditionals: ConditionalHeaders,
422        ) -> Result<PutOutcome, StorageError> {
423            let key = format!("{}/{}", kb.as_str(), path.as_str());
424            let mut objects = self.objects.lock().expect("mutex not poisoned");
425            let existing = objects.get(&key);
426            if let Some(if_match) = conditionals.if_match
427                && existing.is_none_or(|etag| etag != &if_match)
428            {
429                return Err(StorageError::PreconditionFailed);
430            }
431
432            let etag = format!("\"etag-{}\"", objects.len() + 1);
433            objects.insert(key, etag.clone());
434            Ok(PutOutcome { etag: Some(etag) })
435        }
436
437        async fn put_staged_object(
438            &self,
439            kb: &KbSlug,
440            path: &ObjectPath,
441            body: StagedBody,
442            content_type: Option<&str>,
443            conditionals: ConditionalHeaders,
444        ) -> Result<PutOutcome, StorageError> {
445            let bytes =
446                body.memory_bytes()
447                    .cloned()
448                    .ok_or_else(|| StorageError::BackendUnavailable {
449                        message: "file staging is outside this commit unit test".into(),
450                    })?;
451            self.put_object(kb, path, bytes, content_type, conditionals)
452                .await
453        }
454
455        async fn copy_object(
456            &self,
457            kb: &KbSlug,
458            _source: &ObjectPath,
459            destination: &ObjectPath,
460            _options: CopyObjectOptions,
461        ) -> Result<PutOutcome, StorageError> {
462            let key = format!("{}/{}", kb.as_str(), destination.as_str());
463            let mut objects = self.objects.lock().expect("mutex not poisoned");
464            let etag = format!("\"etag-{}\"", objects.len() + 1);
465            objects.insert(key, etag.clone());
466            Ok(PutOutcome { etag: Some(etag) })
467        }
468
469        async fn delete_object(
470            &self,
471            kb: &KbSlug,
472            path: &ObjectPath,
473            _conditionals: ConditionalHeaders,
474        ) -> Result<(), StorageError> {
475            let key = format!("{}/{}", kb.as_str(), path.as_str());
476            self.objects
477                .lock()
478                .expect("mutex not poisoned")
479                .remove(&key);
480            Ok(())
481        }
482
483        async fn list_objects(
484            &self,
485            _kb: &KbSlug,
486            _prefix: Option<&str>,
487            _limit: u32,
488            _cursor: Option<&str>,
489        ) -> Result<ListResponse, StorageError> {
490            unimplemented!()
491        }
492    }
493
494    fn kb() -> KbSlug {
495        KbSlug::try_new("test-kb").expect("valid kb slug")
496    }
497
498    fn path() -> ObjectPath {
499        ObjectPath::try_from_str("test.md").expect("valid path")
500    }
501
502    fn path_named(value: &str) -> ObjectPath {
503        ObjectPath::try_from_str(value).expect("valid path")
504    }
505
506    #[tokio::test]
507    async fn successful_put_enqueues_event() {
508        let storage = TestStorage::default();
509        let kb = kb();
510        let path = path();
511        let (indexer_tx, mut rx) = mpsc::channel(1024);
512
513        let outcome = commit(
514            &storage,
515            &WriteSinks::indexer_only(&indexer_tx),
516            &kb,
517            &path,
518            Bytes::from_static(b"# Test"),
519            Some("text/markdown"),
520            ConditionalHeaders::default(),
521        )
522        .await;
523
524        assert!(outcome.is_ok());
525        assert!(outcome.unwrap().etag.is_some());
526
527        let event = rx.recv().await.expect("event should be enqueued");
528        assert_eq!(event.kb().as_str(), "test-kb");
529        assert_eq!(event.object_key().as_str(), "test.md");
530    }
531
532    #[tokio::test]
533    async fn successful_native_copy_enqueues_destination_event() {
534        let storage = TestStorage::default();
535        let kb = kb();
536        let source = path_named("source.md");
537        let destination = path_named("destination.md");
538        let (indexer_tx, mut rx) = mpsc::channel(1);
539
540        let outcome = commit_copy(
541            &storage,
542            &WriteSinks::indexer_only(&indexer_tx),
543            &kb,
544            &source,
545            &destination,
546            CopyObjectOptions::default(),
547        )
548        .await
549        .expect("copy succeeds");
550
551        assert!(outcome.etag.is_some());
552        let event = rx.recv().await.expect("destination event");
553        assert_eq!(event.object_key(), &destination);
554    }
555
556    #[tokio::test]
557    async fn full_queue_returns_indexer_backpressure() {
558        let storage = TestStorage::default();
559        let kb = kb();
560        let path = path();
561        let (indexer_tx, _rx) = mpsc::channel(1);
562
563        let dummy_event = IndexEvent::Upsert {
564            kb: kb.clone(),
565            object_key: path.clone(),
566            etag: "dummy".to_string(),
567            mtime: 0,
568        };
569        indexer_tx
570            .try_send(dummy_event)
571            .expect("first send should succeed");
572
573        let outcome = commit(
574            &storage,
575            &WriteSinks::indexer_only(&indexer_tx),
576            &kb,
577            &path,
578            Bytes::from_static(b"# Test"),
579            Some("text/markdown"),
580            ConditionalHeaders::default(),
581        )
582        .await;
583
584        let err = outcome.unwrap_err();
585        assert!(
586            matches!(err, WriteError::IndexerBackpressureUpsert),
587            "expected IndexerBackpressureUpsert, got {err:?}"
588        );
589    }
590
591    #[tokio::test]
592    async fn burst_write_returns_backpressure_after_capacity() {
593        let storage = Arc::new(TestStorage::default());
594        let kb = kb();
595        let path_a = path_named("a.md");
596        let path_b = path_named("b.md");
597        let path_c = path_named("c.md");
598        let (indexer_tx, _rx) = mpsc::channel::<IndexEvent>(2);
599
600        let first = commit(
601            storage.as_ref(),
602            &WriteSinks::indexer_only(&indexer_tx),
603            &kb,
604            &path_a,
605            Bytes::from_static(b"# A"),
606            Some("text/markdown"),
607            ConditionalHeaders::default(),
608        )
609        .await;
610        let second = commit(
611            storage.as_ref(),
612            &WriteSinks::indexer_only(&indexer_tx),
613            &kb,
614            &path_b,
615            Bytes::from_static(b"# B"),
616            Some("text/markdown"),
617            ConditionalHeaders::default(),
618        )
619        .await;
620        let third = commit(
621            storage.as_ref(),
622            &WriteSinks::indexer_only(&indexer_tx),
623            &kb,
624            &path_c,
625            Bytes::from_static(b"# C"),
626            Some("text/markdown"),
627            ConditionalHeaders::default(),
628        )
629        .await;
630
631        assert!(first.is_ok(), "first write should fill queue slot one");
632        assert!(second.is_ok(), "second write should fill queue slot two");
633        let err = third.unwrap_err();
634        assert!(
635            matches!(err, WriteError::IndexerBackpressureUpsert),
636            "expected IndexerBackpressureUpsert, got {err:?}"
637        );
638        assert!(
639            storage
640                .objects
641                .lock()
642                .expect("mutex not poisoned")
643                .contains_key("test-kb/c.md"),
644            "stored object should remain after enqueue backpressure"
645        );
646    }
647
648    /// A write cancelled between storing the bytes and enqueuing the `Upsert`
649    /// must still enqueue it. The window is `publish`, which can await for
650    /// seconds against a slow or unreachable broker, and D71's request timeout
651    /// drops the handler where it stands — leaving, without the guard, an
652    /// object that is durable and unsearchable with nothing on `s3` to repair
653    /// it until the next reconciliation pass.
654    #[tokio::test]
655    async fn a_write_cancelled_while_publishing_still_enqueues_its_upsert() {
656        /// A broker that accepts the connection and then never answers.
657        struct NeverAnswers;
658
659        #[async_trait::async_trait]
660        impl EventPublisher for NeverAnswers {
661            async fn publish(&self, _event: ObjectEvent) -> Result<EventId, PublishError> {
662                std::future::pending().await
663            }
664            async fn subscribe(
665                &self,
666                _kb: &KbSlug,
667                _after: Option<EventId>,
668            ) -> Result<EventStream, SubscribeError> {
669                unreachable!("this test never subscribes")
670            }
671            fn ready(&self) -> bool {
672                true
673            }
674            fn backend_name(&self) -> &'static str {
675                "never-answers"
676            }
677        }
678
679        let storage = TestStorage::default();
680        let kb = kb();
681        let path = path_named("cancelled.md");
682        let (indexer_tx, mut rx) = mpsc::channel(4);
683        let events = NeverAnswers;
684        let sinks = WriteSinks {
685            indexer_tx: &indexer_tx,
686            events: Some(&events),
687            index_health: None,
688            source: EventSource::Http,
689        };
690
691        // Given a commit whose publish never returns, when the caller gives up
692        // on it — exactly what the request timeout does.
693        let cancelled = tokio::time::timeout(
694            std::time::Duration::from_millis(50),
695            commit(
696                &storage,
697                &sinks,
698                &kb,
699                &path,
700                Bytes::from_static(b"# Cancelled"),
701                Some("text/markdown"),
702                ConditionalHeaders::default(),
703            ),
704        )
705        .await;
706        assert!(cancelled.is_err(), "the publish should never have returned");
707
708        // Then: the bytes are stored …
709        assert!(
710            storage
711                .head_object(&kb, &path, ConditionalHeaders::default())
712                .await
713                .is_ok(),
714            "the store completed before the publish that was cancelled"
715        );
716        // … and the index was told about them anyway.
717        match rx.try_recv().expect("an Upsert must have been enqueued") {
718            IndexEvent::Upsert { object_key, .. } => assert_eq!(object_key, path),
719            other => panic!("expected an Upsert, got {other:?}"),
720        }
721    }
722
723    #[tokio::test]
724    async fn commit_delete_full_queue_returns_backpressure() {
725        let storage = TestStorage::default();
726        let kb = kb();
727        let path = path_named("delete.md");
728        storage
729            .put_object(
730                &kb,
731                &path,
732                Bytes::from_static(b"# Delete"),
733                Some("text/markdown"),
734                ConditionalHeaders::default(),
735            )
736            .await
737            .expect("prepopulate object");
738        let (indexer_tx, _rx) = mpsc::channel(1);
739        let dummy_event = IndexEvent::Tombstone {
740            kb: kb.clone(),
741            object_key: path.clone(),
742        };
743        indexer_tx
744            .try_send(dummy_event)
745            .expect("first send should succeed");
746
747        let outcome = commit_delete(
748            &storage,
749            &WriteSinks::indexer_only(&indexer_tx),
750            &kb,
751            &path,
752            ConditionalHeaders::default(),
753        )
754        .await;
755
756        let err = outcome.unwrap_err();
757        assert!(
758            matches!(err, WriteError::IndexerBackpressureTombstone),
759            "expected IndexerBackpressureTombstone, got {err:?}"
760        );
761        assert!(
762            !storage
763                .objects
764                .lock()
765                .expect("mutex not poisoned")
766                .contains_key("test-kb/delete.md"),
767            "deleted object should remain deleted after enqueue backpressure"
768        );
769    }
770
771    #[tokio::test]
772    async fn closed_queue_returns_write_success() {
773        let storage = TestStorage::default();
774        let kb = kb();
775        let path = path();
776        let (indexer_tx, rx) = mpsc::channel(1024);
777
778        drop(rx);
779
780        let outcome = commit(
781            &storage,
782            &WriteSinks::indexer_only(&indexer_tx),
783            &kb,
784            &path,
785            Bytes::from_static(b"# Test"),
786            Some("text/markdown"),
787            ConditionalHeaders::default(),
788        )
789        .await;
790
791        assert!(
792            outcome.is_ok(),
793            "write should succeed even if queue is closed"
794        );
795    }
796
797    #[tokio::test]
798    async fn put_failure_returns_error_no_event() {
799        let storage = TestStorage::default();
800        let kb = kb();
801        let path = path();
802        let (indexer_tx, mut rx) = mpsc::channel(1024);
803
804        let conditionals = ConditionalHeaders {
805            if_match: Some("\"wrong-etag\"".to_string()),
806            ..ConditionalHeaders::default()
807        };
808
809        let outcome = commit(
810            &storage,
811            &WriteSinks::indexer_only(&indexer_tx),
812            &kb,
813            &path,
814            Bytes::from_static(b"# Test"),
815            Some("text/markdown"),
816            conditionals,
817        )
818        .await;
819
820        assert!(outcome.is_err(), "put should fail with precondition");
821        assert!(
822            rx.try_recv().is_err(),
823            "no event should be enqueued on put failure"
824        );
825    }
826
827    #[test]
828    fn check_size_over_limit_returns_too_large() {
829        let err = check_size(MAX_UPLOAD_BYTES + 1, MAX_UPLOAD_BYTES).expect_err("too large");
830        assert!(matches!(
831            err,
832            WriteError::TooLarge {
833                size,
834                limit
835            } if size == MAX_UPLOAD_BYTES + 1 && limit == MAX_UPLOAD_BYTES
836        ));
837    }
838
839    #[test]
840    fn check_size_at_limit_returns_ok() {
841        assert!(check_size(MAX_UPLOAD_BYTES, MAX_UPLOAD_BYTES).is_ok());
842    }
843
844    #[test]
845    fn check_size_below_limit_returns_ok() {
846        assert!(check_size(1024, MAX_UPLOAD_BYTES).is_ok());
847    }
848
849    /// Records what was published, or refuses everything.
850    struct RecordingPublisher {
851        published: Mutex<Vec<ObjectEvent>>,
852        refuse: bool,
853    }
854
855    impl RecordingPublisher {
856        fn recording() -> Self {
857            Self {
858                published: Mutex::new(Vec::new()),
859                refuse: false,
860            }
861        }
862
863        fn refusing() -> Self {
864            Self {
865                published: Mutex::new(Vec::new()),
866                refuse: true,
867            }
868        }
869
870        fn events(&self) -> Vec<ObjectEvent> {
871            self.published.lock().expect("mutex not poisoned").clone()
872        }
873    }
874
875    #[async_trait]
876    impl notedthat_core::EventPublisher for RecordingPublisher {
877        async fn publish(
878            &self,
879            event: ObjectEvent,
880        ) -> Result<notedthat_core::EventId, notedthat_core::PublishError> {
881            if self.refuse {
882                return Err(notedthat_core::PublishError::Unavailable {
883                    message: "broker down".into(),
884                });
885            }
886            let mut published = self.published.lock().expect("mutex not poisoned");
887            published.push(event);
888            Ok(notedthat_core::EventId(published.len() as u64))
889        }
890
891        async fn subscribe(
892            &self,
893            _kb: &KbSlug,
894            _after: Option<notedthat_core::EventId>,
895        ) -> Result<notedthat_core::EventStream, notedthat_core::SubscribeError> {
896            unimplemented!()
897        }
898
899        fn ready(&self) -> bool {
900            true
901        }
902
903        fn backend_name(&self) -> &'static str {
904            "recording"
905        }
906    }
907
908    fn sinks<'a>(
909        indexer_tx: &'a mpsc::Sender<IndexEvent>,
910        events: &'a RecordingPublisher,
911        source: notedthat_core::EventSource,
912    ) -> WriteSinks<'a> {
913        WriteSinks {
914            indexer_tx,
915            events: Some(events),
916            index_health: None,
917            source,
918        }
919    }
920
921    #[tokio::test]
922    async fn a_committed_write_is_published_with_its_stamp_and_source() {
923        let storage = TestStorage::default();
924        let events = RecordingPublisher::recording();
925        let (indexer_tx, mut rx) = mpsc::channel(8);
926
927        let outcome = commit(
928            &storage,
929            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Webdav),
930            &kb(),
931            &path(),
932            Bytes::from_static(b"# Test"),
933            Some("text/markdown"),
934            ConditionalHeaders::default(),
935        )
936        .await
937        .expect("write succeeds");
938
939        rx.recv().await.expect("the index event, after the publish");
940        let published = events.events();
941        assert_eq!(published.len(), 1);
942        let event = &published[0];
943        assert_eq!(event.kb.as_str(), "test-kb");
944        assert_eq!(event.object_key.as_str(), "test.md");
945        assert_eq!(event.source, notedthat_core::EventSource::Webdav);
946        assert_eq!(
947            event.kind,
948            notedthat_core::ObjectEventKind::Written {
949                etag: outcome.etag.expect("etag"),
950                size: 6,
951                mime: "text/markdown".into(),
952                mtime: match &event.kind {
953                    notedthat_core::ObjectEventKind::Written { mtime, .. } => *mtime,
954                    other => unreachable!("a write publishes object.written, got {other:?}"),
955                },
956            }
957        );
958    }
959
960    /// The event goes out before the queue is tried, so a write the queue
961    /// refuses has been announced; the client's retry (D38) announces it
962    /// again — at least once, never zero times.
963    #[tokio::test]
964    async fn a_full_queue_still_refuses_the_write_after_its_event_was_published() {
965        let storage = TestStorage::default();
966        let events = RecordingPublisher::recording();
967        let (indexer_tx, _rx) = mpsc::channel(1);
968        indexer_tx
969            .try_send(IndexEvent::Tombstone {
970                kb: kb(),
971                object_key: path(),
972            })
973            .expect("fill the queue");
974
975        let err = commit(
976            &storage,
977            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
978            &kb(),
979            &path(),
980            Bytes::from_static(b"# Test"),
981            Some("text/markdown"),
982            ConditionalHeaders::default(),
983        )
984        .await
985        .unwrap_err();
986
987        assert!(matches!(err, WriteError::IndexerBackpressureUpsert));
988        let published = events.events();
989        assert_eq!(
990            published.len(),
991            1,
992            "announced before the 503: {published:?}"
993        );
994        assert_eq!(published[0].kind.name(), "object.written");
995    }
996
997    #[tokio::test]
998    async fn a_refused_publish_fails_the_write_but_the_object_stays_stored() {
999        let storage = TestStorage::default();
1000        let events = RecordingPublisher::refusing();
1001        let (indexer_tx, mut rx) = mpsc::channel(8);
1002
1003        let err = commit(
1004            &storage,
1005            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
1006            &kb(),
1007            &path(),
1008            Bytes::from_static(b"# Test"),
1009            Some("text/markdown"),
1010            ConditionalHeaders::default(),
1011        )
1012        .await
1013        .unwrap_err();
1014
1015        assert!(
1016            matches!(
1017                err,
1018                WriteError::EventPublishFailed {
1019                    after: WriteEffect::Stored
1020                }
1021            ),
1022            "{err:?}"
1023        );
1024        assert!(
1025            storage
1026                .objects
1027                .lock()
1028                .expect("mutex not poisoned")
1029                .contains_key("test-kb/test.md"),
1030            "the bytes were stored before publishing failed"
1031        );
1032        assert!(
1033            matches!(rx.try_recv(), Ok(IndexEvent::Upsert { .. })),
1034            "the index work was enqueued regardless: the broker is optional, the index is not"
1035        );
1036    }
1037
1038    #[tokio::test]
1039    async fn a_refused_publish_still_enqueues_the_tombstone() {
1040        let storage = TestStorage::default();
1041        let events = RecordingPublisher::refusing();
1042        let (indexer_tx, mut rx) = mpsc::channel(8);
1043
1044        let err = commit_delete(
1045            &storage,
1046            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
1047            &kb(),
1048            &path_named("gone.md"),
1049            ConditionalHeaders::default(),
1050        )
1051        .await
1052        .unwrap_err();
1053
1054        assert!(
1055            matches!(
1056                err,
1057                WriteError::EventPublishFailed {
1058                    after: WriteEffect::Deleted
1059                }
1060            ),
1061            "{err:?}"
1062        );
1063        assert!(
1064            matches!(rx.try_recv(), Ok(IndexEvent::Tombstone { .. })),
1065            "the points come out of the index whether or not the deletion was announced"
1066        );
1067    }
1068
1069    #[tokio::test]
1070    async fn a_delete_publishes_a_deleted_event_and_a_refusal_says_deleted() {
1071        let storage = TestStorage::default();
1072        let (indexer_tx, _rx) = mpsc::channel(8);
1073
1074        let events = RecordingPublisher::recording();
1075        commit_delete(
1076            &storage,
1077            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Mcp),
1078            &kb(),
1079            &path_named("gone.md"),
1080            ConditionalHeaders::default(),
1081        )
1082        .await
1083        .expect("delete succeeds");
1084        let published = events.events();
1085        assert_eq!(published.len(), 1);
1086        assert_eq!(published[0].kind, notedthat_core::ObjectEventKind::Deleted);
1087        assert_eq!(published[0].source, notedthat_core::EventSource::Mcp);
1088
1089        let refusing = RecordingPublisher::refusing();
1090        let err = commit_delete(
1091            &storage,
1092            &sinks(&indexer_tx, &refusing, notedthat_core::EventSource::Http),
1093            &kb(),
1094            &path_named("gone.md"),
1095            ConditionalHeaders::default(),
1096        )
1097        .await
1098        .unwrap_err();
1099        assert!(matches!(
1100            err,
1101            WriteError::EventPublishFailed {
1102                after: WriteEffect::Deleted
1103            }
1104        ));
1105    }
1106
1107    #[tokio::test]
1108    async fn a_copy_heads_the_destination_for_the_stamp_only_when_publishing() {
1109        let storage = TestStorage::default();
1110        let (indexer_tx, _rx) = mpsc::channel(8);
1111        let events = RecordingPublisher::recording();
1112
1113        let outcome = commit_copy(
1114            &storage,
1115            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Webdav),
1116            &kb(),
1117            &path_named("src.md"),
1118            &path_named("dst.md"),
1119            CopyObjectOptions::default(),
1120        )
1121        .await
1122        .expect("copy succeeds");
1123
1124        let published = events.events();
1125        assert_eq!(published.len(), 1);
1126        assert_eq!(published[0].object_key.as_str(), "dst.md");
1127        match &published[0].kind {
1128            notedthat_core::ObjectEventKind::Written {
1129                etag, size, mime, ..
1130            } => {
1131                assert_eq!(Some(etag), outcome.etag.as_ref());
1132                assert_eq!(*size, 7, "size comes from the HEAD after the copy");
1133                assert_eq!(mime, "text/markdown");
1134            }
1135            other => panic!("a copy is a write, got {other:?}"),
1136        }
1137    }
1138}