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("kernel_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("kernel_ask", &result);
64        Ok(recall_view(about, &result))
65    }
66}
67
68impl EmbeddedKernel {
69    /// The kernel accounts for its own recalls.
70    ///
71    /// Observed here, at the contract, rather than left to the consumer:
72    /// a consumer that forgot would silently starve the kernel's quality
73    /// telemetry, and each consumer remembering is the same code written N
74    /// times. The rpc names keep the spelling the telemetry has always
75    /// carried.
76    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            },
84        );
85    }
86}
87
88impl MemoryRecordApi for EmbeddedKernel {
89    fn capabilities(&self) -> ApiCapabilities {
90        ApiCapabilities::new(CONTRACT_VERSION, env!("CARGO_PKG_VERSION"), CAPABILITIES)
91    }
92
93    async fn record(&self, request: MemoryRecordRequest) -> Result<RecordedMemoryView, ApiError> {
94        let outcome = self
95            .service()
96            .ingest(ingest_command(request))
97            .await
98            .map_err(translate_record_error)?;
99        Ok(RecordedMemoryView {
100            about: outcome.about,
101            memory_id: outcome.memory_id,
102            accepted_entries: outcome.accepted.entries,
103            accepted_relations: outcome.accepted.relations,
104            accepted_evidence: outcome.accepted.evidence,
105            read_after_write_ready: outcome.read_after_write_ready,
106            warnings: outcome.warnings,
107        })
108    }
109}
110
111fn ingest_command(request: MemoryRecordRequest) -> MemoryIngestCommand {
112    MemoryIngestCommand {
113        about: request.about,
114        memory: MemoryData {
115            dimensions: request
116                .dimensions
117                .into_iter()
118                .map(|dimension| MemoryDimensionData {
119                    id: dimension.id,
120                    kind: dimension.kind,
121                    title: dimension.title,
122                    metadata: dimension.metadata,
123                })
124                .collect(),
125            entries: request
126                .entries
127                .into_iter()
128                .map(|entry| MemoryEntryData {
129                    id: entry.id,
130                    kind: entry.kind,
131                    text: entry.text,
132                    coordinates: entry
133                        .coordinates
134                        .into_iter()
135                        .map(|coordinate| MemoryCoordinateData {
136                            dimension: coordinate.dimension,
137                            scope_id: coordinate.scope_id,
138                            occurred_at: coordinate.occurred_at,
139                            observed_at: None,
140                            ingested_at: None,
141                            valid_from: None,
142                            valid_until: None,
143                            sequence: coordinate.sequence,
144                            rank: coordinate.rank,
145                            metadata: Default::default(),
146                        })
147                        .collect(),
148                    metadata: entry.metadata,
149                })
150                .collect(),
151            relations: request
152                .relations
153                .into_iter()
154                .map(|relation| kmp_application::MemoryRelationData {
155                    source_ref: relation.from,
156                    target_ref: relation.to,
157                    rel: relation.rel,
158                    semantic_class: relation.semantic_class,
159                    why: relation.why,
160                    evidence: None,
161                    confidence: relation.confidence,
162                    sequence: relation.sequence,
163                    motivation: None,
164                    method: None,
165                    decision_id: None,
166                    caused_by_node_id: None,
167                    coordinate: None,
168                })
169                .collect(),
170            evidence: request
171                .evidence
172                .into_iter()
173                .map(|evidence| MemoryEvidenceData {
174                    id: evidence.id,
175                    supports: evidence.supports,
176                    text: evidence.text,
177                    source: evidence.source,
178                    time: evidence.time,
179                    metadata: evidence.metadata,
180                })
181                .collect(),
182        },
183        provenance: request.provenance.map(|provenance| MemoryProvenanceData {
184            source_kind: provenance.source_kind,
185            source_agent: provenance.source_agent,
186            observed_at: provenance.observed_at,
187            correlation_id: provenance.correlation_id,
188            causation_id: provenance.causation_id,
189        }),
190        idempotency_key: request.idempotency_key,
191        dry_run: false,
192    }
193}
194
195/// The record side's error map.
196///
197/// One case more than the recall side's: a port `Conflict` here is the
198/// idempotency key reused with different content — the kernel looked at the
199/// request and said no, and retrying it unchanged earns the same answer. On
200/// the recall side a conflict cannot mean that, so the shared map's reading
201/// of ports-as-environment stays right there.
202fn translate_record_error(error: ApplicationError) -> ApiError {
203    match error {
204        ApplicationError::Ports(PortError::Conflict(reason)) => ApiError::Refused { reason },
205        other => translate_error(other),
206    }
207}
208
209fn dimensions(kinds: Vec<String>, scoped_to_about: bool) -> DimensionSelection {
210    let selection = if kinds.is_empty() {
211        DimensionSelection::all()
212    } else {
213        DimensionSelection::only(kinds)
214    };
215    if scoped_to_about {
216        selection.with_current_about_scope()
217    } else {
218        selection
219    }
220}
221
222fn tier(tier: MemoryTier) -> ResolutionTier {
223    match tier {
224        MemoryTier::Summary => ResolutionTier::L0Summary,
225        MemoryTier::CausalSpine => ResolutionTier::L1CausalSpine,
226        MemoryTier::EvidencePack => ResolutionTier::L2EvidencePack,
227    }
228}
229
230fn answer_policy(policy: MemoryAnswerPolicy) -> DomainAnswerPolicy {
231    match policy {
232        MemoryAnswerPolicy::EvidenceOrUnknown => DomainAnswerPolicy::EvidenceOrUnknown,
233        MemoryAnswerPolicy::ShowConflicts => DomainAnswerPolicy::ShowConflicts,
234        MemoryAnswerPolicy::BestEffort => DomainAnswerPolicy::BestEffort,
235    }
236}
237
238fn recall_view(about: String, result: &GetContextResult) -> MemoryRecallView {
239    let bundle = &result.bundle;
240    MemoryRecallView {
241        about,
242        revision: bundle.metadata().revision,
243        content_hash: bundle.metadata().content_hash.clone(),
244        root: node_view(bundle.root_node()),
245        neighbors: bundle.neighbor_nodes().iter().map(node_view).collect(),
246        relationships: bundle
247            .relationships()
248            .iter()
249            .map(|relationship| MemoryRelationshipView {
250                source_node_id: relationship.source_node_id().to_string(),
251                target_node_id: relationship.target_node_id().to_string(),
252                relationship_type: relationship.relationship_type().to_string(),
253                why: relationship
254                    .explanation()
255                    .rationale()
256                    .map(ToOwned::to_owned),
257                evidence: relationship.explanation().evidence().map(ToOwned::to_owned),
258            })
259            .collect(),
260        details: bundle
261            .node_details()
262            .iter()
263            .map(|detail| MemoryDetailView {
264                node_id: detail.node_id().to_string(),
265                detail: detail.detail().to_string(),
266                content_hash: detail.content_hash().to_string(),
267                revision: detail.revision(),
268            })
269            .collect(),
270        rendered: RenderedMemoryView {
271            content: result.rendered.content.clone(),
272            content_hash: result.rendered.content_hash.clone(),
273            token_count: result.rendered.token_count,
274            quality: kmp_memory_api::MemoryQualityView {
275                raw_equivalent_tokens: result.rendered.quality.raw_equivalent_tokens(),
276                compression_ratio: result.rendered.quality.compression_ratio(),
277                causal_density: result.rendered.quality.causal_density(),
278                noise_ratio: result.rendered.quality.noise_ratio(),
279                detail_coverage: result.rendered.quality.detail_coverage(),
280            },
281        },
282    }
283}
284
285fn node_view(node: &kmp_domain::BundleNode) -> MemoryNodeView {
286    MemoryNodeView {
287        node_id: node.node_id().to_string(),
288        node_kind: node.node_kind().to_string(),
289        title: node.title().to_string(),
290        summary: node.summary().to_string(),
291        status: node.status().to_string(),
292        labels: node.labels().to_vec(),
293        properties: node.properties().clone(),
294    }
295}
296
297fn translate_error(error: ApplicationError) -> ApiError {
298    match error {
299        ApplicationError::NotFound(what) => ApiError::NotFound { what },
300        ApplicationError::Validation(reason) => ApiError::Refused { reason },
301        refused @ ApplicationError::Domain(_) => ApiError::Refused {
302            reason: refused.to_string(),
303        },
304        // A port failing is the storage or the environment, not the request.
305        unavailable @ ApplicationError::Ports(_) => ApiError::Unavailable {
306            reason: unavailable.to_string(),
307        },
308    }
309}