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;
10pub(crate) mod helpers;
11mod index_health;
12mod index_reconcile;
13mod kbs;
14mod llms;
15mod objects;
16mod well_known;
17
18use crate::bounds::{RequestBounds, bound};
19use crate::middleware::auth_middleware;
20use crate::state::AppState;
21use axum::Router;
22use axum::extract::{DefaultBodyLimit, Request};
23use axum::handler::Handler;
24use axum::http::HeaderName;
25use axum::middleware::from_fn_with_state;
26use axum::routing::get;
27use tower::ServiceBuilder;
28use tower_http::request_id::{
29    MakeRequestId, PropagateRequestIdLayer, RequestId, SetRequestIdLayer,
30};
31use tower_http::trace::TraceLayer;
32use uuid::Uuid;
33
34use browse::{browse_path, browse_root};
35use events::subscribe_events;
36use health::{healthz, readyz};
37use index_health::get_index_health;
38use index_reconcile::post_index_reconcile;
39use kbs::{list_kbs, list_objects};
40use llms::llms_txt;
41use objects::{delete_object, get_object, head_object, patch_object, post_object, put_object};
42use well_known::protected_resource_metadata;
43
44/// The API route table, declared once in two forms.
45///
46/// `ROUTE_*` is the path relative to [`API_V1_PREFIX`], which is what
47/// `build_router` registers under `nest`. `MATCHED_*` is the absolute path,
48/// which is what axum reports as `MatchedPath` and what the public-read match in
49/// [`crate::middleware`] compares against. Both come from one suffix literal, so
50/// the mount point and each route are each written exactly once.
51///
52/// The routes stay nested rather than registered absolutely and merged: a
53/// `.layer()` on an absolutely-routed sub-router also wraps its fallback, and
54/// merging that fallback answers every unrouted path — `/v1/...` included — with
55/// the API's 401 instead of a 404.
56macro_rules! api_routes {
57    ($($route:ident / $matched:ident => $suffix:literal,)+) => {
58        $(
59            pub(crate) const $route: &str = $suffix;
60            pub(crate) const $matched: &str = concat!("/api/v1", $suffix);
61        )+
62    };
63}
64
65/// Mount point of the versioned machine API on the unified listener (D44).
66pub const API_V1_PREFIX: &str = "/api/v1";
67
68/// Mount point of the human-facing browse surface (D52, #100).
69pub const BROWSE_PREFIX: &str = "/browse";
70
71api_routes! {
72    ROUTE_KBS / MATCHED_KBS => "/knowledgebases",
73    ROUTE_KB / MATCHED_KB => "/knowledgebases/{kb_slug}",
74    ROUTE_KB_SEARCH / MATCHED_KB_SEARCH => "/knowledgebases/{kb_slug}/search",
75    ROUTE_KB_EVENTS / MATCHED_KB_EVENTS => "/knowledgebases/{kb_slug}/events",
76    ROUTE_KB_INDEX / MATCHED_KB_INDEX => "/knowledgebases/{kb_slug}/index",
77    ROUTE_KB_INDEX_RECONCILE / MATCHED_KB_INDEX_RECONCILE => "/knowledgebases/{kb_slug}/index/reconcile",
78    ROUTE_KB_OBJECT / MATCHED_KB_OBJECT => "/knowledgebases/{kb_slug}/{*object_path}",
79}
80
81/// Maximum body size for PUT requests: 16 MiB (D35).
82pub const MAX_BODY_BYTES: u64 = 16 * 1024 * 1024;
83
84/// A [`MakeRequestId`] implementation that generates `UUIDv7` request IDs.
85#[derive(Clone, Copy, Default)]
86pub struct MakeRequestUuidV7;
87
88impl MakeRequestId for MakeRequestUuidV7 {
89    fn make_request_id<B>(&mut self, _req: &Request<B>) -> Option<RequestId> {
90        let id = Uuid::now_v7().to_string();
91        let hv = id.parse().ok()?;
92        Some(RequestId::new(hv))
93    }
94}
95
96/// Build the complete axum [`Router`] with all routes and middleware, with no
97/// effective request bounds.
98///
99/// For routers built outside a running server. The server itself calls
100/// [`build_bounded_router`].
101pub fn build_router(state: AppState) -> Router {
102    build_bounded_router(state, &RequestBounds::unbounded())
103}
104
105/// Build the complete axum [`Router`], holding every route that answers once
106/// to `bounds` (D71).
107///
108/// The routes left out are left out by where they are registered, not by a
109/// path list: the events stream, which is meant to stay open for hours, and
110/// the two probes, which answer from memory and must stay reachable to an
111/// orchestrator when the listener is at its cap — a liveness probe refused
112/// `503` under load would restart the process for being busy.
113///
114/// The corollary, since it is the one way this goes wrong quietly: a route
115/// added to the outer `Router::new()` below, beside the probes, is unbounded,
116/// and so is any unmatched path. `bounded_routes::only_the_probes_are_outside`
117/// is what keeps the first of those deliberate.
118pub fn build_bounded_router(state: AppState, bounds: &RequestBounds) -> Router {
119    let request_id_header = HeaderName::from_static("x-request-id");
120    let bounded_api = Router::new()
121        .route(ROUTE_KBS, get(list_kbs))
122        .route(ROUTE_KB, get(list_objects))
123        .route(
124            ROUTE_KB_SEARCH,
125            axum::routing::post(crate::search_route::search_kb).layer(
126                axum::extract::DefaultBodyLimit::max(crate::search_route::SEARCH_BODY_MAX_BYTES),
127            ),
128        )
129        .route(ROUTE_KB_INDEX, get(get_index_health))
130        .route(
131            ROUTE_KB_INDEX_RECONCILE,
132            axum::routing::post(post_index_reconcile),
133        )
134        .route(
135            ROUTE_KB_OBJECT,
136            get(get_object)
137                .head(head_object)
138                .put(put_object)
139                .delete(delete_object)
140                .patch(patch_object.layer(DefaultBodyLimit::disable()))
141                .post(post_object.layer(DefaultBodyLimit::disable())),
142        )
143        .route_layer(from_fn_with_state(bounds.clone(), bound));
144    let streaming_api = Router::new().route(ROUTE_KB_EVENTS, get(subscribe_events));
145    let api_routes = bounded_api
146        .merge(streaming_api)
147        .layer(
148            ServiceBuilder::new()
149                .layer(DefaultBodyLimit::max(helpers::body_limit_usize(
150                    MAX_BODY_BYTES,
151                )))
152                .layer(from_fn_with_state(state.clone(), auth_middleware)),
153        )
154        .with_state(state.clone());
155
156    let bounded_root = Router::new()
157        .route("/llms.txt", get(llms_txt))
158        .route(
159            "/.well-known/oauth-protected-resource",
160            get(protected_resource_metadata),
161        )
162        .route(BROWSE_PREFIX, get(browse_root))
163        .route(&format!("{BROWSE_PREFIX}/"), get(browse_root))
164        .route(&format!("{BROWSE_PREFIX}/{{*path}}"), get(browse_path))
165        .route_layer(from_fn_with_state(bounds.clone(), bound));
166
167    // Request-id generation and tracing wrap every surface this router serves,
168    // including the unauthenticated root routes. Only `auth_middleware` and the
169    // API body limit stay nested on `/api/v1`.
170    Router::new()
171        .route("/healthz", get(healthz))
172        .route("/readyz", get(readyz))
173        .merge(bounded_root)
174        .nest(API_V1_PREFIX, api_routes)
175        .layer(
176            ServiceBuilder::new()
177                .layer(SetRequestIdLayer::new(
178                    request_id_header.clone(),
179                    MakeRequestUuidV7,
180                ))
181                .layer(PropagateRequestIdLayer::new(request_id_header))
182                .layer(TraceLayer::new_for_http()),
183        )
184        // `/browse` is stateful and sits on the outer router, outside
185        // `auth_middleware` — it resolves its own principal with the same rules
186        // (see `notedthat_core::Authenticator`), because nesting it under
187        // `/api/v1` would move the mount point.
188        .with_state(state)
189}
190
191#[cfg(test)]
192mod bounded_routes {
193    use super::{RequestBounds, build_bounded_router};
194    use crate::state::AppState;
195    use crate::testing::InMemoryStorage;
196    use axum::body::{Body, to_bytes};
197    use axum::http::{Method, Request};
198    use notedthat_core::{Authenticator, KbSlug};
199    use std::collections::BTreeMap;
200    use std::sync::Arc;
201    use std::time::Duration;
202    use tokio::sync::Semaphore;
203    use tower::ServiceExt;
204
205    fn state() -> AppState {
206        let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
207        let notes = KbSlug::try_new("notes").expect("valid slug");
208        AppState {
209            storage: Arc::new(InMemoryStorage::default()),
210            declared_kbs: Arc::new(BTreeMap::from([("notes".to_string(), notes)])),
211            access_policies: Arc::new(BTreeMap::new()),
212            kb_details: Arc::new(BTreeMap::new()),
213            authenticator: Arc::new(Authenticator::new("token")),
214            max_body_size: 1024,
215            max_patchable_size: 1024,
216            indexer_tx,
217            searcher: Arc::new(crate::testing::NoopSearcher),
218            events: None,
219            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
220            readiness: crate::testing::ready_receiver(),
221            reconcile: None,
222        }
223    }
224
225    /// Whether the in-flight cap is what answered.
226    async fn refused_by_the_cap(method: Method, path: &str) -> bool {
227        // No permits at all: every bounded route is refused before its
228        // handler runs, and anything that answers otherwise was never bounded.
229        let bounds = RequestBounds::new(
230            Duration::from_secs(60),
231            Duration::from_secs(60),
232            Arc::new(Semaphore::new(0)),
233        );
234        let response = build_bounded_router(state(), &bounds)
235            .oneshot(
236                Request::builder()
237                    .method(method)
238                    .uri(path)
239                    .header("authorization", "Bearer token")
240                    .body(Body::empty())
241                    .expect("request"),
242            )
243            .await
244            .expect("infallible");
245        let body = to_bytes(response.into_body(), usize::MAX)
246            .await
247            .expect("body");
248        String::from_utf8_lossy(&body).contains("limit of requests in flight")
249    }
250
251    #[tokio::test]
252    async fn every_route_that_answers_once_is_bounded() {
253        for (method, path) in [
254            (Method::GET, "/api/v1/knowledgebases"),
255            (Method::GET, "/api/v1/knowledgebases/notes"),
256            (Method::POST, "/api/v1/knowledgebases/notes/search"),
257            (Method::GET, "/api/v1/knowledgebases/notes/index"),
258            (Method::POST, "/api/v1/knowledgebases/notes/index/reconcile"),
259            (Method::GET, "/api/v1/knowledgebases/notes/a.md"),
260            (Method::PUT, "/api/v1/knowledgebases/notes/a.md"),
261            (Method::GET, "/llms.txt"),
262            (Method::GET, "/.well-known/oauth-protected-resource"),
263            (Method::GET, "/browse"),
264            (Method::GET, "/browse/notes/"),
265        ] {
266            assert!(
267                refused_by_the_cap(method.clone(), path).await,
268                "{method} {path} is not bounded"
269            );
270        }
271    }
272
273    /// The outer router is the unbounded one, so what is registered there is
274    /// the one thing that can escape the cap without anybody deciding it.
275    ///
276    /// A source-level check, because axum exposes no way to enumerate a
277    /// router's routes — the same bargain `notedthat-webdav`'s
278    /// `test_router_disables_autoindex` strikes. It fails when a route is
279    /// added beside the probes, which is exactly when someone should be made
280    /// to choose rather than to inherit.
281    #[test]
282    fn only_the_probes_are_outside() {
283        let source = include_str!("mod.rs");
284        // Between the comment introducing the outer router and the merge of
285        // `bounded_root` into it: the routes it registers for itself.
286        let outer = source
287            .split("API body limit stay nested")
288            .nth(1)
289            .and_then(|after| after.split(".merge(bounded_root)").next())
290            .expect("the outer router");
291        let registered: Vec<&str> = outer
292            .lines()
293            .filter(|line| line.trim_start().starts_with(".route(\""))
294            .collect();
295        assert_eq!(
296            registered,
297            [
298                "        .route(\"/healthz\", get(healthz))",
299                "        .route(\"/readyz\", get(readyz))"
300            ],
301            "a route registered outside the bounded sub-routers draws no permit; \
302             add it to `bounded_api` or `bounded_root` instead, or say here why not"
303        );
304    }
305
306    /// The events stream stays open for hours and the probes must answer an
307    /// orchestrator under load, so none of them may draw a permit.
308    #[tokio::test]
309    async fn the_stream_and_the_probes_are_outside_the_bound() {
310        for path in ["/api/v1/knowledgebases/notes/events", "/healthz", "/readyz"] {
311            assert!(
312                !refused_by_the_cap(Method::GET, path).await,
313                "GET {path} is bounded"
314            );
315        }
316    }
317}
318
319#[cfg(test)]
320mod route_constants {
321    use super::{
322        API_V1_PREFIX, MATCHED_KB, MATCHED_KB_EVENTS, MATCHED_KB_INDEX, MATCHED_KB_INDEX_RECONCILE,
323        MATCHED_KB_OBJECT, MATCHED_KB_SEARCH, MATCHED_KBS, ROUTE_KB, ROUTE_KB_EVENTS,
324        ROUTE_KB_INDEX, ROUTE_KB_INDEX_RECONCILE, ROUTE_KB_OBJECT, ROUTE_KB_SEARCH, ROUTE_KBS,
325    };
326
327    /// The router registers the relative form and the middleware matches the
328    /// absolute one. If the two stop agreeing, public-read authorization
329    /// silently stops matching the routes it is meant to guard.
330    #[test]
331    fn matched_paths_are_the_nested_routes_under_the_mount_point() {
332        assert_eq!(API_V1_PREFIX, "/api/v1");
333        for (route, matched) in [
334            (ROUTE_KBS, MATCHED_KBS),
335            (ROUTE_KB, MATCHED_KB),
336            (ROUTE_KB_SEARCH, MATCHED_KB_SEARCH),
337            (ROUTE_KB_EVENTS, MATCHED_KB_EVENTS),
338            (ROUTE_KB_INDEX, MATCHED_KB_INDEX),
339            (ROUTE_KB_INDEX_RECONCILE, MATCHED_KB_INDEX_RECONCILE),
340            (ROUTE_KB_OBJECT, MATCHED_KB_OBJECT),
341        ] {
342            assert_eq!(matched, format!("{API_V1_PREFIX}{route}"));
343        }
344    }
345}
346
347#[cfg(test)]
348mod patch_route {
349    use super::*;
350    use async_trait::async_trait;
351    use axum::body::{Body, to_bytes};
352    use axum::http::StatusCode;
353    use axum::response::Response;
354    use bytes::Bytes;
355    use notedthat_core::{
356        ByteRange, ConditionalHeaders, KbManifest, KbSlug, ListResponse, ObjectMeta, ObjectPath,
357        ObjectRead, PutOutcome, Storage, StorageError,
358    };
359    use notedthat_indexer::IndexEvent;
360    use std::collections::BTreeMap;
361    use std::sync::Arc;
362    use tower::util::ServiceExt;
363
364    const KB: &str = "notes";
365    const OBJECT_PATH: &str = "patch.md";
366    const TOKEN: &str = "test-token-abc";
367
368    async fn router_with_object(
369        initial_body: &'static [u8],
370        max_patchable_size: u64,
371    ) -> (axum::Router, String) {
372        let kb = KbSlug::try_new(KB).unwrap();
373        let object_path = ObjectPath::try_from_str(OBJECT_PATH).unwrap();
374        let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
375        let outcome = storage
376            .put_object(
377                &kb,
378                &object_path,
379                Bytes::from_static(initial_body),
380                Some("text/markdown"),
381                ConditionalHeaders::default(),
382            )
383            .await
384            .unwrap();
385
386        let (indexer_tx, _rx) = tokio::sync::mpsc::channel(16);
387        let mut kbs = BTreeMap::new();
388        kbs.insert(KB.to_string(), kb);
389        let router = build_router(AppState {
390            storage,
391            access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
392            kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
393            declared_kbs: Arc::new(kbs),
394            authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
395            max_body_size: MAX_BODY_BYTES,
396            max_patchable_size,
397            indexer_tx,
398            searcher: Arc::new(crate::testing::NoopSearcher),
399            events: None,
400            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
401            readiness: crate::testing::ready_receiver(),
402            reconcile: None,
403        });
404
405        (router, outcome.etag.unwrap())
406    }
407
408    async fn object_with_etag(
409        storage: &crate::testing::InMemoryStorage,
410        kb: &KbSlug,
411        body: &'static [u8],
412    ) -> String {
413        storage
414            .put_object(
415                kb,
416                &ObjectPath::try_from_str(OBJECT_PATH).unwrap(),
417                Bytes::from_static(body),
418                Some("text/markdown"),
419                ConditionalHeaders::default(),
420            )
421            .await
422            .unwrap()
423            .etag
424            .unwrap()
425    }
426
427    fn router_with_storage(
428        storage: Arc<dyn Storage>,
429        kb: KbSlug,
430        max_patchable_size: u64,
431        indexer_tx: tokio::sync::mpsc::Sender<IndexEvent>,
432    ) -> axum::Router {
433        let mut kbs = BTreeMap::new();
434        kbs.insert(KB.to_string(), kb);
435        build_router(AppState {
436            storage,
437            access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
438            kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
439            declared_kbs: Arc::new(kbs),
440            authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
441            max_body_size: MAX_BODY_BYTES,
442            max_patchable_size,
443            indexer_tx,
444            searcher: Arc::new(crate::testing::NoopSearcher),
445            events: None,
446            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
447            readiness: crate::testing::ready_receiver(),
448            reconcile: None,
449        })
450    }
451
452    async fn patch_request(
453        router: axum::Router,
454        header_name: &'static str,
455        header_value: &str,
456        if_match: Option<&str>,
457        body: Bytes,
458    ) -> Response {
459        let mut builder = Request::builder()
460            .method("PATCH")
461            .uri(format!("/api/v1/knowledgebases/{KB}/{OBJECT_PATH}"))
462            .header("authorization", format!("Bearer {TOKEN}"))
463            .header(header_name, header_value);
464        if let Some(etag) = if_match {
465            builder = builder.header(axum::http::header::IF_MATCH, etag);
466        }
467
468        router
469            .oneshot(builder.body(Body::from(body)).unwrap())
470            .await
471            .unwrap()
472    }
473
474    async fn assert_error_code(response: Response, expected_status: StatusCode, expected: &str) {
475        assert_eq!(response.status(), expected_status);
476        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
477        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
478        assert_eq!(json["error"], expected);
479    }
480
481    #[tokio::test]
482    async fn bytes_content_range_returns_ok_with_etag_and_location() {
483        let (router, etag) = router_with_object(b"0123456789abcdefghij", MAX_BODY_BYTES).await;
484
485        let response = patch_request(
486            router,
487            "content-range",
488            "bytes 0-9/*",
489            Some(&etag),
490            Bytes::from_static(b"ABCDEFGHIJ"),
491        )
492        .await;
493
494        assert_eq!(response.status(), StatusCode::OK);
495        assert!(response.headers().get(axum::http::header::ETAG).is_some());
496        assert_eq!(
497            response
498                .headers()
499                .get(axum::http::header::LOCATION)
500                .unwrap(),
501            &format!("/api/v1/knowledgebases/{KB}/{OBJECT_PATH}")
502        );
503        assert!(
504            response
505                .headers()
506                .get(axum::http::header::CONTENT_RANGE)
507                .is_none()
508        );
509        assert!(response.headers().get("nt-patch-mode").is_none());
510    }
511
512    #[tokio::test]
513    async fn lines_content_range_returns_ok() {
514        let (router, etag) = router_with_object(b"one\ntwo\nthree\nfour\n", MAX_BODY_BYTES).await;
515
516        let response = patch_request(
517            router,
518            "content-range",
519            "lines 2-3/*",
520            Some(&etag),
521            Bytes::from_static(b"TWO\nTHREE\n"),
522        )
523        .await;
524
525        assert_eq!(response.status(), StatusCode::OK);
526    }
527
528    #[tokio::test]
529    async fn append_mode_without_if_match_returns_ok() {
530        let (router, _etag) = router_with_object(b"one\n", MAX_BODY_BYTES).await;
531
532        let response = patch_request(
533            router,
534            "nt-patch-mode",
535            "append",
536            None,
537            Bytes::from_static(b"two\n"),
538        )
539        .await;
540
541        assert_eq!(response.status(), StatusCode::OK);
542    }
543
544    #[tokio::test]
545    async fn append_mode_with_if_match_returns_ok() {
546        let (router, etag) = router_with_object(b"one\n", MAX_BODY_BYTES).await;
547
548        let response = patch_request(
549            router,
550            "nt-patch-mode",
551            "append",
552            Some(&etag),
553            Bytes::from_static(b"two\n"),
554        )
555        .await;
556
557        assert_eq!(response.status(), StatusCode::OK);
558    }
559
560    #[tokio::test]
561    async fn bytes_content_range_without_if_match_returns_invalid_request() {
562        let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
563
564        let response = patch_request(
565            router,
566            "content-range",
567            "bytes 0-1/*",
568            None,
569            Bytes::from_static(b"AB"),
570        )
571        .await;
572
573        assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
574    }
575
576    #[tokio::test]
577    async fn if_match_star_returns_invalid_request() {
578        let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
579
580        let response = patch_request(
581            router,
582            "content-range",
583            "bytes 0-1/*",
584            Some("*"),
585            Bytes::from_static(b"AB"),
586        )
587        .await;
588
589        assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
590    }
591
592    #[tokio::test]
593    async fn multi_value_if_match_returns_invalid_request() {
594        let (router, etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
595
596        let response = patch_request(
597            router,
598            "content-range",
599            "bytes 0-1/*",
600            Some(&format!("{etag}, \"other\"")),
601            Bytes::from_static(b"AB"),
602        )
603        .await;
604
605        assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
606    }
607
608    #[tokio::test]
609    async fn nonexistent_object_returns_not_found() {
610        let (router, etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
611
612        let response = router
613            .oneshot(
614                Request::builder()
615                    .method("PATCH")
616                    .uri(format!("/api/v1/knowledgebases/{KB}/missing.md"))
617                    .header("authorization", format!("Bearer {TOKEN}"))
618                    .header("content-range", "bytes 0-1/*")
619                    .header(axum::http::header::IF_MATCH, etag)
620                    .body(Body::from(Bytes::from_static(b"AB")))
621                    .unwrap(),
622            )
623            .await
624            .unwrap();
625
626        assert_error_code(response, StatusCode::NOT_FOUND, "not_found").await;
627    }
628
629    #[tokio::test]
630    async fn body_larger_than_max_patchable_size_returns_payload_too_large() {
631        let (router, _etag) = router_with_object(b"one\n", 4).await;
632
633        let response = patch_request(
634            router,
635            "nt-patch-mode",
636            "append",
637            None,
638            Bytes::from_static(b"abcde"),
639        )
640        .await;
641
642        assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
643    }
644
645    mod errors {
646        use super::*;
647
648        #[derive(Clone)]
649        struct PutPreconditionFailedStorage {
650            inner: crate::testing::InMemoryStorage,
651        }
652
653        #[async_trait]
654        impl Storage for PutPreconditionFailedStorage {
655            async fn probe(&self, _kb: &KbSlug) -> Result<(), StorageError> {
656                Ok(())
657            }
658
659            async fn ensure_bucket(&self, kb: &KbSlug) -> Result<(), StorageError> {
660                self.inner.ensure_bucket(kb).await
661            }
662
663            async fn read_manifest(&self, kb: &KbSlug) -> Result<KbManifest, StorageError> {
664                self.inner.read_manifest(kb).await
665            }
666
667            async fn write_manifest(
668                &self,
669                kb: &KbSlug,
670                manifest: &KbManifest,
671            ) -> Result<(), StorageError> {
672                self.inner.write_manifest(kb, manifest).await
673            }
674
675            async fn head_object(
676                &self,
677                kb: &KbSlug,
678                path: &ObjectPath,
679                conditionals: ConditionalHeaders,
680            ) -> Result<ObjectMeta, StorageError> {
681                self.inner.head_object(kb, path, conditionals).await
682            }
683
684            async fn get_object(
685                &self,
686                kb: &KbSlug,
687                path: &ObjectPath,
688                range: Option<ByteRange>,
689                conditionals: ConditionalHeaders,
690            ) -> Result<ObjectRead, StorageError> {
691                self.inner.get_object(kb, path, range, conditionals).await
692            }
693
694            async fn get_object_stream(
695                &self,
696                kb: &KbSlug,
697                path: &ObjectPath,
698                range: Option<ByteRange>,
699                conditionals: ConditionalHeaders,
700            ) -> Result<notedthat_core::ObjectStream, StorageError> {
701                self.inner
702                    .get_object_stream(kb, path, range, conditionals)
703                    .await
704            }
705
706            async fn put_object(
707                &self,
708                _kb: &KbSlug,
709                _path: &ObjectPath,
710                _bytes: Bytes,
711                _content_type: Option<&str>,
712                _conditionals: ConditionalHeaders,
713            ) -> Result<PutOutcome, StorageError> {
714                Err(StorageError::PreconditionFailed)
715            }
716
717            async fn put_staged_object(
718                &self,
719                _kb: &KbSlug,
720                _path: &ObjectPath,
721                _body: notedthat_core::StagedBody,
722                _content_type: Option<&str>,
723                _conditionals: ConditionalHeaders,
724            ) -> Result<PutOutcome, StorageError> {
725                Err(StorageError::PreconditionFailed)
726            }
727
728            async fn copy_object(
729                &self,
730                kb: &KbSlug,
731                source: &ObjectPath,
732                destination: &ObjectPath,
733                options: notedthat_core::CopyObjectOptions,
734            ) -> Result<PutOutcome, StorageError> {
735                self.inner
736                    .copy_object(kb, source, destination, options)
737                    .await
738            }
739
740            async fn delete_object(
741                &self,
742                kb: &KbSlug,
743                path: &ObjectPath,
744                conditionals: ConditionalHeaders,
745            ) -> Result<(), StorageError> {
746                self.inner.delete_object(kb, path, conditionals).await
747            }
748
749            async fn list_objects(
750                &self,
751                kb: &KbSlug,
752                prefix: Option<&str>,
753                limit: u32,
754                cursor: Option<&str>,
755            ) -> Result<ListResponse, StorageError> {
756                self.inner.list_objects(kb, prefix, limit, cursor).await
757            }
758        }
759
760        #[tokio::test]
761        async fn missing_if_match_for_bytes_mode_returns_invalid_request() {
762            let (router, _etag) = router_with_object(b"0123456789", MAX_BODY_BYTES).await;
763
764            let response = patch_request(
765                router,
766                "content-range",
767                "bytes 0-9/*",
768                None,
769                Bytes::from_static(b"ABCDEFGHIJ"),
770            )
771            .await;
772
773            assert_error_code(response, StatusCode::BAD_REQUEST, "invalid_request").await;
774        }
775
776        #[tokio::test]
777        async fn body_larger_than_max_patchable_size_returns_payload_too_large() {
778            let (router, _etag) = router_with_object(b"one\n", 10).await;
779
780            let response = patch_request(
781                router,
782                "nt-patch-mode",
783                "append",
784                None,
785                Bytes::from_static(b"more than ten bytes"),
786            )
787            .await;
788
789            assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
790        }
791
792        #[tokio::test]
793        async fn pre_splice_object_larger_than_max_patchable_size_returns_payload_too_large() {
794            let (router, _etag) = router_with_object(b"already too large", 10).await;
795
796            let response = patch_request(
797                router,
798                "nt-patch-mode",
799                "append",
800                None,
801                Bytes::from_static(b"!"),
802            )
803            .await;
804
805            assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
806        }
807
808        #[tokio::test]
809        async fn post_splice_body_larger_than_max_patchable_size_returns_payload_too_large() {
810            let (router, _etag) = router_with_object(b"123456", 10).await;
811
812            let response = patch_request(
813                router,
814                "nt-patch-mode",
815                "append",
816                None,
817                Bytes::from_static(b"78901"),
818            )
819            .await;
820
821            assert_error_code(response, StatusCode::PAYLOAD_TOO_LARGE, "payload_too_large").await;
822        }
823
824        #[tokio::test]
825        async fn line_range_out_of_bounds_returns_dual_416_headers_and_empty_body() {
826            let (router, etag) =
827                router_with_object(b"one\ntwo\nthree\nfour\nfive\n", MAX_BODY_BYTES).await;
828
829            let response = patch_request(
830                router,
831                "content-range",
832                "lines 100-200/*",
833                Some(&etag),
834                Bytes::from_static(b"replacement\n"),
835            )
836            .await;
837
838            assert_eq!(response.status(), StatusCode::RANGE_NOT_SATISFIABLE);
839            assert_eq!(
840                response.headers().get("content-range").unwrap(),
841                "lines */5"
842            );
843            assert_eq!(
844                response.headers().get("x-content-range-bytes").unwrap(),
845                "*/24"
846            );
847            let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
848            assert!(body.is_empty());
849        }
850
851        #[tokio::test]
852        async fn if_match_mismatch_after_retries_returns_precondition_failed_without_content_range()
853        {
854            let kb = KbSlug::try_new(KB).unwrap();
855            let inner = crate::testing::InMemoryStorage::with_kbs([&kb]);
856            let etag = object_with_etag(&inner, &kb, b"0123456789").await;
857            let (indexer_tx, _rx) = tokio::sync::mpsc::channel(16);
858            let router = router_with_storage(
859                Arc::new(PutPreconditionFailedStorage { inner }),
860                kb,
861                MAX_BODY_BYTES,
862                indexer_tx,
863            );
864
865            let response = patch_request(
866                router,
867                "content-range",
868                "bytes 0-1/*",
869                Some(&etag),
870                Bytes::from_static(b"AB"),
871            )
872            .await;
873
874            assert_eq!(response.status(), StatusCode::PRECONDITION_FAILED);
875            assert!(
876                response
877                    .headers()
878                    .get(axum::http::header::CONTENT_RANGE)
879                    .is_none()
880            );
881            let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
882            let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
883            assert_eq!(json["error"], "precondition_failed");
884        }
885
886        #[tokio::test]
887        async fn indexer_queue_full_returns_backend_unavailable_with_retry_after() {
888            let kb = KbSlug::try_new(KB).unwrap();
889            let storage = crate::testing::InMemoryStorage::with_kbs([&kb]);
890            let etag = object_with_etag(&storage, &kb, b"one\n").await;
891            let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
892            indexer_tx
893                .try_send(IndexEvent::Upsert {
894                    kb: kb.clone(),
895                    object_key: ObjectPath::try_from_str("queued.md").unwrap(),
896                    etag: "queued".to_string(),
897                    mtime: 0,
898                })
899                .unwrap();
900            let router = router_with_storage(Arc::new(storage), kb, MAX_BODY_BYTES, indexer_tx);
901
902            let response = patch_request(
903                router,
904                "nt-patch-mode",
905                "append",
906                Some(&etag),
907                Bytes::from_static(b"two\n"),
908            )
909            .await;
910
911            assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
912            assert_eq!(response.headers().get("retry-after").unwrap(), "5");
913            let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
914            let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
915            assert_eq!(json["error"], "backend_unavailable");
916        }
917    }
918}
919
920#[cfg(test)]
921mod line_range_get {
922    use super::*;
923    use axum::body::{Body, to_bytes};
924    use axum::http::StatusCode;
925    use axum::response::Response;
926    use bytes::Bytes;
927    use notedthat_core::{ConditionalHeaders, KbSlug, ObjectPath, Storage};
928    use std::collections::BTreeMap;
929    use std::sync::Arc;
930    use tower::util::ServiceExt;
931
932    const KB: &str = "notes";
933    const TOKEN: &str = "test-token-abc";
934
935    fn twenty_line_markdown() -> String {
936        markdown_lines(1, 20)
937    }
938
939    fn markdown_lines(first: u32, last: u32) -> String {
940        let mut body = String::new();
941        for line in first..=last {
942            std::fmt::Write::write_fmt(&mut body, format_args!("line {line:02}\n")).unwrap();
943        }
944        body
945    }
946
947    async fn router_with_markdown_object(body: String) -> axum::Router {
948        let kb = KbSlug::try_new(KB).unwrap();
949        let object_path = ObjectPath::try_from_str("ranges.md").unwrap();
950        let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
951        storage
952            .put_object(
953                &kb,
954                &object_path,
955                Bytes::from(body),
956                Some("text/markdown"),
957                ConditionalHeaders::default(),
958            )
959            .await
960            .unwrap();
961
962        let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
963        let mut kbs = BTreeMap::new();
964        kbs.insert(KB.to_string(), kb);
965        build_router(AppState {
966            storage,
967            access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
968            kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
969            declared_kbs: Arc::new(kbs),
970            authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
971            max_body_size: MAX_BODY_BYTES,
972            max_patchable_size: MAX_BODY_BYTES,
973            indexer_tx,
974            searcher: Arc::new(crate::testing::NoopSearcher),
975            events: None,
976            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
977            readiness: crate::testing::ready_receiver(),
978            reconcile: None,
979        })
980    }
981
982    async fn get_ranges_md(router: axum::Router, range: &str) -> Response {
983        router
984            .oneshot(
985                Request::builder()
986                    .method("GET")
987                    .uri(format!("/api/v1/knowledgebases/{KB}/ranges.md"))
988                    .header("authorization", format!("Bearer {TOKEN}"))
989                    .header(axum::http::header::RANGE, range)
990                    .body(Body::empty())
991                    .unwrap(),
992            )
993            .await
994            .unwrap()
995    }
996
997    #[tokio::test]
998    async fn returns_first_five_lines_when_closed_range_requested() {
999        let router = router_with_markdown_object(twenty_line_markdown()).await;
1000
1001        let response = get_ranges_md(router, "lines=1-5").await;
1002
1003        assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1004        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1005        assert_eq!(body, Bytes::from(markdown_lines(1, 5)));
1006    }
1007
1008    #[tokio::test]
1009    async fn returns_last_three_lines_when_suffix_range_requested() {
1010        let router = router_with_markdown_object(twenty_line_markdown()).await;
1011
1012        let response = get_ranges_md(router, "lines=-3").await;
1013
1014        assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1015        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1016        assert_eq!(body, Bytes::from(markdown_lines(18, 20)));
1017    }
1018
1019    #[tokio::test]
1020    async fn returns_empty_body_when_insert_range_requested() {
1021        let router = router_with_markdown_object(twenty_line_markdown()).await;
1022
1023        let response = get_ranges_md(router, "lines=5-4").await;
1024
1025        assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1026        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1027        assert!(body.is_empty());
1028    }
1029
1030    #[tokio::test]
1031    async fn returns_full_body_when_unknown_range_unit_requested() {
1032        let body = twenty_line_markdown();
1033        let router = router_with_markdown_object(body.clone()).await;
1034
1035        let response = get_ranges_md(router, "items=0-5").await;
1036
1037        assert_eq!(response.status(), StatusCode::OK);
1038        let actual = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1039        assert_eq!(actual, Bytes::from(body));
1040    }
1041
1042    mod headers {
1043        use super::*;
1044
1045        fn ten_line_markdown() -> String {
1046            markdown_lines(1, 10)
1047        }
1048
1049        #[tokio::test]
1050        async fn returns_line_and_byte_content_ranges_when_closed_range_requested() {
1051            let router = router_with_markdown_object(ten_line_markdown()).await;
1052
1053            let response = get_ranges_md(router, "lines=2-4").await;
1054
1055            assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1056            assert_eq!(
1057                response.headers().get("Content-Range").unwrap(),
1058                "lines 2-4/10"
1059            );
1060            assert_eq!(
1061                response.headers().get("X-Content-Range-Bytes").unwrap(),
1062                "8-31/80"
1063            );
1064        }
1065
1066        #[tokio::test]
1067        async fn returns_slice_content_length_when_closed_range_requested() {
1068            let router = router_with_markdown_object(ten_line_markdown()).await;
1069
1070            let response = get_ranges_md(router, "lines=2-4").await;
1071
1072            assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1073            assert_eq!(response.headers().get("content-length").unwrap(), "24");
1074        }
1075
1076        #[tokio::test]
1077        async fn returns_zero_length_and_empty_byte_range_when_insert_range_requested() {
1078            let router = router_with_markdown_object(ten_line_markdown()).await;
1079
1080            let response = get_ranges_md(router, "lines=5-4").await;
1081
1082            assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1083            assert_eq!(response.headers().get("content-length").unwrap(), "0");
1084            assert_eq!(
1085                response.headers().get("Content-Range").unwrap(),
1086                "lines 5-4/10"
1087            );
1088            assert_eq!(
1089                response.headers().get("X-Content-Range-Bytes").unwrap(),
1090                "32-31/80"
1091            );
1092        }
1093
1094        #[tokio::test]
1095        async fn omits_line_byte_range_header_when_byte_range_requested() {
1096            let router = router_with_markdown_object(ten_line_markdown()).await;
1097
1098            let response = get_ranges_md(router, "bytes=0-9").await;
1099
1100            assert_eq!(response.status(), StatusCode::PARTIAL_CONTENT);
1101            assert!(response.headers().get("X-Content-Range-Bytes").is_none());
1102        }
1103    }
1104}
1105
1106#[cfg(test)]
1107mod tests {
1108    use super::*;
1109    use axum::body::{Body, to_bytes};
1110    use axum::http::StatusCode;
1111    use axum::response::Response;
1112    use bytes::Bytes;
1113    use notedthat_core::{ConditionalHeaders, KbSlug, ObjectPath, Storage, StorageError};
1114    use notedthat_indexer::IndexEvent;
1115    use std::collections::BTreeMap;
1116    use std::sync::Arc;
1117    use tower::util::ServiceExt;
1118
1119    const KB: &str = "notes";
1120    const TOKEN: &str = "test-token-abc";
1121
1122    fn router() -> axum::Router {
1123        let kb = KbSlug::try_new(KB).unwrap();
1124        let mut kbs = BTreeMap::new();
1125        kbs.insert(KB.to_string(), kb);
1126        let (indexer_tx, mut rx) = tokio::sync::mpsc::channel(16);
1127        tokio::spawn(async move { while rx.recv().await.is_some() {} });
1128
1129        build_router(AppState {
1130            storage: Arc::new(crate::testing::InMemoryStorage::with_kbs(kbs.values())),
1131            access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1132            kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1133            declared_kbs: Arc::new(kbs),
1134            authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1135            max_body_size: MAX_BODY_BYTES,
1136            max_patchable_size: MAX_BODY_BYTES,
1137            indexer_tx,
1138            searcher: Arc::new(crate::testing::NoopSearcher),
1139            events: None,
1140            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1141            readiness: crate::testing::ready_receiver(),
1142            reconcile: None,
1143        })
1144    }
1145
1146    async fn put_object(router: axum::Router, path: &str, body: &'static [u8]) -> Response {
1147        router
1148            .oneshot(
1149                Request::builder()
1150                    .method("PUT")
1151                    .uri(format!("/api/v1/knowledgebases/{KB}/{path}"))
1152                    .header("authorization", format!("Bearer {TOKEN}"))
1153                    .header(axum::http::header::CONTENT_TYPE, "text/markdown")
1154                    .body(Body::from(Bytes::from_static(body)))
1155                    .unwrap(),
1156            )
1157            .await
1158            .unwrap()
1159    }
1160
1161    async fn put_object_etag(router: axum::Router, path: &str, body: &'static [u8]) -> String {
1162        let response = put_object(router, path, body).await;
1163        assert_eq!(response.status(), StatusCode::CREATED);
1164        response
1165            .headers()
1166            .get(axum::http::header::ETAG)
1167            .unwrap()
1168            .to_str()
1169            .unwrap()
1170            .to_string()
1171    }
1172
1173    async fn get_object(router: axum::Router, path: &str) -> Response {
1174        router
1175            .oneshot(
1176                Request::builder()
1177                    .method("GET")
1178                    .uri(format!("/api/v1/knowledgebases/{KB}/{path}"))
1179                    .header("authorization", format!("Bearer {TOKEN}"))
1180                    .body(Body::empty())
1181                    .unwrap(),
1182            )
1183            .await
1184            .unwrap()
1185    }
1186
1187    async fn post_replace(
1188        router: axum::Router,
1189        path: &str,
1190        if_match: &str,
1191        body: &'static [u8],
1192    ) -> Response {
1193        router
1194            .oneshot(
1195                Request::builder()
1196                    .method("POST")
1197                    .uri(format!("/api/v1/knowledgebases/{KB}/replace/{path}"))
1198                    .header("authorization", format!("Bearer {TOKEN}"))
1199                    .header(axum::http::header::CONTENT_TYPE, "application/json")
1200                    .header(axum::http::header::IF_MATCH, if_match)
1201                    .body(Body::from(Bytes::from_static(body)))
1202                    .unwrap(),
1203            )
1204            .await
1205            .unwrap()
1206    }
1207
1208    #[tokio::test]
1209    async fn get_on_replace_prefixed_path_still_reads_object_via_catch_all() {
1210        let router = router();
1211        put_object_etag(router.clone(), "replace/foo.md", b"hi").await;
1212
1213        let response = get_object(router, "replace/foo.md").await;
1214
1215        assert_eq!(response.status(), StatusCode::OK);
1216        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1217        assert_eq!(&body[..], b"hi");
1218    }
1219
1220    #[tokio::test]
1221    async fn patch_on_replace_prefixed_path_still_reaches_patch_object() {
1222        let router = router();
1223        let etag = put_object_etag(router.clone(), "replace/bar.md", b"old\n").await;
1224
1225        let response = router
1226            .oneshot(
1227                Request::builder()
1228                    .method("PATCH")
1229                    .uri(format!("/api/v1/knowledgebases/{KB}/replace/bar.md"))
1230                    .header("authorization", format!("Bearer {TOKEN}"))
1231                    .header(axum::http::header::CONTENT_RANGE, "lines 1-1/*")
1232                    .header(axum::http::header::IF_MATCH, etag)
1233                    .body(Body::from(Bytes::from_static(b"new\n")))
1234                    .unwrap(),
1235            )
1236            .await
1237            .unwrap();
1238
1239        assert_eq!(response.status(), StatusCode::OK);
1240    }
1241
1242    #[tokio::test]
1243    async fn put_and_delete_on_replace_prefixed_path_still_work() {
1244        let router = router();
1245        let put = put_object(router.clone(), "replace/delete.md", b"gone").await;
1246
1247        assert_eq!(put.status(), StatusCode::CREATED);
1248        let delete = router
1249            .oneshot(
1250                Request::builder()
1251                    .method("DELETE")
1252                    .uri(format!("/api/v1/knowledgebases/{KB}/replace/delete.md"))
1253                    .header("authorization", format!("Bearer {TOKEN}"))
1254                    .body(Body::empty())
1255                    .unwrap(),
1256            )
1257            .await
1258            .unwrap();
1259        assert_eq!(delete.status(), StatusCode::NO_CONTENT);
1260    }
1261
1262    #[tokio::test]
1263    async fn post_on_non_replace_path_returns_404_not_found() {
1264        let response = router()
1265            .oneshot(
1266                Request::builder()
1267                    .method("POST")
1268                    .uri(format!("/api/v1/knowledgebases/{KB}/foo.md"))
1269                    .header("authorization", format!("Bearer {TOKEN}"))
1270                    .body(Body::empty())
1271                    .unwrap(),
1272            )
1273            .await
1274            .unwrap();
1275
1276        assert_eq!(response.status(), StatusCode::NOT_FOUND);
1277        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1278        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1279        assert_eq!(json["error"], "not_found");
1280        assert!(
1281            json["message"]
1282                .as_str()
1283                .unwrap()
1284                .contains("supported actions: 'replace/<path>'")
1285        );
1286    }
1287
1288    #[tokio::test]
1289    async fn post_on_replace_prefixed_path_dispatches_to_replace_handler() {
1290        let router = router();
1291        let etag = put_object_etag(router.clone(), "target.md", b"hello world").await;
1292
1293        let response = post_replace(
1294            router,
1295            "target.md",
1296            &etag,
1297            br#"{"old_string":"world","new_string":"planet"}"#,
1298        )
1299        .await;
1300
1301        assert_eq!(response.status(), StatusCode::OK);
1302        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1303        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1304        assert_eq!(json["match_count"], 1);
1305    }
1306
1307    #[tokio::test]
1308    async fn post_on_replace_replace_path_targets_the_replace_prefixed_object() {
1309        let router = router();
1310        let etag = put_object_etag(router.clone(), "replace/nested.md", b"foo bar").await;
1311
1312        let response = post_replace(
1313            router.clone(),
1314            "replace/nested.md",
1315            &etag,
1316            br#"{"old_string":"bar","new_string":"baz"}"#,
1317        )
1318        .await;
1319
1320        assert_eq!(response.status(), StatusCode::OK);
1321        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1322        let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
1323        assert_eq!(json["match_count"], 1);
1324        let get = get_object(router, "replace/nested.md").await;
1325        let body = to_bytes(get.into_body(), 64 * 1024).await.unwrap();
1326        assert_eq!(&body[..], b"foo baz");
1327    }
1328
1329    #[tokio::test]
1330    async fn test_conditional_put_503_then_naive_retry_412_keeps_object_stored() {
1331        let kb = KbSlug::try_new(KB).unwrap();
1332        let object_path = ObjectPath::try_from_str("cond.md").unwrap();
1333        let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
1334
1335        let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
1336        indexer_tx
1337            .try_send(IndexEvent::Upsert {
1338                kb: kb.clone(),
1339                object_key: ObjectPath::try_from_str("queued.md").unwrap(),
1340                etag: "etag".to_string(),
1341                mtime: 0,
1342            })
1343            .unwrap();
1344
1345        let mut kbs = BTreeMap::new();
1346        kbs.insert(KB.to_string(), kb.clone());
1347        let state = AppState {
1348            storage: storage.clone(),
1349            access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1350            kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1351            declared_kbs: Arc::new(kbs),
1352            authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1353            max_body_size: MAX_BODY_BYTES,
1354            max_patchable_size: MAX_BODY_BYTES,
1355            indexer_tx,
1356            searcher: Arc::new(crate::testing::NoopSearcher),
1357            events: None,
1358            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1359            readiness: crate::testing::ready_receiver(),
1360            reconcile: None,
1361        };
1362        let router = build_router(state);
1363
1364        let response = router
1365            .clone()
1366            .oneshot(
1367                Request::builder()
1368                    .method("PUT")
1369                    .uri(format!("/api/v1/knowledgebases/{KB}/cond.md"))
1370                    .header("authorization", format!("Bearer {TOKEN}"))
1371                    .header("if-none-match", "*")
1372                    .body(Body::from(Bytes::from_static(b"first content")))
1373                    .unwrap(),
1374            )
1375            .await
1376            .unwrap();
1377
1378        assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
1379        assert_eq!(response.headers().get("retry-after").unwrap(), "5");
1380        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1381        let body = String::from_utf8(body.to_vec()).unwrap();
1382        assert!(body.contains("\"error\":\"backend_unavailable\""));
1383        assert!(body.contains("object stored; indexer queue full — retry to re-enqueue"));
1384
1385        let stored = storage
1386            .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1387            .await
1388            .unwrap();
1389        assert_eq!(stored.bytes, Bytes::from_static(b"first content"));
1390
1391        let retry = router
1392            .oneshot(
1393                Request::builder()
1394                    .method("PUT")
1395                    .uri(format!("/api/v1/knowledgebases/{KB}/cond.md"))
1396                    .header("authorization", format!("Bearer {TOKEN}"))
1397                    .header("if-none-match", "*")
1398                    .body(Body::from(Bytes::from_static(b"second content")))
1399                    .unwrap(),
1400            )
1401            .await
1402            .unwrap();
1403
1404        assert_eq!(retry.status(), StatusCode::PRECONDITION_FAILED);
1405        assert!(retry.headers().get("retry-after").is_none());
1406
1407        let stored = storage
1408            .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1409            .await
1410            .unwrap();
1411        assert_eq!(stored.bytes, Bytes::from_static(b"first content"));
1412    }
1413
1414    #[tokio::test]
1415    async fn test_delete_returns_delete_specific_503_body_when_indexer_backpressure() {
1416        let kb = KbSlug::try_new(KB).unwrap();
1417        let object_path = ObjectPath::try_from_str("to-delete.md").unwrap();
1418        let storage = Arc::new(crate::testing::InMemoryStorage::with_kbs([&kb]));
1419        storage
1420            .put_object(
1421                &kb,
1422                &object_path,
1423                Bytes::from_static(b"content"),
1424                Some("text/plain"),
1425                ConditionalHeaders::default(),
1426            )
1427            .await
1428            .unwrap();
1429
1430        let (indexer_tx, _rx) = tokio::sync::mpsc::channel(1);
1431        indexer_tx
1432            .try_send(IndexEvent::Upsert {
1433                kb: kb.clone(),
1434                object_key: ObjectPath::try_from_str("queued.md").unwrap(),
1435                etag: "etag".to_string(),
1436                mtime: 0,
1437            })
1438            .unwrap();
1439
1440        let mut kbs = BTreeMap::new();
1441        kbs.insert(KB.to_string(), kb.clone());
1442        let state = AppState {
1443            storage: storage.clone(),
1444            access_policies: Arc::new(notedthat_core::signed_in_policies(&kbs)),
1445            kb_details: Arc::new(notedthat_core::slug_kb_details(&kbs)),
1446            declared_kbs: Arc::new(kbs),
1447            authenticator: Arc::new(notedthat_core::Authenticator::new(TOKEN)),
1448            max_body_size: MAX_BODY_BYTES,
1449            max_patchable_size: MAX_BODY_BYTES,
1450            indexer_tx,
1451            searcher: Arc::new(crate::testing::NoopSearcher),
1452            events: None,
1453            index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
1454            readiness: crate::testing::ready_receiver(),
1455            reconcile: None,
1456        };
1457        let router = build_router(state);
1458
1459        let response = router
1460            .oneshot(
1461                Request::builder()
1462                    .method("DELETE")
1463                    .uri(format!("/api/v1/knowledgebases/{KB}/to-delete.md"))
1464                    .header("authorization", format!("Bearer {TOKEN}"))
1465                    .body(Body::empty())
1466                    .unwrap(),
1467            )
1468            .await
1469            .unwrap();
1470
1471        assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
1472        assert_eq!(response.headers().get("retry-after").unwrap(), "5");
1473        let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap();
1474        let body = String::from_utf8(body.to_vec()).unwrap();
1475        assert!(body.contains("\"error\":\"backend_unavailable\""));
1476        assert!(
1477            body.contains("\"message\":\"deleted from storage; retry to clear from search index\"")
1478        );
1479        assert!(!body.contains("object stored; indexer queue full — retry to re-enqueue"));
1480
1481        let deleted = storage
1482            .get_object(&kb, &object_path, None, ConditionalHeaders::default())
1483            .await;
1484        assert!(matches!(deleted, Err(StorageError::NotFound { .. })));
1485    }
1486}