mod budget;
pub mod chunk;
mod classify;
mod dedup;
pub mod estimator;
#[cfg(not(target_arch = "wasm32"))]
pub mod ingest;
pub mod insights;
mod log_normalize;
pub(crate) mod media;
pub mod model;
pub mod model_windows;
pub(crate) mod provenance;
mod relevance;
pub mod segment;
pub mod transcript_bridge;
pub mod wire;
pub use chunk::{chunk_text, ChunkBoundary, ChunkPolicy, TextChunk};
pub use estimator::{DynTokenEstimator, HeuristicEstimator, TokenEstimator};
#[cfg(not(target_arch = "wasm32"))]
pub use ingest::IngestRoots;
pub use insights::{CompilationInsights, ModelPricing, PricingTable};
pub use model::{
CompilePolicy, CompileRequest, CompiledContext, CompiledSection, ContextAction,
ContextDecision, ContextDecisionRef, ContextFact, ContextFragment, ContextSavings,
ContextSource, ContextWarning, FidelityRisk, ImportanceWeights, LoadedWorkingContext, MediaRef,
MemoryScope, RetrievalHandle, SectionKind, SourceReference, WorkingContext,
WorkingContextIndex, WorkingContextSession,
};
pub use model_windows::{model_window, suggest_token_budget, SuggestedBudget};
pub use relevance::DeterministicReranker;
pub use segment::{
segment_transcript, SegmentFormat, SegmentKind, SegmentationOutcome, SegmentationPolicy,
TranscriptSegment,
};
pub use transcript_bridge::{
build_transcript_compile_request, SegmentInfo, SegmentationReport, TranscriptCompileInput,
};
use std::collections::BTreeMap;
use crate::error::MemoryError;
use crate::id::stable_id;
use crate::limits;
use budget::PackItem;
use classify::RuleMatch;
use dedup::{DupKind, Duplicate};
#[must_use]
pub fn fragment_id(content: &str) -> u64 {
stable_id(content)
}
pub struct ContextCompiler {
policy: CompilePolicy,
estimator: DynTokenEstimator,
pricing: Option<PricingTable>,
}
impl ContextCompiler {
#[must_use]
pub fn new(policy: CompilePolicy) -> Self {
Self {
policy,
estimator: Box::new(HeuristicEstimator),
pricing: None,
}
}
#[must_use]
pub fn with_estimator(mut self, estimator: DynTokenEstimator) -> Self {
self.estimator = estimator;
self
}
#[must_use]
pub fn with_pricing(mut self, pricing: PricingTable) -> Self {
self.pricing = Some(pricing);
self
}
pub(crate) fn effective_policy<'a>(&'a self, request: &'a CompileRequest) -> &'a CompilePolicy {
request.policy.as_ref().unwrap_or(&self.policy)
}
pub fn compile(&self, request: &CompileRequest) -> Result<CompiledContext, MemoryError> {
let compiled = self.compile_raw(request)?;
Ok(apply_slim(compiled, self.effective_policy(request)))
}
pub(crate) fn compile_raw(
&self,
request: &CompileRequest,
) -> Result<CompiledContext, MemoryError> {
let policy = self.effective_policy(request);
let usable = validate(request, policy)?;
let analyses = analyze(request, policy, self.estimator.as_ref());
let items = pack_items(&analyses, policy, usable, self.estimator.as_ref());
let taken = budget::pack(&items, usable, &self.estimator);
let emissions = emissions(&items, &taken);
Ok(self.finish(request, &analyses, &emissions))
}
fn finish(
&self,
request: &CompileRequest,
analyses: &[Analysis],
emissions: &BTreeMap<usize, Emission>,
) -> CompiledContext {
let sections = sections(analyses, emissions);
let content = sections
.iter()
.map(|section| section.content.as_str())
.collect::<Vec<_>>()
.join(budget::JOINER);
let decisions: Vec<ContextDecision> = analyses
.iter()
.map(|analysis| decision(analysis, analyses, emissions))
.collect();
let insights = self.insights(request, analyses, &decisions, emissions, &content);
let warnings = warnings_for(&decisions);
CompiledContext {
retrieval_handles: retrieval_handles(analyses, &decisions),
sources: analyses
.iter()
.filter(|analysis| analysis.dup.is_none())
.map(|analysis| {
provenance::source_for(analysis.fragment_id, analysis.handle_hash())
})
.collect(),
risk: decisions
.iter()
.map(|decision| decision.risk)
.max()
.unwrap_or_default(),
content,
sections,
decisions,
insights,
warnings,
}
}
fn insights(
&self,
request: &CompileRequest,
analyses: &[Analysis],
decisions: &[ContextDecision],
emissions: &BTreeMap<usize, Emission>,
content: &str,
) -> CompilationInsights {
let estimator = self.estimator.as_ref();
let tokens_in: u64 = analyses
.iter()
.map(|analysis| analysis.tokens)
.fold(0, u64::saturating_add);
let media_tokens_out: u64 = analyses
.iter()
.filter(|analysis| emissions.contains_key(&analysis.seq))
.filter_map(|analysis| analysis.media.as_ref())
.map(|media| media.image_tokens)
.fold(0, u64::saturating_add);
let tokens_out = estimator.estimate(content).saturating_add(media_tokens_out);
let tokens_saved = tokens_in.saturating_sub(tokens_out);
let mut insights = CompilationInsights {
tokens_in,
tokens_out,
tokens_saved,
tokens_saved_by_rule: saved_by_rule(analyses, decisions, emissions, estimator),
..CompilationInsights::default()
};
let cost = request.target_model.as_deref().and_then(|model| {
let pricing = self
.effective_policy(request)
.pricing
.as_ref()
.or(self.pricing.as_ref())?;
let micros = pricing.cost_micros(model, tokens_saved)?;
Some((micros, pricing.currency.clone(), pricing.version.clone()))
});
if let Some((micros, currency, version)) = cost {
insights.estimated_cost_saved_micros = Some(micros);
insights.currency = Some(currency);
insights.pricing_version = Some(version);
}
insights
}
}
struct Analysis<'a> {
seq: usize,
fragment_id: u64,
content_hash: u64,
original: &'a str,
tokens: u64,
rule: RuleMatch,
relevance: f32,
priority: u8,
dup: Option<Duplicate>,
abstract_collapse: Option<(String, bool)>,
media: Option<media::MediaAnalysis>,
superseded: bool,
}
impl Analysis<'_> {
fn handle_hash(&self) -> u64 {
self.media
.as_ref()
.map_or(self.content_hash, |media| media.raw_hash)
}
}
fn media_fragment_tokens(
media: &media::MediaAnalysis,
caption: &str,
estimator: &dyn TokenEstimator,
) -> u64 {
media
.image_tokens
.saturating_add(estimator.estimate(caption))
}
struct Emission {
text: String,
taken: usize,
total: usize,
}
impl Emission {
fn is_full(&self) -> bool {
self.taken == self.total
}
}
fn validate(request: &CompileRequest, policy: &CompilePolicy) -> Result<u64, MemoryError> {
if request.fragments.iter().any(|f| f.path.is_some()) {
return Err(MemoryError::IngestDisabled);
}
if request.fragments.len() > limits::MAX_FRAGMENTS {
return Err(MemoryError::ContextOverLimit(format!(
"{} fragments exceed the cap of {}",
request.fragments.len(),
limits::MAX_FRAGMENTS
)));
}
if let Some(oversized) = request
.fragments
.iter()
.find(|fragment| fragment.content.len() > limits::MAX_FRAGMENT_BYTES)
{
return Err(MemoryError::ContextOverLimit(format!(
"a fragment of {} bytes exceeds the cap of {} bytes",
oversized.content.len(),
limits::MAX_FRAGMENT_BYTES
)));
}
for fragment in &request.fragments {
let Some(metadata) = fragment.metadata.as_ref() else {
continue;
};
let bytes = limits::metadata_bytes(metadata);
if bytes > limits::MAX_METADATA_BYTES {
return Err(MemoryError::MetadataTooLarge {
bytes,
max: limits::MAX_METADATA_BYTES,
});
}
}
validate_media(&request.fragments)?;
let budget = limits::clamp_token_budget(request.token_budget);
let usable = budget.saturating_sub(policy.response_reserve_tokens);
if usable == 0 {
return Err(MemoryError::ContextBudget {
budget,
reserve: policy.response_reserve_tokens,
});
}
Ok(usable)
}
fn validate_media(fragments: &[ContextFragment]) -> Result<(), MemoryError> {
let mut total_media_bytes: usize = 0;
for (seq, fragment) in fragments.iter().enumerate() {
let Some(media_ref) = &fragment.media else {
continue;
};
total_media_bytes = total_media_bytes.saturating_add(media_ref.bytes_b64.len());
if total_media_bytes > limits::MAX_TOTAL_MEDIA_BYTES {
return Err(MemoryError::ContextOverLimit(format!(
"total media payload exceeds the request cap of {} base64 bytes",
limits::MAX_TOTAL_MEDIA_BYTES
)));
}
if media_ref.bytes_b64.len() > limits::MAX_MEDIA_BYTES {
return Err(MemoryError::ContextOverLimit(format!(
"fragment #{seq} media payload of {} base64 bytes exceeds the cap of {} bytes",
media_ref.bytes_b64.len(),
limits::MAX_MEDIA_BYTES
)));
}
if !media::is_valid_base64(&media_ref.bytes_b64) {
return Err(MemoryError::ContextOverLimit(format!(
"fragment #{seq} media payload is not valid base64"
)));
}
}
Ok(())
}
fn analyze<'a>(
request: &'a CompileRequest,
policy: &CompilePolicy,
estimator: &dyn TokenEstimator,
) -> Vec<Analysis<'a>> {
let contents: Vec<&str> = request
.fragments
.iter()
.map(|fragment| fragment.content.as_str())
.collect();
let media_analyses: Vec<Option<media::MediaAnalysis>> = request
.fragments
.iter()
.map(|fragment| fragment.media.as_ref().map(media::analyze))
.collect();
let media_hashes: Vec<Option<u64>> = media_analyses
.iter()
.map(|analysis| analysis.as_ref().map(|analysis| analysis.raw_hash))
.collect();
let superseded_flags = classify::screenshot_supersession(&request.fragments);
let duplicates = dedup::find_duplicates(
&contents,
policy.near_dup_dedup,
&media_hashes,
&superseded_flags,
);
let query_terms = relevance::terms(&request.query);
let mut analyses: Vec<Analysis<'a>> = request
.fragments
.iter()
.zip(duplicates)
.zip(media_analyses)
.enumerate()
.map(|(seq, ((fragment, dup), media_analysis))| {
let content_hash = stable_id(&fragment.content);
let rule = classify::classify(fragment, policy);
let abstract_collapse = (rule.action == ContextAction::Abstract).then(|| {
classify::collapse_repeated_lines(
&fragment.content,
policy.normalize_log_timestamps,
)
});
let tokens = media_analysis.as_ref().map_or_else(
|| estimator.estimate(&fragment.content),
|media| media_fragment_tokens(media, &fragment.content, estimator),
);
let superseded = superseded_flags[seq]
&& !policy
.disabled_rules
.iter()
.any(|disabled| disabled == classify::SCREENSHOT_SUPERSEDED_RULE_ID);
Analysis {
seq,
fragment_id: fragment.id.unwrap_or(content_hash),
content_hash,
original: &fragment.content,
tokens,
rule,
relevance: relevance::lexical_relevance(&query_terms, &fragment.content),
priority: fragment.priority.unwrap_or(0),
dup,
abstract_collapse,
media: media_analysis,
superseded,
}
})
.collect();
retain_safe_duplicates(&mut analyses);
analyses
}
fn retain_safe_duplicates(analyses: &mut [Analysis<'_>]) {
for index in 0..analyses.len() {
let Some(dup) = analyses[index].dup else {
continue;
};
let twin_verbatim = matches!(
analyses[dup.kept_seq].rule.action,
ContextAction::Preserve | ContextAction::Cache
);
let critical_near = dup.kind == DupKind::Near && analyses[index].rule.critical;
if !twin_verbatim || critical_near {
analyses[index].dup = None;
}
}
}
fn pack_items(
analyses: &[Analysis],
policy: &CompilePolicy,
usable: u64,
estimator: &dyn TokenEstimator,
) -> Vec<PackItem> {
let chunk_policy = effective_chunk_policy(policy, usable, estimator);
analyses
.iter()
.filter(|analysis| analysis.dup.is_none() && !analysis.superseded)
.map(|analysis| PackItem {
seq: analysis.seq,
critical: analysis.rule.critical,
priority: analysis.priority,
relevance: analysis.relevance,
cache: analysis.rule.action == ContextAction::Cache,
pieces: pieces(analysis, &chunk_policy, estimator),
})
.collect()
}
fn pieces(
analysis: &Analysis,
chunk_policy: &ChunkPolicy,
estimator: &dyn TokenEstimator,
) -> Vec<budget::Piece> {
if let Some(media) = &analysis.media {
let cost = media_fragment_tokens(media, analysis.original, estimator);
return vec![budget::Piece {
text: analysis.original.to_owned(),
cost: Some(cost),
}];
}
if let Some((collapsed, _normalized)) = &analysis.abstract_collapse {
return vec![budget::Piece {
text: collapsed.clone(),
cost: None,
}];
}
chunk_text(analysis.original, chunk_policy)
.into_iter()
.map(|chunk| budget::Piece {
text: chunk.text,
cost: None,
})
.collect()
}
const MIN_CHUNK_BYTES: usize = 256;
fn effective_chunk_policy(
policy: &CompilePolicy,
usable: u64,
estimator: &dyn TokenEstimator,
) -> ChunkPolicy {
let budget_bytes = usize::try_from(usable.saturating_mul(estimator.bytes_per_token_hint()))
.unwrap_or(usize::MAX);
ChunkPolicy {
max_chunk_bytes: policy
.chunk
.max_chunk_bytes
.min(budget_bytes)
.max(MIN_CHUNK_BYTES),
overlap_bytes: 0,
boundary: policy.chunk.boundary,
}
}
fn emissions(items: &[PackItem], taken: &[usize]) -> BTreeMap<usize, Emission> {
items
.iter()
.zip(taken.iter().copied())
.filter(|&(item, count)| count > 0 || item.pieces.is_empty())
.map(|(item, count)| {
(
item.seq,
Emission {
text: item.pieces[..count]
.iter()
.map(|piece| piece.text.as_str())
.collect(),
taken: count,
total: item.pieces.len(),
},
)
})
.collect()
}
fn sections(analyses: &[Analysis], emissions: &BTreeMap<usize, Emission>) -> Vec<CompiledSection> {
let mut result = Vec::new();
for kind in [SectionKind::Cache, SectionKind::Body] {
let mut blocks: Vec<&str> = Vec::new();
let mut ids: Vec<u64> = Vec::new();
for analysis in analyses {
let cache = analysis.rule.action == ContextAction::Cache;
let wanted = (kind == SectionKind::Cache) == cache;
if let Some(emission) = emissions
.get(&analysis.seq)
.filter(|emission| wanted && !emission.text.is_empty())
{
blocks.push(&emission.text);
ids.push(analysis.fragment_id);
}
}
if !blocks.is_empty() {
result.push(CompiledSection {
kind,
content: blocks.join(budget::JOINER),
fragment_ids: ids,
});
}
}
result
}
fn decision(
analysis: &Analysis,
all: &[Analysis],
emissions: &BTreeMap<usize, Emission>,
) -> ContextDecision {
let emission = emissions.get(&analysis.seq);
let (action, rule_id, risk, reason, handle) = match (&analysis.dup, emission) {
(Some(dup), _) => dup_verdict(analysis, *dup, &all[dup.kept_seq], emissions),
(None, _) if analysis.superseded => superseded_screenshot_verdict(analysis),
(None, Some(emission)) if emission.is_full() => full_verdict(analysis),
(None, Some(emission)) => partial_verdict(analysis, emission),
(None, None) => externalized_verdict(analysis),
};
ContextDecision {
fragment_id: analysis.fragment_id,
content_hash: analysis.content_hash,
action,
rule_id,
relevance: analysis.relevance,
risk,
reason,
memory_id: None,
handle,
}
}
type Verdict = (ContextAction, String, FidelityRisk, String, Option<String>);
fn critical_risk(critical: bool) -> FidelityRisk {
if critical {
FidelityRisk::High
} else {
FidelityRisk::Medium
}
}
fn dup_verdict(
analysis: &Analysis,
dup: Duplicate,
twin: &Analysis,
emissions: &BTreeMap<usize, Emission>,
) -> Verdict {
let (rule_id, variant) = match dup.kind {
DupKind::Exact => ("drop.duplicate", "exact duplicate"),
DupKind::Near => ("drop.near_duplicate", "near-duplicate"),
};
let twin_full = emissions.get(&twin.seq).is_some_and(Emission::is_full);
if twin_full {
let caption_diverges = analysis.media.is_some() && analysis.original != twin.original;
let reason = if caption_diverges {
format!(
"{variant} of fragment #{} — image survives through it; this fragment's differing caption does not",
dup.kept_seq
)
} else {
format!(
"{variant} of fragment #{} — content survives through it",
dup.kept_seq
)
};
return (
ContextAction::Drop,
rule_id.to_owned(),
FidelityRisk::Low,
reason,
Some(provenance::handle_for(analysis.handle_hash())),
);
}
(
ContextAction::Drop,
rule_id.to_owned(),
critical_risk(analysis.rule.critical),
format!(
"{variant} of fragment #{} — but that twin was not fully emitted — recover via the handle",
dup.kept_seq
),
Some(provenance::handle_for(analysis.handle_hash())),
)
}
fn superseded_screenshot_verdict(analysis: &Analysis) -> Verdict {
(
ContextAction::Retrieve,
classify::SCREENSHOT_SUPERSEDED_RULE_ID.to_owned(),
FidelityRisk::Medium,
classify::SCREENSHOT_SUPERSEDED_REASON.to_owned(),
Some(provenance::handle_for(analysis.handle_hash())),
)
}
fn full_verdict(analysis: &Analysis) -> Verdict {
let risk = if analysis.rule.action == ContextAction::Abstract {
FidelityRisk::Medium
} else {
FidelityRisk::Low
};
(
analysis.rule.action,
analysis.rule.id.to_owned(),
risk,
reason_with_normalization(analysis),
None,
)
}
fn partial_verdict(analysis: &Analysis, emission: &Emission) -> Verdict {
(
analysis.rule.action,
analysis.rule.id.to_owned(),
critical_risk(analysis.rule.critical),
format!(
"{} — packed {}/{} chunks, the rest stays retrievable",
reason_with_normalization(analysis),
emission.taken,
emission.total
),
Some(provenance::handle_for(analysis.handle_hash())),
)
}
fn reason_with_normalization(analysis: &Analysis) -> String {
match &analysis.abstract_collapse {
Some((_, true)) => {
format!(
"{} — timestamps normalized before collapsing",
analysis.rule.reason
)
}
_ => analysis.rule.reason.to_owned(),
}
}
fn externalized_verdict(analysis: &Analysis) -> Verdict {
(
ContextAction::Retrieve,
"budget.externalize".to_owned(),
critical_risk(analysis.rule.critical),
format!(
"did not fit the budget ({}); retrievable via its handle",
analysis.rule.reason
),
Some(provenance::handle_for(analysis.handle_hash())),
)
}
const WARNING_RELEVANCE_THRESHOLD: f32 = 0.35;
pub(crate) fn warnings_for(decisions: &[ContextDecision]) -> Vec<ContextWarning> {
decisions
.iter()
.filter(|decision| {
decision.action == ContextAction::Retrieve
&& decision.relevance >= WARNING_RELEVANCE_THRESHOLD
})
.map(|decision| ContextWarning {
fragment_id: decision.fragment_id,
action: decision.action,
relevance: decision.relevance,
reason: decision.reason.clone(),
})
.collect()
}
pub(crate) fn apply_slim(mut compiled: CompiledContext, policy: &CompilePolicy) -> CompiledContext {
if policy.slim_response {
compiled.sections.clear();
compiled.decisions.clear();
}
compiled
}
fn retrieval_handles(analyses: &[Analysis], decisions: &[ContextDecision]) -> Vec<RetrievalHandle> {
analyses
.iter()
.zip(decisions)
.filter(|(_, decision)| decision.action == ContextAction::Retrieve)
.map(|(analysis, _)| RetrievalHandle {
handle: provenance::handle_for(analysis.handle_hash()),
fragment_id: analysis.fragment_id,
estimated_tokens: analysis.tokens,
})
.collect()
}
fn emitted_tokens(
analysis: &Analysis,
emissions: &BTreeMap<usize, Emission>,
estimator: &dyn TokenEstimator,
) -> u64 {
let Some(emission) = emissions.get(&analysis.seq) else {
return 0;
};
if analysis.media.is_some() {
analysis.tokens
} else {
estimator.estimate(&emission.text)
}
}
fn saved_by_rule(
analyses: &[Analysis],
decisions: &[ContextDecision],
emissions: &BTreeMap<usize, Emission>,
estimator: &dyn TokenEstimator,
) -> BTreeMap<String, u64> {
let mut by_rule = BTreeMap::new();
for (analysis, decision) in analyses.iter().zip(decisions) {
let emitted = emitted_tokens(analysis, emissions, estimator);
let saved = analysis.tokens.saturating_sub(emitted);
if saved > 0 {
*by_rule.entry(decision.rule_id.clone()).or_insert(0) += saved;
}
}
by_rule
}
#[cfg(test)]
#[path = "context/media_pipeline_tests.rs"]
mod media_pipeline_tests;
#[cfg(test)]
#[path = "chunk_policy_tests.rs"]
mod chunk_policy_tests;
#[cfg(test)]
#[path = "warning_completeness_tests.rs"]
mod warning_completeness_tests;