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