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, SqliteQualityTelemetryReader,
7    SqliteQualityTelemetryWriter, StorageEngine,
8};
9use kmp_application::{
10    CommandApplicationService, KernelMemoryApplicationService, QueryApplicationService,
11    UpdateContextUseCase,
12};
13use kmp_domain::{PortError, QualityMetricsObserver};
14use kmp_observability::{BufferedQualityMetricsObserver, EmbeddedTelemetryGuard};
15
16const GENERATOR_VERSION: &str = env!("CARGO_PKG_VERSION");
17
18/// The KMP memory facade composed over one stamped embedded store: every port
19/// shares the selected engine, and ingest projects synchronously in-process,
20/// which is what makes `read_after_write_ready` unconditionally true.
21pub type EmbeddedMemoryService = KernelMemoryApplicationService<
22    EmbeddedKernelStore,
23    EmbeddedKernelStore,
24    EmbeddedKernelStore,
25    EmbeddedKernelStore,
26    EmbeddedKernelStore,
27>;
28
29/// One opened embedded kernel: the composed service plus the store handle
30/// for operational tooling (replay, stats, compaction).
31pub struct EmbeddedKernel {
32    data_dir: PathBuf,
33    engine: StorageEngine,
34    store: EmbeddedKernelStore,
35    service: Arc<EmbeddedMemoryService>,
36    quality_observer: Arc<BufferedQualityMetricsObserver>,
37    telemetry_guard: Option<EmbeddedTelemetryGuard>,
38    telemetry_writer: Option<Arc<SqliteQualityTelemetryWriter>>,
39    quality_telemetry_error: Option<String>,
40}
41
42impl EmbeddedKernel {
43    /// Opens the store at `data_dir` (fail-fast per ADR-012) and composes the
44    /// kernel. An existing directory opens with the engine it was created
45    /// with; a fresh one gets the caller's resolved default. SQLite permits
46    /// several sessions to share the same data directory.
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?;
65
66        let graph = Arc::new(store.clone());
67        let detail = Arc::new(store.clone());
68        let query_application = Arc::new(
69            QueryApplicationService::new(
70                Arc::clone(&graph),
71                Arc::clone(&detail),
72                Arc::new(store.clone()),
73                GENERATOR_VERSION,
74            )
75            .with_read_snapshots(Arc::new(store.clone())),
76        );
77        let update_context = Arc::new(UpdateContextUseCase::new_with_projection_writer(
78            Arc::new(store.clone()),
79            store.clone(),
80            GENERATOR_VERSION,
81        ));
82        let service = Arc::new(KernelMemoryApplicationService::new(
83            query_application,
84            Arc::new(CommandApplicationService::new(update_context)),
85        ));
86        let (quality_observer, telemetry_guard, telemetry_writer, quality_telemetry_error) =
87            compose_quality_telemetry(data_dir);
88
89        // Read back rather than trusted from the request: what the directory
90        // says it is, is what we are on.
91        let engine = EmbeddedKernelStore::engine_of(data_dir)?;
92
93        Ok(Self {
94            data_dir: data_dir.to_path_buf(),
95            engine,
96            store,
97            service,
98            quality_observer,
99            telemetry_guard,
100            telemetry_writer,
101            quality_telemetry_error,
102        })
103    }
104
105    /// The engine behind this kernel's store, as its data directory records
106    /// it. Surfaced at startup and by the doctor so a user can see which one
107    /// they are on.
108    pub fn engine(&self) -> StorageEngine {
109        self.engine
110    }
111
112    pub fn data_dir(&self) -> &Path {
113        &self.data_dir
114    }
115
116    pub fn store(&self) -> &EmbeddedKernelStore {
117        &self.store
118    }
119
120    pub fn service(&self) -> Arc<EmbeddedMemoryService> {
121        Arc::clone(&self.service)
122    }
123
124    pub fn quality_observer(&self) -> Arc<dyn QualityMetricsObserver> {
125        self.quality_observer.clone()
126    }
127
128    pub fn quality_telemetry_dropped_observations(&self) -> u64 {
129        self.quality_observer.dropped_observations()
130    }
131
132    pub fn quality_telemetry_write_failures(&self) -> u64 {
133        self.telemetry_writer
134            .as_ref()
135            .map_or(0, |writer| writer.write_failures())
136    }
137
138    pub fn quality_telemetry_error(&self) -> Option<&str> {
139        self.quality_telemetry_error.as_deref()
140    }
141
142    pub fn quality_telemetry_active(&self) -> bool {
143        self.telemetry_guard.is_some()
144    }
145
146    /// Live query side for the same shareable SQLite journal used by the
147    /// quality observer.
148    pub fn quality_telemetry_reader(&self) -> Option<SqliteQualityTelemetryReader> {
149        self.telemetry_writer.as_ref().map(|writer| writer.reader())
150    }
151}
152
153fn compose_quality_telemetry(
154    data_dir: &Path,
155) -> (
156    Arc<BufferedQualityMetricsObserver>,
157    Option<EmbeddedTelemetryGuard>,
158    Option<Arc<SqliteQualityTelemetryWriter>>,
159    Option<String>,
160) {
161    let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(1_024);
162    let observer = Arc::new(observer);
163    let writer =
164        match SqliteQualityTelemetryWriter::open(data_dir, QualityTelemetryRetention::default()) {
165            Ok(writer) => Arc::new(writer),
166            Err(error) => {
167                drop(receiver);
168                let raw = error.to_string();
169                let reason = if raw.contains("Cannot acquire lock")
170                    || raw.to_ascii_lowercase().contains("already open")
171                {
172                    format!("the store's quality telemetry is held by another process ({raw})")
173                } else {
174                    raw
175                };
176                return (observer, None, None, Some(reason));
177            }
178        };
179    let batch_writer = Arc::clone(&writer);
180    let final_writer = Arc::clone(&writer);
181    let guard = match EmbeddedTelemetryGuard::try_spawn(
182        receiver,
183        64,
184        Duration::from_millis(250),
185        move |batch| {
186            let _ = batch_writer.write_batch(&batch);
187        },
188        move || {
189            let _ = final_writer.flush_durable();
190        },
191    ) {
192        Ok(guard) => guard,
193        Err(error) => {
194            return (
195                observer,
196                None,
197                Some(writer),
198                Some(format!("quality telemetry worker could not start: {error}")),
199            );
200        }
201    };
202    (observer, Some(guard), Some(writer), None)
203}
204
205#[cfg(test)]
206mod tests {
207    use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
208
209    use super::EmbeddedKernel;
210
211    #[test]
212    fn telemetry_startup_failure_never_prevents_the_kernel_from_opening() {
213        let data_dir = tempfile::tempdir().expect("temp data dir");
214        std::fs::write(
215            data_dir.path().join("telemetry"),
216            b"blocks directory creation",
217        )
218        .expect("blocking file");
219
220        let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel remains available");
221        assert!(!kernel.quality_telemetry_active());
222        assert!(kernel.quality_telemetry_error().is_some());
223        let metrics = BundleQualityMetrics::new(1, 1.0, 0.0, 0.0, 0.0).expect("valid metrics");
224        kernel.quality_observer.observe(
225            &metrics,
226            &QualityObservationContext {
227                rpc: "kmp_wake".to_string(),
228                root_node_id: "question:fail-open".to_string(),
229                role: "resumer".to_string(),
230                revision: Some(1),
231            },
232        );
233        assert_eq!(kernel.quality_telemetry_dropped_observations(), 1);
234        assert_eq!(kernel.quality_telemetry_write_failures(), 0);
235    }
236
237    #[test]
238    fn live_quality_reader_shares_the_kernel_journal() {
239        let data_dir = tempfile::tempdir().expect("temp data dir");
240        let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel");
241        let reader = kernel
242            .quality_telemetry_reader()
243            .expect("active telemetry exposes its read side");
244
245        assert_eq!(reader.count().expect("shared journal is readable"), 0);
246    }
247
248    #[test]
249    fn two_telemetry_compositions_share_the_sqlite_journal() {
250        let data_dir = tempfile::tempdir().expect("temp data dir");
251        let first = super::compose_quality_telemetry(data_dir.path());
252        assert!(first.2.is_some(), "the first writer opens the journal");
253        let second = super::compose_quality_telemetry(data_dir.path());
254        assert!(
255            second.2.is_some(),
256            "the second writer opens the same journal"
257        );
258        assert!(
259            second.3.is_none(),
260            "SQLite does not report process ownership"
261        );
262    }
263}