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