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::error::TrySendError;
9
10use crate::WriteError;
11use crate::error::WriteEffect;
12use crate::mime::sniff_content_type;
13use crate::sinks::WriteSinks;
14
15/// Maximum upload size accepted by shared write paths: 5 GiB.
16pub const MAX_UPLOAD_BYTES: u64 = 5 * 1024 * 1024 * 1024;
17
18/// Validate an upload size against a byte limit.
19pub fn check_size(size: u64, limit: u64) -> Result<(), WriteError> {
20    if size > limit {
21        Err(WriteError::TooLarge { size, limit })
22    } else {
23        Ok(())
24    }
25}
26
27/// Store an object, enqueue a best-effort index upsert event and publish the
28/// change.
29pub async fn commit<B: Into<StagedBody>>(
30    storage: &dyn Storage,
31    sinks: &WriteSinks<'_>,
32    kb: &KbSlug,
33    path: &ObjectPath,
34    body: B,
35    caller_content_type: Option<&str>,
36    conditionals: ConditionalHeaders,
37) -> Result<PutOutcome, WriteError> {
38    let body = body.into();
39    let size = body.len();
40    check_size(size, MAX_UPLOAD_BYTES)?;
41    let _prefix = body.prefix(512).await.map_err(|source| {
42        WriteError::Storage(StorageError::Other {
43            source: Box::new(source),
44        })
45    })?;
46    let mime = sniff_content_type(caller_content_type, path);
47    let outcome = storage
48        .put_staged_object(kb, path, body, Some(&mime), conditionals)
49        .await?;
50
51    after_write(sinks, kb, path, &outcome, size, &mime).await?;
52    Ok(outcome)
53}
54
55/// Copy an object natively, enqueue its destination for indexing and publish
56/// the change.
57pub async fn commit_copy(
58    storage: &dyn Storage,
59    sinks: &WriteSinks<'_>,
60    kb: &KbSlug,
61    source: &ObjectPath,
62    destination: &ObjectPath,
63    options: CopyObjectOptions,
64) -> Result<PutOutcome, WriteError> {
65    let content_type = options.content_type.clone();
66    let outcome = storage
67        .copy_object(kb, source, destination, options)
68        .await?;
69    // A copy returns only the new ETag. The event wants the size too, and only
70    // a subscriber cares, so the extra HEAD is paid only when one can exist.
71    // The copy is already durable by now, so a HEAD that fails — a racing
72    // delete, a transient error — degrades the stamp rather than the copy.
73    let (size, mime) = if sinks.events.is_some() {
74        match storage
75            .head_object(kb, destination, ConditionalHeaders::default())
76            .await
77        {
78            Ok(meta) => (meta.size, meta.content_type.or(content_type)),
79            Err(error) => {
80                tracing::warn!(
81                    target: "notedthat::events",
82                    kb = %kb, path = %destination, %error,
83                    "could not HEAD the copy destination for its event; publishing without a size"
84                );
85                (0, content_type)
86            }
87        }
88    } else {
89        (0, content_type)
90    };
91    after_write(
92        sinks,
93        kb,
94        destination,
95        &outcome,
96        size,
97        mime.as_deref().unwrap_or_default(),
98    )
99    .await?;
100    Ok(outcome)
101}
102
103/// Report a durable write: index first, then publish.
104///
105/// The order matters for D38's promise that a 503 makes the client the retry
106/// mechanism. Either failure returns an error after the bytes are stored, and
107/// a retried write re-runs both — an event may be published twice, never
108/// zero times.
109pub(crate) async fn after_write(
110    sinks: &WriteSinks<'_>,
111    kb: &KbSlug,
112    path: &ObjectPath,
113    outcome: &PutOutcome,
114    size: u64,
115    mime: &str,
116) -> Result<(), WriteError> {
117    let event = IndexEvent::Upsert {
118        kb: kb.clone(),
119        object_key: path.clone(),
120        etag: outcome.etag.clone().unwrap_or_default(),
121        mtime: current_unix_seconds(),
122    };
123    match sinks.indexer_tx.try_send(event) {
124        Ok(()) => {}
125        Err(TrySendError::Full(ev)) => {
126            tracing::warn!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_FULL");
127            let _ = ev;
128            return Err(WriteError::IndexerBackpressureUpsert);
129        }
130        Err(TrySendError::Closed(ev)) => {
131            tracing::error!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_CLOSED");
132            let _ = ev;
133            // Closed = indexer worker ended (shutdown OR panic). v1 preserves success-with-error-log
134            // until post-v1 worker liveness detection is added.
135        }
136    }
137
138    let Some(events) = sinks.events else {
139        return Ok(());
140    };
141    let event = ObjectEvent::written(
142        kb.clone(),
143        path.clone(),
144        outcome.etag.clone().unwrap_or_default(),
145        size,
146        mime.to_string(),
147        current_unix_seconds(),
148        sinks.source,
149    );
150    publish(events, event, WriteEffect::Stored).await
151}
152
153/// Delete an object idempotently, enqueue a best-effort tombstone event and
154/// publish the change.
155pub async fn commit_delete(
156    storage: &dyn Storage,
157    sinks: &WriteSinks<'_>,
158    kb: &KbSlug,
159    path: &ObjectPath,
160    conditionals: ConditionalHeaders,
161) -> Result<(), WriteError> {
162    match storage.delete_object(kb, path, conditionals).await {
163        Ok(()) | Err(StorageError::NotFound { .. }) => {}
164        Err(e) => return Err(WriteError::Storage(e)),
165    }
166
167    let event = IndexEvent::Tombstone {
168        kb: kb.clone(),
169        object_key: path.clone(),
170    };
171    match sinks.indexer_tx.try_send(event) {
172        Ok(()) => {}
173        Err(TrySendError::Full(ev)) => {
174            tracing::warn!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_FULL");
175            let _ = ev;
176            return Err(WriteError::IndexerBackpressureTombstone);
177        }
178        // 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, and post-v1 worker liveness detection will revisit this.
179        Err(TrySendError::Closed(ev)) => {
180            tracing::error!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_CLOSED");
181            let _ = ev;
182        }
183    }
184
185    let Some(events) = sinks.events else {
186        return Ok(());
187    };
188    let event = ObjectEvent::deleted(kb.clone(), path.clone(), sinks.source);
189    publish(events, event, WriteEffect::Deleted).await
190}
191
192async fn publish(
193    events: &dyn notedthat_core::EventPublisher,
194    event: ObjectEvent,
195    after: WriteEffect,
196) -> Result<(), WriteError> {
197    let (kb, path) = (event.kb.clone(), event.object_key.clone());
198    match events.publish(event).await {
199        Ok(_) => Ok(()),
200        Err(error) => {
201            tracing::error!(target: "notedthat::events", kb = %kb, path = %path, %error, "EVENT_PUBLISH_FAILED");
202            Err(WriteError::EventPublishFailed { after })
203        }
204    }
205}
206
207pub(crate) fn current_unix_seconds() -> i64 {
208    std::time::SystemTime::now()
209        .duration_since(std::time::UNIX_EPOCH)
210        .map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
211}
212
213#[cfg(test)]
214mod tests {
215    use super::*;
216    use crate::WriteEffect;
217    use async_trait::async_trait;
218    use bytes::Bytes;
219    use notedthat_core::{KbManifest, ListResponse, ObjectMeta, ObjectRead};
220    use std::collections::HashMap;
221    use std::sync::{Arc, Mutex};
222    use tokio::sync::mpsc;
223
224    #[derive(Default)]
225    struct TestStorage {
226        objects: Mutex<HashMap<String, String>>,
227    }
228
229    #[async_trait]
230    impl Storage for TestStorage {
231        async fn ensure_bucket(&self, _kb: &KbSlug) -> Result<(), StorageError> {
232            unimplemented!()
233        }
234
235        async fn read_manifest(&self, _kb: &KbSlug) -> Result<KbManifest, StorageError> {
236            unimplemented!()
237        }
238
239        async fn write_manifest(
240            &self,
241            _kb: &KbSlug,
242            _manifest: &KbManifest,
243        ) -> Result<(), StorageError> {
244            unimplemented!()
245        }
246
247        async fn head_object(
248            &self,
249            kb: &KbSlug,
250            path: &ObjectPath,
251            _conditionals: ConditionalHeaders,
252        ) -> Result<ObjectMeta, StorageError> {
253            let key = format!("{}/{}", kb.as_str(), path.as_str());
254            let objects = self.objects.lock().expect("mutex not poisoned");
255            let etag = objects
256                .get(&key)
257                .ok_or_else(|| StorageError::NotFound { key: key.clone() })?;
258            Ok(ObjectMeta {
259                key: path.as_str().to_string(),
260                size: 7,
261                last_modified: Some(1_700_000_000),
262                content_type: Some("text/markdown".into()),
263                etag: Some(etag.clone()),
264            })
265        }
266
267        async fn get_object(
268            &self,
269            _kb: &KbSlug,
270            _path: &ObjectPath,
271            _range: Option<Vec<notedthat_core::ByteRange>>,
272            _conditionals: ConditionalHeaders,
273        ) -> Result<ObjectRead, StorageError> {
274            unimplemented!()
275        }
276
277        async fn get_object_stream(
278            &self,
279            _kb: &KbSlug,
280            _path: &ObjectPath,
281            _range: Option<Vec<notedthat_core::ByteRange>>,
282            _conditionals: ConditionalHeaders,
283        ) -> Result<notedthat_core::ObjectStream, StorageError> {
284            unimplemented!()
285        }
286
287        async fn put_object(
288            &self,
289            kb: &KbSlug,
290            path: &ObjectPath,
291            _bytes: Bytes,
292            _content_type: Option<&str>,
293            conditionals: ConditionalHeaders,
294        ) -> Result<PutOutcome, StorageError> {
295            let key = format!("{}/{}", kb.as_str(), path.as_str());
296            let mut objects = self.objects.lock().expect("mutex not poisoned");
297            let existing = objects.get(&key);
298            if let Some(if_match) = conditionals.if_match
299                && existing.is_none_or(|etag| etag != &if_match)
300            {
301                return Err(StorageError::PreconditionFailed);
302            }
303
304            let etag = format!("\"etag-{}\"", objects.len() + 1);
305            objects.insert(key, etag.clone());
306            Ok(PutOutcome { etag: Some(etag) })
307        }
308
309        async fn put_staged_object(
310            &self,
311            kb: &KbSlug,
312            path: &ObjectPath,
313            body: StagedBody,
314            content_type: Option<&str>,
315            conditionals: ConditionalHeaders,
316        ) -> Result<PutOutcome, StorageError> {
317            let bytes =
318                body.memory_bytes()
319                    .cloned()
320                    .ok_or_else(|| StorageError::BackendUnavailable {
321                        message: "file staging is outside this commit unit test".into(),
322                    })?;
323            self.put_object(kb, path, bytes, content_type, conditionals)
324                .await
325        }
326
327        async fn copy_object(
328            &self,
329            kb: &KbSlug,
330            _source: &ObjectPath,
331            destination: &ObjectPath,
332            _options: CopyObjectOptions,
333        ) -> Result<PutOutcome, StorageError> {
334            let key = format!("{}/{}", kb.as_str(), destination.as_str());
335            let mut objects = self.objects.lock().expect("mutex not poisoned");
336            let etag = format!("\"etag-{}\"", objects.len() + 1);
337            objects.insert(key, etag.clone());
338            Ok(PutOutcome { etag: Some(etag) })
339        }
340
341        async fn delete_object(
342            &self,
343            kb: &KbSlug,
344            path: &ObjectPath,
345            _conditionals: ConditionalHeaders,
346        ) -> Result<(), StorageError> {
347            let key = format!("{}/{}", kb.as_str(), path.as_str());
348            self.objects
349                .lock()
350                .expect("mutex not poisoned")
351                .remove(&key);
352            Ok(())
353        }
354
355        async fn list_objects(
356            &self,
357            _kb: &KbSlug,
358            _prefix: Option<&str>,
359            _limit: u32,
360            _cursor: Option<&str>,
361        ) -> Result<ListResponse, StorageError> {
362            unimplemented!()
363        }
364    }
365
366    fn kb() -> KbSlug {
367        KbSlug::try_new("test-kb").expect("valid kb slug")
368    }
369
370    fn path() -> ObjectPath {
371        ObjectPath::try_from_str("test.md").expect("valid path")
372    }
373
374    fn path_named(value: &str) -> ObjectPath {
375        ObjectPath::try_from_str(value).expect("valid path")
376    }
377
378    #[tokio::test]
379    async fn successful_put_enqueues_event() {
380        let storage = TestStorage::default();
381        let kb = kb();
382        let path = path();
383        let (indexer_tx, mut rx) = mpsc::channel(1024);
384
385        let outcome = commit(
386            &storage,
387            &WriteSinks::indexer_only(&indexer_tx),
388            &kb,
389            &path,
390            Bytes::from_static(b"# Test"),
391            Some("text/markdown"),
392            ConditionalHeaders::default(),
393        )
394        .await;
395
396        assert!(outcome.is_ok());
397        assert!(outcome.unwrap().etag.is_some());
398
399        let event = rx.recv().await.expect("event should be enqueued");
400        assert_eq!(event.kb().as_str(), "test-kb");
401        assert_eq!(event.object_key().as_str(), "test.md");
402    }
403
404    #[tokio::test]
405    async fn successful_native_copy_enqueues_destination_event() {
406        let storage = TestStorage::default();
407        let kb = kb();
408        let source = path_named("source.md");
409        let destination = path_named("destination.md");
410        let (indexer_tx, mut rx) = mpsc::channel(1);
411
412        let outcome = commit_copy(
413            &storage,
414            &WriteSinks::indexer_only(&indexer_tx),
415            &kb,
416            &source,
417            &destination,
418            CopyObjectOptions::default(),
419        )
420        .await
421        .expect("copy succeeds");
422
423        assert!(outcome.etag.is_some());
424        let event = rx.recv().await.expect("destination event");
425        assert_eq!(event.object_key(), &destination);
426    }
427
428    #[tokio::test]
429    async fn full_queue_returns_indexer_backpressure() {
430        let storage = TestStorage::default();
431        let kb = kb();
432        let path = path();
433        let (indexer_tx, _rx) = mpsc::channel(1);
434
435        let dummy_event = IndexEvent::Upsert {
436            kb: kb.clone(),
437            object_key: path.clone(),
438            etag: "dummy".to_string(),
439            mtime: 0,
440        };
441        indexer_tx
442            .try_send(dummy_event)
443            .expect("first send should succeed");
444
445        let outcome = commit(
446            &storage,
447            &WriteSinks::indexer_only(&indexer_tx),
448            &kb,
449            &path,
450            Bytes::from_static(b"# Test"),
451            Some("text/markdown"),
452            ConditionalHeaders::default(),
453        )
454        .await;
455
456        let err = outcome.unwrap_err();
457        assert!(
458            matches!(err, WriteError::IndexerBackpressureUpsert),
459            "expected IndexerBackpressureUpsert, got {err:?}"
460        );
461    }
462
463    #[tokio::test]
464    async fn burst_write_returns_backpressure_after_capacity() {
465        let storage = Arc::new(TestStorage::default());
466        let kb = kb();
467        let path_a = path_named("a.md");
468        let path_b = path_named("b.md");
469        let path_c = path_named("c.md");
470        let (indexer_tx, _rx) = mpsc::channel::<IndexEvent>(2);
471
472        let first = commit(
473            storage.as_ref(),
474            &WriteSinks::indexer_only(&indexer_tx),
475            &kb,
476            &path_a,
477            Bytes::from_static(b"# A"),
478            Some("text/markdown"),
479            ConditionalHeaders::default(),
480        )
481        .await;
482        let second = commit(
483            storage.as_ref(),
484            &WriteSinks::indexer_only(&indexer_tx),
485            &kb,
486            &path_b,
487            Bytes::from_static(b"# B"),
488            Some("text/markdown"),
489            ConditionalHeaders::default(),
490        )
491        .await;
492        let third = commit(
493            storage.as_ref(),
494            &WriteSinks::indexer_only(&indexer_tx),
495            &kb,
496            &path_c,
497            Bytes::from_static(b"# C"),
498            Some("text/markdown"),
499            ConditionalHeaders::default(),
500        )
501        .await;
502
503        assert!(first.is_ok(), "first write should fill queue slot one");
504        assert!(second.is_ok(), "second write should fill queue slot two");
505        let err = third.unwrap_err();
506        assert!(
507            matches!(err, WriteError::IndexerBackpressureUpsert),
508            "expected IndexerBackpressureUpsert, got {err:?}"
509        );
510        assert!(
511            storage
512                .objects
513                .lock()
514                .expect("mutex not poisoned")
515                .contains_key("test-kb/c.md"),
516            "stored object should remain after enqueue backpressure"
517        );
518    }
519
520    #[tokio::test]
521    async fn commit_delete_full_queue_returns_backpressure() {
522        let storage = TestStorage::default();
523        let kb = kb();
524        let path = path_named("delete.md");
525        storage
526            .put_object(
527                &kb,
528                &path,
529                Bytes::from_static(b"# Delete"),
530                Some("text/markdown"),
531                ConditionalHeaders::default(),
532            )
533            .await
534            .expect("prepopulate object");
535        let (indexer_tx, _rx) = mpsc::channel(1);
536        let dummy_event = IndexEvent::Tombstone {
537            kb: kb.clone(),
538            object_key: path.clone(),
539        };
540        indexer_tx
541            .try_send(dummy_event)
542            .expect("first send should succeed");
543
544        let outcome = commit_delete(
545            &storage,
546            &WriteSinks::indexer_only(&indexer_tx),
547            &kb,
548            &path,
549            ConditionalHeaders::default(),
550        )
551        .await;
552
553        let err = outcome.unwrap_err();
554        assert!(
555            matches!(err, WriteError::IndexerBackpressureTombstone),
556            "expected IndexerBackpressureTombstone, got {err:?}"
557        );
558        assert!(
559            !storage
560                .objects
561                .lock()
562                .expect("mutex not poisoned")
563                .contains_key("test-kb/delete.md"),
564            "deleted object should remain deleted after enqueue backpressure"
565        );
566    }
567
568    #[tokio::test]
569    async fn closed_queue_returns_write_success() {
570        let storage = TestStorage::default();
571        let kb = kb();
572        let path = path();
573        let (indexer_tx, rx) = mpsc::channel(1024);
574
575        drop(rx);
576
577        let outcome = commit(
578            &storage,
579            &WriteSinks::indexer_only(&indexer_tx),
580            &kb,
581            &path,
582            Bytes::from_static(b"# Test"),
583            Some("text/markdown"),
584            ConditionalHeaders::default(),
585        )
586        .await;
587
588        assert!(
589            outcome.is_ok(),
590            "write should succeed even if queue is closed"
591        );
592    }
593
594    #[tokio::test]
595    async fn put_failure_returns_error_no_event() {
596        let storage = TestStorage::default();
597        let kb = kb();
598        let path = path();
599        let (indexer_tx, mut rx) = mpsc::channel(1024);
600
601        let conditionals = ConditionalHeaders {
602            if_match: Some("\"wrong-etag\"".to_string()),
603            ..ConditionalHeaders::default()
604        };
605
606        let outcome = commit(
607            &storage,
608            &WriteSinks::indexer_only(&indexer_tx),
609            &kb,
610            &path,
611            Bytes::from_static(b"# Test"),
612            Some("text/markdown"),
613            conditionals,
614        )
615        .await;
616
617        assert!(outcome.is_err(), "put should fail with precondition");
618        assert!(
619            rx.try_recv().is_err(),
620            "no event should be enqueued on put failure"
621        );
622    }
623
624    #[test]
625    fn check_size_over_limit_returns_too_large() {
626        let err = check_size(MAX_UPLOAD_BYTES + 1, MAX_UPLOAD_BYTES).expect_err("too large");
627        assert!(matches!(
628            err,
629            WriteError::TooLarge {
630                size,
631                limit
632            } if size == MAX_UPLOAD_BYTES + 1 && limit == MAX_UPLOAD_BYTES
633        ));
634    }
635
636    #[test]
637    fn check_size_at_limit_returns_ok() {
638        assert!(check_size(MAX_UPLOAD_BYTES, MAX_UPLOAD_BYTES).is_ok());
639    }
640
641    #[test]
642    fn check_size_below_limit_returns_ok() {
643        assert!(check_size(1024, MAX_UPLOAD_BYTES).is_ok());
644    }
645
646    /// Records what was published, or refuses everything.
647    struct RecordingPublisher {
648        published: Mutex<Vec<ObjectEvent>>,
649        refuse: bool,
650    }
651
652    impl RecordingPublisher {
653        fn recording() -> Self {
654            Self {
655                published: Mutex::new(Vec::new()),
656                refuse: false,
657            }
658        }
659
660        fn refusing() -> Self {
661            Self {
662                published: Mutex::new(Vec::new()),
663                refuse: true,
664            }
665        }
666
667        fn events(&self) -> Vec<ObjectEvent> {
668            self.published.lock().expect("mutex not poisoned").clone()
669        }
670    }
671
672    #[async_trait]
673    impl notedthat_core::EventPublisher for RecordingPublisher {
674        async fn publish(
675            &self,
676            event: ObjectEvent,
677        ) -> Result<notedthat_core::EventId, notedthat_core::PublishError> {
678            if self.refuse {
679                return Err(notedthat_core::PublishError::Unavailable {
680                    message: "broker down".into(),
681                });
682            }
683            let mut published = self.published.lock().expect("mutex not poisoned");
684            published.push(event);
685            Ok(notedthat_core::EventId(published.len() as u64))
686        }
687
688        async fn subscribe(
689            &self,
690            _kb: &KbSlug,
691            _after: Option<notedthat_core::EventId>,
692        ) -> Result<notedthat_core::EventStream, notedthat_core::SubscribeError> {
693            unimplemented!()
694        }
695
696        fn ready(&self) -> bool {
697            true
698        }
699
700        fn backend_name(&self) -> &'static str {
701            "recording"
702        }
703    }
704
705    fn sinks<'a>(
706        indexer_tx: &'a mpsc::Sender<IndexEvent>,
707        events: &'a RecordingPublisher,
708        source: notedthat_core::EventSource,
709    ) -> WriteSinks<'a> {
710        WriteSinks::new(indexer_tx, Some(events), source)
711    }
712
713    #[tokio::test]
714    async fn a_committed_write_is_published_with_its_stamp_and_source() {
715        let storage = TestStorage::default();
716        let events = RecordingPublisher::recording();
717        let (indexer_tx, mut rx) = mpsc::channel(8);
718
719        let outcome = commit(
720            &storage,
721            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Webdav),
722            &kb(),
723            &path(),
724            Bytes::from_static(b"# Test"),
725            Some("text/markdown"),
726            ConditionalHeaders::default(),
727        )
728        .await
729        .expect("write succeeds");
730
731        rx.recv().await.expect("index event first");
732        let published = events.events();
733        assert_eq!(published.len(), 1);
734        let event = &published[0];
735        assert_eq!(event.kb.as_str(), "test-kb");
736        assert_eq!(event.object_key.as_str(), "test.md");
737        assert_eq!(event.source, notedthat_core::EventSource::Webdav);
738        assert_eq!(
739            event.kind,
740            notedthat_core::ObjectEventKind::Written {
741                etag: outcome.etag.expect("etag"),
742                size: 6,
743                mime: "text/markdown".into(),
744                mtime: match &event.kind {
745                    notedthat_core::ObjectEventKind::Written { mtime, .. } => *mtime,
746                    notedthat_core::ObjectEventKind::Deleted => unreachable!(),
747                },
748            }
749        );
750    }
751
752    #[tokio::test]
753    async fn indexer_backpressure_wins_and_nothing_is_published() {
754        let storage = TestStorage::default();
755        let events = RecordingPublisher::recording();
756        let (indexer_tx, _rx) = mpsc::channel(1);
757        indexer_tx
758            .try_send(IndexEvent::Tombstone {
759                kb: kb(),
760                object_key: path(),
761            })
762            .expect("fill the queue");
763
764        let err = commit(
765            &storage,
766            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
767            &kb(),
768            &path(),
769            Bytes::from_static(b"# Test"),
770            Some("text/markdown"),
771            ConditionalHeaders::default(),
772        )
773        .await
774        .unwrap_err();
775
776        assert!(matches!(err, WriteError::IndexerBackpressureUpsert));
777        assert!(events.events().is_empty(), "no event behind a 503");
778    }
779
780    #[tokio::test]
781    async fn a_refused_publish_fails_the_write_but_the_object_stays_stored() {
782        let storage = TestStorage::default();
783        let events = RecordingPublisher::refusing();
784        let (indexer_tx, mut rx) = mpsc::channel(8);
785
786        let err = commit(
787            &storage,
788            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
789            &kb(),
790            &path(),
791            Bytes::from_static(b"# Test"),
792            Some("text/markdown"),
793            ConditionalHeaders::default(),
794        )
795        .await
796        .unwrap_err();
797
798        assert!(
799            matches!(
800                err,
801                WriteError::EventPublishFailed {
802                    after: WriteEffect::Stored
803                }
804            ),
805            "{err:?}"
806        );
807        assert!(
808            storage
809                .objects
810                .lock()
811                .expect("mutex not poisoned")
812                .contains_key("test-kb/test.md"),
813            "the bytes were stored before publishing failed"
814        );
815        rx.recv().await.expect("the index event was still enqueued");
816    }
817
818    #[tokio::test]
819    async fn a_delete_publishes_a_deleted_event_and_a_refusal_says_deleted() {
820        let storage = TestStorage::default();
821        let (indexer_tx, _rx) = mpsc::channel(8);
822
823        let events = RecordingPublisher::recording();
824        commit_delete(
825            &storage,
826            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Mcp),
827            &kb(),
828            &path_named("gone.md"),
829            ConditionalHeaders::default(),
830        )
831        .await
832        .expect("delete succeeds");
833        let published = events.events();
834        assert_eq!(published.len(), 1);
835        assert_eq!(published[0].kind, notedthat_core::ObjectEventKind::Deleted);
836        assert_eq!(published[0].source, notedthat_core::EventSource::Mcp);
837
838        let refusing = RecordingPublisher::refusing();
839        let err = commit_delete(
840            &storage,
841            &sinks(&indexer_tx, &refusing, notedthat_core::EventSource::Http),
842            &kb(),
843            &path_named("gone.md"),
844            ConditionalHeaders::default(),
845        )
846        .await
847        .unwrap_err();
848        assert!(matches!(
849            err,
850            WriteError::EventPublishFailed {
851                after: WriteEffect::Deleted
852            }
853        ));
854    }
855
856    #[tokio::test]
857    async fn a_copy_heads_the_destination_for_the_stamp_only_when_publishing() {
858        let storage = TestStorage::default();
859        let (indexer_tx, _rx) = mpsc::channel(8);
860        let events = RecordingPublisher::recording();
861
862        let outcome = commit_copy(
863            &storage,
864            &sinks(&indexer_tx, &events, notedthat_core::EventSource::Webdav),
865            &kb(),
866            &path_named("src.md"),
867            &path_named("dst.md"),
868            CopyObjectOptions::default(),
869        )
870        .await
871        .expect("copy succeeds");
872
873        let published = events.events();
874        assert_eq!(published.len(), 1);
875        assert_eq!(published[0].object_key.as_str(), "dst.md");
876        match &published[0].kind {
877            notedthat_core::ObjectEventKind::Written {
878                etag, size, mime, ..
879            } => {
880                assert_eq!(Some(etag), outcome.etag.as_ref());
881                assert_eq!(*size, 7, "size comes from the HEAD after the copy");
882                assert_eq!(mime, "text/markdown");
883            }
884            notedthat_core::ObjectEventKind::Deleted => panic!("a copy is a write"),
885        }
886    }
887}