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