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
18pub type EmbeddedMemoryService = KernelMemoryApplicationService<
22 EmbeddedKernelStore,
23 EmbeddedKernelStore,
24 EmbeddedKernelStore,
25 EmbeddedKernelStore,
26 EmbeddedKernelStore,
27>;
28
29pub 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 pub fn open(data_dir: &Path) -> Result<Self, PortError> {
48 Self::open_with_engine(data_dir, None)
49 }
50
51 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 .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 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 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 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}