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