notedthat_api_http/state.rs
1//! Shared application state for the axum router.
2
3use notedthat_core::{
4 AccessPolicy, Authenticator, EventPublisher, EventSource, KbDetails, KbSlug, Storage,
5};
6use notedthat_indexer::{IndexHealth, Searcher};
7
8use crate::readiness::ReadinessReceiver;
9use std::collections::BTreeMap;
10use std::sync::Arc;
11
12/// Application state shared across all axum handlers.
13///
14/// This is cloned cheaply for each request (all fields are behind [`Arc`]).
15#[derive(Clone)]
16pub struct AppState {
17 /// The backing storage implementation (injected at startup).
18 pub storage: Arc<dyn Storage>,
19 /// Canonical map of slug string → [`KbSlug`] for declared knowledge bases.
20 pub declared_kbs: Arc<BTreeMap<String, KbSlug>>,
21 /// Startup snapshot of manifest access policies, keyed by KB slug.
22 ///
23 /// One `Arc` per policy so a request-scoped authorization handle can hold
24 /// its knowledge base's policy cheaply instead of borrowing the whole map.
25 pub access_policies: Arc<BTreeMap<String, Arc<AccessPolicy>>>,
26 /// Startup snapshot of each manifest's display name and description,
27 /// keyed by KB slug: what `GET /api/v1/knowledgebases` shows (#98).
28 pub kb_details: Arc<BTreeMap<String, KbDetails>>,
29 /// The credential rules every request's principal is resolved through.
30 pub authenticator: Arc<Authenticator>,
31 /// Maximum accepted PUT body size in bytes (16 MiB in M2).
32 pub max_body_size: u64,
33 /// Maximum object size eligible for patch operations, in bytes.
34 pub max_patchable_size: u64,
35 /// Sender half of the async indexing queue.
36 pub indexer_tx: tokio::sync::mpsc::Sender<notedthat_indexer::IndexEvent>,
37 /// The search implementation (injected at startup).
38 pub searcher: Arc<dyn Searcher>,
39 /// The object change event log, when `NOTEDTHAT_EVENTS_BACKEND` selects one.
40 /// `None` is the default: writes are not announced and the events route
41 /// answers 404.
42 pub events: Option<Arc<dyn EventPublisher>>,
43 /// The per-knowledge-base index health record, shared with the indexer
44 /// worker and every write path (#97).
45 pub index_health: Arc<IndexHealth>,
46 /// The latest background probe of the storage and vector backends.
47 /// `/readyz` reads it and never probes inline.
48 pub readiness: ReadinessReceiver,
49 /// What `POST …/index/reconcile` starts (D67): a comparison of one
50 /// knowledge base's storage against the index, run by the server. `None`
51 /// where the backend has no on-demand pass — the `fs` backend today, whose
52 /// watcher covers it — and the route answers `404`.
53 pub reconcile: Option<Arc<dyn ReconcileTrigger>>,
54}
55
56/// A knowledge base already has a pass running; a second one now would not see
57/// what the first is past, so the caller retries once `last_reconcile` moves.
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
59pub struct ReconcileBusy;
60
61/// Starts a reconciliation pass for one knowledge base (D67).
62///
63/// Synchronous on purpose: an implementation only claims the knowledge base's
64/// slot and spawns the pass, so the route answers before any bucket is listed.
65pub trait ReconcileTrigger: Send + Sync {
66 /// Start a pass, or report that one is already running.
67 ///
68 /// # Errors
69 ///
70 /// [`ReconcileBusy`] when the knowledge base's previous pass has not finished.
71 fn trigger(&self, kb: &KbSlug) -> Result<(), ReconcileBusy>;
72}
73
74impl AppState {
75 /// Where a write made through this surface is reported, attributed to
76 /// `source` — the HTTP API itself, or MCP when the request says so.
77 pub(crate) fn sinks(&self, source: EventSource) -> notedthat_write::WriteSinks<'_> {
78 notedthat_write::WriteSinks::new(
79 &self.indexer_tx,
80 self.events.as_deref(),
81 &self.index_health,
82 source,
83 )
84 }
85}
86
87#[cfg(test)]
88mod tests {
89 use super::*;
90 use crate::testing::InMemoryStorage;
91 use std::collections::BTreeMap;
92 use std::sync::Arc;
93
94 fn minimal_state(tx: tokio::sync::mpsc::Sender<notedthat_indexer::IndexEvent>) -> AppState {
95 AppState {
96 storage: Arc::new(InMemoryStorage::default()),
97 declared_kbs: Arc::new(BTreeMap::new()),
98 access_policies: Arc::new(BTreeMap::new()),
99 kb_details: Arc::new(BTreeMap::new()),
100 authenticator: Arc::new(Authenticator::new("token")),
101 max_body_size: 1024,
102 max_patchable_size: 1024,
103 indexer_tx: tx,
104 searcher: Arc::new(crate::testing::NoopSearcher),
105 events: None,
106 index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
107 readiness: crate::testing::ready_receiver(),
108 reconcile: None,
109 }
110 }
111
112 #[tokio::test]
113 async fn clone_shares_same_channel() {
114 let (tx, mut rx) = tokio::sync::mpsc::channel(10);
115 let state = minimal_state(tx);
116 let cloned = state.clone();
117
118 let event = notedthat_indexer::IndexEvent::Tombstone {
119 kb: notedthat_core::KbSlug::try_new("test").expect("valid kb slug"),
120 object_key: notedthat_core::ObjectPath::try_from("a.md").expect("valid path"),
121 };
122 cloned
123 .indexer_tx
124 .send(event.clone())
125 .await
126 .expect("send on cloned tx");
127
128 let received = rx.recv().await.expect("receive on original rx");
129 assert_eq!(
130 received, event,
131 "cloned Sender must share the same underlying channel"
132 );
133 }
134}