use notedthat_core::{
AccessPolicy, Authenticator, EventPublisher, EventSource, KbDetails, KbSlug, Storage,
};
use notedthat_indexer::{IndexHealth, IndexQueueSender, Searcher};
use crate::readiness::ReadinessReceiver;
use std::collections::BTreeMap;
use std::sync::Arc;
#[derive(Clone)]
pub struct AppState {
pub storage: Arc<dyn Storage>,
pub declared_kbs: Arc<BTreeMap<String, KbSlug>>,
pub access_policies: Arc<BTreeMap<String, Arc<AccessPolicy>>>,
pub kb_details: Arc<BTreeMap<String, KbDetails>>,
pub authenticator: Arc<Authenticator>,
pub max_body_size: u64,
pub max_patchable_size: u64,
pub indexer_tx: IndexQueueSender,
pub searcher: Arc<dyn Searcher>,
pub events: Option<Arc<dyn EventPublisher>>,
pub index_health: Arc<IndexHealth>,
pub readiness: ReadinessReceiver,
pub reconcile: Option<Arc<dyn ReconcileTrigger>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ReconcileBusy;
pub trait ReconcileTrigger: Send + Sync {
fn trigger(&self, kb: &KbSlug) -> Result<(), ReconcileBusy>;
}
impl AppState {
pub(crate) fn sinks(&self, source: EventSource) -> notedthat_write::WriteSinks<'_> {
notedthat_write::WriteSinks::new(
&self.indexer_tx,
self.events.as_deref(),
&self.index_health,
source,
)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::testing::InMemoryStorage;
use std::collections::BTreeMap;
use std::sync::Arc;
#[allow(clippy::needless_pass_by_value)]
fn minimal_state(tx: tokio::sync::mpsc::Sender<notedthat_indexer::IndexEvent>) -> AppState {
AppState {
storage: Arc::new(InMemoryStorage::default()),
declared_kbs: Arc::new(BTreeMap::new()),
access_policies: Arc::new(BTreeMap::new()),
kb_details: Arc::new(BTreeMap::new()),
authenticator: Arc::new(Authenticator::new("token")),
max_body_size: 1024,
max_patchable_size: 1024,
indexer_tx: (&tx).into(),
searcher: Arc::new(crate::testing::NoopSearcher),
events: None,
index_health: Arc::new(notedthat_indexer::IndexHealth::new()),
readiness: crate::testing::ready_receiver(),
reconcile: None,
}
}
#[tokio::test]
async fn clone_shares_same_channel() {
let (tx, mut rx) = tokio::sync::mpsc::channel(10);
let state = minimal_state(tx);
let cloned = state.clone();
let event = notedthat_indexer::IndexEvent::Tombstone {
kb: notedthat_core::KbSlug::try_new("test").expect("valid kb slug"),
object_key: notedthat_core::ObjectPath::try_from("a.md").expect("valid path"),
};
cloned
.indexer_tx
.send(event.clone())
.await
.expect("send on cloned tx");
let received = rx.recv().await.expect("receive on original rx");
assert_eq!(
received, event,
"cloned Sender must share the same underlying channel"
);
}
}