Skip to main content

notedthat_api_http/router/
mod.rs

1//! Axum router builder and HTTP handlers for the `NotedThat` API.
2
3// `lookup_kb` is the single definition of "is this knowledge base declared";
4// `crate::authz::KbAccess::resolve` is its only caller.
5pub(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
41/// The API route table, declared once in two forms.
42///
43/// `ROUTE_*` is the path relative to [`API_V1_PREFIX`], which is what
44/// `build_router` registers under `nest`. `MATCHED_*` is the absolute path,
45/// which is what axum reports as `MatchedPath` and what the public-read match in
46/// [`crate::middleware`] compares against. Both come from one suffix literal, so
47/// the mount point and each route are each written exactly once.
48///
49/// The routes stay nested rather than registered absolutely and merged: a
50/// `.layer()` on an absolutely-routed sub-router also wraps its fallback, and
51/// merging that fallback answers every unrouted path — `/v1/...` included — with
52/// the API's 401 instead of a 404.
53macro_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
62/// Mount point of the versioned machine API on the unified listener (D44).
63pub const API_V1_PREFIX: &str = "/api/v1";
64
65/// Mount point of the human-facing browse surface (D52, #100).
66pub 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
77/// Maximum body size for PUT requests: 16 MiB (D35).
78pub const MAX_BODY_BYTES: u64 = 16 * 1024 * 1024;
79
80/// A [`MakeRequestId`] implementation that generates `UUIDv7` request IDs.
81#[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
92/// Build the complete axum [`Router`] with all routes and middleware.
93pub 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    // Request-id generation and tracing wrap every surface this router serves,
125    // including the unauthenticated root routes. Only `auth_middleware` and the
126    // API body limit stay nested on `/api/v1`.
127    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        // `/browse` is stateful and sits on the outer router, outside
149        // `auth_middleware` — it resolves its own principal with the same rules
150        // (see `notedthat_core::Authenticator`), because nesting it under
151        // `/api/v1` would move the mount point.
152        .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    /// The router registers the relative form and the middleware matches the
164    /// absolute one. If the two stop agreeing, public-read authorization
165    /// silently stops matching the routes it is meant to guard.
166    #[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}