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