Skip to main content

deepstrike_core/context/
execution.rs

1//! Canonical Context execution contracts.
2//!
3//! Context has three deliberately different representations:
4//!
5//! * [`ContextState`] is the semantic state from which a context is selected;
6//! * [`ContextPlan`] is the admitted runtime optimisation decision;
7//! * [`ContextExecutionInput`] is the frozen identity of one provider execution input.
8//!
9//! Provider request bytes remain host evidence. None of these objects is a provider request
10//! serializer, and none of them carries token projections as semantic message fields.
11
12use serde::{Deserialize, Serialize};
13use std::collections::HashSet;
14use thiserror::Error;
15
16use crate::evolution::ContentDigest;
17use crate::runtime::kernel::wire::record::canonical_bytes;
18use crate::types::message::CoreMessage;
19
20use super::partitions::ContextPartitions;
21use super::renderer::InternalRenderedContext;
22
23pub const CONTEXT_SCHEMA: &str = "context/v1";
24
25#[derive(Debug, Error, Clone, PartialEq, Eq)]
26pub enum ContextContractError {
27    #[error("context canonical value could not be serialized: {0}")]
28    Canonical(String),
29    #[error("context {kind} digest mismatch: expected {expected}, got {actual}")]
30    DigestMismatch {
31        kind: &'static str,
32        expected: ContentDigest,
33        actual: ContentDigest,
34    },
35    #[error("context plan references an entry that is not in state: {0}")]
36    UnknownEntry(String),
37    #[error("context plan state generation mismatch: expected {expected}, got {actual}")]
38    GenerationMismatch { expected: u64, actual: u64 },
39    #[error("context execution input field mismatch: {field}")]
40    FieldMismatch { field: &'static str },
41    #[error("unsupported context schema: {0}")]
42    UnsupportedSchema(String),
43    #[error("duplicate context entry or selection: {0}")]
44    DuplicateEntry(String),
45}
46
47fn digest<T: Serialize>(value: &T) -> Result<ContentDigest, ContextContractError> {
48    let bytes = canonical_bytes(value)
49        .map_err(|error| ContextContractError::Canonical(error.to_string()))?;
50    Ok(ContentDigest::from_bytes(bytes.as_slice()))
51}
52
53pub(crate) fn message_digest(
54    message: &CoreMessage,
55    handles: &crate::mm::handle::HandleTable,
56) -> Result<ContentDigest, ContextContractError> {
57    digest(&message_material(message, handles))
58}
59
60pub(crate) fn message_material(
61    message: &CoreMessage,
62    handles: &crate::mm::handle::HandleTable,
63) -> serde_json::Value {
64    use crate::types::durable_content::DurableContent;
65    use crate::types::message::{Content, ContentPart};
66    // Checkpoints may rebuild text results with explicit durable blocks, and page-out replaces
67    // their resident body with a preview. These are representations of the same semantic object.
68    let blocks = match &message.content {
69        Content::Text(text) => vec![serde_json::json!({"type": "text", "text": text})],
70        Content::Parts(parts) => parts.iter().map(|part| match part {
71            ContentPart::ToolResult { call_id, output, is_error, durable_content } => {
72                let reference = handles.all().iter().find(|h| h.source.as_deref() == Some(call_id.as_str()))
73                    .filter(|h| h.residency.digest().is_some());
74                match reference {
75                    Some(handle) => serde_json::json!({
76                        "type": "tool_result", "call_id": call_id, "is_error": is_error,
77                        "reference": {"payload_ref": handle.residency.payload_ref(), "digest": handle.residency.digest()},
78                    }),
79                    None => serde_json::json!({
80                        "type": "tool_result", "call_id": call_id, "is_error": is_error,
81                        "content": durable_content.clone().unwrap_or_else(|| DurableContent::text(output)),
82                    }),
83                }
84            }
85            _ => serde_json::to_value(part).expect("content part is serializable"),
86        }).collect(),
87    };
88    serde_json::json!([&message.role, blocks, &message.tool_calls])
89}
90
91fn verify_schema(schema: &str) -> Result<(), ContextContractError> {
92    if schema != CONTEXT_SCHEMA {
93        return Err(ContextContractError::UnsupportedSchema(schema.to_string()));
94    }
95    Ok(())
96}
97
98/// The semantic area from which an entry came. This is provenance, not a provider role.
99#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
100#[serde(rename_all = "snake_case")]
101pub enum ContextEntrySource {
102    System,
103    Knowledge,
104    History,
105    State,
106    Signal,
107}
108
109/// A stable reference to one addressable semantic item in ContextState.
110#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
111#[serde(deny_unknown_fields)]
112pub struct ContextEntryRef {
113    pub entry_id: String,
114    pub content_digest: ContentDigest,
115    pub source: ContextEntrySource,
116    pub ordinal: u32,
117}
118
119/// Canonical Context state used to prepare one execution input.
120///
121/// Bodies remain in the kernel's existing semantic state/checkpoint surfaces. This object is the
122/// explicit, content-addressed index of that state; measurements and provider rendering are not
123/// part of it.
124#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
125#[serde(deny_unknown_fields)]
126pub struct ContextState {
127    pub schema: String,
128    pub generation: u64,
129    pub system: Vec<ContextEntryRef>,
130    pub knowledge: Vec<ContextEntryRef>,
131    pub history: Vec<ContextEntryRef>,
132    pub state: Vec<ContextEntryRef>,
133    pub task_state: ContentDigest,
134    pub signals: Vec<ContentDigest>,
135    pub digest: ContentDigest,
136}
137
138#[derive(Debug, Serialize)]
139struct ContextStateBody<'a> {
140    schema: &'a str,
141    generation: u64,
142    system: &'a [ContextEntryRef],
143    knowledge: &'a [ContextEntryRef],
144    history: &'a [ContextEntryRef],
145    state: &'a [ContextEntryRef],
146    task_state: &'a ContentDigest,
147    signals: &'a [ContentDigest],
148}
149
150impl ContextState {
151    pub fn from_partitions(
152        partitions: &ContextPartitions,
153        generation: u64,
154    ) -> Result<Self, ContextContractError> {
155        Self::from_partitions_with_handles(
156            partitions,
157            generation,
158            &crate::mm::handle::HandleTable::new(),
159        )
160    }
161
162    pub fn from_partitions_with_handles(
163        partitions: &ContextPartitions,
164        generation: u64,
165        handles: &crate::mm::handle::HandleTable,
166    ) -> Result<Self, ContextContractError> {
167        let refs_for_messages = |source: ContextEntrySource, messages: &[CoreMessage]| {
168            messages
169                .iter()
170                .enumerate()
171                .map(|(ordinal, message)| {
172                    let content_digest = message_digest(message, handles)?;
173                    Ok(ContextEntryRef {
174                        entry_id: format!("{}:{ordinal}:{}", source_label(source), content_digest),
175                        content_digest,
176                        source,
177                        ordinal: ordinal as u32,
178                    })
179                })
180                .collect::<Result<Vec<_>, ContextContractError>>()
181        };
182
183        let system = refs_for_messages(ContextEntrySource::System, &partitions.system.messages)?;
184        let knowledge = partitions
185            .knowledge
186            .entries
187            .iter()
188            .enumerate()
189            .map(|(ordinal, entry)| {
190                let content_digest = message_digest(&entry.message, handles)?;
191                Ok(ContextEntryRef {
192                    entry_id: format!(
193                        "knowledge:{ordinal}:{}:{}",
194                        entry.key.as_deref().unwrap_or("unkeyed"),
195                        content_digest
196                    ),
197                    content_digest,
198                    source: ContextEntrySource::Knowledge,
199                    ordinal: ordinal as u32,
200                })
201            })
202            .collect::<Result<Vec<_>, ContextContractError>>()?;
203        let history = refs_for_messages(ContextEntrySource::History, &partitions.history.messages)?;
204        let task_state = digest(&partitions.task_state)?;
205        let mut state = Vec::with_capacity(1 + partitions.signals.len());
206        state.push(ContextEntryRef {
207            entry_id: format!("state:task_state:{task_state}"),
208            content_digest: task_state.clone(),
209            source: ContextEntrySource::State,
210            ordinal: 0,
211        });
212        let signals = partitions
213            .signals
214            .iter()
215            .enumerate()
216            .map(|(ordinal, signal)| {
217                let content_digest = digest(signal)?;
218                state.push(ContextEntryRef {
219                    entry_id: format!("signal:{ordinal}:{content_digest}"),
220                    content_digest: content_digest.clone(),
221                    source: ContextEntrySource::Signal,
222                    ordinal: ordinal as u32,
223                });
224                Ok(content_digest)
225            })
226            .collect::<Result<Vec<_>, ContextContractError>>()?;
227        let unsigned = Self {
228            schema: CONTEXT_SCHEMA.to_string(),
229            generation,
230            system,
231            knowledge,
232            history,
233            state,
234            task_state,
235            signals,
236            digest: ContentDigest::from_bytes(b"context-state-placeholder"),
237        };
238        let digest = digest(&ContextStateBody::from(&unsigned))?;
239        Ok(Self { digest, ..unsigned })
240    }
241
242    pub fn verify_digest(&self) -> Result<(), ContextContractError> {
243        verify_schema(&self.schema)?;
244        let mut ids = HashSet::new();
245        for entry in self
246            .system
247            .iter()
248            .chain(&self.knowledge)
249            .chain(&self.history)
250            .chain(&self.state)
251        {
252            if !ids.insert(&entry.entry_id) {
253                return Err(ContextContractError::DuplicateEntry(entry.entry_id.clone()));
254            }
255        }
256        let expected = digest(&ContextStateBody::from(self))?;
257        if expected != self.digest {
258            return Err(ContextContractError::DigestMismatch {
259                kind: "state",
260                expected,
261                actual: self.digest.clone(),
262            });
263        }
264        Ok(())
265    }
266
267    pub fn contains_entry(&self, entry_id: &str) -> bool {
268        self.system
269            .iter()
270            .chain(self.knowledge.iter())
271            .chain(self.history.iter())
272            .chain(self.state.iter())
273            .any(|entry| entry.entry_id == entry_id)
274    }
275}
276
277impl<'a> From<&'a ContextState> for ContextStateBody<'a> {
278    fn from(value: &'a ContextState) -> Self {
279        Self {
280            schema: &value.schema,
281            generation: value.generation,
282            system: &value.system,
283            knowledge: &value.knowledge,
284            history: &value.history,
285            state: &value.state,
286            task_state: &value.task_state,
287            signals: &value.signals,
288        }
289    }
290}
291
292fn source_label(source: ContextEntrySource) -> &'static str {
293    match source {
294        ContextEntrySource::System => "system",
295        ContextEntrySource::Knowledge => "knowledge",
296        ContextEntrySource::History => "history",
297        ContextEntrySource::State => "state",
298        ContextEntrySource::Signal => "signal",
299    }
300}
301
302/// The action the runtime admitted for one Context entry.
303#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
304#[serde(rename_all = "snake_case")]
305pub enum ContextPlanAction {
306    Include,
307    Excerpt,
308    Collapse,
309    PageOut,
310    Omit,
311}
312
313#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
314#[serde(deny_unknown_fields)]
315pub struct ContextSelection {
316    pub entry_id: String,
317    pub action: ContextPlanAction,
318    pub reason: String,
319}
320
321#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
322#[serde(deny_unknown_fields)]
323pub struct CachePrefixBoundary {
324    pub digest: ContentDigest,
325    pub entries: u32,
326}
327
328/// A deterministic runtime decision. It is not a mutable copy of ContextState.
329#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
330#[serde(deny_unknown_fields)]
331pub struct ContextPlan {
332    pub schema: String,
333    pub plan_id: ContentDigest,
334    pub operation_id: String,
335    pub step_id: String,
336    pub state_digest: ContentDigest,
337    pub state_generation: u64,
338    pub runtime_inputs: ContentDigest,
339    pub policy_digest: ContentDigest,
340    pub provider_profile_digest: ContentDigest,
341    pub measurement_fingerprints: Vec<ContentDigest>,
342    pub selections: Vec<ContextSelection>,
343    pub input_budget_tokens: u32,
344    pub projected_tokens: u32,
345    pub pressure_ppm: u32,
346    pub cache_prefix: Option<CachePrefixBoundary>,
347}
348
349#[derive(Debug, Serialize)]
350struct ContextPlanBody<'a> {
351    schema: &'a str,
352    operation_id: &'a str,
353    step_id: &'a str,
354    state_digest: &'a ContentDigest,
355    state_generation: u64,
356    runtime_inputs: &'a ContentDigest,
357    policy_digest: &'a ContentDigest,
358    provider_profile_digest: &'a ContentDigest,
359    measurement_fingerprints: &'a [ContentDigest],
360    selections: &'a [ContextSelection],
361    input_budget_tokens: u32,
362    projected_tokens: u32,
363    pressure_ppm: u32,
364    cache_prefix: Option<&'a CachePrefixBoundary>,
365}
366
367impl ContextPlan {
368    #[allow(clippy::too_many_arguments)]
369    pub fn new(
370        operation_id: impl Into<String>,
371        step_id: impl Into<String>,
372        state: &ContextState,
373        policy_digest: ContentDigest,
374        provider_profile_digest: ContentDigest,
375        measurement_fingerprints: Vec<ContentDigest>,
376        selections: Vec<ContextSelection>,
377        input_budget_tokens: u32,
378        projected_tokens: u32,
379        pressure_ppm: u32,
380        cache_prefix: Option<CachePrefixBoundary>,
381        runtime_inputs: ContentDigest,
382    ) -> Result<Self, ContextContractError> {
383        let unsigned = Self {
384            schema: CONTEXT_SCHEMA.to_string(),
385            plan_id: ContentDigest::from_bytes(b"context-plan-placeholder"),
386            operation_id: operation_id.into(),
387            step_id: step_id.into(),
388            state_digest: state.digest.clone(),
389            state_generation: state.generation,
390            runtime_inputs,
391            policy_digest,
392            provider_profile_digest,
393            measurement_fingerprints,
394            selections,
395            input_budget_tokens,
396            projected_tokens,
397            pressure_ppm,
398            cache_prefix,
399        };
400        let plan_id = digest(&ContextPlanBody::from(&unsigned))?;
401        let plan = Self {
402            plan_id,
403            ..unsigned
404        };
405        plan.verify(state)?;
406        Ok(plan)
407    }
408
409    pub fn verify(&self, state: &ContextState) -> Result<(), ContextContractError> {
410        state.verify_digest()?;
411        if self.state_generation != state.generation {
412            return Err(ContextContractError::GenerationMismatch {
413                expected: state.generation,
414                actual: self.state_generation,
415            });
416        }
417        if self.state_digest != state.digest {
418            return Err(ContextContractError::DigestMismatch {
419                kind: "plan state",
420                expected: state.digest.clone(),
421                actual: self.state_digest.clone(),
422            });
423        }
424        let state_entries: HashSet<&str> = state
425            .system
426            .iter()
427            .chain(&state.knowledge)
428            .chain(&state.history)
429            .chain(&state.state)
430            .map(|entry| entry.entry_id.as_str())
431            .collect();
432        for selection in &self.selections {
433            if !state_entries.contains(selection.entry_id.as_str()) {
434                return Err(ContextContractError::UnknownEntry(
435                    selection.entry_id.clone(),
436                ));
437            }
438        }
439        self.verify_digest()
440    }
441
442    pub fn verify_digest(&self) -> Result<(), ContextContractError> {
443        verify_schema(&self.schema)?;
444        if self.operation_id.is_empty() || self.step_id.is_empty() || self.pressure_ppm > 1_000_000
445        {
446            return Err(ContextContractError::FieldMismatch {
447                field: "plan identity or pressure",
448            });
449        }
450        let mut ids = HashSet::new();
451        for selection in &self.selections {
452            if !ids.insert(&selection.entry_id) {
453                return Err(ContextContractError::DuplicateEntry(
454                    selection.entry_id.clone(),
455                ));
456            }
457        }
458        let expected = digest(&ContextPlanBody::from(self))?;
459        if expected != self.plan_id {
460            return Err(ContextContractError::DigestMismatch {
461                kind: "plan",
462                expected,
463                actual: self.plan_id.clone(),
464            });
465        }
466        Ok(())
467    }
468}
469
470impl<'a> From<&'a ContextPlan> for ContextPlanBody<'a> {
471    fn from(value: &'a ContextPlan) -> Self {
472        Self {
473            schema: &value.schema,
474            operation_id: &value.operation_id,
475            step_id: &value.step_id,
476            state_digest: &value.state_digest,
477            state_generation: value.state_generation,
478            runtime_inputs: &value.runtime_inputs,
479            policy_digest: &value.policy_digest,
480            provider_profile_digest: &value.provider_profile_digest,
481            measurement_fingerprints: &value.measurement_fingerprints,
482            selections: &value.selections,
483            input_budget_tokens: value.input_budget_tokens,
484            projected_tokens: value.projected_tokens,
485            pressure_ppm: value.pressure_ppm,
486            cache_prefix: value.cache_prefix.as_ref(),
487        }
488    }
489}
490
491/// The frozen identity of one provider execution input.
492#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
493#[serde(deny_unknown_fields)]
494pub struct ContextExecutionInput {
495    pub schema: String,
496    pub input_digest: ContentDigest,
497    pub operation_id: String,
498    pub step_id: String,
499    pub input_sequence: u64,
500    pub state_digest: ContentDigest,
501    pub policy_digest: ContentDigest,
502    pub plan_digest: ContentDigest,
503    pub rendered_snapshot: ContentDigest,
504    pub prompt_measurement: ContentDigest,
505    pub provider_route: ContentDigest,
506    pub cache_prefix: Option<CachePrefixBoundary>,
507}
508
509#[derive(Debug, Serialize)]
510struct ContextExecutionInputBody<'a> {
511    schema: &'a str,
512    operation_id: &'a str,
513    step_id: &'a str,
514    input_sequence: u64,
515    state_digest: &'a ContentDigest,
516    policy_digest: &'a ContentDigest,
517    plan_digest: &'a ContentDigest,
518    rendered_snapshot: &'a ContentDigest,
519    prompt_measurement: &'a ContentDigest,
520    provider_route: &'a ContentDigest,
521    cache_prefix: Option<&'a CachePrefixBoundary>,
522}
523
524impl ContextExecutionInput {
525    #[allow(clippy::too_many_arguments)]
526    pub fn new(
527        operation_id: impl Into<String>,
528        step_id: impl Into<String>,
529        input_sequence: u64,
530        state: &ContextState,
531        plan: &ContextPlan,
532        rendered_snapshot: ContentDigest,
533        prompt_measurement: ContentDigest,
534        provider_route: ContentDigest,
535    ) -> Result<Self, ContextContractError> {
536        plan.verify(state)?;
537        let operation_id = operation_id.into();
538        let step_id = step_id.into();
539        if operation_id != plan.operation_id {
540            return Err(ContextContractError::FieldMismatch {
541                field: "operation_id",
542            });
543        }
544        if step_id != plan.step_id {
545            return Err(ContextContractError::FieldMismatch { field: "step_id" });
546        }
547        let unsigned = Self {
548            schema: CONTEXT_SCHEMA.to_string(),
549            input_digest: ContentDigest::from_bytes(b"context-input-placeholder"),
550            operation_id,
551            step_id,
552            input_sequence,
553            state_digest: state.digest.clone(),
554            policy_digest: plan.policy_digest.clone(),
555            plan_digest: plan.plan_id.clone(),
556            rendered_snapshot,
557            prompt_measurement,
558            provider_route,
559            cache_prefix: plan.cache_prefix.clone(),
560        };
561        let input_digest = digest(&ContextExecutionInputBody::from(&unsigned))?;
562        let input = Self {
563            input_digest,
564            ..unsigned
565        };
566        input.verify(plan)?;
567        Ok(input)
568    }
569
570    pub fn verify(&self, plan: &ContextPlan) -> Result<(), ContextContractError> {
571        verify_schema(&self.schema)?;
572        plan.verify_digest()?;
573        if self.provider_route != plan.provider_profile_digest {
574            return Err(ContextContractError::FieldMismatch {
575                field: "provider_route",
576            });
577        }
578        if !plan
579            .measurement_fingerprints
580            .contains(&self.prompt_measurement)
581        {
582            return Err(ContextContractError::FieldMismatch {
583                field: "prompt_measurement",
584            });
585        }
586        if self.operation_id != plan.operation_id {
587            return Err(ContextContractError::FieldMismatch {
588                field: "operation_id",
589            });
590        }
591        if self.step_id != plan.step_id {
592            return Err(ContextContractError::FieldMismatch { field: "step_id" });
593        }
594        if self.state_digest != plan.state_digest {
595            return Err(ContextContractError::FieldMismatch {
596                field: "state_digest",
597            });
598        }
599        if self.policy_digest != plan.policy_digest {
600            return Err(ContextContractError::FieldMismatch {
601                field: "policy_digest",
602            });
603        }
604        if self.cache_prefix != plan.cache_prefix {
605            return Err(ContextContractError::FieldMismatch {
606                field: "cache_prefix",
607            });
608        }
609        if self.plan_digest != plan.plan_id {
610            return Err(ContextContractError::DigestMismatch {
611                kind: "input plan",
612                expected: plan.plan_id.clone(),
613                actual: self.plan_digest.clone(),
614            });
615        }
616        let expected = digest(&ContextExecutionInputBody::from(self))?;
617        if expected != self.input_digest {
618            return Err(ContextContractError::DigestMismatch {
619                kind: "execution input",
620                expected,
621                actual: self.input_digest.clone(),
622            });
623        }
624        Ok(())
625    }
626}
627
628impl<'a> From<&'a ContextExecutionInput> for ContextExecutionInputBody<'a> {
629    fn from(value: &'a ContextExecutionInput) -> Self {
630        Self {
631            schema: &value.schema,
632            operation_id: &value.operation_id,
633            step_id: &value.step_id,
634            input_sequence: value.input_sequence,
635            state_digest: &value.state_digest,
636            policy_digest: &value.policy_digest,
637            plan_digest: &value.plan_digest,
638            rendered_snapshot: &value.rendered_snapshot,
639            prompt_measurement: &value.prompt_measurement,
640            provider_route: &value.provider_route,
641            cache_prefix: value.cache_prefix.as_ref(),
642        }
643    }
644}
645
646/// The non-authoritative result returned by the preparation boundary. The provider adapter keeps
647/// the projection transient and records the execution-input identity alongside host evidence.
648#[derive(Debug, Clone)]
649pub struct ContextPreparation {
650    pub execution_input: ContextExecutionInput,
651    pub plan: ContextPlan,
652    pub rendered_projection: InternalRenderedContext,
653}
654
655/// Facts supplied by the runtime at the preparation boundary. The renderer remains a kernel
656/// projection; provider adapters add their protocol-specific evidence after this call.
657#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
658#[serde(deny_unknown_fields)]
659pub struct ContextPreparationRequest {
660    pub operation_id: String,
661    pub step_id: String,
662    pub input_sequence: u64,
663    pub policy_digest: ContentDigest,
664    pub prompt_measurement: ContentDigest,
665    pub provider_route: ContentDigest,
666}
667
668/// Kernel-owned facts frozen in a provider effect. Host route and native count do not exist
669/// when the kernel emits this candidate; the host binds them before dispatch through core.
670#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
671#[serde(deny_unknown_fields)]
672pub struct ContextCandidate {
673    pub schema: String,
674    pub operation_id: String,
675    pub step_id: String,
676    pub input_sequence: u64,
677    pub state: ContextState,
678    pub runtime_inputs: ContentDigest,
679    pub policy_digest: ContentDigest,
680    pub rendered_snapshot: ContentDigest,
681    pub selections: Vec<ContextSelection>,
682    pub input_budget_tokens: u32,
683    pub projected_tokens: u32,
684    pub pressure_ppm: u32,
685    pub cache_prefix: Option<CachePrefixBoundary>,
686}
687
688impl ContextCandidate {
689    pub fn bind(
690        &self,
691        prompt_measurement: ContentDigest,
692        provider_route: ContentDigest,
693    ) -> Result<(ContextPlan, ContextExecutionInput), ContextContractError> {
694        verify_schema(&self.schema)?;
695        let plan = ContextPlan::new(
696            &self.operation_id,
697            &self.step_id,
698            &self.state,
699            self.policy_digest.clone(),
700            provider_route.clone(),
701            vec![prompt_measurement.clone()],
702            self.selections.clone(),
703            self.input_budget_tokens,
704            self.projected_tokens,
705            self.pressure_ppm,
706            self.cache_prefix.clone(),
707            self.runtime_inputs.clone(),
708        )?;
709        let input = ContextExecutionInput::new(
710            &self.operation_id,
711            &self.step_id,
712            self.input_sequence,
713            &self.state,
714            &plan,
715            self.rendered_snapshot.clone(),
716            prompt_measurement,
717            provider_route,
718        )?;
719        Ok((plan, input))
720    }
721}
722
723/// The host count names fingerprinted request material under the route's declared scope.
724#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
725#[serde(rename_all = "camelCase", deny_unknown_fields)]
726pub struct ContextPromptMeasurement {
727    pub request_fingerprint: ContentDigest,
728    pub input_tokens: u64,
729    pub source: super::measurement::MeasurementSource,
730    pub confidence: super::measurement::MeasurementConfidence,
731}
732
733#[derive(Debug, Clone, Serialize, Deserialize)]
734#[serde(deny_unknown_fields)]
735pub struct ContextDispatchRequest {
736    pub effect: crate::runtime::kernel::wire::effect::CallProviderEffect,
737    pub request_fingerprint: ContentDigest,
738    pub provider_route: serde_json::Value,
739    pub prompt_measurement: ContextPromptMeasurement,
740}
741
742/// Evidence is returned with its canonical identity so storage adapters can persist every
743/// referenced body. This is host evidence; the kernel's immutable candidate remains in the effect.
744#[derive(Debug, Clone, Serialize, Deserialize)]
745pub struct ContextDispatchPreparation {
746    pub execution_input: ContextExecutionInput,
747    pub plan: ContextPlan,
748    pub binding: crate::evolution::EvaluationContextBinding,
749    pub state: ContextState,
750    pub provider_route: serde_json::Value,
751    pub prompt_measurement: ContextPromptMeasurement,
752}
753
754pub fn prepare_context_dispatch(
755    request: &ContextDispatchRequest,
756) -> Result<ContextDispatchPreparation, ContextContractError> {
757    let candidate = &request.effect.context_candidate;
758    let projection = digest(&(&request.effect.context, &request.effect.tools))?;
759    if projection != candidate.rendered_snapshot {
760        return Err(ContextContractError::DigestMismatch {
761            kind: "provider projection",
762            expected: candidate.rendered_snapshot.clone(),
763            actual: projection,
764        });
765    }
766    if request.prompt_measurement.request_fingerprint != request.request_fingerprint {
767        return Err(ContextContractError::FieldMismatch {
768            field: "measurement request fingerprint",
769        });
770    }
771    let context = &request.effect.context;
772    let prefix = match context.frozen_prefix_len {
773        Some(entries) => {
774            let turns = context.turns.get(..entries as usize).ok_or(
775                ContextContractError::FieldMismatch {
776                    field: "cache prefix boundary",
777                },
778            )?;
779            Some(CachePrefixBoundary {
780                entries,
781                digest: digest(&(&context.system_stable, &context.system_knowledge, turns))?,
782            })
783        }
784        None => None,
785    };
786    if prefix != candidate.cache_prefix {
787        return Err(ContextContractError::FieldMismatch {
788            field: "cache_prefix",
789        });
790    }
791    if !request.provider_route.is_object()
792        || request
793            .provider_route
794            .as_object()
795            .is_some_and(|route| route.is_empty())
796    {
797        return Err(ContextContractError::FieldMismatch {
798            field: "provider_route",
799        });
800    }
801    let (plan, execution_input) = candidate.bind(
802        digest(&request.prompt_measurement)?,
803        digest(&request.provider_route)?,
804    )?;
805    let binding =
806        crate::evolution::EvaluationContextBinding::from_execution_input(&execution_input);
807    Ok(ContextDispatchPreparation {
808        execution_input,
809        plan,
810        binding,
811        state: candidate.state.clone(),
812        provider_route: request.provider_route.clone(),
813        prompt_measurement: request.prompt_measurement.clone(),
814    })
815}
816
817/// Single canonical bridge shared by all SDK mirrors; no SDK owns selection or digest rules.
818pub fn prepare_context_dispatch_json(request: &str) -> Result<String, String> {
819    let request: ContextDispatchRequest =
820        serde_json::from_str(request).map_err(|e| e.to_string())?;
821    let preparation = prepare_context_dispatch(&request).map_err(|e| e.to_string())?;
822    serde_json::to_string(&preparation).map_err(|e| e.to_string())
823}
824
825/// Rebind recorded host evidence to the replayed kernel effect and compare all derived objects.
826/// A self-consistent record from a different effect, route or request is not equivalent.
827pub fn verify_context_dispatch(
828    effect: &crate::runtime::kernel::wire::effect::CallProviderEffect,
829    preparation: &ContextDispatchPreparation,
830) -> Result<(), ContextContractError> {
831    let expected = prepare_context_dispatch(&ContextDispatchRequest {
832        effect: effect.clone(),
833        request_fingerprint: preparation.prompt_measurement.request_fingerprint.clone(),
834        provider_route: preparation.provider_route.clone(),
835        prompt_measurement: preparation.prompt_measurement.clone(),
836    })?;
837    let expected_digest = digest(&expected)?;
838    let actual = digest(preparation)?;
839    if expected_digest != actual {
840        return Err(ContextContractError::DigestMismatch {
841            kind: "replayed context execution",
842            expected: expected_digest,
843            actual,
844        });
845    }
846    Ok(())
847}
848
849pub fn verify_context_dispatch_json(request: &str) -> Result<String, String> {
850    #[derive(Deserialize)]
851    #[serde(deny_unknown_fields)]
852    struct Request {
853        effect: crate::runtime::kernel::wire::effect::CallProviderEffect,
854        preparation: ContextDispatchPreparation,
855    }
856    let request: Request = serde_json::from_str(request).map_err(|e| e.to_string())?;
857    verify_context_dispatch(&request.effect, &request.preparation).map_err(|e| e.to_string())?;
858    Ok("true".to_string())
859}
860
861#[cfg(test)]
862mod tests {
863    use super::*;
864    use crate::context::config::ContextConfig;
865    use crate::types::message::CoreMessage;
866
867    fn digest(text: &str) -> ContentDigest {
868        ContentDigest::from_bytes(text.as_bytes())
869    }
870
871    fn state() -> ContextState {
872        let mut partitions = ContextPartitions::new(&ContextConfig::default());
873        partitions.system.push(CoreMessage::system("rules"), 1);
874        partitions.history.push(CoreMessage::user("hello"), 1);
875        ContextState::from_partitions(&partitions, 3).unwrap()
876    }
877
878    fn dispatch_request() -> ContextDispatchRequest {
879        use crate::runtime::kernel::wire::effect::{
880            CallProviderEffect, RenderedContext as WireRenderedContext,
881        };
882        let manager = crate::context::manager::ContextManager::new(100_000);
883        let (mut candidate, _) = manager
884            .prepare_candidate("op".to_string(), "step".to_string(), 1, digest("policy"))
885            .unwrap();
886        let context = WireRenderedContext::default();
887        let tools = vec![];
888        candidate.rendered_snapshot = super::digest(&(&context, &tools)).unwrap();
889        ContextDispatchRequest {
890            effect: CallProviderEffect {
891                context,
892                tools,
893                context_candidate: Box::new(candidate),
894            },
895            request_fingerprint: digest("actual request"),
896            provider_route: serde_json::json!({"protocol":"test", "model":"model"}),
897            prompt_measurement: ContextPromptMeasurement {
898                request_fingerprint: digest("actual request"),
899                input_tokens: 10,
900                source: super::super::measurement::MeasurementSource::Heuristic,
901                confidence: super::super::measurement::MeasurementConfidence::LowConfidence,
902            },
903        }
904    }
905
906    #[test]
907    fn dispatch_replay_recomputes_the_same_input_and_rejects_changed_evidence() {
908        let request = dispatch_request();
909        let prepared = prepare_context_dispatch(&request).unwrap();
910        verify_context_dispatch(&request.effect, &prepared).unwrap();
911        let restored: ContextDispatchPreparation =
912            serde_json::from_slice(&serde_json::to_vec(&prepared).unwrap()).unwrap();
913        verify_context_dispatch(&request.effect, &restored).unwrap();
914        let mut changed = restored.clone();
915        changed.provider_route["model"] = "another-model".into();
916        assert!(verify_context_dispatch(&request.effect, &changed).is_err());
917        let mut changed = restored.clone();
918        changed.prompt_measurement.input_tokens += 1;
919        assert!(verify_context_dispatch(&request.effect, &changed).is_err());
920        let mut changed = restored;
921        changed.plan.selections[0].reason = "changed".to_string();
922        assert!(changed.execution_input.verify(&changed.plan).is_err());
923    }
924
925    #[test]
926    fn dispatch_rejects_projection_measurement_and_cache_mismatch() {
927        let request = dispatch_request();
928        let mut changed = request.clone();
929        changed.effect.context.system_stable = "altered prompt".into();
930        assert!(prepare_context_dispatch(&changed).is_err());
931        let mut changed = request.clone();
932        changed.prompt_measurement.request_fingerprint = digest("another request");
933        assert!(prepare_context_dispatch(&changed).is_err());
934        let mut changed = request;
935        changed.effect.context_candidate.cache_prefix = Some(CachePrefixBoundary {
936            digest: digest("forged cache"),
937            entries: 99,
938        });
939        assert!(prepare_context_dispatch(&changed).is_err());
940    }
941
942    #[test]
943    fn identical_unkeyed_knowledge_entries_have_distinct_identity() {
944        let mut partitions = ContextPartitions::new(&ContextConfig::default());
945        partitions.knowledge.push(CoreMessage::system("same"), 1);
946        partitions.knowledge.push(CoreMessage::system("same"), 1);
947        let state = ContextState::from_partitions(&partitions, 0).unwrap();
948        assert_ne!(state.knowledge[0].entry_id, state.knowledge[1].entry_id);
949        state.verify_digest().unwrap();
950    }
951
952    #[test]
953    fn context_execution_shared_sdk_fixture() {
954        let request = dispatch_request();
955        let prepared = prepare_context_dispatch(&request).unwrap();
956        verify_context_dispatch(&request.effect, &prepared).unwrap();
957        let produced = serde_json::json!({
958            "id": "context-execution", "domain": "context_execution",
959            "input": { "request": request },
960            "expected": { "canonical": {
961                "input_digest": prepared.execution_input.input_digest,
962                "plan_digest": prepared.plan.plan_id, "verified": true,
963            } },
964        });
965        let path = std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
966            .join("../../tests/fixtures/sdk-conformance/canonical/context-execution.json");
967        if std::env::var("BLESS_CONTEXT_FIXTURES").as_deref() == Ok("1") {
968            std::fs::write(
969                &path,
970                format!("{}\n", serde_json::to_string_pretty(&produced).unwrap()),
971            )
972            .unwrap();
973        }
974        let fixture: serde_json::Value =
975            serde_json::from_slice(&std::fs::read(path).unwrap()).unwrap();
976        assert_eq!(produced, fixture, "shared SDK Context fixture drifted");
977    }
978
979    #[test]
980    fn state_digest_covers_semantic_entries() {
981        let state = state();
982        state.verify_digest().unwrap();
983        let mut changed = state.clone();
984        changed.history[0].content_digest = digest("changed");
985        assert!(changed.verify_digest().is_err());
986    }
987
988    #[test]
989    fn plan_rejects_unknown_entries_and_tampering() {
990        let state = state();
991        let unknown = ContextSelection {
992            entry_id: "history:missing".to_string(),
993            action: ContextPlanAction::Include,
994            reason: "test".to_string(),
995        };
996        assert!(
997            ContextPlan::new(
998                "op",
999                "step",
1000                &state,
1001                digest("policy"),
1002                digest("route"),
1003                vec![],
1004                vec![unknown],
1005                100,
1006                10,
1007                50_000,
1008                None,
1009                digest("runtime-inputs"),
1010            )
1011            .is_err()
1012        );
1013
1014        let selection = ContextSelection {
1015            entry_id: state.history[0].entry_id.clone(),
1016            action: ContextPlanAction::Include,
1017            reason: "fits".to_string(),
1018        };
1019        let mut plan = ContextPlan::new(
1020            "op",
1021            "step",
1022            &state,
1023            digest("policy"),
1024            digest("route"),
1025            vec![digest("measurement")],
1026            vec![selection],
1027            100,
1028            10,
1029            50_000,
1030            None,
1031            digest("runtime-inputs"),
1032        )
1033        .unwrap();
1034        plan.projected_tokens = 11;
1035        assert!(plan.verify(&state).is_err());
1036    }
1037
1038    #[test]
1039    fn large_plan_validates_all_partitions_and_rejects_a_missing_tail_entry() {
1040        let mut partitions = ContextPartitions::new(&ContextConfig::default());
1041        partitions.system.push(CoreMessage::system("rules"), 1);
1042        partitions
1043            .knowledge
1044            .push(CoreMessage::system("reference"), 1);
1045        partitions.signals.push("current signal".to_string());
1046        for ordinal in 0..4096 {
1047            partitions
1048                .history
1049                .push(CoreMessage::user(format!("turn {ordinal}")), 1);
1050        }
1051        let state = ContextState::from_partitions(&partitions, 7).unwrap();
1052        let selections = state
1053            .system
1054            .iter()
1055            .chain(&state.knowledge)
1056            .chain(&state.history)
1057            .chain(&state.state)
1058            .rev()
1059            .map(|entry| ContextSelection {
1060                entry_id: entry.entry_id.clone(),
1061                action: ContextPlanAction::Include,
1062                reason: "selected".to_string(),
1063            })
1064            .collect::<Vec<_>>();
1065        let plan = ContextPlan::new(
1066            "op",
1067            "step",
1068            &state,
1069            digest("policy"),
1070            digest("route"),
1071            vec![digest("measurement")],
1072            selections.clone(),
1073            100_000,
1074            5000,
1075            50_000,
1076            None,
1077            digest("runtime"),
1078        )
1079        .unwrap();
1080        plan.verify(&state).unwrap();
1081        let mut invalid = selections;
1082        invalid.last_mut().unwrap().entry_id = "absent-entry".to_string();
1083        let rejected = ContextPlan::new(
1084            "op",
1085            "step",
1086            &state,
1087            digest("policy"),
1088            digest("route"),
1089            vec![digest("measurement")],
1090            invalid,
1091            100_000,
1092            5000,
1093            50_000,
1094            None,
1095            digest("runtime"),
1096        )
1097        .unwrap_err();
1098        assert_eq!(
1099            rejected,
1100            ContextContractError::UnknownEntry("absent-entry".to_string())
1101        );
1102    }
1103
1104    #[test]
1105    fn execution_input_is_bound_to_plan_and_rejects_tampering() {
1106        let state = state();
1107        let selection = ContextSelection {
1108            entry_id: state.history[0].entry_id.clone(),
1109            action: ContextPlanAction::Include,
1110            reason: "fits".to_string(),
1111        };
1112        let plan = ContextPlan::new(
1113            "op",
1114            "step",
1115            &state,
1116            digest("policy"),
1117            digest("route"),
1118            vec![digest("measurement")],
1119            vec![selection],
1120            100,
1121            10,
1122            50_000,
1123            None,
1124            digest("runtime-inputs"),
1125        )
1126        .unwrap();
1127        let mut input = ContextExecutionInput::new(
1128            "op",
1129            "step",
1130            1,
1131            &state,
1132            &plan,
1133            digest("render"),
1134            digest("measurement"),
1135            digest("route"),
1136        )
1137        .unwrap();
1138        input.verify(&plan).unwrap();
1139        input.rendered_snapshot = digest("tampered-render");
1140        assert!(input.verify(&plan).is_err());
1141    }
1142
1143    #[test]
1144    fn execution_input_cannot_relabel_a_plan() {
1145        let state = state();
1146        let plan = ContextPlan::new(
1147            "op",
1148            "step",
1149            &state,
1150            digest("policy"),
1151            digest("route"),
1152            vec![],
1153            vec![],
1154            100,
1155            10,
1156            50_000,
1157            None,
1158            digest("runtime-inputs"),
1159        )
1160        .unwrap();
1161        let error = ContextExecutionInput::new(
1162            "other-op",
1163            "step",
1164            1,
1165            &state,
1166            &plan,
1167            digest("render"),
1168            digest("measurement"),
1169            digest("route"),
1170        )
1171        .unwrap_err();
1172        assert_eq!(
1173            error,
1174            ContextContractError::FieldMismatch {
1175                field: "operation_id"
1176            }
1177        );
1178    }
1179}