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