1use std::path::{Path, PathBuf};
2use std::sync::Arc;
3use std::time::Duration;
4
5use kmp_adapter_embedded::{
6 EmbeddedKernelStore, QualityTelemetryRetention, RedbQualityTelemetryWriter, StorageEngine,
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
17pub type EmbeddedMemoryService = KernelMemoryApplicationService<
21 EmbeddedKernelStore,
22 EmbeddedKernelStore,
23 EmbeddedKernelStore,
24 EmbeddedKernelStore,
25 RoutingProjectionWriter<Arc<EmbeddedKernelStore>, Arc<EmbeddedKernelStore>>,
26>;
27
28pub struct EmbeddedKernel {
31 data_dir: PathBuf,
32 engine: StorageEngine,
33 store: EmbeddedKernelStore,
34 service: Arc<EmbeddedMemoryService>,
35 quality_observer: Arc<BufferedQualityMetricsObserver>,
36 telemetry_guard: Option<EmbeddedTelemetryGuard>,
37 telemetry_writer: Option<Arc<RedbQualityTelemetryWriter>>,
38 quality_telemetry_error: Option<String>,
39}
40
41impl EmbeddedKernel {
42 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.map_err(|error| match error {
65 PortError::Unavailable(message) if message.contains("Cannot acquire lock") => {
68 PortError::Unavailable(format!(
69 "{message}; if another agent session is using this data dir, close it first \
70 (the redb engine is single-writer per ADR-011 — to share one store between \
71 hosts, migrate it to the sqlite engine: `kmp-mcp migrate <this-dir> \
72 <new-dir> --engine sqlite`)"
73 ))
74 }
75 other => other,
76 })?;
77
78 let graph = Arc::new(store.clone());
79 let detail = Arc::new(store.clone());
80 let query_application = Arc::new(QueryApplicationService::new(
81 Arc::clone(&graph),
82 Arc::clone(&detail),
83 Arc::new(store.clone()),
84 GENERATOR_VERSION,
85 ));
86 let update_context = Arc::new(UpdateContextUseCase::new_with_projection_writer(
87 Arc::new(store.clone()),
88 RoutingProjectionWriter::new(graph, detail),
89 GENERATOR_VERSION,
90 ));
91 let service = Arc::new(KernelMemoryApplicationService::new(
92 query_application,
93 Arc::new(CommandApplicationService::new(update_context)),
94 ));
95 let (quality_observer, telemetry_guard, telemetry_writer, quality_telemetry_error) =
96 compose_quality_telemetry(data_dir);
97
98 let engine = EmbeddedKernelStore::engine_of(data_dir)?;
101
102 Ok(Self {
103 data_dir: data_dir.to_path_buf(),
104 engine,
105 store,
106 service,
107 quality_observer,
108 telemetry_guard,
109 telemetry_writer,
110 quality_telemetry_error,
111 })
112 }
113
114 pub fn engine(&self) -> StorageEngine {
118 self.engine
119 }
120
121 pub fn data_dir(&self) -> &Path {
122 &self.data_dir
123 }
124
125 pub fn store(&self) -> &EmbeddedKernelStore {
126 &self.store
127 }
128
129 pub fn service(&self) -> Arc<EmbeddedMemoryService> {
130 Arc::clone(&self.service)
131 }
132
133 pub fn quality_observer(&self) -> Arc<dyn QualityMetricsObserver> {
134 self.quality_observer.clone()
135 }
136
137 pub fn quality_telemetry_dropped_observations(&self) -> u64 {
138 self.quality_observer.dropped_observations()
139 }
140
141 pub fn quality_telemetry_write_failures(&self) -> u64 {
142 self.telemetry_writer
143 .as_ref()
144 .map_or(0, |writer| writer.write_failures())
145 }
146
147 pub fn quality_telemetry_error(&self) -> Option<&str> {
148 self.quality_telemetry_error.as_deref()
149 }
150
151 pub fn quality_telemetry_active(&self) -> bool {
152 self.telemetry_guard.is_some()
153 }
154}
155
156fn compose_quality_telemetry(
157 data_dir: &Path,
158) -> (
159 Arc<BufferedQualityMetricsObserver>,
160 Option<EmbeddedTelemetryGuard>,
161 Option<Arc<RedbQualityTelemetryWriter>>,
162 Option<String>,
163) {
164 let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(1_024);
165 let observer = Arc::new(observer);
166 let writer =
167 match RedbQualityTelemetryWriter::open(data_dir, QualityTelemetryRetention::default()) {
168 Ok(writer) => Arc::new(writer),
169 Err(error) => {
170 drop(receiver);
171 return (observer, None, None, Some(error.to_string()));
172 }
173 };
174 let batch_writer = Arc::clone(&writer);
175 let final_writer = Arc::clone(&writer);
176 let guard = match EmbeddedTelemetryGuard::try_spawn(
177 receiver,
178 64,
179 Duration::from_millis(250),
180 move |batch| {
181 let _ = batch_writer.write_batch(&batch);
182 },
183 move || {
184 let _ = final_writer.flush_durable();
185 },
186 ) {
187 Ok(guard) => guard,
188 Err(error) => {
189 return (
190 observer,
191 None,
192 Some(writer),
193 Some(format!("quality telemetry worker could not start: {error}")),
194 );
195 }
196 };
197 (observer, Some(guard), Some(writer), None)
198}
199
200#[cfg(test)]
201mod tests {
202 use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
203
204 use super::EmbeddedKernel;
205
206 #[test]
207 fn telemetry_startup_failure_never_prevents_the_kernel_from_opening() {
208 let data_dir = tempfile::tempdir().expect("temp data dir");
209 std::fs::write(
210 data_dir.path().join("telemetry"),
211 b"blocks directory creation",
212 )
213 .expect("blocking file");
214
215 let kernel = EmbeddedKernel::open(data_dir.path()).expect("kernel remains available");
216 assert!(!kernel.quality_telemetry_active());
217 assert!(kernel.quality_telemetry_error().is_some());
218 let metrics = BundleQualityMetrics::new(1, 1.0, 0.0, 0.0, 0.0).expect("valid metrics");
219 kernel.quality_observer.observe(
220 &metrics,
221 &QualityObservationContext {
222 rpc: "kernel_wake".to_string(),
223 root_node_id: "question:fail-open".to_string(),
224 role: "resumer".to_string(),
225 },
226 );
227 assert_eq!(kernel.quality_telemetry_dropped_observations(), 1);
228 assert_eq!(kernel.quality_telemetry_write_failures(), 0);
229 }
230}