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 receipt_context: None,
122 about: request.about,
123 memory: MemoryData {
124 dimensions: request
125 .dimensions
126 .into_iter()
127 .map(|dimension| MemoryDimensionData {
128 id: dimension.id,
129 kind: dimension.kind,
130 title: dimension.title,
131 metadata: dimension.metadata,
132 })
133 .collect(),
134 entries: request
135 .entries
136 .into_iter()
137 .map(|entry| MemoryEntryData {
138 id: entry.id,
139 kind: entry.kind,
140 text: entry.text,
141 coordinates: entry
142 .coordinates
143 .into_iter()
144 .map(|coordinate| MemoryCoordinateData {
145 dimension: coordinate.dimension,
146 scope_id: coordinate.scope_id,
147 occurred_at: coordinate.occurred_at,
148 observed_at: None,
149 ingested_at: None,
150 valid_from: None,
151 valid_until: None,
152 sequence: coordinate.sequence,
153 rank: coordinate.rank,
154 metadata: Default::default(),
155 })
156 .collect(),
157 metadata: entry.metadata,
158 })
159 .collect(),
160 relations: request
161 .relations
162 .into_iter()
163 .map(|relation| kmp_application::MemoryRelationData {
164 source_ref: relation.from,
165 target_ref: relation.to,
166 rel: relation.rel,
167 semantic_class: relation.semantic_class,
168 why: relation.why,
169 evidence: None,
170 confidence: relation.confidence,
171 sequence: relation.sequence,
172 motivation: None,
173 method: None,
174 decision_id: None,
175 caused_by_node_id: None,
176 coordinate: None,
177 })
178 .collect(),
179 evidence: request
180 .evidence
181 .into_iter()
182 .map(|evidence| MemoryEvidenceData {
183 id: evidence.id,
184 supports: evidence.supports,
185 text: evidence.text,
186 source: evidence.source,
187 time: evidence.time,
188 metadata: evidence.metadata,
189 })
190 .collect(),
191 },
192 provenance: request.provenance.map(|provenance| MemoryProvenanceData {
193 source_kind: provenance.source_kind,
194 source_agent: provenance.source_agent,
195 observed_at: provenance.observed_at,
196 correlation_id: provenance.correlation_id,
197 causation_id: provenance.causation_id,
198 }),
199 idempotency_key: request.idempotency_key,
200 dry_run: false,
201 label_policy: Default::default(),
202 }
203}
204
205fn translate_record_error(error: ApplicationError) -> ApiError {
212 match error {
213 ApplicationError::Ports(PortError::Conflict(reason)) => ApiError::Refused { reason },
214 ApplicationError::RetryableConflict(reason) => ApiError::Refused {
215 reason: format!(
216 "retryable write conflict; rebase and retry the same logical write with the same idempotency_key — replay is safe: {reason}"
217 ),
218 },
219 other => translate_error(other),
220 }
221}
222
223fn dimensions(kinds: Vec<String>, scoped_to_about: bool) -> DimensionSelection {
224 let selection = if kinds.is_empty() {
225 DimensionSelection::all()
226 } else {
227 DimensionSelection::only(kinds)
228 };
229 if scoped_to_about {
230 selection.with_current_about_scope()
231 } else {
232 selection
233 }
234}
235
236fn tier(tier: MemoryTier) -> ResolutionTier {
237 match tier {
238 MemoryTier::Summary => ResolutionTier::L0Summary,
239 MemoryTier::CausalSpine => ResolutionTier::L1CausalSpine,
240 MemoryTier::EvidencePack => ResolutionTier::L2EvidencePack,
241 }
242}
243
244fn answer_policy(policy: MemoryAnswerPolicy) -> DomainAnswerPolicy {
245 match policy {
246 MemoryAnswerPolicy::EvidenceOrUnknown => DomainAnswerPolicy::EvidenceOrUnknown,
247 MemoryAnswerPolicy::ShowConflicts => DomainAnswerPolicy::ShowConflicts,
248 MemoryAnswerPolicy::BestEffort => DomainAnswerPolicy::BestEffort,
249 }
250}
251
252fn recall_view(about: String, result: &GetContextResult) -> MemoryRecallView {
253 let bundle = &result.bundle;
254 MemoryRecallView {
255 about,
256 revision: bundle.metadata().revision,
257 content_hash: bundle.metadata().content_hash.clone(),
258 root: node_view(bundle.root_node()),
259 neighbors: bundle.neighbor_nodes().iter().map(node_view).collect(),
260 relationships: bundle
261 .relationships()
262 .iter()
263 .map(|relationship| MemoryRelationshipView {
264 source_node_id: relationship.source_node_id().to_string(),
265 target_node_id: relationship.target_node_id().to_string(),
266 relationship_type: relationship.relationship_type().to_string(),
267 why: relationship
268 .explanation()
269 .rationale()
270 .map(ToOwned::to_owned),
271 evidence: relationship.explanation().evidence().map(ToOwned::to_owned),
272 })
273 .collect(),
274 details: bundle
275 .node_details()
276 .iter()
277 .map(|detail| MemoryDetailView {
278 node_id: detail.node_id().to_string(),
279 detail: detail.detail().to_string(),
280 content_hash: detail.content_hash().to_string(),
281 revision: detail.revision(),
282 })
283 .collect(),
284 rendered: RenderedMemoryView {
285 content: result.rendered.content.clone(),
286 content_hash: result.rendered.content_hash.clone(),
287 token_count: result.rendered.token_count,
288 quality: kmp_memory_api::MemoryQualityView {
289 raw_equivalent_tokens: result.rendered.quality.raw_equivalent_tokens(),
290 compression_ratio: result.rendered.quality.compression_ratio(),
291 causal_density: result.rendered.quality.causal_density(),
292 noise_ratio: result.rendered.quality.noise_ratio(),
293 detail_coverage: result.rendered.quality.detail_coverage(),
294 },
295 },
296 }
297}
298
299fn node_view(node: &kmp_domain::BundleNode) -> MemoryNodeView {
300 MemoryNodeView {
301 node_id: node.node_id().to_string(),
302 node_kind: node.node_kind().to_string(),
303 title: node.title().to_string(),
304 summary: node.summary().to_string(),
305 status: node.status().to_string(),
306 labels: node.labels().to_vec(),
307 properties: node.properties().clone(),
308 }
309}
310
311fn translate_error(error: ApplicationError) -> ApiError {
312 match error {
313 ApplicationError::NotFound(what) => ApiError::NotFound { what },
314 ApplicationError::Validation(reason) => ApiError::Refused { reason },
315 ApplicationError::RetryableConflict(reason) => ApiError::Refused {
316 reason: format!(
317 "retryable write conflict; retry the same logical write with the same idempotency_key: {reason}"
318 ),
319 },
320 refused @ ApplicationError::Domain(_) => ApiError::Refused {
321 reason: refused.to_string(),
322 },
323 unavailable @ ApplicationError::Ports(_) => ApiError::Unavailable {
325 reason: unavailable.to_string(),
326 },
327 }
328}