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,
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    store: EmbeddedKernelStore,
33    service: Arc<EmbeddedMemoryService>,
34    quality_observer: Arc<BufferedQualityMetricsObserver>,
35    telemetry_guard: Option<EmbeddedTelemetryGuard>,
36    telemetry_writer: Option<Arc<RedbQualityTelemetryWriter>>,
37    quality_telemetry_error: Option<String>,
38}
39
40impl EmbeddedKernel {
41    /// Opens the store at `data_dir` (fail-fast per ADR-012) and composes the
42    /// kernel. A second session on the same data dir fails here with an
43    /// explicit single-writer error (ADR-011; the engine holds the lock).
44    pub fn open(data_dir: &Path) -> Result<Self, PortError> {
45        let store = EmbeddedKernelStore::open(data_dir).map_err(|error| match error {
46            PortError::Unavailable(message) if message.contains("could not open") => {
47                PortError::Unavailable(format!(
48                    "{message}; if another agent session is using this data dir, close it first \
49                     (the embedded store is single-writer per ADR-011)"
50                ))
51            }
52            other => other,
53        })?;
54
55        let graph = Arc::new(store.clone());
56        let detail = Arc::new(store.clone());
57        let query_application = Arc::new(QueryApplicationService::new(
58            Arc::clone(&graph),
59            Arc::clone(&detail),
60            Arc::new(store.clone()),
61            GENERATOR_VERSION,
62        ));
63        let update_context = Arc::new(UpdateContextUseCase::new_with_projection_writer(
64            Arc::new(store.clone()),
65            RoutingProjectionWriter::new(graph, detail),
66            GENERATOR_VERSION,
67        ));
68        let service = Arc::new(KernelMemoryApplicationService::new(
69            query_application,
70            Arc::new(CommandApplicationService::new(update_context)),
71        ));
72        let (quality_observer, telemetry_guard, telemetry_writer, quality_telemetry_error) =
73            compose_quality_telemetry(data_dir);
74
75        Ok(Self {
76            data_dir: data_dir.to_path_buf(),
77            store,
78            service,
79            quality_observer,
80            telemetry_guard,
81            telemetry_writer,
82            quality_telemetry_error,
83        })
84    }
85
86    pub fn data_dir(&self) -> &Path {
87        &self.data_dir
88    }
89
90    pub fn store(&self) -> &EmbeddedKernelStore {
91        &self.store
92    }
93
94    pub fn service(&self) -> Arc<EmbeddedMemoryService> {
95        Arc::clone(&self.service)
96    }
97
98    pub fn quality_observer(&self) -> Arc<dyn QualityMetricsObserver> {
99        self.quality_observer.clone()
100    }
101
102    pub fn quality_telemetry_dropped_observations(&self) -> u64 {
103        self.quality_observer.dropped_observations()
104    }
105
106    pub fn quality_telemetry_write_failures(&self) -> u64 {
107        self.telemetry_writer
108            .as_ref()
109            .map_or(0, |writer| writer.write_failures())
110    }
111
112    pub fn quality_telemetry_error(&self) -> Option<&str> {
113        self.quality_telemetry_error.as_deref()
114    }
115
116    pub fn quality_telemetry_active(&self) -> bool {
117        self.telemetry_guard.is_some()
118    }
119}
120
121fn compose_quality_telemetry(
122    data_dir: &Path,
123) -> (
124    Arc<BufferedQualityMetricsObserver>,
125    Option<EmbeddedTelemetryGuard>,
126    Option<Arc<RedbQualityTelemetryWriter>>,
127    Option<String>,
128) {
129    let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(1_024);
130    let observer = Arc::new(observer);
131    let writer =
132        match RedbQualityTelemetryWriter::open(data_dir, QualityTelemetryRetention::default()) {
133            Ok(writer) => Arc::new(writer),
134            Err(error) => {
135                drop(receiver);
136                return (observer, None, None, Some(error.to_string()));
137            }
138        };
139    let batch_writer = Arc::clone(&writer);
140    let final_writer = Arc::clone(&writer);
141    let guard = match EmbeddedTelemetryGuard::try_spawn(
142        receiver,
143        64,
144        Duration::from_millis(250),
145        move |batch| {
146            let _ = batch_writer.write_batch(&batch);
147        },
148        move || {
149            let _ = final_writer.flush_durable();
150        },
151    ) {
152        Ok(guard) => guard,
153        Err(error) => {
154            return (
155                observer,
156                None,
157                Some(writer),
158                Some(format!("quality telemetry worker could not start: {error}")),
159            );
160        }
161    };
162    (observer, Some(guard), Some(writer), None)
163}
164
165#[cfg(test)]
166mod tests {
167    use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
168
169    use super::EmbeddedKernel;
170
171    #[test]
172    fn telemetry_startup_failure_never_prevents_the_kernel_from_opening() {
173        let data_dir = tempfile::tempdir().expect("temp data dir");
174        std::fs::write(
175            data_dir.path().join("telemetry"),
176            b"blocks directory creation",
177        )
178        .expect("blocking file");
179
180        let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel remains available");
181        assert!(!kernel.quality_telemetry_active());
182        assert!(kernel.quality_telemetry_error().is_some());
183        let metrics = BundleQualityMetrics::new(1, 1.0, 0.0, 0.0, 0.0).expect("valid metrics");
184        kernel.quality_observer.observe(
185            &metrics,
186            &QualityObservationContext {
187                rpc: "kernel_wake".to_string(),
188                root_node_id: "question:fail-open".to_string(),
189                role: "resumer".to_string(),
190            },
191        );
192        assert_eq!(kernel.quality_telemetry_dropped_observations(), 1);
193        assert_eq!(kernel.quality_telemetry_write_failures(), 0);
194    }
195}