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