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