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(
83            KernelMemoryApplicationService::new(
84                query_application,
85                Arc::new(CommandApplicationService::new(update_context)),
86            )
87            // The derived card view lives on the same store as the bodies it
88            // describes, so a card is written and read inside one snapshot.
89            .with_node_cards(Arc::new(store.clone())),
90        );
91        let (quality_observer, telemetry_guard, telemetry_writer, quality_telemetry_error) =
92            compose_quality_telemetry(data_dir);
93
94        // Read back rather than trusted from the request: what the directory
95        // says it is, is what we are on.
96        let engine = EmbeddedKernelStore::engine_of(data_dir)?;
97
98        Ok(Self {
99            data_dir: data_dir.to_path_buf(),
100            engine,
101            store,
102            service,
103            quality_observer,
104            telemetry_guard,
105            telemetry_writer,
106            quality_telemetry_error,
107        })
108    }
109
110    /// The engine behind this kernel's store, as its data directory records
111    /// it. Surfaced at startup and by the doctor so a user can see which one
112    /// they are on.
113    pub fn engine(&self) -> StorageEngine {
114        self.engine
115    }
116
117    pub fn data_dir(&self) -> &Path {
118        &self.data_dir
119    }
120
121    pub fn store(&self) -> &EmbeddedKernelStore {
122        &self.store
123    }
124
125    pub fn service(&self) -> Arc<EmbeddedMemoryService> {
126        Arc::clone(&self.service)
127    }
128
129    pub fn quality_observer(&self) -> Arc<dyn QualityMetricsObserver> {
130        self.quality_observer.clone()
131    }
132
133    pub fn quality_telemetry_dropped_observations(&self) -> u64 {
134        self.quality_observer.dropped_observations()
135    }
136
137    pub fn quality_telemetry_write_failures(&self) -> u64 {
138        self.telemetry_writer
139            .as_ref()
140            .map_or(0, |writer| writer.write_failures())
141    }
142
143    pub fn quality_telemetry_error(&self) -> Option<&str> {
144        self.quality_telemetry_error.as_deref()
145    }
146
147    pub fn quality_telemetry_active(&self) -> bool {
148        self.telemetry_guard.is_some()
149    }
150
151    /// Live query side for the same shareable SQLite journal used by the
152    /// quality observer.
153    pub fn quality_telemetry_reader(&self) -> Option<SqliteQualityTelemetryReader> {
154        self.telemetry_writer.as_ref().map(|writer| writer.reader())
155    }
156}
157
158fn compose_quality_telemetry(
159    data_dir: &Path,
160) -> (
161    Arc<BufferedQualityMetricsObserver>,
162    Option<EmbeddedTelemetryGuard>,
163    Option<Arc<SqliteQualityTelemetryWriter>>,
164    Option<String>,
165) {
166    let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(1_024);
167    let observer = Arc::new(observer);
168    let writer =
169        match SqliteQualityTelemetryWriter::open(data_dir, QualityTelemetryRetention::default()) {
170            Ok(writer) => Arc::new(writer),
171            Err(error) => {
172                drop(receiver);
173                let raw = error.to_string();
174                let reason = if raw.contains("Cannot acquire lock")
175                    || raw.to_ascii_lowercase().contains("already open")
176                {
177                    format!("the store's quality telemetry is held by another process ({raw})")
178                } else {
179                    raw
180                };
181                return (observer, None, None, Some(reason));
182            }
183        };
184    let batch_writer = Arc::clone(&writer);
185    let final_writer = Arc::clone(&writer);
186    let guard = match EmbeddedTelemetryGuard::try_spawn(
187        receiver,
188        64,
189        Duration::from_millis(250),
190        move |batch| {
191            let _ = batch_writer.write_batch(&batch);
192        },
193        move || {
194            let _ = final_writer.flush_durable();
195        },
196    ) {
197        Ok(guard) => guard,
198        Err(error) => {
199            return (
200                observer,
201                None,
202                Some(writer),
203                Some(format!("quality telemetry worker could not start: {error}")),
204            );
205        }
206    };
207    (observer, Some(guard), Some(writer), None)
208}
209
210#[cfg(test)]
211mod tests {
212    use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
213
214    use super::EmbeddedKernel;
215
216    #[test]
217    fn telemetry_startup_failure_never_prevents_the_kernel_from_opening() {
218        let data_dir = tempfile::tempdir().expect("temp data dir");
219        std::fs::write(
220            data_dir.path().join("telemetry"),
221            b"blocks directory creation",
222        )
223        .expect("blocking file");
224
225        let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel remains available");
226        assert!(!kernel.quality_telemetry_active());
227        assert!(kernel.quality_telemetry_error().is_some());
228        let metrics = BundleQualityMetrics::new(1, 1.0, 0.0, 0.0, 0.0).expect("valid metrics");
229        kernel.quality_observer.observe(
230            &metrics,
231            &QualityObservationContext {
232                rpc: "kmp_wake".to_string(),
233                root_node_id: "question:fail-open".to_string(),
234                role: "resumer".to_string(),
235                revision: Some(1),
236            },
237        );
238        assert_eq!(kernel.quality_telemetry_dropped_observations(), 1);
239        assert_eq!(kernel.quality_telemetry_write_failures(), 0);
240    }
241
242    #[test]
243    fn live_quality_reader_shares_the_kernel_journal() {
244        let data_dir = tempfile::tempdir().expect("temp data dir");
245        let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel");
246        let reader = kernel
247            .quality_telemetry_reader()
248            .expect("active telemetry exposes its read side");
249
250        assert_eq!(reader.count().expect("shared journal is readable"), 0);
251    }
252
253    #[test]
254    fn two_telemetry_compositions_share_the_sqlite_journal() {
255        let data_dir = tempfile::tempdir().expect("temp data dir");
256        let first = super::compose_quality_telemetry(data_dir.path());
257        assert!(first.2.is_some(), "the first writer opens the journal");
258        let second = super::compose_quality_telemetry(data_dir.path());
259        assert!(
260            second.2.is_some(),
261            "the second writer opens the same journal"
262        );
263        assert!(
264            second.3.is_none(),
265            "SQLite does not report process ownership"
266        );
267    }
268}