1pub(crate) use helpers::{kb_not_found, lookup_kb};
6
7mod browse;
8mod events;
9mod health;
10pub(crate) mod helpers;
11mod index_health;
12mod index_reconcile;
13mod kbs;
14mod llms;
15mod objects;
16mod well_known;
17
18use crate::bounds::{RequestBounds, bound};
19use crate::middleware::auth_middleware;
20use crate::state::AppState;
21use axum::Router;
22use axum::extract::{DefaultBodyLimit, Request};
23use axum::handler::Handler;
24use axum::http::HeaderName;
25use axum::middleware::from_fn_with_state;
26use axum::routing::get;
27use tower::ServiceBuilder;
28use tower_http::request_id::{
29 MakeRequestId, PropagateRequestIdLayer, RequestId, SetRequestIdLayer,
30};
31use tower_http::trace::TraceLayer;
32use uuid::Uuid;
33
34use browse::{browse_path, browse_root};
35use events::subscribe_events;
36use health::{healthz, readyz};
37use index_health::get_index_health;
38use index_reconcile::post_index_reconcile;
39use kbs::{list_kbs, list_objects};
40use llms::llms_txt;
41use objects::{delete_object, get_object, head_object, patch_object, post_object, put_object};
42use well_known::protected_resource_metadata;
43
44macro_rules! api_routes {
57 ($($route:ident / $matched:ident => $suffix:literal,)+) => {
58 $(
59 pub(crate) const $route: &str = $suffix;
60 pub(crate) const $matched: &str = concat!("/api/v1", $suffix);
61 )+
62 };
63}
64
65pub const API_V1_PREFIX: &str = "/api/v1";
67
68pub const BROWSE_PREFIX: &str = "/browse";
70
71api_routes! {
72 ROUTE_KBS / MATCHED_KBS => "/knowledgebases",
73 ROUTE_KB / MATCHED_KB => "/knowledgebases/{kb_slug}",
74 ROUTE_KB_SEARCH / MATCHED_KB_SEARCH => "/knowledgebases/{kb_slug}/search",
75 ROUTE_KB_EVENTS / MATCHED_KB_EVENTS => "/knowledgebases/{kb_slug}/events",
76 ROUTE_KB_INDEX / MATCHED_KB_INDEX => "/knowledgebases/{kb_slug}/index",
77 ROUTE_KB_INDEX_RECONCILE / MATCHED_KB_INDEX_RECONCILE => "/knowledgebases/{kb_slug}/index/reconcile",
78 ROUTE_KB_OBJECT / MATCHED_KB_OBJECT => "/knowledgebases/{kb_slug}/{*object_path}",
79}
80
81pub const MAX_BODY_BYTES: u64 = 16 * 1024 * 1024;
83
84#[derive(Clone, Copy, Default)]
86pub struct MakeRequestUuidV7;
87
88impl MakeRequestId for MakeRequestUuidV7 {
89 fn make_request_id<B>(&mut self, _req: &Request<B>) -> Option<RequestId> {
90 let id = Uuid::now_v7().to_string();
91 let hv = id.parse().ok()?;
92 Some(RequestId::new(hv))
93 }
94}
95
96pub fn build_router(state: AppState) -> Router {
102 build_bounded_router(state, &RequestBounds::unbounded())
103}
104
105pub fn build_bounded_router(state: AppState, bounds: &RequestBounds) -> Router {
119 let request_id_header = HeaderName::from_static("x-request-id");
120 let bounded_api = Router::new()
121 .route(ROUTE_KBS, get(list_kbs))
122 .route(ROUTE_KB, get(list_objects))
123 .route(
124 ROUTE_KB_SEARCH,
125 axum::routing::post(crate::search_route::search_kb).layer(
126 axum::extract::DefaultBodyLimit::max(crate::search_route::SEARCH_BODY_MAX_BYTES),
127 ),
128 )
129 .route(ROUTE_KB_INDEX, get(get_index_health))
130 .route(
131 ROUTE_KB_INDEX_RECONCILE,
132 axum::routing::post(post_index_reconcile),
133 )
134 .route(
135 ROUTE_KB_OBJECT,
136 get(get_object)
137 .head(head_object)
138 .put(put_object)
139 .delete(delete_object)
140 .patch(patch_object.layer(DefaultBodyLimit::disable()))
141 .post(post_object.layer(DefaultBodyLimit::disable())),
142 )
143 .route_layer(from_fn_with_state(bounds.clone(), bound));
144 let streaming_api = Router::new().route(ROUTE_KB_EVENTS, get(subscribe_events));
145 let api_routes = bounded_api
146 .merge(streaming_api)
147 .layer(
148 ServiceBuilder::new()
149 .layer(DefaultBodyLimit::max(helpers::body_limit_usize(
150 MAX_BODY_BYTES,
151 )))
152 .layer(from_fn_with_state(state.clone(), auth_middleware)),
153 )
154 .with_state(state.clone());
155
156 let bounded_root = Router::new()
157 .route("/llms.txt", get(llms_txt))
158 .route(
159 "/.well-known/oauth-protected-resource",
160 get(protected_resource_metadata),
161 )
162 .route(BROWSE_PREFIX, get(browse_root))
163 .route(&format!("{BROWSE_PREFIX}/"), get(browse_root))
164 .route(&format!("{BROWSE_PREFIX}/{{*path}}"), get(browse_path))
165 .route_layer(from_fn_with_state(bounds.clone(), bound));
166
167 Router::new()
171 .route("/healthz", get(healthz))
172 .route("/readyz", get(readyz))
173 .merge(bounded_root)
174 .nest(API_V1_PREFIX, api_routes)
175 .layer(
176 ServiceBuilder::new()
177 .layer(SetRequestIdLayer::new(
178 request_id_header.clone(),
179 MakeRequestUuidV7,
180 ))
181 .layer(PropagateRequestIdLayer::new(request_id_header))
182 .layer(TraceLayer::new_for_http()),
183 )
184 .with_state(state)
189}
190
191#[cfg(test)]
192mod bounded_routes {
193 use super::{RequestBounds, build_bounded_router};
194 use crate::state::AppState;
195 use crate::testing::InMemoryStorage;
196 use axum::body::{Body, to_bytes};
197 use axum::http::{Method, Request};
198 use notedthat_core::{Authenticator, KbSlug};
199 use std::collections::BTreeMap;
200 use std::sync::Arc;
201 use std::time::Duration;
202 use tokio::sync::Semaphore;
203 use tower::ServiceExt;
204
205 fn state() -> AppState {
206 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
207 let notes = KbSlug::try_new("notes").expect("valid slug");
208 AppState {
209 storage: Arc::new(InMemoryStorage::default()),
210 declared_kbs: Arc::new(BTreeMap::from([("notes".to_string(), notes)])),
211 access_policies: Arc::new(BTreeMap::new()),
212 kb_details: Arc::new(BTreeMap::new()),
213 authenticator: Arc::new(Authenticator::new("token")),
214 max_body_size: 1024,
215 max_patchable_size: 1024,
216 indexer_tx: (&indexer_tx).into(),
217 searcher: Arc::new(crate::testing::NoopSearcher),
218 events: None,
219 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
220 readiness: crate::testing::ready_receiver(),
221 reconcile: None,
222 }
223 }
224
225 async fn refused_by_the_cap(method: Method, path: &str) -> bool {
227 let bounds = RequestBounds::new(
230 Duration::from_secs(60),
231 Duration::from_secs(60),
232 Arc::new(Semaphore::new(0)),
233 );
234 let response = build_bounded_router(state(), &bounds)
235 .oneshot(
236 Request::builder()
237 .method(method)
238 .uri(path)
239 .header("authorization", "Bearer token")
240 .body(Body::empty())
241 .expect("request"),
242 )
243 .await
244 .expect("infallible");
245 let body = to_bytes(response.into_body(), usize::MAX)
246 .await
247 .expect("body");
248 String::from_utf8_lossy(&body).contains("limit of requests in flight")
249 }
250
251 #[tokio::test]
252 async fn every_route_that_answers_once_is_bounded() {
253 for (method, path) in [
254 (Method::GET, "/api/v1/knowledgebases"),
255 (Method::GET, "/api/v1/knowledgebases/notes"),
256 (Method::POST, "/api/v1/knowledgebases/notes/search"),
257 (Method::GET, "/api/v1/knowledgebases/notes/index"),
258 (Method::POST, "/api/v1/knowledgebases/notes/index/reconcile"),
259 (Method::GET, "/api/v1/knowledgebases/notes/a.md"),
260 (Method::PUT, "/api/v1/knowledgebases/notes/a.md"),
261 (Method::GET, "/llms.txt"),
262 (Method::GET, "/.well-known/oauth-protected-resource"),
263 (Method::GET, "/browse"),
264 (Method::GET, "/browse/notes/"),
265 ] {
266 assert!(
267 refused_by_the_cap(method.clone(), path).await,
268 "{method} {path} is not bounded"
269 );
270 }
271 }
272
273 #[test]
282 fn only_the_probes_are_outside() {
283 let source = include_str!("mod.rs");
284 let outer = source
287 .split("API body limit stay nested")
288 .nth(1)
289 .and_then(|after| after.split(".merge(bounded_root)").next())
290 .expect("the outer router");
291 let registered: Vec<&str> = outer
292 .lines()
293 .filter(|line| line.trim_start().starts_with(".route(\""))
294 .collect();
295 assert_eq!(
296 registered,
297 [
298 " .route(\"/healthz\", get(healthz))",
299 " .route(\"/readyz\", get(readyz))"
300 ],
301 "a route registered outside the bounded sub-routers draws no permit; \
302 add it to `bounded_api` or `bounded_root` instead, or say here why not"
303 );
304 }
305
306 #[tokio::test]
309 async fn the_stream_and_the_probes_are_outside_the_bound() {
310 for path in ["/api/v1/knowledgebases/notes/events", "/healthz", "/readyz"] {
311 assert!(
312 !refused_by_the_cap(Method::GET, path).await,
313 "GET {path} is bounded"
314 );
315 }
316 }
317}
318
319#[cfg(test)]
320mod route_constants {
321 use super::{
322 API_V1_PREFIX, MATCHED_KB, MATCHED_KB_EVENTS, MATCHED_KB_INDEX, MATCHED_KB_INDEX_RECONCILE,
323 MATCHED_KB_OBJECT, MATCHED_KB_SEARCH, MATCHED_KBS, ROUTE_KB, ROUTE_KB_EVENTS,
324 ROUTE_KB_INDEX, ROUTE_KB_INDEX_RECONCILE, ROUTE_KB_OBJECT, ROUTE_KB_SEARCH, ROUTE_KBS,
325 };
326
327 #[test]
331 fn matched_paths_are_the_nested_routes_under_the_mount_point() {
332 assert_eq!(API_V1_PREFIX, "/api/v1");
333 for (route, matched) in [
334 (ROUTE_KBS, MATCHED_KBS),
335 (ROUTE_KB, MATCHED_KB),
336 (ROUTE_KB_SEARCH, MATCHED_KB_SEARCH),
337 (ROUTE_KB_EVENTS, MATCHED_KB_EVENTS),
338 (ROUTE_KB_INDEX, MATCHED_KB_INDEX),
339 (ROUTE_KB_INDEX_RECONCILE, MATCHED_KB_INDEX_RECONCILE),
340 (ROUTE_KB_OBJECT, MATCHED_KB_OBJECT),
341 ] {
342 assert_eq!(matched, format!("{API_V1_PREFIX}{route}"));
343 }
344 }
345}
346
347#[cfg(test)]
348mod patch_route {
349 use super::*;
350 use async_trait::async_trait;
351 use axum::body::{Body, to_bytes};
352 use axum::http::StatusCode;
353 use axum::response::Response;
354 use bytes::Bytes;
355 use notedthat_core::{
356 ByteRange, ConditionalHeaders, KbManifest, KbSlug, ListResponse, ObjectMeta, ObjectPath,
357 ObjectRead, PutOutcome, Storage, StorageError,
358 };
359 use notedthat_indexer::IndexEvent;
360 use std::collections::BTreeMap;
361 use std::sync::Arc;
362 use tower::util::ServiceExt;
363
364 const KB: &str = "notes";
365 const OBJECT_PATH: &str = "patch.md";
366 const TOKEN: &str = "test-token-abc";
367
368 async fn router_with_object(
369 initial_body: &'static [u8],
370 max_patchable_size: u64,
371 ) -> (axum::Router, String) {
372 let kb = KbSlug::try_new(KB).unwrap();
373 let object_path = ObjectPath::try_from_str(OBJECT_PATH).unwrap();
374 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
375 let outcome = storage
376 .put_object(
377 &kb,
378 &object_path,
379 Bytes::from_static(initial_body),
380 Some("text/markdown"),
381 ConditionalHeaders::default(),
382 )
383 .await
384 .unwrap();
385
386 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(16);
387 let mut kbs = BTreeMap::new();
388 kbs.insert(KB.to_string(), kb);
389 let router = build_router(AppState {
390 storage,
391 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
392 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
393 declared_kbs: Arc::new(kbs),
394 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
395 max_body_size: MAX_BODY_BYTES,
396 max_patchable_size,
397 indexer_tx: (&indexer_tx).into(),
398 searcher: Arc::new(crate::testing::NoopSearcher),
399 events: None,
400 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
401 readiness: crate::testing::ready_receiver(),
402 reconcile: None,
403 });
404
405 (router, outcome.etag.unwrap())
406 }
407
408 async fn object_with_etag(
409 storage: &crate::testing::InMemoryStorage,
410 kb: &KbSlug,
411 body: &'static [u8],
412 ) -> String {
413 storage
414 .put_object(
415 kb,
416 &ObjectPath::try_from_str(OBJECT_PATH).unwrap(),
417 Bytes::from_static(body),
418 Some("text/markdown"),
419 ConditionalHeaders::default(),
420 )
421 .await
422 .unwrap()
423 .etag
424 .unwrap()
425 }
426
427 #[allow(clippy::needless_pass_by_value)]
428 fn router_with_storage(
429 storage: Arc<dyn Storage>,
430 kb: KbSlug,
431 max_patchable_size: u64,
432 indexer_tx: tokio::sync::mpsc::Sender<IndexEvent>,
433 ) -> axum::Router {
434 let mut kbs = BTreeMap::new();
435 kbs.insert(KB.to_string(), kb);
436 build_router(AppState {
437 storage,
438 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
439 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
440 declared_kbs: Arc::new(kbs),
441 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
442 max_body_size: MAX_BODY_BYTES,
443 max_patchable_size,
444 indexer_tx: (&indexer_tx).into(),
445 searcher: Arc::new(crate::testing::NoopSearcher),
446 events: None,
447 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
448 readiness: crate::testing::ready_receiver(),
449 reconcile: None,
450 })
451 }
452
453 async fn patch_request(
454 router: axum::Router,
455 header_name: &'static str,
456 header_value: &str,
457 if_match: Option<&str>,
458 body: Bytes,
459 ) -> Response {
460 let mut builder = Request::builder()
461 .method("PATCH")
462 .uri(format!("/api/v1/knowledgebases/{KB}/{OBJECT_PATH}"))
463 .header("authorization", format!("Bearer {TOKEN}"))
464 .header(header_name, header_value);
465 if let Some(etag) = if_match {
466 builder = builder.header(axum::http::header::IF_MATCH, etag);
467 }
468
469 router
470 .oneshot(builder.body(Body::from(body)).unwrap())
471 .await
472 .unwrap()
473 }
474
475 async fn assert_error_code(response: Response, expected_status: StatusCode, expected: &str) {
476 assert_eq!(response.status(), expected_status);
477 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
478 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
479 assert_eq!(json["error"], expected);
480 }
481
482 #[tokio::test]
483 async fn bytes_content_range_returns_ok_with_etag_and_location() {
484 let (router, etag) = router_with_object(b"0123456789abcdefghij", MAX_BODY_BYTES).await;
485
486 let response = patch_request(
487 router,
488 "content-range",
489 "bytes 0-9/*",
490 Some(&etag),
491 Bytes::from_static(b"ABCDEFGHIJ"),
492 )
493 .await;
494
495 assert_eq!(response.status(), StatusCode::OK);
496 assert!(response.headers().get(axum::http::header::ETAG).is_some());
497 assert_eq!(
498 response
499 .headers()
500 .get(axum::http::header::LOCATION)
501 .unwrap(),
502 &format!("/api/v1/knowledgebases/{KB}/{OBJECT_PATH}")
503 );
504 assert!(
505 response
506 .headers()
507 .get(axum::http::header::CONTENT_RANGE)
508 .is_none()
509 );
510 assert!(response.headers().get("nt-patch-mode").is_none());
511 }
512
513 #[tokio::test]
514 async fn lines_content_range_returns_ok() {
515 let (router, etag) = router_with_object(b"one\ntwo\nthree\nfour\n", MAX_BODY_BYTES).await;
516
517 let response = patch_request(
518 router,
519 "content-range",
520 "lines 2-3/*",
521 Some(&etag),
522 Bytes::from_static(b"TWO\nTHREE\n"),
523 )
524 .await;
525
526 assert_eq!(response.status(), StatusCode::OK);
527 }
528
529 #[tokio::test]
530 async fn append_mode_without_if_match_returns_ok() {
531 let (router, _etag) = router_with_object(b"one\n", MAX_BODY_BYTES).await;
532
533 let response = patch_request(
534 router,
535 "nt-patch-mode",
536 "append",
537 None,
538 Bytes::from_static(b"two\n"),
539 )
540 .await;
541
542 assert_eq!(response.status(), StatusCode::OK);
543 }
544
545 #[tokio::test]
546 async fn append_mode_with_if_match_returns_ok() {
547 let (router, etag) = router_with_object(b"one\n", MAX_BODY_BYTES).await;
548
549 let response = patch_request(
550 router,
551 "nt-patch-mode",
552 "append",
553 Some(&etag),
554 Bytes::from_static(b"two\n"),
555 )
556 .await;
557
558 assert_eq!(response.status(), StatusCode::OK);
559 }
560
561 #[tokio::test]
562 async fn bytes_content_range_without_if_match_returns_invalid_request() {
563 let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
564
565 let response = patch_request(
566 router,
567 "content-range",
568 "bytes 0-1/*",
569 None,
570 Bytes::from_static(b"AB"),
571 )
572 .await;
573
574 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
575 }
576
577 #[tokio::test]
578 async fn if_match_star_returns_invalid_request() {
579 let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
580
581 let response = patch_request(
582 router,
583 "content-range",
584 "bytes 0-1/*",
585 Some("*"),
586 Bytes::from_static(b"AB"),
587 )
588 .await;
589
590 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
591 }
592
593 #[tokio::test]
594 async fn multi_value_if_match_returns_invalid_request() {
595 let (router, etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
596
597 let response = patch_request(
598 router,
599 "content-range",
600 "bytes 0-1/*",
601 Some(&format!("{etag}, \"other\"")),
602 Bytes::from_static(b"AB"),
603 )
604 .await;
605
606 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
607 }
608
609 #[tokio::test]
610 async fn nonexistent_object_returns_not_found() {
611 let (router, etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
612
613 let response = router
614 .oneshot(
615 Request::builder()
616 .method("PATCH")
617 .uri(format!("/api/v1/knowledgebases/{KB}/missing.md"))
618 .header("authorization", format!("Bearer {TOKEN}"))
619 .header("content-range", "bytes 0-1/*")
620 .header(axum::http::header::IF_MATCH, etag)
621 .body(Body::from(Bytes::from_static(b"AB")))
622 .unwrap(),
623 )
624 .await
625 .unwrap();
626
627 assert_error_code(response, StatusCode::NOT_FOUND, "not_found").await;
628 }
629
630 #[tokio::test]
631 async fn body_larger_than_max_patchable_size_returns_payload_too_large() {
632 let (router, _etag) = router_with_object(b"one\n", 4).await;
633
634 let response = patch_request(
635 router,
636 "nt-patch-mode",
637 "append",
638 None,
639 Bytes::from_static(b"abcde"),
640 )
641 .await;
642
643 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
644 }
645
646 mod errors {
647 use super::*;
648
649 #[derive(Clone)]
650 struct PutPreconditionFailedStorage {
651 inner: crate::testing::InMemoryStorage,
652 }
653
654 #[async_trait]
655 impl Storage for PutPreconditionFailedStorage {
656 async fn probe(&self, _kb: &KbSlug) -> Result<(), StorageError> {
657 Ok(())
658 }
659
660 async fn ensure_bucket(&self, kb: &KbSlug) -> Result<(), StorageError> {
661 self.inner.ensure_bucket(kb).await
662 }
663
664 async fn read_manifest(&self, kb: &KbSlug) -> Result<KbManifest, StorageError> {
665 self.inner.read_manifest(kb).await
666 }
667
668 async fn write_manifest(
669 &self,
670 kb: &KbSlug,
671 manifest: &KbManifest,
672 ) -> Result<(), StorageError> {
673 self.inner.write_manifest(kb, manifest).await
674 }
675
676 async fn head_object(
677 &self,
678 kb: &KbSlug,
679 path: &ObjectPath,
680 conditionals: ConditionalHeaders,
681 ) -> Result<ObjectMeta, StorageError> {
682 self.inner.head_object(kb, path, conditionals).await
683 }
684
685 async fn get_object(
686 &self,
687 kb: &KbSlug,
688 path: &ObjectPath,
689 range: Option<ByteRange>,
690 conditionals: ConditionalHeaders,
691 ) -> Result<ObjectRead, StorageError> {
692 self.inner.get_object(kb, path, range, conditionals).await
693 }
694
695 async fn get_object_stream(
696 &self,
697 kb: &KbSlug,
698 path: &ObjectPath,
699 range: Option<ByteRange>,
700 conditionals: ConditionalHeaders,
701 ) -> Result<notedthat_core::ObjectStream, StorageError> {
702 self.inner
703 .get_object_stream(kb, path, range, conditionals)
704 .await
705 }
706
707 async fn put_object(
708 &self,
709 _kb: &KbSlug,
710 _path: &ObjectPath,
711 _bytes: Bytes,
712 _content_type: Option<&str>,
713 _conditionals: ConditionalHeaders,
714 ) -> Result<PutOutcome, StorageError> {
715 Err(StorageError::PreconditionFailed)
716 }
717
718 async fn put_staged_object(
719 &self,
720 _kb: &KbSlug,
721 _path: &ObjectPath,
722 _body: notedthat_core::StagedBody,
723 _content_type: Option<&str>,
724 _conditionals: ConditionalHeaders,
725 ) -> Result<PutOutcome, StorageError> {
726 Err(StorageError::PreconditionFailed)
727 }
728
729 async fn copy_object(
730 &self,
731 kb: &KbSlug,
732 source: &ObjectPath,
733 destination: &ObjectPath,
734 options: notedthat_core::CopyObjectOptions,
735 ) -> Result<PutOutcome, StorageError> {
736 self.inner
737 .copy_object(kb, source, destination, options)
738 .await
739 }
740
741 async fn delete_object(
742 &self,
743 kb: &KbSlug,
744 path: &ObjectPath,
745 conditionals: ConditionalHeaders,
746 ) -> Result<(), StorageError> {
747 self.inner.delete_object(kb, path, conditionals).await
748 }
749
750 async fn list_objects(
751 &self,
752 kb: &KbSlug,
753 prefix: Option<&str>,
754 limit: u32,
755 cursor: Option<&str>,
756 ) -> Result<ListResponse, StorageError> {
757 self.inner.list_objects(kb, prefix, limit, cursor).await
758 }
759 }
760
761 #[tokio::test]
762 async fn missing_if_match_for_bytes_mode_returns_invalid_request() {
763 let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
764
765 let response = patch_request(
766 router,
767 "content-range",
768 "bytes 0-9/*",
769 None,
770 Bytes::from_static(b"ABCDEFGHIJ"),
771 )
772 .await;
773
774 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
775 }
776
777 #[tokio::test]
778 async fn body_larger_than_max_patchable_size_returns_payload_too_large() {
779 let (router, _etag) = router_with_object(b"one\n", 10).await;
780
781 let response = patch_request(
782 router,
783 "nt-patch-mode",
784 "append",
785 None,
786 Bytes::from_static(b"more than ten bytes"),
787 )
788 .await;
789
790 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
791 }
792
793 #[tokio::test]
794 async fn pre_splice_object_larger_than_max_patchable_size_returns_payload_too_large() {
795 let (router, _etag) = router_with_object(b"already too large", 10).await;
796
797 let response = patch_request(
798 router,
799 "nt-patch-mode",
800 "append",
801 None,
802 Bytes::from_static(b"!"),
803 )
804 .await;
805
806 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
807 }
808
809 #[tokio::test]
810 async fn post_splice_body_larger_than_max_patchable_size_returns_payload_too_large() {
811 let (router, _etag) = router_with_object(b"123456", 10).await;
812
813 let response = patch_request(
814 router,
815 "nt-patch-mode",
816 "append",
817 None,
818 Bytes::from_static(b"78901"),
819 )
820 .await;
821
822 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
823 }
824
825 #[tokio::test]
826 async fn line_range_out_of_bounds_returns_dual_416_headers_and_empty_body() {
827 let (router, etag) =
828 router_with_object(b"one\ntwo\nthree\nfour\nfive\n", MAX_BODY_BYTES).await;
829
830 let response = patch_request(
831 router,
832 "content-range",
833 "lines 100-200/*",
834 Some(&etag),
835 Bytes::from_static(b"replacement\n"),
836 )
837 .await;
838
839 assert_eq!(response.status(), StatusCode::RANGE_NOT_SATISFIABLE);
840 assert_eq!(
841 response.headers().get("content-range").unwrap(),
842 "lines */5"
843 );
844 assert_eq!(
845 response.headers().get("x-content-range-bytes").unwrap(),
846 "*/24"
847 );
848 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
849 assert!(body.is_empty());
850 }
851
852 #[tokio::test]
853 async fn if_match_mismatch_after_retries_returns_precondition_failed_without_content_range()
854 {
855 let kb = KbSlug::try_new(KB).unwrap();
856 let inner = crate::testing::InMemoryStorage::with_kbs([&kb]);
857 let etag = object_with_etag(&inner, &kb, b"0123456789").await;
858 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(16);
859 let router = router_with_storage(
860 Arc::new(PutPreconditionFailedStorage { inner }),
861 kb,
862 MAX_BODY_BYTES,
863 indexer_tx,
864 );
865
866 let response = patch_request(
867 router,
868 "content-range",
869 "bytes 0-1/*",
870 Some(&etag),
871 Bytes::from_static(b"AB"),
872 )
873 .await;
874
875 assert_eq!(response.status(), StatusCode::PRECONDITION_FAILED);
876 assert!(
877 response
878 .headers()
879 .get(axum::http::header::CONTENT_RANGE)
880 .is_none()
881 );
882 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
883 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
884 assert_eq!(json["error"], "precondition_failed");
885 }
886
887 #[tokio::test]
888 async fn indexer_queue_full_returns_backend_unavailable_with_retry_after() {
889 let kb = KbSlug::try_new(KB).unwrap();
890 let storage = crate::testing::InMemoryStorage::with_kbs([&kb]);
891 let etag = object_with_etag(&storage, &kb, b"one\n").await;
892 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
893 indexer_tx
894 .try_send(IndexEvent::Upsert {
895 kb: kb.clone(),
896 object_key: ObjectPath::try_from_str("queued.md").unwrap(),
897 etag: "queued".to_string(),
898 mtime: 0,
899 })
900 .unwrap();
901 let router = router_with_storage(Arc::new(storage), kb, MAX_BODY_BYTES, indexer_tx);
902
903 let response = patch_request(
904 router,
905 "nt-patch-mode",
906 "append",
907 Some(&etag),
908 Bytes::from_static(b"two\n"),
909 )
910 .await;
911
912 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
913 assert_eq!(response.headers().get("retry-after").unwrap(), "5");
914 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
915 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
916 assert_eq!(json["error"], "backend_unavailable");
917 }
918 }
919}
920
921#[cfg(test)]
922mod line_range_get {
923 use super::*;
924 use axum::body::{Body, to_bytes};
925 use axum::http::StatusCode;
926 use axum::response::Response;
927 use bytes::Bytes;
928 use notedthat_core::{ConditionalHeaders, KbSlug, ObjectPath, Storage};
929 use std::collections::BTreeMap;
930 use std::sync::Arc;
931 use tower::util::ServiceExt;
932
933 const KB: &str = "notes";
934 const TOKEN: &str = "test-token-abc";
935
936 fn twenty_line_markdown() -> String {
937 markdown_lines(1, 20)
938 }
939
940 fn markdown_lines(first: u32, last: u32) -> String {
941 let mut body = String::new();
942 for line in first..=last {
943 std::fmt::Write::write_fmt(&mut body, format_args!("line {line:02}\n")).unwrap();
944 }
945 body
946 }
947
948 async fn router_with_markdown_object(body: String) -> axum::Router {
949 let kb = KbSlug::try_new(KB).unwrap();
950 let object_path = ObjectPath::try_from_str("ranges.md").unwrap();
951 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
952 storage
953 .put_object(
954 &kb,
955 &object_path,
956 Bytes::from(body),
957 Some("text/markdown"),
958 ConditionalHeaders::default(),
959 )
960 .await
961 .unwrap();
962
963 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
964 let mut kbs = BTreeMap::new();
965 kbs.insert(KB.to_string(), kb);
966 build_router(AppState {
967 storage,
968 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
969 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
970 declared_kbs: Arc::new(kbs),
971 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
972 max_body_size: MAX_BODY_BYTES,
973 max_patchable_size: MAX_BODY_BYTES,
974 indexer_tx: (&indexer_tx).into(),
975 searcher: Arc::new(crate::testing::NoopSearcher),
976 events: None,
977 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
978 readiness: crate::testing::ready_receiver(),
979 reconcile: None,
980 })
981 }
982
983 async fn get_ranges_md(router: axum::Router, range: &str) -> Response {
984 router
985 .oneshot(
986 Request::builder()
987 .method("GET")
988 .uri(format!("/api/v1/knowledgebases/{KB}/ranges.md"))
989 .header("authorization", format!("Bearer {TOKEN}"))
990 .header(axum::http::header::RANGE, range)
991 .body(Body::empty())
992 .unwrap(),
993 )
994 .await
995 .unwrap()
996 }
997
998 #[tokio::test]
999 async fn returns_first_five_lines_when_closed_range_requested() {
1000 let router = router_with_markdown_object(twenty_line_markdown()).await;
1001
1002 let response = get_ranges_md(router, "lines=1-5").await;
1003
1004 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1005 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1006 assert_eq!(body, Bytes::from(markdown_lines(1, 5)));
1007 }
1008
1009 #[tokio::test]
1010 async fn returns_last_three_lines_when_suffix_range_requested() {
1011 let router = router_with_markdown_object(twenty_line_markdown()).await;
1012
1013 let response = get_ranges_md(router, "lines=-3").await;
1014
1015 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1016 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1017 assert_eq!(body, Bytes::from(markdown_lines(18, 20)));
1018 }
1019
1020 #[tokio::test]
1021 async fn returns_empty_body_when_insert_range_requested() {
1022 let router = router_with_markdown_object(twenty_line_markdown()).await;
1023
1024 let response = get_ranges_md(router, "lines=5-4").await;
1025
1026 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1027 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1028 assert!(body.is_empty());
1029 }
1030
1031 #[tokio::test]
1032 async fn returns_full_body_when_unknown_range_unit_requested() {
1033 let body = twenty_line_markdown();
1034 let router = router_with_markdown_object(body.clone()).await;
1035
1036 let response = get_ranges_md(router, "items=0-5").await;
1037
1038 assert_eq!(response.status(), StatusCode::OK);
1039 let actual = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1040 assert_eq!(actual, Bytes::from(body));
1041 }
1042
1043 mod headers {
1044 use super::*;
1045
1046 fn ten_line_markdown() -> String {
1047 markdown_lines(1, 10)
1048 }
1049
1050 #[tokio::test]
1051 async fn returns_line_and_byte_content_ranges_when_closed_range_requested() {
1052 let router = router_with_markdown_object(ten_line_markdown()).await;
1053
1054 let response = get_ranges_md(router, "lines=2-4").await;
1055
1056 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1057 assert_eq!(
1058 response.headers().get("Content-Range").unwrap(),
1059 "lines 2-4/10"
1060 );
1061 assert_eq!(
1062 response.headers().get("X-Content-Range-Bytes").unwrap(),
1063 "8-31/80"
1064 );
1065 }
1066
1067 #[tokio::test]
1068 async fn returns_slice_content_length_when_closed_range_requested() {
1069 let router = router_with_markdown_object(ten_line_markdown()).await;
1070
1071 let response = get_ranges_md(router, "lines=2-4").await;
1072
1073 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1074 assert_eq!(response.headers().get("content-length").unwrap(), "24");
1075 }
1076
1077 #[tokio::test]
1078 async fn returns_zero_length_and_empty_byte_range_when_insert_range_requested() {
1079 let router = router_with_markdown_object(ten_line_markdown()).await;
1080
1081 let response = get_ranges_md(router, "lines=5-4").await;
1082
1083 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1084 assert_eq!(response.headers().get("content-length").unwrap(), "0");
1085 assert_eq!(
1086 response.headers().get("Content-Range").unwrap(),
1087 "lines 5-4/10"
1088 );
1089 assert_eq!(
1090 response.headers().get("X-Content-Range-Bytes").unwrap(),
1091 "32-31/80"
1092 );
1093 }
1094
1095 #[tokio::test]
1096 async fn omits_line_byte_range_header_when_byte_range_requested() {
1097 let router = router_with_markdown_object(ten_line_markdown()).await;
1098
1099 let response = get_ranges_md(router, "bytes=0-9").await;
1100
1101 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1102 assert!(response.headers().get("X-Content-Range-Bytes").is_none());
1103 }
1104 }
1105}
1106
1107#[cfg(test)]
1108mod tests {
1109 use super::*;
1110 use axum::body::{Body, to_bytes};
1111 use axum::http::StatusCode;
1112 use axum::response::Response;
1113 use bytes::Bytes;
1114 use notedthat_core::{ConditionalHeaders, KbSlug, ObjectPath, Storage, StorageError};
1115 use notedthat_indexer::IndexEvent;
1116 use std::collections::BTreeMap;
1117 use std::sync::Arc;
1118 use tower::util::ServiceExt;
1119
1120 const KB: &str = "notes";
1121 const TOKEN: &str = "test-token-abc";
1122
1123 fn router() -> axum::Router {
1124 let kb = KbSlug::try_new(KB).unwrap();
1125 let mut kbs = BTreeMap::new();
1126 kbs.insert(KB.to_string(), kb);
1127 let (indexer_tx, mut rx) = tokio::sync::mpsc::channel(16);
1128 tokio::spawn(async move { while rx.recv().await.is_some() {} });
1129
1130 build_router(AppState {
1131 storage: Arc::new(crate::testing::InMemoryStorage::with_kbs(kbs.values())),
1132 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1133 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1134 declared_kbs: Arc::new(kbs),
1135 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1136 max_body_size: MAX_BODY_BYTES,
1137 max_patchable_size: MAX_BODY_BYTES,
1138 indexer_tx: (&indexer_tx).into(),
1139 searcher: Arc::new(crate::testing::NoopSearcher),
1140 events: None,
1141 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1142 readiness: crate::testing::ready_receiver(),
1143 reconcile: None,
1144 })
1145 }
1146
1147 async fn put_object(router: axum::Router, path: &str, body: &'static [u8]) -> Response {
1148 router
1149 .oneshot(
1150 Request::builder()
1151 .method("PUT")
1152 .uri(format!("/api/v1/knowledgebases/{KB}/{path}"))
1153 .header("authorization", format!("Bearer {TOKEN}"))
1154 .header(axum::http::header::CONTENT_TYPE, "text/markdown")
1155 .body(Body::from(Bytes::from_static(body)))
1156 .unwrap(),
1157 )
1158 .await
1159 .unwrap()
1160 }
1161
1162 async fn put_object_etag(router: axum::Router, path: &str, body: &'static [u8]) -> String {
1163 let response = put_object(router, path, body).await;
1164 assert_eq!(response.status(), StatusCode::CREATED);
1165 response
1166 .headers()
1167 .get(axum::http::header::ETAG)
1168 .unwrap()
1169 .to_str()
1170 .unwrap()
1171 .to_string()
1172 }
1173
1174 async fn get_object(router: axum::Router, path: &str) -> Response {
1175 router
1176 .oneshot(
1177 Request::builder()
1178 .method("GET")
1179 .uri(format!("/api/v1/knowledgebases/{KB}/{path}"))
1180 .header("authorization", format!("Bearer {TOKEN}"))
1181 .body(Body::empty())
1182 .unwrap(),
1183 )
1184 .await
1185 .unwrap()
1186 }
1187
1188 async fn post_replace(
1189 router: axum::Router,
1190 path: &str,
1191 if_match: &str,
1192 body: &'static [u8],
1193 ) -> Response {
1194 router
1195 .oneshot(
1196 Request::builder()
1197 .method("POST")
1198 .uri(format!("/api/v1/knowledgebases/{KB}/replace/{path}"))
1199 .header("authorization", format!("Bearer {TOKEN}"))
1200 .header(axum::http::header::CONTENT_TYPE, "application/json")
1201 .header(axum::http::header::IF_MATCH, if_match)
1202 .body(Body::from(Bytes::from_static(body)))
1203 .unwrap(),
1204 )
1205 .await
1206 .unwrap()
1207 }
1208
1209 #[tokio::test]
1210 async fn get_on_replace_prefixed_path_still_reads_object_via_catch_all() {
1211 let router = router();
1212 put_object_etag(router.clone(), "replace/foo.md", b"hi").await;
1213
1214 let response = get_object(router, "replace/foo.md").await;
1215
1216 assert_eq!(response.status(), StatusCode::OK);
1217 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1218 assert_eq!(&body[..], b"hi");
1219 }
1220
1221 #[tokio::test]
1222 async fn patch_on_replace_prefixed_path_still_reaches_patch_object() {
1223 let router = router();
1224 let etag = put_object_etag(router.clone(), "replace/bar.md", b"old\n").await;
1225
1226 let response = router
1227 .oneshot(
1228 Request::builder()
1229 .method("PATCH")
1230 .uri(format!("/api/v1/knowledgebases/{KB}/replace/bar.md"))
1231 .header("authorization", format!("Bearer {TOKEN}"))
1232 .header(axum::http::header::CONTENT_RANGE, "lines 1-1/*")
1233 .header(axum::http::header::IF_MATCH, etag)
1234 .body(Body::from(Bytes::from_static(b"new\n")))
1235 .unwrap(),
1236 )
1237 .await
1238 .unwrap();
1239
1240 assert_eq!(response.status(), StatusCode::OK);
1241 }
1242
1243 #[tokio::test]
1244 async fn put_and_delete_on_replace_prefixed_path_still_work() {
1245 let router = router();
1246 let put = put_object(router.clone(), "replace/delete.md", b"gone").await;
1247
1248 assert_eq!(put.status(), StatusCode::CREATED);
1249 let delete = router
1250 .oneshot(
1251 Request::builder()
1252 .method("DELETE")
1253 .uri(format!("/api/v1/knowledgebases/{KB}/replace/delete.md"))
1254 .header("authorization", format!("Bearer {TOKEN}"))
1255 .body(Body::empty())
1256 .unwrap(),
1257 )
1258 .await
1259 .unwrap();
1260 assert_eq!(delete.status(), StatusCode::NO_CONTENT);
1261 }
1262
1263 #[tokio::test]
1264 async fn post_on_non_replace_path_returns_404_not_found() {
1265 let response = router()
1266 .oneshot(
1267 Request::builder()
1268 .method("POST")
1269 .uri(format!("/api/v1/knowledgebases/{KB}/foo.md"))
1270 .header("authorization", format!("Bearer {TOKEN}"))
1271 .body(Body::empty())
1272 .unwrap(),
1273 )
1274 .await
1275 .unwrap();
1276
1277 assert_eq!(response.status(), StatusCode::NOT_FOUND);
1278 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1279 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1280 assert_eq!(json["error"], "not_found");
1281 assert!(
1282 json["message"]
1283 .as_str()
1284 .unwrap()
1285 .contains("supported actions: 'replace/<path>'")
1286 );
1287 }
1288
1289 #[tokio::test]
1290 async fn post_on_replace_prefixed_path_dispatches_to_replace_handler() {
1291 let router = router();
1292 let etag = put_object_etag(router.clone(), "target.md", b"hello world").await;
1293
1294 let response = post_replace(
1295 router,
1296 "target.md",
1297 &etag,
1298 br#"{"old_string":"world","new_string":"planet"}"#,
1299 )
1300 .await;
1301
1302 assert_eq!(response.status(), StatusCode::OK);
1303 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1304 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1305 assert_eq!(json["match_count"], 1);
1306 }
1307
1308 #[tokio::test]
1309 async fn post_on_replace_replace_path_targets_the_replace_prefixed_object() {
1310 let router = router();
1311 let etag = put_object_etag(router.clone(), "replace/nested.md", b"foo bar").await;
1312
1313 let response = post_replace(
1314 router.clone(),
1315 "replace/nested.md",
1316 &etag,
1317 br#"{"old_string":"bar","new_string":"baz"}"#,
1318 )
1319 .await;
1320
1321 assert_eq!(response.status(), StatusCode::OK);
1322 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1323 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1324 assert_eq!(json["match_count"], 1);
1325 let get = get_object(router, "replace/nested.md").await;
1326 let body = to_bytes(get.into_body(), 64 * 1024).await.unwrap();
1327 assert_eq!(&body[..], b"foo baz");
1328 }
1329
1330 #[tokio::test]
1331 async fn test_conditional_put_503_then_naive_retry_412_keeps_object_stored() {
1332 let kb = KbSlug::try_new(KB).unwrap();
1333 let object_path = ObjectPath::try_from_str("cond.md").unwrap();
1334 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
1335
1336 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
1337 indexer_tx
1338 .try_send(IndexEvent::Upsert {
1339 kb: kb.clone(),
1340 object_key: ObjectPath::try_from_str("queued.md").unwrap(),
1341 etag: "etag".to_string(),
1342 mtime: 0,
1343 })
1344 .unwrap();
1345
1346 let mut kbs = BTreeMap::new();
1347 kbs.insert(KB.to_string(), kb.clone());
1348 let state = AppState {
1349 storage: storage.clone(),
1350 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1351 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1352 declared_kbs: Arc::new(kbs),
1353 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1354 max_body_size: MAX_BODY_BYTES,
1355 max_patchable_size: MAX_BODY_BYTES,
1356 indexer_tx: (&indexer_tx).into(),
1357 searcher: Arc::new(crate::testing::NoopSearcher),
1358 events: None,
1359 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1360 readiness: crate::testing::ready_receiver(),
1361 reconcile: None,
1362 };
1363 let router = build_router(state);
1364
1365 let response = router
1366 .clone()
1367 .oneshot(
1368 Request::builder()
1369 .method("PUT")
1370 .uri(format!("/api/v1/knowledgebases/{KB}/cond.md"))
1371 .header("authorization", format!("Bearer {TOKEN}"))
1372 .header("if-none-match", "*")
1373 .body(Body::from(Bytes::from_static(b"first content")))
1374 .unwrap(),
1375 )
1376 .await
1377 .unwrap();
1378
1379 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
1380 assert_eq!(response.headers().get("retry-after").unwrap(), "5");
1381 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1382 let body = String::from_utf8(body.to_vec()).unwrap();
1383 assert!(body.contains("\"error\":\"backend_unavailable\""));
1384 assert!(body.contains("object stored; indexer queue full — retry to re-enqueue"));
1385
1386 let stored = storage
1387 .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1388 .await
1389 .unwrap();
1390 assert_eq!(stored.bytes, Bytes::from_static(b"first content"));
1391
1392 let retry = router
1393 .oneshot(
1394 Request::builder()
1395 .method("PUT")
1396 .uri(format!("/api/v1/knowledgebases/{KB}/cond.md"))
1397 .header("authorization", format!("Bearer {TOKEN}"))
1398 .header("if-none-match", "*")
1399 .body(Body::from(Bytes::from_static(b"second content")))
1400 .unwrap(),
1401 )
1402 .await
1403 .unwrap();
1404
1405 assert_eq!(retry.status(), StatusCode::PRECONDITION_FAILED);
1406 assert!(retry.headers().get("retry-after").is_none());
1407
1408 let stored = storage
1409 .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1410 .await
1411 .unwrap();
1412 assert_eq!(stored.bytes, Bytes::from_static(b"first content"));
1413 }
1414
1415 #[tokio::test]
1416 async fn test_delete_returns_delete_specific_503_body_when_indexer_backpressure() {
1417 let kb = KbSlug::try_new(KB).unwrap();
1418 let object_path = ObjectPath::try_from_str("to-delete.md").unwrap();
1419 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
1420 storage
1421 .put_object(
1422 &kb,
1423 &object_path,
1424 Bytes::from_static(b"content"),
1425 Some("text/plain"),
1426 ConditionalHeaders::default(),
1427 )
1428 .await
1429 .unwrap();
1430
1431 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
1432 indexer_tx
1433 .try_send(IndexEvent::Upsert {
1434 kb: kb.clone(),
1435 object_key: ObjectPath::try_from_str("queued.md").unwrap(),
1436 etag: "etag".to_string(),
1437 mtime: 0,
1438 })
1439 .unwrap();
1440
1441 let mut kbs = BTreeMap::new();
1442 kbs.insert(KB.to_string(), kb.clone());
1443 let state = AppState {
1444 storage: storage.clone(),
1445 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1446 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1447 declared_kbs: Arc::new(kbs),
1448 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1449 max_body_size: MAX_BODY_BYTES,
1450 max_patchable_size: MAX_BODY_BYTES,
1451 indexer_tx: (&indexer_tx).into(),
1452 searcher: Arc::new(crate::testing::NoopSearcher),
1453 events: None,
1454 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1455 readiness: crate::testing::ready_receiver(),
1456 reconcile: None,
1457 };
1458 let router = build_router(state);
1459
1460 let response = router
1461 .oneshot(
1462 Request::builder()
1463 .method("DELETE")
1464 .uri(format!("/api/v1/knowledgebases/{KB}/to-delete.md"))
1465 .header("authorization", format!("Bearer {TOKEN}"))
1466 .body(Body::empty())
1467 .unwrap(),
1468 )
1469 .await
1470 .unwrap();
1471
1472 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
1473 assert_eq!(response.headers().get("retry-after").unwrap(), "5");
1474 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1475 let body = String::from_utf8(body.to_vec()).unwrap();
1476 assert!(body.contains("\"error\":\"backend_unavailable\""));
1477 assert!(
1478 body.contains("\"message\":\"deleted from storage; retry to clear from search index\"")
1479 );
1480 assert!(!body.contains("object stored; indexer queue full — retry to re-enqueue"));
1481
1482 let deleted = storage
1483 .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1484 .await;
1485 assert!(matches!(deleted, Err(StorageError::NotFound { .. })));
1486 }
1487}