1use 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 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#[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#[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#[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#[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#[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#[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#[derive(Debug, Clone)]
649pub struct ContextPreparation {
650 pub execution_input: ContextExecutionInput,
651 pub plan: ContextPlan,
652 pub rendered_projection: InternalRenderedContext,
653}
654
655#[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#[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#[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#[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
817pub 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
825pub 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}