1use kmp_application::{
8 ApplicationError, AskMemoryQuery, GetContextResult, MemoryAnswerPolicy as DomainAnswerPolicy,
9 MemoryCoordinateData, MemoryData, MemoryDimensionData, MemoryEntryData, MemoryEvidenceData,
10 MemoryIngestCommand, MemoryProvenanceData, WakeMemoryQuery,
11};
12use kmp_domain::{DimensionSelection, PortError, ResolutionTier};
13use kmp_memory_api::{
14 ApiCapabilities, ApiError, CONTRACT_VERSION, MemoryAnswerPolicy, MemoryAskRequest,
15 MemoryDetailView, MemoryNodeView, MemoryRecallApi, MemoryRecallView, MemoryRecordApi,
16 MemoryRecordRequest, MemoryRelationshipView, MemoryTier, MemoryWakeRequest, RecordedMemoryView,
17 RenderedMemoryView,
18};
19
20use crate::EmbeddedKernel;
21
22const CAPABILITIES: [&str; 3] = ["wake", "ask", "record"];
27
28impl MemoryRecallApi for EmbeddedKernel {
29 fn capabilities(&self) -> ApiCapabilities {
30 ApiCapabilities::new(CONTRACT_VERSION, env!("CARGO_PKG_VERSION"), CAPABILITIES)
31 }
32
33 async fn wake(&self, request: MemoryWakeRequest) -> Result<MemoryRecallView, ApiError> {
34 let about = request.about.clone();
35 let query = WakeMemoryQuery {
36 about: request.about,
37 role: request.role,
38 intent: request.intent,
39 dimensions: dimensions(request.dimension_kinds, request.scoped_to_about),
40 token_budget: request.token_budget,
41 depth: request.depth,
42 max_tier: request.max_tier.map(tier),
43 max_entries: request.max_entries.map(|entries| entries as usize),
44 temporal: kmp_domain::TemporalSelection::Frontier,
47 };
48 let result = self.service().wake(query).await.map_err(translate_error)?;
49 self.observe_recall_quality("kmp_wake", &result);
50 Ok(recall_view(about, &result))
51 }
52
53 async fn ask(&self, request: MemoryAskRequest) -> Result<MemoryRecallView, ApiError> {
54 let about = request.about.clone();
55 let query = AskMemoryQuery {
56 about: request.about,
57 question: request.question,
58 asked_as: None,
61 answer_policy: answer_policy(request.answer_policy),
62 dimensions: dimensions(request.dimension_kinds, request.scoped_to_about),
63 token_budget: request.token_budget,
64 depth: request.depth,
65 max_tier: request.max_tier.map(tier),
66 max_entries: None,
67 temporal: kmp_domain::TemporalSelection::Frontier,
68 };
69 let result = self.service().ask(query).await.map_err(translate_error)?;
70 self.observe_recall_quality("kmp_ask", &result);
71 Ok(recall_view(about, &result))
72 }
73}
74
75impl EmbeddedKernel {
76 fn observe_recall_quality(&self, rpc: &str, result: &GetContextResult) {
84 self.quality_observer().observe(
85 &result.rendered.quality,
86 &kmp_domain::QualityObservationContext {
87 rpc: rpc.to_owned(),
88 root_node_id: result.bundle.root_node_id().as_str().to_owned(),
89 role: result.bundle.role().as_str().to_owned(),
90 revision: Some(result.bundle.metadata().revision),
91 },
92 );
93 }
94}
95
96impl MemoryRecordApi for EmbeddedKernel {
97 fn capabilities(&self) -> ApiCapabilities {
98 ApiCapabilities::new(CONTRACT_VERSION, env!("CARGO_PKG_VERSION"), CAPABILITIES)
99 }
100
101 async fn record(&self, request: MemoryRecordRequest) -> Result<RecordedMemoryView, ApiError> {
102 let outcome = self
103 .service()
104 .ingest(ingest_command(request))
105 .await
106 .map_err(translate_record_error)?;
107 Ok(RecordedMemoryView {
108 about: outcome.about,
109 memory_id: outcome.memory_id,
110 accepted_entries: outcome.accepted.entries,
111 accepted_relations: outcome.accepted.relations,
112 accepted_evidence: outcome.accepted.evidence,
113 read_after_write_ready: outcome.read_after_write_ready,
114 warnings: outcome.warnings,
115 })
116 }
117}
118
119fn ingest_command(request: MemoryRecordRequest) -> MemoryIngestCommand {
120 MemoryIngestCommand {
121 about: request.about,
122 memory: MemoryData {
123 dimensions: request
124 .dimensions
125 .into_iter()
126 .map(|dimension| MemoryDimensionData {
127 id: dimension.id,
128 kind: dimension.kind,
129 title: dimension.title,
130 metadata: dimension.metadata,
131 })
132 .collect(),
133 entries: request
134 .entries
135 .into_iter()
136 .map(|entry| MemoryEntryData {
137 id: entry.id,
138 kind: entry.kind,
139 text: entry.text,
140 coordinates: entry
141 .coordinates
142 .into_iter()
143 .map(|coordinate| MemoryCoordinateData {
144 dimension: coordinate.dimension,
145 scope_id: coordinate.scope_id,
146 occurred_at: coordinate.occurred_at,
147 observed_at: None,
148 ingested_at: None,
149 valid_from: None,
150 valid_until: None,
151 sequence: coordinate.sequence,
152 rank: coordinate.rank,
153 metadata: Default::default(),
154 })
155 .collect(),
156 metadata: entry.metadata,
157 })
158 .collect(),
159 relations: request
160 .relations
161 .into_iter()
162 .map(|relation| kmp_application::MemoryRelationData {
163 source_ref: relation.from,
164 target_ref: relation.to,
165 rel: relation.rel,
166 semantic_class: relation.semantic_class,
167 why: relation.why,
168 evidence: None,
169 confidence: relation.confidence,
170 sequence: relation.sequence,
171 motivation: None,
172 method: None,
173 decision_id: None,
174 caused_by_node_id: None,
175 coordinate: None,
176 })
177 .collect(),
178 evidence: request
179 .evidence
180 .into_iter()
181 .map(|evidence| MemoryEvidenceData {
182 id: evidence.id,
183 supports: evidence.supports,
184 text: evidence.text,
185 source: evidence.source,
186 time: evidence.time,
187 metadata: evidence.metadata,
188 })
189 .collect(),
190 },
191 provenance: request.provenance.map(|provenance| MemoryProvenanceData {
192 source_kind: provenance.source_kind,
193 source_agent: provenance.source_agent,
194 observed_at: provenance.observed_at,
195 correlation_id: provenance.correlation_id,
196 causation_id: provenance.causation_id,
197 }),
198 idempotency_key: request.idempotency_key,
199 dry_run: false,
200 label_policy: Default::default(),
201 }
202}
203
204fn translate_record_error(error: ApplicationError) -> ApiError {
211 match error {
212 ApplicationError::Ports(PortError::Conflict(reason)) => ApiError::Refused { reason },
213 ApplicationError::RetryableConflict(reason) => ApiError::Refused {
214 reason: format!(
215 "retryable write conflict; rebase and retry the same logical write with the same idempotency_key — replay is safe: {reason}"
216 ),
217 },
218 other => translate_error(other),
219 }
220}
221
222fn dimensions(kinds: Vec<String>, scoped_to_about: bool) -> DimensionSelection {
223 let selection = if kinds.is_empty() {
224 DimensionSelection::all()
225 } else {
226 DimensionSelection::only(kinds)
227 };
228 if scoped_to_about {
229 selection.with_current_about_scope()
230 } else {
231 selection
232 }
233}
234
235fn tier(tier: MemoryTier) -> ResolutionTier {
236 match tier {
237 MemoryTier::Summary => ResolutionTier::L0Summary,
238 MemoryTier::CausalSpine => ResolutionTier::L1CausalSpine,
239 MemoryTier::EvidencePack => ResolutionTier::L2EvidencePack,
240 }
241}
242
243fn answer_policy(policy: MemoryAnswerPolicy) -> DomainAnswerPolicy {
244 match policy {
245 MemoryAnswerPolicy::EvidenceOrUnknown => DomainAnswerPolicy::EvidenceOrUnknown,
246 MemoryAnswerPolicy::ShowConflicts => DomainAnswerPolicy::ShowConflicts,
247 MemoryAnswerPolicy::BestEffort => DomainAnswerPolicy::BestEffort,
248 }
249}
250
251fn recall_view(about: String, result: &GetContextResult) -> MemoryRecallView {
252 let bundle = &result.bundle;
253 MemoryRecallView {
254 about,
255 revision: bundle.metadata().revision,
256 content_hash: bundle.metadata().content_hash.clone(),
257 root: node_view(bundle.root_node()),
258 neighbors: bundle.neighbor_nodes().iter().map(node_view).collect(),
259 relationships: bundle
260 .relationships()
261 .iter()
262 .map(|relationship| MemoryRelationshipView {
263 source_node_id: relationship.source_node_id().to_string(),
264 target_node_id: relationship.target_node_id().to_string(),
265 relationship_type: relationship.relationship_type().to_string(),
266 why: relationship
267 .explanation()
268 .rationale()
269 .map(ToOwned::to_owned),
270 evidence: relationship.explanation().evidence().map(ToOwned::to_owned),
271 })
272 .collect(),
273 details: bundle
274 .node_details()
275 .iter()
276 .map(|detail| MemoryDetailView {
277 node_id: detail.node_id().to_string(),
278 detail: detail.detail().to_string(),
279 content_hash: detail.content_hash().to_string(),
280 revision: detail.revision(),
281 })
282 .collect(),
283 rendered: RenderedMemoryView {
284 content: result.rendered.content.clone(),
285 content_hash: result.rendered.content_hash.clone(),
286 token_count: result.rendered.token_count,
287 quality: kmp_memory_api::MemoryQualityView {
288 raw_equivalent_tokens: result.rendered.quality.raw_equivalent_tokens(),
289 compression_ratio: result.rendered.quality.compression_ratio(),
290 causal_density: result.rendered.quality.causal_density(),
291 noise_ratio: result.rendered.quality.noise_ratio(),
292 detail_coverage: result.rendered.quality.detail_coverage(),
293 },
294 },
295 }
296}
297
298fn node_view(node: &kmp_domain::BundleNode) -> MemoryNodeView {
299 MemoryNodeView {
300 node_id: node.node_id().to_string(),
301 node_kind: node.node_kind().to_string(),
302 title: node.title().to_string(),
303 summary: node.summary().to_string(),
304 status: node.status().to_string(),
305 labels: node.labels().to_vec(),
306 properties: node.properties().clone(),
307 }
308}
309
310fn translate_error(error: ApplicationError) -> ApiError {
311 match error {
312 ApplicationError::NotFound(what) => ApiError::NotFound { what },
313 ApplicationError::Validation(reason) => ApiError::Refused { reason },
314 ApplicationError::RetryableConflict(reason) => ApiError::Refused {
315 reason: format!(
316 "retryable write conflict; retry the same logical write with the same idempotency_key: {reason}"
317 ),
318 },
319 refused @ ApplicationError::Domain(_) => ApiError::Refused {
320 reason: refused.to_string(),
321 },
322 unavailable @ ApplicationError::Ports(_) => ApiError::Unavailable {
324 reason: unavailable.to_string(),
325 },
326 }
327}