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
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 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 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}