1use notedthat_core::{
4 ConditionalHeaders, CopyObjectOptions, KbSlug, ObjectEvent, ObjectPath, PutOutcome, StagedBody,
5 Storage, StorageError,
6};
7use notedthat_indexer::IndexEvent;
8use tokio::sync::mpsc::Sender;
9use tokio::sync::mpsc::error::TrySendError;
10
11use crate::WriteError;
12use crate::error::WriteEffect;
13use crate::mime::sniff_content_type;
14use crate::sinks::WriteSinks;
15
16pub const MAX_UPLOAD_BYTES: u64 = 5 * 1024 * 1024 * 1024;
18
19pub fn check_size(size: u64, limit: u64) -> Result<(), WriteError> {
21 if size > limit {
22 Err(WriteError::TooLarge { size, limit })
23 } else {
24 Ok(())
25 }
26}
27
28pub async fn commit<B: Into<StagedBody>>(
31 storage: &dyn Storage,
32 sinks: &WriteSinks<'_>,
33 kb: &KbSlug,
34 path: &ObjectPath,
35 body: B,
36 caller_content_type: Option<&str>,
37 conditionals: ConditionalHeaders,
38) -> Result<PutOutcome, WriteError> {
39 let body = body.into();
40 let size = body.len();
41 check_size(size, MAX_UPLOAD_BYTES)?;
42 crate::manifest::check_staged_manifest(kb, path, &body).await?;
43 let _prefix = body.prefix(512).await.map_err(|source| {
44 WriteError::Storage(StorageError::Other {
45 source: Box::new(source),
46 })
47 })?;
48 let mime = sniff_content_type(caller_content_type, path);
49 let outcome = storage
50 .put_staged_object(kb, path, body, Some(&mime), conditionals)
51 .await?;
52
53 after_write(sinks, kb, path, &outcome, size, &mime).await?;
54 Ok(outcome)
55}
56
57pub async fn commit_copy(
60 storage: &dyn Storage,
61 sinks: &WriteSinks<'_>,
62 kb: &KbSlug,
63 source: &ObjectPath,
64 destination: &ObjectPath,
65 options: CopyObjectOptions,
66) -> Result<PutOutcome, WriteError> {
67 let content_type = options.content_type.clone();
68 if crate::manifest::is_manifest(destination) {
69 let meta = storage
75 .head_object(kb, source, ConditionalHeaders::default())
76 .await?;
77 if meta.size > crate::manifest::MANIFEST_MAX_BYTES {
78 return Err(WriteError::InvalidManifest {
79 message: format!(
80 "manifest is {} bytes; a manifest is at most {}",
81 meta.size,
82 crate::manifest::MANIFEST_MAX_BYTES
83 ),
84 });
85 }
86 let read = storage
87 .get_object(kb, source, None, ConditionalHeaders::default())
88 .await?;
89 crate::manifest::check_manifest_bytes(kb, destination, &read.bytes)?;
90 }
91 let outcome = storage
92 .copy_object(kb, source, destination, options)
93 .await?;
94 let (size, mime) = if sinks.events.is_some() {
99 match storage
100 .head_object(kb, destination, ConditionalHeaders::default())
101 .await
102 {
103 Ok(meta) => (meta.size, meta.content_type.or(content_type)),
104 Err(error) => {
105 tracing::warn!(
106 target: "notedthat::events",
107 kb = %kb, path = %destination, %error,
108 "could not HEAD the copy destination for its event; publishing without a size"
109 );
110 (0, content_type)
111 }
112 }
113 } else {
114 (0, content_type)
115 };
116 after_write(
117 sinks,
118 kb,
119 destination,
120 &outcome,
121 size,
122 mime.as_deref().unwrap_or_default(),
123 )
124 .await?;
125 Ok(outcome)
126}
127
128struct EnqueueOnDrop {
137 indexer_tx: Sender<IndexEvent>,
138 event: Option<IndexEvent>,
139}
140
141impl EnqueueOnDrop {
142 fn disarm(&mut self) -> IndexEvent {
144 self.event.take().expect("disarmed once")
145 }
146}
147
148impl Drop for EnqueueOnDrop {
149 fn drop(&mut self) {
150 let Some(event) = self.event.take() else {
151 return;
152 };
153 if self.indexer_tx.try_send(event).is_err() {
157 tracing::warn!(
158 target: "notedthat::indexing",
159 "INDEX_ENQUEUE_LOST_ON_CANCEL: a write was cancelled after its bytes were stored \
160 and its Upsert could not be queued; reconciliation will repair it"
161 );
162 }
163 }
164}
165
166pub(crate) async fn after_write(
177 sinks: &WriteSinks<'_>,
178 kb: &KbSlug,
179 path: &ObjectPath,
180 outcome: &PutOutcome,
181 size: u64,
182 mime: &str,
183) -> Result<(), WriteError> {
184 let mut enqueue = EnqueueOnDrop {
204 indexer_tx: sinks.indexer_tx.clone(),
205 event: Some(IndexEvent::Upsert {
206 kb: kb.clone(),
207 object_key: path.clone(),
208 etag: outcome.etag.clone().unwrap_or_default(),
209 mtime: current_unix_seconds(),
210 }),
211 };
212
213 let published = match sinks.events {
214 Some(events) => {
215 let event = ObjectEvent::written(
216 kb.clone(),
217 path.clone(),
218 outcome.etag.clone().unwrap_or_default(),
219 size,
220 mime.to_string(),
221 current_unix_seconds(),
222 sinks.source,
223 );
224 publish(events, event, WriteEffect::Stored).await
225 }
226 None => Ok(()),
227 };
228
229 let event = enqueue.disarm();
230 match sinks.indexer_tx.try_send(event) {
231 Ok(()) => {
232 if let Some(health) = sinks.index_health {
233 health.enqueued(kb.as_str());
234 }
235 }
236 Err(TrySendError::Full(ev)) => {
237 tracing::warn!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_FULL");
238 let _ = ev;
239 if let Some(health) = sinks.index_health {
240 health.backpressured(kb.as_str());
241 }
242 return Err(WriteError::IndexerBackpressureUpsert);
243 }
244 Err(TrySendError::Closed(ev)) => {
245 tracing::error!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_CLOSED");
246 let _ = ev;
247 if let Some(health) = sinks.index_health {
248 health.worker_stopped();
249 }
250 }
253 }
254 published
255}
256
257pub async fn commit_delete(
261 storage: &dyn Storage,
262 sinks: &WriteSinks<'_>,
263 kb: &KbSlug,
264 path: &ObjectPath,
265 conditionals: ConditionalHeaders,
266) -> Result<(), WriteError> {
267 match storage.delete_object(kb, path, conditionals).await {
268 Ok(()) | Err(StorageError::NotFound { .. }) => {}
269 Err(e) => return Err(WriteError::Storage(e)),
270 }
271
272 let published = match sinks.events {
276 Some(events) => {
277 let event = ObjectEvent::deleted(kb.clone(), path.clone(), sinks.source);
278 publish(events, event, WriteEffect::Deleted).await
279 }
280 None => Ok(()),
281 };
282
283 let event = IndexEvent::Tombstone {
284 kb: kb.clone(),
285 object_key: path.clone(),
286 };
287 match sinks.indexer_tx.try_send(event) {
288 Ok(()) => {
289 if let Some(health) = sinks.index_health {
290 health.enqueued(kb.as_str());
291 }
292 }
293 Err(TrySendError::Full(ev)) => {
294 tracing::warn!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_FULL");
295 let _ = ev;
296 if let Some(health) = sinks.index_health {
297 health.backpressured(kb.as_str());
298 }
299 return Err(WriteError::IndexerBackpressureTombstone);
300 }
301 Err(TrySendError::Closed(ev)) => {
303 tracing::error!(target: "notedthat::indexing", kb = %kb, path = %path, "INDEX_QUEUE_CLOSED");
304 let _ = ev;
305 if let Some(health) = sinks.index_health {
306 health.worker_stopped();
307 }
308 }
309 }
310 published
311}
312
313async fn publish(
314 events: &dyn notedthat_core::EventPublisher,
315 event: ObjectEvent,
316 after: WriteEffect,
317) -> Result<(), WriteError> {
318 let (kb, path) = (event.kb.clone(), event.object_key.clone());
319 match events.publish(event).await {
320 Ok(_) => Ok(()),
321 Err(error) => {
322 tracing::error!(target: "notedthat::events", kb = %kb, path = %path, %error, "EVENT_PUBLISH_FAILED");
323 Err(WriteError::EventPublishFailed { after })
324 }
325 }
326}
327
328pub(crate) fn current_unix_seconds() -> i64 {
329 std::time::SystemTime::now()
330 .duration_since(std::time::UNIX_EPOCH)
331 .map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
332}
333
334#[cfg(test)]
335mod tests {
336 use super::*;
337 use crate::WriteEffect;
338 use async_trait::async_trait;
339 use bytes::Bytes;
340 use notedthat_core::{
341 EventId, EventPublisher, EventSource, EventStream, KbManifest, ListResponse, ObjectEvent,
342 ObjectMeta, ObjectRead, PublishError, SubscribeError,
343 };
344 use std::collections::HashMap;
345 use std::sync::{Arc, Mutex};
346 use tokio::sync::mpsc;
347
348 #[derive(Default)]
349 struct TestStorage {
350 objects: Mutex<HashMap<String, String>>,
351 }
352
353 #[async_trait]
354 impl Storage for TestStorage {
355 async fn probe(&self, _kb: &KbSlug) -> Result<(), StorageError> {
356 Ok(())
357 }
358
359 async fn ensure_bucket(&self, _kb: &KbSlug) -> Result<(), StorageError> {
360 unimplemented!()
361 }
362
363 async fn read_manifest(&self, _kb: &KbSlug) -> Result<KbManifest, StorageError> {
364 unimplemented!()
365 }
366
367 async fn write_manifest(
368 &self,
369 _kb: &KbSlug,
370 _manifest: &KbManifest,
371 ) -> Result<(), StorageError> {
372 unimplemented!()
373 }
374
375 async fn head_object(
376 &self,
377 kb: &KbSlug,
378 path: &ObjectPath,
379 _conditionals: ConditionalHeaders,
380 ) -> Result<ObjectMeta, StorageError> {
381 let key = format!("{}/{}", kb.as_str(), path.as_str());
382 let objects = self.objects.lock().expect("mutex not poisoned");
383 let etag = objects
384 .get(&key)
385 .ok_or_else(|| StorageError::NotFound { key: key.clone() })?;
386 Ok(ObjectMeta {
387 key: path.as_str().to_string(),
388 size: 7,
389 last_modified: Some(1_700_000_000),
390 content_type: Some("text/markdown".into()),
391 etag: Some(etag.clone()),
392 })
393 }
394
395 async fn get_object(
396 &self,
397 _kb: &KbSlug,
398 _path: &ObjectPath,
399 _range: Option<notedthat_core::ByteRange>,
400 _conditionals: ConditionalHeaders,
401 ) -> Result<ObjectRead, StorageError> {
402 unimplemented!()
403 }
404
405 async fn get_object_stream(
406 &self,
407 _kb: &KbSlug,
408 _path: &ObjectPath,
409 _range: Option<notedthat_core::ByteRange>,
410 _conditionals: ConditionalHeaders,
411 ) -> Result<notedthat_core::ObjectStream, StorageError> {
412 unimplemented!()
413 }
414
415 async fn put_object(
416 &self,
417 kb: &KbSlug,
418 path: &ObjectPath,
419 _bytes: Bytes,
420 _content_type: Option<&str>,
421 conditionals: ConditionalHeaders,
422 ) -> Result<PutOutcome, StorageError> {
423 let key = format!("{}/{}", kb.as_str(), path.as_str());
424 let mut objects = self.objects.lock().expect("mutex not poisoned");
425 let existing = objects.get(&key);
426 if let Some(if_match) = conditionals.if_match
427 && existing.is_none_or(|etag| etag != &if_match)
428 {
429 return Err(StorageError::PreconditionFailed);
430 }
431
432 let etag = format!("\"etag-{}\"", objects.len() + 1);
433 objects.insert(key, etag.clone());
434 Ok(PutOutcome { etag: Some(etag) })
435 }
436
437 async fn put_staged_object(
438 &self,
439 kb: &KbSlug,
440 path: &ObjectPath,
441 body: StagedBody,
442 content_type: Option<&str>,
443 conditionals: ConditionalHeaders,
444 ) -> Result<PutOutcome, StorageError> {
445 let bytes =
446 body.memory_bytes()
447 .cloned()
448 .ok_or_else(|| StorageError::BackendUnavailable {
449 message: "file staging is outside this commit unit test".into(),
450 })?;
451 self.put_object(kb, path, bytes, content_type, conditionals)
452 .await
453 }
454
455 async fn copy_object(
456 &self,
457 kb: &KbSlug,
458 _source: &ObjectPath,
459 destination: &ObjectPath,
460 _options: CopyObjectOptions,
461 ) -> Result<PutOutcome, StorageError> {
462 let key = format!("{}/{}", kb.as_str(), destination.as_str());
463 let mut objects = self.objects.lock().expect("mutex not poisoned");
464 let etag = format!("\"etag-{}\"", objects.len() + 1);
465 objects.insert(key, etag.clone());
466 Ok(PutOutcome { etag: Some(etag) })
467 }
468
469 async fn delete_object(
470 &self,
471 kb: &KbSlug,
472 path: &ObjectPath,
473 _conditionals: ConditionalHeaders,
474 ) -> Result<(), StorageError> {
475 let key = format!("{}/{}", kb.as_str(), path.as_str());
476 self.objects
477 .lock()
478 .expect("mutex not poisoned")
479 .remove(&key);
480 Ok(())
481 }
482
483 async fn list_objects(
484 &self,
485 _kb: &KbSlug,
486 _prefix: Option<&str>,
487 _limit: u32,
488 _cursor: Option<&str>,
489 ) -> Result<ListResponse, StorageError> {
490 unimplemented!()
491 }
492 }
493
494 fn kb() -> KbSlug {
495 KbSlug::try_new("test-kb").expect("valid kb slug")
496 }
497
498 fn path() -> ObjectPath {
499 ObjectPath::try_from_str("test.md").expect("valid path")
500 }
501
502 fn path_named(value: &str) -> ObjectPath {
503 ObjectPath::try_from_str(value).expect("valid path")
504 }
505
506 #[tokio::test]
507 async fn successful_put_enqueues_event() {
508 let storage = TestStorage::default();
509 let kb = kb();
510 let path = path();
511 let (indexer_tx, mut rx) = mpsc::channel(1024);
512
513 let outcome = commit(
514 &storage,
515 &WriteSinks::indexer_only(&indexer_tx),
516 &kb,
517 &path,
518 Bytes::from_static(b"# Test"),
519 Some("text/markdown"),
520 ConditionalHeaders::default(),
521 )
522 .await;
523
524 assert!(outcome.is_ok());
525 assert!(outcome.unwrap().etag.is_some());
526
527 let event = rx.recv().await.expect("event should be enqueued");
528 assert_eq!(event.kb().as_str(), "test-kb");
529 assert_eq!(event.object_key().as_str(), "test.md");
530 }
531
532 #[tokio::test]
533 async fn successful_native_copy_enqueues_destination_event() {
534 let storage = TestStorage::default();
535 let kb = kb();
536 let source = path_named("source.md");
537 let destination = path_named("destination.md");
538 let (indexer_tx, mut rx) = mpsc::channel(1);
539
540 let outcome = commit_copy(
541 &storage,
542 &WriteSinks::indexer_only(&indexer_tx),
543 &kb,
544 &source,
545 &destination,
546 CopyObjectOptions::default(),
547 )
548 .await
549 .expect("copy succeeds");
550
551 assert!(outcome.etag.is_some());
552 let event = rx.recv().await.expect("destination event");
553 assert_eq!(event.object_key(), &destination);
554 }
555
556 #[tokio::test]
557 async fn full_queue_returns_indexer_backpressure() {
558 let storage = TestStorage::default();
559 let kb = kb();
560 let path = path();
561 let (indexer_tx, _rx) = mpsc::channel(1);
562
563 let dummy_event = IndexEvent::Upsert {
564 kb: kb.clone(),
565 object_key: path.clone(),
566 etag: "dummy".to_string(),
567 mtime: 0,
568 };
569 indexer_tx
570 .try_send(dummy_event)
571 .expect("first send should succeed");
572
573 let outcome = commit(
574 &storage,
575 &WriteSinks::indexer_only(&indexer_tx),
576 &kb,
577 &path,
578 Bytes::from_static(b"# Test"),
579 Some("text/markdown"),
580 ConditionalHeaders::default(),
581 )
582 .await;
583
584 let err = outcome.unwrap_err();
585 assert!(
586 matches!(err, WriteError::IndexerBackpressureUpsert),
587 "expected IndexerBackpressureUpsert, got {err:?}"
588 );
589 }
590
591 #[tokio::test]
592 async fn burst_write_returns_backpressure_after_capacity() {
593 let storage = Arc::new(TestStorage::default());
594 let kb = kb();
595 let path_a = path_named("a.md");
596 let path_b = path_named("b.md");
597 let path_c = path_named("c.md");
598 let (indexer_tx, _rx) = mpsc::channel::<IndexEvent>(2);
599
600 let first = commit(
601 storage.as_ref(),
602 &WriteSinks::indexer_only(&indexer_tx),
603 &kb,
604 &path_a,
605 Bytes::from_static(b"# A"),
606 Some("text/markdown"),
607 ConditionalHeaders::default(),
608 )
609 .await;
610 let second = commit(
611 storage.as_ref(),
612 &WriteSinks::indexer_only(&indexer_tx),
613 &kb,
614 &path_b,
615 Bytes::from_static(b"# B"),
616 Some("text/markdown"),
617 ConditionalHeaders::default(),
618 )
619 .await;
620 let third = commit(
621 storage.as_ref(),
622 &WriteSinks::indexer_only(&indexer_tx),
623 &kb,
624 &path_c,
625 Bytes::from_static(b"# C"),
626 Some("text/markdown"),
627 ConditionalHeaders::default(),
628 )
629 .await;
630
631 assert!(first.is_ok(), "first write should fill queue slot one");
632 assert!(second.is_ok(), "second write should fill queue slot two");
633 let err = third.unwrap_err();
634 assert!(
635 matches!(err, WriteError::IndexerBackpressureUpsert),
636 "expected IndexerBackpressureUpsert, got {err:?}"
637 );
638 assert!(
639 storage
640 .objects
641 .lock()
642 .expect("mutex not poisoned")
643 .contains_key("test-kb/c.md"),
644 "stored object should remain after enqueue backpressure"
645 );
646 }
647
648 #[tokio::test]
655 async fn a_write_cancelled_while_publishing_still_enqueues_its_upsert() {
656 struct NeverAnswers;
658
659 #[async_trait::async_trait]
660 impl EventPublisher for NeverAnswers {
661 async fn publish(&self, _event: ObjectEvent) -> Result<EventId, PublishError> {
662 std::future::pending().await
663 }
664 async fn subscribe(
665 &self,
666 _kb: &KbSlug,
667 _after: Option<EventId>,
668 ) -> Result<EventStream, SubscribeError> {
669 unreachable!("this test never subscribes")
670 }
671 fn ready(&self) -> bool {
672 true
673 }
674 fn backend_name(&self) -> &'static str {
675 "never-answers"
676 }
677 }
678
679 let storage = TestStorage::default();
680 let kb = kb();
681 let path = path_named("cancelled.md");
682 let (indexer_tx, mut rx) = mpsc::channel(4);
683 let events = NeverAnswers;
684 let sinks = WriteSinks {
685 indexer_tx: &indexer_tx,
686 events: Some(&events),
687 index_health: None,
688 source: EventSource::Http,
689 };
690
691 let cancelled = tokio::time::timeout(
694 std::time::Duration::from_millis(50),
695 commit(
696 &storage,
697 &sinks,
698 &kb,
699 &path,
700 Bytes::from_static(b"# Cancelled"),
701 Some("text/markdown"),
702 ConditionalHeaders::default(),
703 ),
704 )
705 .await;
706 assert!(cancelled.is_err(), "the publish should never have returned");
707
708 assert!(
710 storage
711 .head_object(&kb, &path, ConditionalHeaders::default())
712 .await
713 .is_ok(),
714 "the store completed before the publish that was cancelled"
715 );
716 match rx.try_recv().expect("an Upsert must have been enqueued") {
718 IndexEvent::Upsert { object_key, .. } => assert_eq!(object_key, path),
719 other => panic!("expected an Upsert, got {other:?}"),
720 }
721 }
722
723 #[tokio::test]
724 async fn commit_delete_full_queue_returns_backpressure() {
725 let storage = TestStorage::default();
726 let kb = kb();
727 let path = path_named("delete.md");
728 storage
729 .put_object(
730 &kb,
731 &path,
732 Bytes::from_static(b"# Delete"),
733 Some("text/markdown"),
734 ConditionalHeaders::default(),
735 )
736 .await
737 .expect("prepopulate object");
738 let (indexer_tx, _rx) = mpsc::channel(1);
739 let dummy_event = IndexEvent::Tombstone {
740 kb: kb.clone(),
741 object_key: path.clone(),
742 };
743 indexer_tx
744 .try_send(dummy_event)
745 .expect("first send should succeed");
746
747 let outcome = commit_delete(
748 &storage,
749 &WriteSinks::indexer_only(&indexer_tx),
750 &kb,
751 &path,
752 ConditionalHeaders::default(),
753 )
754 .await;
755
756 let err = outcome.unwrap_err();
757 assert!(
758 matches!(err, WriteError::IndexerBackpressureTombstone),
759 "expected IndexerBackpressureTombstone, got {err:?}"
760 );
761 assert!(
762 !storage
763 .objects
764 .lock()
765 .expect("mutex not poisoned")
766 .contains_key("test-kb/delete.md"),
767 "deleted object should remain deleted after enqueue backpressure"
768 );
769 }
770
771 #[tokio::test]
772 async fn closed_queue_returns_write_success() {
773 let storage = TestStorage::default();
774 let kb = kb();
775 let path = path();
776 let (indexer_tx, rx) = mpsc::channel(1024);
777
778 drop(rx);
779
780 let outcome = commit(
781 &storage,
782 &WriteSinks::indexer_only(&indexer_tx),
783 &kb,
784 &path,
785 Bytes::from_static(b"# Test"),
786 Some("text/markdown"),
787 ConditionalHeaders::default(),
788 )
789 .await;
790
791 assert!(
792 outcome.is_ok(),
793 "write should succeed even if queue is closed"
794 );
795 }
796
797 #[tokio::test]
798 async fn put_failure_returns_error_no_event() {
799 let storage = TestStorage::default();
800 let kb = kb();
801 let path = path();
802 let (indexer_tx, mut rx) = mpsc::channel(1024);
803
804 let conditionals = ConditionalHeaders {
805 if_match: Some("\"wrong-etag\"".to_string()),
806 ..ConditionalHeaders::default()
807 };
808
809 let outcome = commit(
810 &storage,
811 &WriteSinks::indexer_only(&indexer_tx),
812 &kb,
813 &path,
814 Bytes::from_static(b"# Test"),
815 Some("text/markdown"),
816 conditionals,
817 )
818 .await;
819
820 assert!(outcome.is_err(), "put should fail with precondition");
821 assert!(
822 rx.try_recv().is_err(),
823 "no event should be enqueued on put failure"
824 );
825 }
826
827 #[test]
828 fn check_size_over_limit_returns_too_large() {
829 let err = check_size(MAX_UPLOAD_BYTES + 1, MAX_UPLOAD_BYTES).expect_err("too large");
830 assert!(matches!(
831 err,
832 WriteError::TooLarge {
833 size,
834 limit
835 } if size == MAX_UPLOAD_BYTES + 1 && limit == MAX_UPLOAD_BYTES
836 ));
837 }
838
839 #[test]
840 fn check_size_at_limit_returns_ok() {
841 assert!(check_size(MAX_UPLOAD_BYTES, MAX_UPLOAD_BYTES).is_ok());
842 }
843
844 #[test]
845 fn check_size_below_limit_returns_ok() {
846 assert!(check_size(1024, MAX_UPLOAD_BYTES).is_ok());
847 }
848
849 struct RecordingPublisher {
851 published: Mutex<Vec<ObjectEvent>>,
852 refuse: bool,
853 }
854
855 impl RecordingPublisher {
856 fn recording() -> Self {
857 Self {
858 published: Mutex::new(Vec::new()),
859 refuse: false,
860 }
861 }
862
863 fn refusing() -> Self {
864 Self {
865 published: Mutex::new(Vec::new()),
866 refuse: true,
867 }
868 }
869
870 fn events(&self) -> Vec<ObjectEvent> {
871 self.published.lock().expect("mutex not poisoned").clone()
872 }
873 }
874
875 #[async_trait]
876 impl notedthat_core::EventPublisher for RecordingPublisher {
877 async fn publish(
878 &self,
879 event: ObjectEvent,
880 ) -> Result<notedthat_core::EventId, notedthat_core::PublishError> {
881 if self.refuse {
882 return Err(notedthat_core::PublishError::Unavailable {
883 message: "broker down".into(),
884 });
885 }
886 let mut published = self.published.lock().expect("mutex not poisoned");
887 published.push(event);
888 Ok(notedthat_core::EventId(published.len() as u64))
889 }
890
891 async fn subscribe(
892 &self,
893 _kb: &KbSlug,
894 _after: Option<notedthat_core::EventId>,
895 ) -> Result<notedthat_core::EventStream, notedthat_core::SubscribeError> {
896 unimplemented!()
897 }
898
899 fn ready(&self) -> bool {
900 true
901 }
902
903 fn backend_name(&self) -> &'static str {
904 "recording"
905 }
906 }
907
908 fn sinks<'a>(
909 indexer_tx: &'a mpsc::Sender<IndexEvent>,
910 events: &'a RecordingPublisher,
911 source: notedthat_core::EventSource,
912 ) -> WriteSinks<'a> {
913 WriteSinks {
914 indexer_tx,
915 events: Some(events),
916 index_health: None,
917 source,
918 }
919 }
920
921 #[tokio::test]
922 async fn a_committed_write_is_published_with_its_stamp_and_source() {
923 let storage = TestStorage::default();
924 let events = RecordingPublisher::recording();
925 let (indexer_tx, mut rx) = mpsc::channel(8);
926
927 let outcome = commit(
928 &storage,
929 &sinks(&indexer_tx, &events, notedthat_core::EventSource::Webdav),
930 &kb(),
931 &path(),
932 Bytes::from_static(b"# Test"),
933 Some("text/markdown"),
934 ConditionalHeaders::default(),
935 )
936 .await
937 .expect("write succeeds");
938
939 rx.recv().await.expect("the index event, after the publish");
940 let published = events.events();
941 assert_eq!(published.len(), 1);
942 let event = &published[0];
943 assert_eq!(event.kb.as_str(), "test-kb");
944 assert_eq!(event.object_key.as_str(), "test.md");
945 assert_eq!(event.source, notedthat_core::EventSource::Webdav);
946 assert_eq!(
947 event.kind,
948 notedthat_core::ObjectEventKind::Written {
949 etag: outcome.etag.expect("etag"),
950 size: 6,
951 mime: "text/markdown".into(),
952 mtime: match &event.kind {
953 notedthat_core::ObjectEventKind::Written { mtime, .. } => *mtime,
954 other => unreachable!("a write publishes object.written, got {other:?}"),
955 },
956 }
957 );
958 }
959
960 #[tokio::test]
964 async fn a_full_queue_still_refuses_the_write_after_its_event_was_published() {
965 let storage = TestStorage::default();
966 let events = RecordingPublisher::recording();
967 let (indexer_tx, _rx) = mpsc::channel(1);
968 indexer_tx
969 .try_send(IndexEvent::Tombstone {
970 kb: kb(),
971 object_key: path(),
972 })
973 .expect("fill the queue");
974
975 let err = commit(
976 &storage,
977 &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
978 &kb(),
979 &path(),
980 Bytes::from_static(b"# Test"),
981 Some("text/markdown"),
982 ConditionalHeaders::default(),
983 )
984 .await
985 .unwrap_err();
986
987 assert!(matches!(err, WriteError::IndexerBackpressureUpsert));
988 let published = events.events();
989 assert_eq!(
990 published.len(),
991 1,
992 "announced before the 503: {published:?}"
993 );
994 assert_eq!(published[0].kind.name(), "object.written");
995 }
996
997 #[tokio::test]
998 async fn a_refused_publish_fails_the_write_but_the_object_stays_stored() {
999 let storage = TestStorage::default();
1000 let events = RecordingPublisher::refusing();
1001 let (indexer_tx, mut rx) = mpsc::channel(8);
1002
1003 let err = commit(
1004 &storage,
1005 &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
1006 &kb(),
1007 &path(),
1008 Bytes::from_static(b"# Test"),
1009 Some("text/markdown"),
1010 ConditionalHeaders::default(),
1011 )
1012 .await
1013 .unwrap_err();
1014
1015 assert!(
1016 matches!(
1017 err,
1018 WriteError::EventPublishFailed {
1019 after: WriteEffect::Stored
1020 }
1021 ),
1022 "{err:?}"
1023 );
1024 assert!(
1025 storage
1026 .objects
1027 .lock()
1028 .expect("mutex not poisoned")
1029 .contains_key("test-kb/test.md"),
1030 "the bytes were stored before publishing failed"
1031 );
1032 assert!(
1033 matches!(rx.try_recv(), Ok(IndexEvent::Upsert { .. })),
1034 "the index work was enqueued regardless: the broker is optional, the index is not"
1035 );
1036 }
1037
1038 #[tokio::test]
1039 async fn a_refused_publish_still_enqueues_the_tombstone() {
1040 let storage = TestStorage::default();
1041 let events = RecordingPublisher::refusing();
1042 let (indexer_tx, mut rx) = mpsc::channel(8);
1043
1044 let err = commit_delete(
1045 &storage,
1046 &sinks(&indexer_tx, &events, notedthat_core::EventSource::Http),
1047 &kb(),
1048 &path_named("gone.md"),
1049 ConditionalHeaders::default(),
1050 )
1051 .await
1052 .unwrap_err();
1053
1054 assert!(
1055 matches!(
1056 err,
1057 WriteError::EventPublishFailed {
1058 after: WriteEffect::Deleted
1059 }
1060 ),
1061 "{err:?}"
1062 );
1063 assert!(
1064 matches!(rx.try_recv(), Ok(IndexEvent::Tombstone { .. })),
1065 "the points come out of the index whether or not the deletion was announced"
1066 );
1067 }
1068
1069 #[tokio::test]
1070 async fn a_delete_publishes_a_deleted_event_and_a_refusal_says_deleted() {
1071 let storage = TestStorage::default();
1072 let (indexer_tx, _rx) = mpsc::channel(8);
1073
1074 let events = RecordingPublisher::recording();
1075 commit_delete(
1076 &storage,
1077 &sinks(&indexer_tx, &events, notedthat_core::EventSource::Mcp),
1078 &kb(),
1079 &path_named("gone.md"),
1080 ConditionalHeaders::default(),
1081 )
1082 .await
1083 .expect("delete succeeds");
1084 let published = events.events();
1085 assert_eq!(published.len(), 1);
1086 assert_eq!(published[0].kind, notedthat_core::ObjectEventKind::Deleted);
1087 assert_eq!(published[0].source, notedthat_core::EventSource::Mcp);
1088
1089 let refusing = RecordingPublisher::refusing();
1090 let err = commit_delete(
1091 &storage,
1092 &sinks(&indexer_tx, &refusing, notedthat_core::EventSource::Http),
1093 &kb(),
1094 &path_named("gone.md"),
1095 ConditionalHeaders::default(),
1096 )
1097 .await
1098 .unwrap_err();
1099 assert!(matches!(
1100 err,
1101 WriteError::EventPublishFailed {
1102 after: WriteEffect::Deleted
1103 }
1104 ));
1105 }
1106
1107 #[tokio::test]
1108 async fn a_copy_heads_the_destination_for_the_stamp_only_when_publishing() {
1109 let storage = TestStorage::default();
1110 let (indexer_tx, _rx) = mpsc::channel(8);
1111 let events = RecordingPublisher::recording();
1112
1113 let outcome = commit_copy(
1114 &storage,
1115 &sinks(&indexer_tx, &events, notedthat_core::EventSource::Webdav),
1116 &kb(),
1117 &path_named("src.md"),
1118 &path_named("dst.md"),
1119 CopyObjectOptions::default(),
1120 )
1121 .await
1122 .expect("copy succeeds");
1123
1124 let published = events.events();
1125 assert_eq!(published.len(), 1);
1126 assert_eq!(published[0].object_key.as_str(), "dst.md");
1127 match &published[0].kind {
1128 notedthat_core::ObjectEventKind::Written {
1129 etag, size, mime, ..
1130 } => {
1131 assert_eq!(Some(etag), outcome.etag.as_ref());
1132 assert_eq!(*size, 7, "size comes from the HEAD after the copy");
1133 assert_eq!(mime, "text/markdown");
1134 }
1135 other => panic!("a copy is a write, got {other:?}"),
1136 }
1137 }
1138}