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