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