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            // The versioned memory API names no instant yet; the packet
45            // stands on the memory's frontier.
46            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            // The versioned memory API does not carry the user's words yet;
59            // the kernel searches the question as given and warns on nothing.
60            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    /// The kernel accounts for its own recalls.
77    ///
78    /// Observed here, at the contract, rather than left to the consumer:
79    /// a consumer that forgot would silently starve the kernel's quality
80    /// telemetry, and each consumer remembering is the same code written N
81    /// times. The rpc names keep the spelling the telemetry has always
82    /// carried.
83    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
209/// The record side's error map.
210///
211/// Two cases more than the recall side's. A port `Conflict` here is an
212/// idempotency key reused with different content and will not change on
213/// retry. A `RetryableConflict` is optimistic concurrency: this attempt did
214/// not land, and replaying the same logical write under the same key is safe.
215fn 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        // A port failing is the storage or the environment, not the request.
328        unavailable @ ApplicationError::Ports(_) => ApiError::Unavailable {
329            reason: unavailable.to_string(),
330        },
331    }
332}