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