Skip to main content

kmp_embedded/
memory_api.rs

1//! [`EmbeddedKernel`] as an implementation of the published consumer contract.
2//!
3//! The conversions in this module are the whole of the coupling a consumer is
4//! allowed: domain aggregate in, plain view out. Nothing of the domain crosses
5//! the trait.
6
7use 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
22/// What this build can do, by name.
23///
24/// Listed next to the implementation, so that adding a method to the trait
25/// without adding its name is a diff a reviewer sees in one place.
26const 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            // The versioned memory API does not carry the user's words yet;
56            // the kernel searches the question as given and warns on nothing.
57            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    /// The kernel accounts for its own recalls.
73    ///
74    /// Observed here, at the contract, rather than left to the consumer:
75    /// a consumer that forgot would silently starve the kernel's quality
76    /// telemetry, and each consumer remembering is the same code written N
77    /// times. The rpc names keep the spelling the telemetry has always
78    /// carried.
79    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
199/// The record side's error map.
200///
201/// Two cases more than the recall side's. A port `Conflict` here is an
202/// idempotency key reused with different content and will not change on
203/// retry. A `RetryableConflict` is optimistic concurrency: this attempt did
204/// not land, and replaying the same logical write under the same key is safe.
205fn 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        // A port failing is the storage or the environment, not the request.
318        unavailable @ ApplicationError::Ports(_) => ApiError::Unavailable {
319            reason: unavailable.to_string(),
320        },
321    }
322}