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