Skip to main content

kmp_embedded/
kernel.rs

1use std::path::{Path, PathBuf};
2use std::sync::Arc;
3use std::time::Duration;
4
5use kmp_adapter_embedded::{
6    EmbeddedKernelStore, QualityTelemetryRetention, RedbQualityTelemetryWriter, StorageEngine,
7};
8use kmp_application::{
9    CommandApplicationService, KernelMemoryApplicationService, QueryApplicationService,
10    RoutingProjectionWriter, UpdateContextUseCase,
11};
12use kmp_domain::{PortError, QualityMetricsObserver};
13use kmp_observability::{BufferedQualityMetricsObserver, EmbeddedTelemetryGuard};
14
15const GENERATOR_VERSION: &str = env!("CARGO_PKG_VERSION");
16
17/// The KMP memory facade composed over the embedded store: every port is the
18/// same single-file redb store, and ingest projects synchronously in-process,
19/// which is what makes `read_after_write_ready` unconditionally true.
20pub type EmbeddedMemoryService = KernelMemoryApplicationService<
21    EmbeddedKernelStore,
22    EmbeddedKernelStore,
23    EmbeddedKernelStore,
24    EmbeddedKernelStore,
25    RoutingProjectionWriter<Arc<EmbeddedKernelStore>, Arc<EmbeddedKernelStore>>,
26>;
27
28/// One opened embedded kernel: the composed service plus the store handle
29/// for operational tooling (replay, stats, compaction).
30pub struct EmbeddedKernel {
31    data_dir: PathBuf,
32    engine: StorageEngine,
33    store: EmbeddedKernelStore,
34    service: Arc<EmbeddedMemoryService>,
35    quality_observer: Arc<BufferedQualityMetricsObserver>,
36    telemetry_guard: Option<EmbeddedTelemetryGuard>,
37    telemetry_writer: Option<Arc<RedbQualityTelemetryWriter>>,
38    quality_telemetry_error: Option<String>,
39}
40
41impl EmbeddedKernel {
42    /// Opens the store at `data_dir` (fail-fast per ADR-012) and composes the
43    /// kernel. An existing directory opens with the engine it was created
44    /// with; a fresh one gets the caller's resolved default. On redb a second
45    /// session on the same data dir fails here with an explicit single-writer
46    /// error (ADR-011; the engine holds the lock); on SQLite it does not.
47    pub fn open(data_dir: &Path) -> Result<Self, PortError> {
48        Self::open_with_engine(data_dir, None)
49    }
50
51    /// [`open`](Self::open) with a say in the engine (ADR-018): a fresh
52    /// directory is created for `engine`; an existing one must already be
53    /// `engine`, and the mismatch is refused by name rather than quietly
54    /// opened as whatever it is — that is how a user ends up on the wrong
55    /// engine without knowing. `None` means no preference.
56    pub fn open_with_engine(
57        data_dir: &Path,
58        engine: Option<StorageEngine>,
59    ) -> Result<Self, PortError> {
60        let opened = match engine {
61            Some(engine) => EmbeddedKernelStore::open_with_engine(data_dir, engine),
62            None => EmbeddedKernelStore::open(data_dir),
63        };
64        let store = opened.map_err(|error| match error {
65            // redb's own words for "another process has it". Only redb says
66            // this: SQLite waits for the commit lock instead of refusing.
67            PortError::Unavailable(message) if message.contains("Cannot acquire lock") => {
68                PortError::Unavailable(format!(
69                    "{message}; if another agent session is using this data dir, close it first \
70                     (the redb engine is single-writer per ADR-011 — to share one store between \
71                     hosts, migrate it to the sqlite engine: `kmp-mcp migrate <this-dir> \
72                     <new-dir> --engine sqlite`)"
73                ))
74            }
75            other => other,
76        })?;
77
78        let graph = Arc::new(store.clone());
79        let detail = Arc::new(store.clone());
80        let query_application = Arc::new(QueryApplicationService::new(
81            Arc::clone(&graph),
82            Arc::clone(&detail),
83            Arc::new(store.clone()),
84            GENERATOR_VERSION,
85        ));
86        let update_context = Arc::new(UpdateContextUseCase::new_with_projection_writer(
87            Arc::new(store.clone()),
88            RoutingProjectionWriter::new(graph, detail),
89            GENERATOR_VERSION,
90        ));
91        let service = Arc::new(KernelMemoryApplicationService::new(
92            query_application,
93            Arc::new(CommandApplicationService::new(update_context)),
94        ));
95        let (quality_observer, telemetry_guard, telemetry_writer, quality_telemetry_error) =
96            compose_quality_telemetry(data_dir);
97
98        // Read back rather than trusted from the request: what the directory
99        // says it is, is what we are on.
100        let engine = EmbeddedKernelStore::engine_of(data_dir)?;
101
102        Ok(Self {
103            data_dir: data_dir.to_path_buf(),
104            engine,
105            store,
106            service,
107            quality_observer,
108            telemetry_guard,
109            telemetry_writer,
110            quality_telemetry_error,
111        })
112    }
113
114    /// The engine behind this kernel's store, as its data directory records
115    /// it. Surfaced at startup and by the doctor so a user can see which one
116    /// they are on.
117    pub fn engine(&self) -> StorageEngine {
118        self.engine
119    }
120
121    pub fn data_dir(&self) -> &Path {
122        &self.data_dir
123    }
124
125    pub fn store(&self) -> &EmbeddedKernelStore {
126        &self.store
127    }
128
129    pub fn service(&self) -> Arc<EmbeddedMemoryService> {
130        Arc::clone(&self.service)
131    }
132
133    pub fn quality_observer(&self) -> Arc<dyn QualityMetricsObserver> {
134        self.quality_observer.clone()
135    }
136
137    pub fn quality_telemetry_dropped_observations(&self) -> u64 {
138        self.quality_observer.dropped_observations()
139    }
140
141    pub fn quality_telemetry_write_failures(&self) -> u64 {
142        self.telemetry_writer
143            .as_ref()
144            .map_or(0, |writer| writer.write_failures())
145    }
146
147    pub fn quality_telemetry_error(&self) -> Option<&str> {
148        self.quality_telemetry_error.as_deref()
149    }
150
151    pub fn quality_telemetry_active(&self) -> bool {
152        self.telemetry_guard.is_some()
153    }
154}
155
156fn compose_quality_telemetry(
157    data_dir: &Path,
158) -> (
159    Arc<BufferedQualityMetricsObserver>,
160    Option<EmbeddedTelemetryGuard>,
161    Option<Arc<RedbQualityTelemetryWriter>>,
162    Option<String>,
163) {
164    let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(1_024);
165    let observer = Arc::new(observer);
166    let writer =
167        match RedbQualityTelemetryWriter::open(data_dir, QualityTelemetryRetention::default()) {
168            Ok(writer) => Arc::new(writer),
169            Err(error) => {
170                drop(receiver);
171                return (observer, None, None, Some(error.to_string()));
172            }
173        };
174    let batch_writer = Arc::clone(&writer);
175    let final_writer = Arc::clone(&writer);
176    let guard = match EmbeddedTelemetryGuard::try_spawn(
177        receiver,
178        64,
179        Duration::from_millis(250),
180        move |batch| {
181            let _ = batch_writer.write_batch(&batch);
182        },
183        move || {
184            let _ = final_writer.flush_durable();
185        },
186    ) {
187        Ok(guard) => guard,
188        Err(error) => {
189            return (
190                observer,
191                None,
192                Some(writer),
193                Some(format!("quality telemetry worker could not start: {error}")),
194            );
195        }
196    };
197    (observer, Some(guard), Some(writer), None)
198}
199
200#[cfg(test)]
201mod tests {
202    use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
203
204    use super::EmbeddedKernel;
205
206    #[test]
207    fn telemetry_startup_failure_never_prevents_the_kernel_from_opening() {
208        let data_dir = tempfile::tempdir().expect("temp data dir");
209        std::fs::write(
210            data_dir.path().join("telemetry"),
211            b"blocks directory creation",
212        )
213        .expect("blocking file");
214
215        let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel remains available");
216        assert!(!kernel.quality_telemetry_active());
217        assert!(kernel.quality_telemetry_error().is_some());
218        let metrics = BundleQualityMetrics::new(1, 1.0, 0.0, 0.0, 0.0).expect("valid metrics");
219        kernel.quality_observer.observe(
220            &metrics,
221            &QualityObservationContext {
222                rpc: "kernel_wake".to_string(),
223                root_node_id: "question:fail-open".to_string(),
224                role: "resumer".to_string(),
225            },
226        );
227        assert_eq!(kernel.quality_telemetry_dropped_observations(), 1);
228        assert_eq!(kernel.quality_telemetry_write_failures(), 0);
229    }
230}