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