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,
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,
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 fn router_with_storage(
428 storage: Arc<dyn Storage>,
429 kb: KbSlug,
430 max_patchable_size: u64,
431 indexer_tx: tokio::sync::mpsc::Sender<IndexEvent>,
432 ) -> axum::Router {
433 let mut kbs = BTreeMap::new();
434 kbs.insert(KB.to_string(), kb);
435 build_router(AppState {
436 storage,
437 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
438 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
439 declared_kbs: Arc::new(kbs),
440 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
441 max_body_size: MAX_BODY_BYTES,
442 max_patchable_size,
443 indexer_tx,
444 searcher: Arc::new(crate::testing::NoopSearcher),
445 events: None,
446 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
447 readiness: crate::testing::ready_receiver(),
448 reconcile: None,
449 })
450 }
451
452 async fn patch_request(
453 router: axum::Router,
454 header_name: &'static str,
455 header_value: &str,
456 if_match: Option<&str>,
457 body: Bytes,
458 ) -> Response {
459 let mut builder = Request::builder()
460 .method("PATCH")
461 .uri(format!("/api/v1/knowledgebases/{KB}/{OBJECT_PATH}"))
462 .header("authorization", format!("Bearer {TOKEN}"))
463 .header(header_name, header_value);
464 if let Some(etag) = if_match {
465 builder = builder.header(axum::http::header::IF_MATCH, etag);
466 }
467
468 router
469 .oneshot(builder.body(Body::from(body)).unwrap())
470 .await
471 .unwrap()
472 }
473
474 async fn assert_error_code(response: Response, expected_status: StatusCode, expected: &str) {
475 assert_eq!(response.status(), expected_status);
476 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
477 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
478 assert_eq!(json["error"], expected);
479 }
480
481 #[tokio::test]
482 async fn bytes_content_range_returns_ok_with_etag_and_location() {
483 let (router, etag) = router_with_object(b"0123456789abcdefghij", MAX_BODY_BYTES).await;
484
485 let response = patch_request(
486 router,
487 "content-range",
488 "bytes 0-9/*",
489 Some(&etag),
490 Bytes::from_static(b"ABCDEFGHIJ"),
491 )
492 .await;
493
494 assert_eq!(response.status(), StatusCode::OK);
495 assert!(response.headers().get(axum::http::header::ETAG).is_some());
496 assert_eq!(
497 response
498 .headers()
499 .get(axum::http::header::LOCATION)
500 .unwrap(),
501 &format!("/api/v1/knowledgebases/{KB}/{OBJECT_PATH}")
502 );
503 assert!(
504 response
505 .headers()
506 .get(axum::http::header::CONTENT_RANGE)
507 .is_none()
508 );
509 assert!(response.headers().get("nt-patch-mode").is_none());
510 }
511
512 #[tokio::test]
513 async fn lines_content_range_returns_ok() {
514 let (router, etag) = router_with_object(b"one\ntwo\nthree\nfour\n", MAX_BODY_BYTES).await;
515
516 let response = patch_request(
517 router,
518 "content-range",
519 "lines 2-3/*",
520 Some(&etag),
521 Bytes::from_static(b"TWO\nTHREE\n"),
522 )
523 .await;
524
525 assert_eq!(response.status(), StatusCode::OK);
526 }
527
528 #[tokio::test]
529 async fn append_mode_without_if_match_returns_ok() {
530 let (router, _etag) = router_with_object(b"one\n", MAX_BODY_BYTES).await;
531
532 let response = patch_request(
533 router,
534 "nt-patch-mode",
535 "append",
536 None,
537 Bytes::from_static(b"two\n"),
538 )
539 .await;
540
541 assert_eq!(response.status(), StatusCode::OK);
542 }
543
544 #[tokio::test]
545 async fn append_mode_with_if_match_returns_ok() {
546 let (router, etag) = router_with_object(b"one\n", MAX_BODY_BYTES).await;
547
548 let response = patch_request(
549 router,
550 "nt-patch-mode",
551 "append",
552 Some(&etag),
553 Bytes::from_static(b"two\n"),
554 )
555 .await;
556
557 assert_eq!(response.status(), StatusCode::OK);
558 }
559
560 #[tokio::test]
561 async fn bytes_content_range_without_if_match_returns_invalid_request() {
562 let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
563
564 let response = patch_request(
565 router,
566 "content-range",
567 "bytes 0-1/*",
568 None,
569 Bytes::from_static(b"AB"),
570 )
571 .await;
572
573 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
574 }
575
576 #[tokio::test]
577 async fn if_match_star_returns_invalid_request() {
578 let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
579
580 let response = patch_request(
581 router,
582 "content-range",
583 "bytes 0-1/*",
584 Some("*"),
585 Bytes::from_static(b"AB"),
586 )
587 .await;
588
589 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
590 }
591
592 #[tokio::test]
593 async fn multi_value_if_match_returns_invalid_request() {
594 let (router, etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
595
596 let response = patch_request(
597 router,
598 "content-range",
599 "bytes 0-1/*",
600 Some(&format!("{etag}, \"other\"")),
601 Bytes::from_static(b"AB"),
602 )
603 .await;
604
605 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
606 }
607
608 #[tokio::test]
609 async fn nonexistent_object_returns_not_found() {
610 let (router, etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
611
612 let response = router
613 .oneshot(
614 Request::builder()
615 .method("PATCH")
616 .uri(format!("/api/v1/knowledgebases/{KB}/missing.md"))
617 .header("authorization", format!("Bearer {TOKEN}"))
618 .header("content-range", "bytes 0-1/*")
619 .header(axum::http::header::IF_MATCH, etag)
620 .body(Body::from(Bytes::from_static(b"AB")))
621 .unwrap(),
622 )
623 .await
624 .unwrap();
625
626 assert_error_code(response, StatusCode::NOT_FOUND, "not_found").await;
627 }
628
629 #[tokio::test]
630 async fn body_larger_than_max_patchable_size_returns_payload_too_large() {
631 let (router, _etag) = router_with_object(b"one\n", 4).await;
632
633 let response = patch_request(
634 router,
635 "nt-patch-mode",
636 "append",
637 None,
638 Bytes::from_static(b"abcde"),
639 )
640 .await;
641
642 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
643 }
644
645 mod errors {
646 use super::*;
647
648 #[derive(Clone)]
649 struct PutPreconditionFailedStorage {
650 inner: crate::testing::InMemoryStorage,
651 }
652
653 #[async_trait]
654 impl Storage for PutPreconditionFailedStorage {
655 async fn probe(&self, _kb: &KbSlug) -> Result<(), StorageError> {
656 Ok(())
657 }
658
659 async fn ensure_bucket(&self, kb: &KbSlug) -> Result<(), StorageError> {
660 self.inner.ensure_bucket(kb).await
661 }
662
663 async fn read_manifest(&self, kb: &KbSlug) -> Result<KbManifest, StorageError> {
664 self.inner.read_manifest(kb).await
665 }
666
667 async fn write_manifest(
668 &self,
669 kb: &KbSlug,
670 manifest: &KbManifest,
671 ) -> Result<(), StorageError> {
672 self.inner.write_manifest(kb, manifest).await
673 }
674
675 async fn head_object(
676 &self,
677 kb: &KbSlug,
678 path: &ObjectPath,
679 conditionals: ConditionalHeaders,
680 ) -> Result<ObjectMeta, StorageError> {
681 self.inner.head_object(kb, path, conditionals).await
682 }
683
684 async fn get_object(
685 &self,
686 kb: &KbSlug,
687 path: &ObjectPath,
688 range: Option<ByteRange>,
689 conditionals: ConditionalHeaders,
690 ) -> Result<ObjectRead, StorageError> {
691 self.inner.get_object(kb, path, range, conditionals).await
692 }
693
694 async fn get_object_stream(
695 &self,
696 kb: &KbSlug,
697 path: &ObjectPath,
698 range: Option<ByteRange>,
699 conditionals: ConditionalHeaders,
700 ) -> Result<notedthat_core::ObjectStream, StorageError> {
701 self.inner
702 .get_object_stream(kb, path, range, conditionals)
703 .await
704 }
705
706 async fn put_object(
707 &self,
708 _kb: &KbSlug,
709 _path: &ObjectPath,
710 _bytes: Bytes,
711 _content_type: Option<&str>,
712 _conditionals: ConditionalHeaders,
713 ) -> Result<PutOutcome, StorageError> {
714 Err(StorageError::PreconditionFailed)
715 }
716
717 async fn put_staged_object(
718 &self,
719 _kb: &KbSlug,
720 _path: &ObjectPath,
721 _body: notedthat_core::StagedBody,
722 _content_type: Option<&str>,
723 _conditionals: ConditionalHeaders,
724 ) -> Result<PutOutcome, StorageError> {
725 Err(StorageError::PreconditionFailed)
726 }
727
728 async fn copy_object(
729 &self,
730 kb: &KbSlug,
731 source: &ObjectPath,
732 destination: &ObjectPath,
733 options: notedthat_core::CopyObjectOptions,
734 ) -> Result<PutOutcome, StorageError> {
735 self.inner
736 .copy_object(kb, source, destination, options)
737 .await
738 }
739
740 async fn delete_object(
741 &self,
742 kb: &KbSlug,
743 path: &ObjectPath,
744 conditionals: ConditionalHeaders,
745 ) -> Result<(), StorageError> {
746 self.inner.delete_object(kb, path, conditionals).await
747 }
748
749 async fn list_objects(
750 &self,
751 kb: &KbSlug,
752 prefix: Option<&str>,
753 limit: u32,
754 cursor: Option<&str>,
755 ) -> Result<ListResponse, StorageError> {
756 self.inner.list_objects(kb, prefix, limit, cursor).await
757 }
758 }
759
760 #[tokio::test]
761 async fn missing_if_match_for_bytes_mode_returns_invalid_request() {
762 let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
763
764 let response = patch_request(
765 router,
766 "content-range",
767 "bytes 0-9/*",
768 None,
769 Bytes::from_static(b"ABCDEFGHIJ"),
770 )
771 .await;
772
773 assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
774 }
775
776 #[tokio::test]
777 async fn body_larger_than_max_patchable_size_returns_payload_too_large() {
778 let (router, _etag) = router_with_object(b"one\n", 10).await;
779
780 let response = patch_request(
781 router,
782 "nt-patch-mode",
783 "append",
784 None,
785 Bytes::from_static(b"more than ten bytes"),
786 )
787 .await;
788
789 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
790 }
791
792 #[tokio::test]
793 async fn pre_splice_object_larger_than_max_patchable_size_returns_payload_too_large() {
794 let (router, _etag) = router_with_object(b"already too large", 10).await;
795
796 let response = patch_request(
797 router,
798 "nt-patch-mode",
799 "append",
800 None,
801 Bytes::from_static(b"!"),
802 )
803 .await;
804
805 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
806 }
807
808 #[tokio::test]
809 async fn post_splice_body_larger_than_max_patchable_size_returns_payload_too_large() {
810 let (router, _etag) = router_with_object(b"123456", 10).await;
811
812 let response = patch_request(
813 router,
814 "nt-patch-mode",
815 "append",
816 None,
817 Bytes::from_static(b"78901"),
818 )
819 .await;
820
821 assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
822 }
823
824 #[tokio::test]
825 async fn line_range_out_of_bounds_returns_dual_416_headers_and_empty_body() {
826 let (router, etag) =
827 router_with_object(b"one\ntwo\nthree\nfour\nfive\n", MAX_BODY_BYTES).await;
828
829 let response = patch_request(
830 router,
831 "content-range",
832 "lines 100-200/*",
833 Some(&etag),
834 Bytes::from_static(b"replacement\n"),
835 )
836 .await;
837
838 assert_eq!(response.status(), StatusCode::RANGE_NOT_SATISFIABLE);
839 assert_eq!(
840 response.headers().get("content-range").unwrap(),
841 "lines */5"
842 );
843 assert_eq!(
844 response.headers().get("x-content-range-bytes").unwrap(),
845 "*/24"
846 );
847 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
848 assert!(body.is_empty());
849 }
850
851 #[tokio::test]
852 async fn if_match_mismatch_after_retries_returns_precondition_failed_without_content_range()
853 {
854 let kb = KbSlug::try_new(KB).unwrap();
855 let inner = crate::testing::InMemoryStorage::with_kbs([&kb]);
856 let etag = object_with_etag(&inner, &kb, b"0123456789").await;
857 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(16);
858 let router = router_with_storage(
859 Arc::new(PutPreconditionFailedStorage { inner }),
860 kb,
861 MAX_BODY_BYTES,
862 indexer_tx,
863 );
864
865 let response = patch_request(
866 router,
867 "content-range",
868 "bytes 0-1/*",
869 Some(&etag),
870 Bytes::from_static(b"AB"),
871 )
872 .await;
873
874 assert_eq!(response.status(), StatusCode::PRECONDITION_FAILED);
875 assert!(
876 response
877 .headers()
878 .get(axum::http::header::CONTENT_RANGE)
879 .is_none()
880 );
881 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
882 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
883 assert_eq!(json["error"], "precondition_failed");
884 }
885
886 #[tokio::test]
887 async fn indexer_queue_full_returns_backend_unavailable_with_retry_after() {
888 let kb = KbSlug::try_new(KB).unwrap();
889 let storage = crate::testing::InMemoryStorage::with_kbs([&kb]);
890 let etag = object_with_etag(&storage, &kb, b"one\n").await;
891 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
892 indexer_tx
893 .try_send(IndexEvent::Upsert {
894 kb: kb.clone(),
895 object_key: ObjectPath::try_from_str("queued.md").unwrap(),
896 etag: "queued".to_string(),
897 mtime: 0,
898 })
899 .unwrap();
900 let router = router_with_storage(Arc::new(storage), kb, MAX_BODY_BYTES, indexer_tx);
901
902 let response = patch_request(
903 router,
904 "nt-patch-mode",
905 "append",
906 Some(&etag),
907 Bytes::from_static(b"two\n"),
908 )
909 .await;
910
911 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
912 assert_eq!(response.headers().get("retry-after").unwrap(), "5");
913 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
914 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
915 assert_eq!(json["error"], "backend_unavailable");
916 }
917 }
918}
919
920#[cfg(test)]
921mod line_range_get {
922 use super::*;
923 use axum::body::{Body, to_bytes};
924 use axum::http::StatusCode;
925 use axum::response::Response;
926 use bytes::Bytes;
927 use notedthat_core::{ConditionalHeaders, KbSlug, ObjectPath, Storage};
928 use std::collections::BTreeMap;
929 use std::sync::Arc;
930 use tower::util::ServiceExt;
931
932 const KB: &str = "notes";
933 const TOKEN: &str = "test-token-abc";
934
935 fn twenty_line_markdown() -> String {
936 markdown_lines(1, 20)
937 }
938
939 fn markdown_lines(first: u32, last: u32) -> String {
940 let mut body = String::new();
941 for line in first..=last {
942 std::fmt::Write::write_fmt(&mut body, format_args!("line {line:02}\n")).unwrap();
943 }
944 body
945 }
946
947 async fn router_with_markdown_object(body: String) -> axum::Router {
948 let kb = KbSlug::try_new(KB).unwrap();
949 let object_path = ObjectPath::try_from_str("ranges.md").unwrap();
950 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
951 storage
952 .put_object(
953 &kb,
954 &object_path,
955 Bytes::from(body),
956 Some("text/markdown"),
957 ConditionalHeaders::default(),
958 )
959 .await
960 .unwrap();
961
962 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
963 let mut kbs = BTreeMap::new();
964 kbs.insert(KB.to_string(), kb);
965 build_router(AppState {
966 storage,
967 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
968 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
969 declared_kbs: Arc::new(kbs),
970 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
971 max_body_size: MAX_BODY_BYTES,
972 max_patchable_size: MAX_BODY_BYTES,
973 indexer_tx,
974 searcher: Arc::new(crate::testing::NoopSearcher),
975 events: None,
976 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
977 readiness: crate::testing::ready_receiver(),
978 reconcile: None,
979 })
980 }
981
982 async fn get_ranges_md(router: axum::Router, range: &str) -> Response {
983 router
984 .oneshot(
985 Request::builder()
986 .method("GET")
987 .uri(format!("/api/v1/knowledgebases/{KB}/ranges.md"))
988 .header("authorization", format!("Bearer {TOKEN}"))
989 .header(axum::http::header::RANGE, range)
990 .body(Body::empty())
991 .unwrap(),
992 )
993 .await
994 .unwrap()
995 }
996
997 #[tokio::test]
998 async fn returns_first_five_lines_when_closed_range_requested() {
999 let router = router_with_markdown_object(twenty_line_markdown()).await;
1000
1001 let response = get_ranges_md(router, "lines=1-5").await;
1002
1003 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1004 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1005 assert_eq!(body, Bytes::from(markdown_lines(1, 5)));
1006 }
1007
1008 #[tokio::test]
1009 async fn returns_last_three_lines_when_suffix_range_requested() {
1010 let router = router_with_markdown_object(twenty_line_markdown()).await;
1011
1012 let response = get_ranges_md(router, "lines=-3").await;
1013
1014 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1015 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1016 assert_eq!(body, Bytes::from(markdown_lines(18, 20)));
1017 }
1018
1019 #[tokio::test]
1020 async fn returns_empty_body_when_insert_range_requested() {
1021 let router = router_with_markdown_object(twenty_line_markdown()).await;
1022
1023 let response = get_ranges_md(router, "lines=5-4").await;
1024
1025 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1026 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1027 assert!(body.is_empty());
1028 }
1029
1030 #[tokio::test]
1031 async fn returns_full_body_when_unknown_range_unit_requested() {
1032 let body = twenty_line_markdown();
1033 let router = router_with_markdown_object(body.clone()).await;
1034
1035 let response = get_ranges_md(router, "items=0-5").await;
1036
1037 assert_eq!(response.status(), StatusCode::OK);
1038 let actual = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1039 assert_eq!(actual, Bytes::from(body));
1040 }
1041
1042 mod headers {
1043 use super::*;
1044
1045 fn ten_line_markdown() -> String {
1046 markdown_lines(1, 10)
1047 }
1048
1049 #[tokio::test]
1050 async fn returns_line_and_byte_content_ranges_when_closed_range_requested() {
1051 let router = router_with_markdown_object(ten_line_markdown()).await;
1052
1053 let response = get_ranges_md(router, "lines=2-4").await;
1054
1055 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1056 assert_eq!(
1057 response.headers().get("Content-Range").unwrap(),
1058 "lines 2-4/10"
1059 );
1060 assert_eq!(
1061 response.headers().get("X-Content-Range-Bytes").unwrap(),
1062 "8-31/80"
1063 );
1064 }
1065
1066 #[tokio::test]
1067 async fn returns_slice_content_length_when_closed_range_requested() {
1068 let router = router_with_markdown_object(ten_line_markdown()).await;
1069
1070 let response = get_ranges_md(router, "lines=2-4").await;
1071
1072 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1073 assert_eq!(response.headers().get("content-length").unwrap(), "24");
1074 }
1075
1076 #[tokio::test]
1077 async fn returns_zero_length_and_empty_byte_range_when_insert_range_requested() {
1078 let router = router_with_markdown_object(ten_line_markdown()).await;
1079
1080 let response = get_ranges_md(router, "lines=5-4").await;
1081
1082 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1083 assert_eq!(response.headers().get("content-length").unwrap(), "0");
1084 assert_eq!(
1085 response.headers().get("Content-Range").unwrap(),
1086 "lines 5-4/10"
1087 );
1088 assert_eq!(
1089 response.headers().get("X-Content-Range-Bytes").unwrap(),
1090 "32-31/80"
1091 );
1092 }
1093
1094 #[tokio::test]
1095 async fn omits_line_byte_range_header_when_byte_range_requested() {
1096 let router = router_with_markdown_object(ten_line_markdown()).await;
1097
1098 let response = get_ranges_md(router, "bytes=0-9").await;
1099
1100 assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1101 assert!(response.headers().get("X-Content-Range-Bytes").is_none());
1102 }
1103 }
1104}
1105
1106#[cfg(test)]
1107mod tests {
1108 use super::*;
1109 use axum::body::{Body, to_bytes};
1110 use axum::http::StatusCode;
1111 use axum::response::Response;
1112 use bytes::Bytes;
1113 use notedthat_core::{ConditionalHeaders, KbSlug, ObjectPath, Storage, StorageError};
1114 use notedthat_indexer::IndexEvent;
1115 use std::collections::BTreeMap;
1116 use std::sync::Arc;
1117 use tower::util::ServiceExt;
1118
1119 const KB: &str = "notes";
1120 const TOKEN: &str = "test-token-abc";
1121
1122 fn router() -> axum::Router {
1123 let kb = KbSlug::try_new(KB).unwrap();
1124 let mut kbs = BTreeMap::new();
1125 kbs.insert(KB.to_string(), kb);
1126 let (indexer_tx, mut rx) = tokio::sync::mpsc::channel(16);
1127 tokio::spawn(async move { while rx.recv().await.is_some() {} });
1128
1129 build_router(AppState {
1130 storage: Arc::new(crate::testing::InMemoryStorage::with_kbs(kbs.values())),
1131 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1132 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1133 declared_kbs: Arc::new(kbs),
1134 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1135 max_body_size: MAX_BODY_BYTES,
1136 max_patchable_size: MAX_BODY_BYTES,
1137 indexer_tx,
1138 searcher: Arc::new(crate::testing::NoopSearcher),
1139 events: None,
1140 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1141 readiness: crate::testing::ready_receiver(),
1142 reconcile: None,
1143 })
1144 }
1145
1146 async fn put_object(router: axum::Router, path: &str, body: &'static [u8]) -> Response {
1147 router
1148 .oneshot(
1149 Request::builder()
1150 .method("PUT")
1151 .uri(format!("/api/v1/knowledgebases/{KB}/{path}"))
1152 .header("authorization", format!("Bearer {TOKEN}"))
1153 .header(axum::http::header::CONTENT_TYPE, "text/markdown")
1154 .body(Body::from(Bytes::from_static(body)))
1155 .unwrap(),
1156 )
1157 .await
1158 .unwrap()
1159 }
1160
1161 async fn put_object_etag(router: axum::Router, path: &str, body: &'static [u8]) -> String {
1162 let response = put_object(router, path, body).await;
1163 assert_eq!(response.status(), StatusCode::CREATED);
1164 response
1165 .headers()
1166 .get(axum::http::header::ETAG)
1167 .unwrap()
1168 .to_str()
1169 .unwrap()
1170 .to_string()
1171 }
1172
1173 async fn get_object(router: axum::Router, path: &str) -> Response {
1174 router
1175 .oneshot(
1176 Request::builder()
1177 .method("GET")
1178 .uri(format!("/api/v1/knowledgebases/{KB}/{path}"))
1179 .header("authorization", format!("Bearer {TOKEN}"))
1180 .body(Body::empty())
1181 .unwrap(),
1182 )
1183 .await
1184 .unwrap()
1185 }
1186
1187 async fn post_replace(
1188 router: axum::Router,
1189 path: &str,
1190 if_match: &str,
1191 body: &'static [u8],
1192 ) -> Response {
1193 router
1194 .oneshot(
1195 Request::builder()
1196 .method("POST")
1197 .uri(format!("/api/v1/knowledgebases/{KB}/replace/{path}"))
1198 .header("authorization", format!("Bearer {TOKEN}"))
1199 .header(axum::http::header::CONTENT_TYPE, "application/json")
1200 .header(axum::http::header::IF_MATCH, if_match)
1201 .body(Body::from(Bytes::from_static(body)))
1202 .unwrap(),
1203 )
1204 .await
1205 .unwrap()
1206 }
1207
1208 #[tokio::test]
1209 async fn get_on_replace_prefixed_path_still_reads_object_via_catch_all() {
1210 let router = router();
1211 put_object_etag(router.clone(), "replace/foo.md", b"hi").await;
1212
1213 let response = get_object(router, "replace/foo.md").await;
1214
1215 assert_eq!(response.status(), StatusCode::OK);
1216 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1217 assert_eq!(&body[..], b"hi");
1218 }
1219
1220 #[tokio::test]
1221 async fn patch_on_replace_prefixed_path_still_reaches_patch_object() {
1222 let router = router();
1223 let etag = put_object_etag(router.clone(), "replace/bar.md", b"old\n").await;
1224
1225 let response = router
1226 .oneshot(
1227 Request::builder()
1228 .method("PATCH")
1229 .uri(format!("/api/v1/knowledgebases/{KB}/replace/bar.md"))
1230 .header("authorization", format!("Bearer {TOKEN}"))
1231 .header(axum::http::header::CONTENT_RANGE, "lines 1-1/*")
1232 .header(axum::http::header::IF_MATCH, etag)
1233 .body(Body::from(Bytes::from_static(b"new\n")))
1234 .unwrap(),
1235 )
1236 .await
1237 .unwrap();
1238
1239 assert_eq!(response.status(), StatusCode::OK);
1240 }
1241
1242 #[tokio::test]
1243 async fn put_and_delete_on_replace_prefixed_path_still_work() {
1244 let router = router();
1245 let put = put_object(router.clone(), "replace/delete.md", b"gone").await;
1246
1247 assert_eq!(put.status(), StatusCode::CREATED);
1248 let delete = router
1249 .oneshot(
1250 Request::builder()
1251 .method("DELETE")
1252 .uri(format!("/api/v1/knowledgebases/{KB}/replace/delete.md"))
1253 .header("authorization", format!("Bearer {TOKEN}"))
1254 .body(Body::empty())
1255 .unwrap(),
1256 )
1257 .await
1258 .unwrap();
1259 assert_eq!(delete.status(), StatusCode::NO_CONTENT);
1260 }
1261
1262 #[tokio::test]
1263 async fn post_on_non_replace_path_returns_404_not_found() {
1264 let response = router()
1265 .oneshot(
1266 Request::builder()
1267 .method("POST")
1268 .uri(format!("/api/v1/knowledgebases/{KB}/foo.md"))
1269 .header("authorization", format!("Bearer {TOKEN}"))
1270 .body(Body::empty())
1271 .unwrap(),
1272 )
1273 .await
1274 .unwrap();
1275
1276 assert_eq!(response.status(), StatusCode::NOT_FOUND);
1277 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1278 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1279 assert_eq!(json["error"], "not_found");
1280 assert!(
1281 json["message"]
1282 .as_str()
1283 .unwrap()
1284 .contains("supported actions: 'replace/<path>'")
1285 );
1286 }
1287
1288 #[tokio::test]
1289 async fn post_on_replace_prefixed_path_dispatches_to_replace_handler() {
1290 let router = router();
1291 let etag = put_object_etag(router.clone(), "target.md", b"hello world").await;
1292
1293 let response = post_replace(
1294 router,
1295 "target.md",
1296 &etag,
1297 br#"{"old_string":"world","new_string":"planet"}"#,
1298 )
1299 .await;
1300
1301 assert_eq!(response.status(), StatusCode::OK);
1302 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1303 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1304 assert_eq!(json["match_count"], 1);
1305 }
1306
1307 #[tokio::test]
1308 async fn post_on_replace_replace_path_targets_the_replace_prefixed_object() {
1309 let router = router();
1310 let etag = put_object_etag(router.clone(), "replace/nested.md", b"foo bar").await;
1311
1312 let response = post_replace(
1313 router.clone(),
1314 "replace/nested.md",
1315 &etag,
1316 br#"{"old_string":"bar","new_string":"baz"}"#,
1317 )
1318 .await;
1319
1320 assert_eq!(response.status(), StatusCode::OK);
1321 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1322 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1323 assert_eq!(json["match_count"], 1);
1324 let get = get_object(router, "replace/nested.md").await;
1325 let body = to_bytes(get.into_body(), 64 * 1024).await.unwrap();
1326 assert_eq!(&body[..], b"foo baz");
1327 }
1328
1329 #[tokio::test]
1330 async fn test_conditional_put_503_then_naive_retry_412_keeps_object_stored() {
1331 let kb = KbSlug::try_new(KB).unwrap();
1332 let object_path = ObjectPath::try_from_str("cond.md").unwrap();
1333 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
1334
1335 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
1336 indexer_tx
1337 .try_send(IndexEvent::Upsert {
1338 kb: kb.clone(),
1339 object_key: ObjectPath::try_from_str("queued.md").unwrap(),
1340 etag: "etag".to_string(),
1341 mtime: 0,
1342 })
1343 .unwrap();
1344
1345 let mut kbs = BTreeMap::new();
1346 kbs.insert(KB.to_string(), kb.clone());
1347 let state = AppState {
1348 storage: storage.clone(),
1349 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1350 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1351 declared_kbs: Arc::new(kbs),
1352 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1353 max_body_size: MAX_BODY_BYTES,
1354 max_patchable_size: MAX_BODY_BYTES,
1355 indexer_tx,
1356 searcher: Arc::new(crate::testing::NoopSearcher),
1357 events: None,
1358 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1359 readiness: crate::testing::ready_receiver(),
1360 reconcile: None,
1361 };
1362 let router = build_router(state);
1363
1364 let response = router
1365 .clone()
1366 .oneshot(
1367 Request::builder()
1368 .method("PUT")
1369 .uri(format!("/api/v1/knowledgebases/{KB}/cond.md"))
1370 .header("authorization", format!("Bearer {TOKEN}"))
1371 .header("if-none-match", "*")
1372 .body(Body::from(Bytes::from_static(b"first content")))
1373 .unwrap(),
1374 )
1375 .await
1376 .unwrap();
1377
1378 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
1379 assert_eq!(response.headers().get("retry-after").unwrap(), "5");
1380 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1381 let body = String::from_utf8(body.to_vec()).unwrap();
1382 assert!(body.contains("\"error\":\"backend_unavailable\""));
1383 assert!(body.contains("object stored; indexer queue full — retry to re-enqueue"));
1384
1385 let stored = storage
1386 .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1387 .await
1388 .unwrap();
1389 assert_eq!(stored.bytes, Bytes::from_static(b"first content"));
1390
1391 let retry = router
1392 .oneshot(
1393 Request::builder()
1394 .method("PUT")
1395 .uri(format!("/api/v1/knowledgebases/{KB}/cond.md"))
1396 .header("authorization", format!("Bearer {TOKEN}"))
1397 .header("if-none-match", "*")
1398 .body(Body::from(Bytes::from_static(b"second content")))
1399 .unwrap(),
1400 )
1401 .await
1402 .unwrap();
1403
1404 assert_eq!(retry.status(), StatusCode::PRECONDITION_FAILED);
1405 assert!(retry.headers().get("retry-after").is_none());
1406
1407 let stored = storage
1408 .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1409 .await
1410 .unwrap();
1411 assert_eq!(stored.bytes, Bytes::from_static(b"first content"));
1412 }
1413
1414 #[tokio::test]
1415 async fn test_delete_returns_delete_specific_503_body_when_indexer_backpressure() {
1416 let kb = KbSlug::try_new(KB).unwrap();
1417 let object_path = ObjectPath::try_from_str("to-delete.md").unwrap();
1418 let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
1419 storage
1420 .put_object(
1421 &kb,
1422 &object_path,
1423 Bytes::from_static(b"content"),
1424 Some("text/plain"),
1425 ConditionalHeaders::default(),
1426 )
1427 .await
1428 .unwrap();
1429
1430 let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
1431 indexer_tx
1432 .try_send(IndexEvent::Upsert {
1433 kb: kb.clone(),
1434 object_key: ObjectPath::try_from_str("queued.md").unwrap(),
1435 etag: "etag".to_string(),
1436 mtime: 0,
1437 })
1438 .unwrap();
1439
1440 let mut kbs = BTreeMap::new();
1441 kbs.insert(KB.to_string(), kb.clone());
1442 let state = AppState {
1443 storage: storage.clone(),
1444 access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1445 kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1446 declared_kbs: Arc::new(kbs),
1447 authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1448 max_body_size: MAX_BODY_BYTES,
1449 max_patchable_size: MAX_BODY_BYTES,
1450 indexer_tx,
1451 searcher: Arc::new(crate::testing::NoopSearcher),
1452 events: None,
1453 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1454 readiness: crate::testing::ready_receiver(),
1455 reconcile: None,
1456 };
1457 let router = build_router(state);
1458
1459 let response = router
1460 .oneshot(
1461 Request::builder()
1462 .method("DELETE")
1463 .uri(format!("/api/v1/knowledgebases/{KB}/to-delete.md"))
1464 .header("authorization", format!("Bearer {TOKEN}"))
1465 .body(Body::empty())
1466 .unwrap(),
1467 )
1468 .await
1469 .unwrap();
1470
1471 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
1472 assert_eq!(response.headers().get("retry-after").unwrap(), "5");
1473 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1474 let body = String::from_utf8(body.to_vec()).unwrap();
1475 assert!(body.contains("\"error\":\"backend_unavailable\""));
1476 assert!(
1477 body.contains("\"message\":\"deleted from storage; retry to clear from search index\"")
1478 );
1479 assert!(!body.contains("object stored; indexer queue full — retry to re-enqueue"));
1480
1481 let deleted = storage
1482 .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1483 .await;
1484 assert!(matches!(deleted, Err(StorageError::NotFound { .. })));
1485 }
1486}