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 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
56pub 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 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 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
127pub(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 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 }
199 }
200 published
201}
202
203pub 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 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 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 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 #[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}