use super::*;
#[derive(Component, Clone)]
pub struct CompactionSettings(pub leviath_core::CompactionConfig);
#[derive(Component, Debug, Clone, Copy, PartialEq, Eq)]
pub struct AwaitingCompaction;
#[derive(Resource)]
pub struct CompactionResults(pub UnboundedReceiver<CompactionOutcome>);
pub(crate) const EVICTION_THRESHOLD: f32 = 0.9;
#[allow(clippy::type_complexity)]
pub fn dispatch_compaction(
mut agents: Query<
(Entity, &AgentState, &mut ContextWindow, &CompactionSettings),
(With<ReadyToInfer>, Without<AwaitingCompaction>),
>,
stage: Res<InferenceStage>,
providers: Res<Providers>,
mut commands: Commands,
) {
crate::tick_scope::clear();
for (entity, state, mut window, settings) in agents.iter_mut() {
crate::tick_scope::enter(entity);
if state.status != AgentStatus::Active {
continue; }
if !window.needs_eviction(EVICTION_THRESHOLD) {
continue; }
let target_free = window.max_tokens / 10;
let Ok(eviction) = window.try_evict(target_free) else {
continue; };
let config = &settings.0;
let mut requests = Vec::new();
for region_name in &eviction.needs_compaction {
let region = window
.get_region(region_name)
.expect("needs_compaction region present: named by try_evict's own scan");
let content: String = region
.content
.iter()
.map(|e| e.content.as_str())
.collect::<Vec<_>>()
.join("\n\n");
if content.is_empty() {
continue; }
requests.push((
region_name.clone(),
compaction_request(config, &content, region_name),
));
}
if requests.is_empty() {
continue; }
let Some(provider) = providers.0.get(&config.provider) else {
continue; };
let Some(permit) = stage.pools.try_acquire(&config.model) else {
continue; };
stage.runtime.spawn(run_compaction_job(
CompactionJob {
entity,
provider,
requests,
permit,
},
std::time::Duration::from_secs(leviath_providers::DEFAULT_INFERENCE_TIMEOUT_SECS),
stage.compaction_outcomes.clone(),
stage.wake.clone(),
));
commands
.entity(entity)
.remove::<ReadyToInfer>()
.insert(AwaitingCompaction);
}
}
pub fn collect_compaction(
mut results: ResMut<CompactionResults>,
mut agents: Query<
(
&mut ContextWindow,
Option<&mut crate::telemetry::StageActivity>,
),
With<AwaitingCompaction>,
>,
mut commands: Commands,
) {
crate::tick_scope::clear();
while let Ok(outcome) = results.0.try_recv() {
let Ok((mut window, activity)) = agents.get_mut(outcome.entity) else {
continue; };
crate::tick_scope::enter(outcome.entity);
if let Some(mut activity) = activity {
activity
.0
.push(crate::telemetry::ActivityRecord::Compaction {
success: outcome.result.is_ok(),
});
}
if let Ok(summaries) = outcome.result {
for (region_name, summary) in summaries {
let summary_tokens = leviath_core::estimate_tokens(&summary);
let history = window
.regions
.iter()
.find(|r| {
matches!(&r.kind, leviath_core::RegionKind::CompactHistory { source_region }
if source_region == ®ion_name)
})
.map(|r| r.name.clone());
if let Some(history_name) = history {
let _ = window.add_to_region(&history_name, summary, summary_tokens);
}
if let Some(region) = window.get_region_mut(®ion_name) {
region.clear();
}
}
window.current_tokens = window.calculate_tokens();
}
commands
.entity(outcome.entity)
.remove::<AwaitingCompaction>()
.insert(ReadyToInfer);
}
}
pub(crate) fn compaction_request(
config: &leviath_core::CompactionConfig,
content: &str,
region_name: &str,
) -> InferenceRequest {
InferenceRequest {
system: vec![],
messages: vec![
leviath_providers::Message {
role: "system".to_string(),
content: config.system_prompt().to_string().into(),
cache_breakpoint: false,
},
leviath_providers::Message {
role: "user".to_string(),
content: config.user_prompt(content, region_name).into(),
cache_breakpoint: false,
},
],
model: config.model.clone(),
max_tokens: config.max_summary_tokens,
temperature: config.temperature,
tools: Vec::new(),
extra: serde_json::Value::Null,
request_timeout_secs: None,
}
}
#[derive(Component, Debug, Clone)]
pub struct PendingEdgeCompact(pub Vec<String>);
pub(crate) fn is_stage_specific(kind: &leviath_core::RegionKind) -> bool {
!matches!(
kind,
leviath_core::RegionKind::Pinned
| leviath_core::RegionKind::CompactHistory { .. }
| leviath_core::RegionKind::HashMap { .. }
| leviath_core::RegionKind::Custom {
persistent: true,
..
}
)
}
pub(crate) fn apply_edge_transform(
window: &mut ContextWindow,
transform: &leviath_core::blueprint::EdgeTransform,
) -> Vec<String> {
use leviath_core::blueprint::EdgeTransform;
match transform {
EdgeTransform::Direct => Vec::new(),
EdgeTransform::Clear => {
window
.regions
.iter_mut()
.filter(|r| is_stage_specific(&r.kind))
.for_each(|r| r.clear());
window.current_tokens = window.calculate_tokens();
Vec::new()
}
EdgeTransform::Compact { .. } => window
.regions
.iter()
.filter(|r| is_stage_specific(&r.kind) && !r.content.is_empty())
.map(|r| r.name.clone())
.collect(),
EdgeTransform::Custom {
carry,
compact,
clear,
..
} => {
clear
.iter()
.filter(|n| !carry.contains(n))
.for_each(|name| {
window
.get_region_mut(name)
.into_iter()
.for_each(|r| r.clear());
});
window.current_tokens = window.calculate_tokens();
compact
.iter()
.filter(|n| !carry.contains(n))
.filter(|n| window.get_region(n).is_some_and(|r| !r.content.is_empty()))
.cloned()
.collect()
}
}
}
#[allow(clippy::type_complexity)]
pub fn dispatch_edge_compact(
mut agents: Query<
(
Entity,
&AgentState,
&ContextWindow,
&PendingEdgeCompact,
Option<&CompactionSettings>,
),
(With<ReadyToInfer>, Without<AwaitingCompaction>),
>,
stage: Res<InferenceStage>,
providers: Res<Providers>,
mut commands: Commands,
) {
crate::tick_scope::clear();
for (entity, state, window, pending, settings) in agents.iter_mut() {
crate::tick_scope::enter(entity);
if state.status != AgentStatus::Active {
continue; }
let started = settings
.and_then(|s| {
let config = &s.0;
let requests = build_edge_compact_requests(window, &pending.0, config)?;
let provider = providers.0.get(&config.provider)?;
let permit = stage.pools.try_acquire(&config.model)?;
stage.runtime.spawn(run_compaction_job(
CompactionJob {
entity,
provider,
requests,
permit,
},
std::time::Duration::from_secs(
leviath_providers::DEFAULT_INFERENCE_TIMEOUT_SECS,
),
stage.compaction_outcomes.clone(),
stage.wake.clone(),
));
Some(())
})
.is_some();
let mut ec = commands.entity(entity);
ec.remove::<PendingEdgeCompact>();
if started {
ec.remove::<ReadyToInfer>().insert(AwaitingCompaction);
}
}
}
pub(crate) fn build_edge_compact_requests(
window: &ContextWindow,
regions: &[String],
config: &leviath_core::CompactionConfig,
) -> Option<Vec<(String, InferenceRequest)>> {
let requests: Vec<(String, InferenceRequest)> = regions
.iter()
.filter_map(|name| {
let region = window.get_region(name)?;
let content = region
.content
.iter()
.map(|e| e.content.as_str())
.collect::<Vec<_>>()
.join("\n\n");
(!content.is_empty())
.then(|| (name.clone(), compaction_request(config, &content, name)))
})
.collect();
(!requests.is_empty()).then_some(requests)
}