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