Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 56 additions & 5 deletions crates/arcan/arcan-core/src/context_compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,21 +6,35 @@ pub enum ContextBlockKind {
Persona,
Rules,
Memory,
/// Compressed summary of earlier-in-session conversation turns (BRO-425).
///
/// Produced by [`crate::summarization::SummarizationMiddleware`] when a
/// long session's context exceeds the token threshold: older turns are
/// summarized to "summary + key decisions" while the most recent N turns
/// stay at full fidelity. The full history is never lost — it remains in
/// Lago's event journal; only the in-context copy is compressed.
Compressed,
Retrieval,
Workspace,
Task,
}

impl ContextBlockKind {
/// Fixed assembly order index: Persona=0, Rules=1, Memory=2, Retrieval=3, Workspace=4, Task=5
/// Fixed assembly order index:
/// Persona=0, Rules=1, Memory=2, Compressed=3, Retrieval=4, Workspace=5, Task=6.
///
/// Compressed session history sits just after long-term Memory and before
/// Retrieval — grouping the "what happened before" context ahead of the
/// task-specific blocks.
fn order(self) -> u8 {
match self {
Self::Persona => 0,
Self::Rules => 1,
Self::Memory => 2,
Self::Retrieval => 3,
Self::Workspace => 4,
Self::Task => 5,
Self::Compressed => 3,
Self::Retrieval => 4,
Self::Workspace => 5,
Self::Task => 6,
}
}
}
Expand Down Expand Up @@ -53,6 +67,7 @@ impl Default for ContextCompilerConfig {
(ContextBlockKind::Persona, 2_000),
(ContextBlockKind::Rules, 5_000),
(ContextBlockKind::Memory, 8_000),
(ContextBlockKind::Compressed, 6_000),
(ContextBlockKind::Retrieval, 6_000),
(ContextBlockKind::Workspace, 5_000),
(ContextBlockKind::Task, 4_000),
Expand Down Expand Up @@ -338,7 +353,43 @@ mod tests {
fn default_config_reasonable() {
let config = ContextCompilerConfig::default();
assert_eq!(config.total_budget, 30_000);
assert_eq!(config.block_budgets.len(), 6);
assert_eq!(config.block_budgets.len(), 7);
}

#[test]
fn compressed_block_assembles_after_memory_before_retrieval() {
let blocks = vec![
make_block(ContextBlockKind::Task, "current task", 50),
make_block(ContextBlockKind::Persona, "I am an AI", 255),
make_block(ContextBlockKind::Retrieval, "retrieved docs", 90),
make_block(
ContextBlockKind::Compressed,
"summary of earlier turns",
120,
),
make_block(ContextBlockKind::Memory, "long-term memory", 100),
];
let config = ContextCompilerConfig {
total_budget: 100_000,
block_budgets: Vec::new(),
};
let result = compile_context(&blocks, &config);
assert_eq!(result.system_messages.len(), 5);
// Persona, Memory, Compressed, Retrieval, Task
assert!(result.system_messages[0].content.contains("I am an AI"));
assert!(
result.system_messages[1]
.content
.contains("long-term memory")
);
assert!(
result.system_messages[2]
.content
.contains("summary of earlier turns")
);
assert!(result.system_messages[3].content.contains("retrieved docs"));
assert!(result.system_messages[4].content.contains("current task"));
assert!(result.dropped_blocks.is_empty());
}

#[test]
Expand Down
5 changes: 5 additions & 0 deletions crates/arcan/arcan-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ pub mod protocol_bridge;
pub mod queue;
pub mod runtime;
pub mod state;
pub mod summarization;

pub use context::{CompactionResult, ContextConfig, compact_messages, estimate_tokens};
pub use context_compiler::{
Expand All @@ -21,3 +22,7 @@ pub use lifecycle::LifecycleHook;
pub use protocol::*;
pub use runtime::*;
pub use state::{AppState, StateError};
pub use summarization::{
COMPRESSED_SUMMARY_HEADER, CompressionOutcome, HeuristicSummarizer, SummarizationConfig,
SummarizationMiddleware, Summarizer, compressed_block,
};
115 changes: 115 additions & 0 deletions crates/arcan/arcan-core/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1229,6 +1229,121 @@ mod tests {
assert_eq!(output.state.data["last_echo"], "rewritten input");
}

#[test]
fn summarization_middleware_compresses_provider_request() {
use crate::summarization::{
COMPRESSED_SUMMARY_HEADER, SummarizationConfig, SummarizationMiddleware,
};

// A provider that records the messages of the call it receives.
struct RecordingProvider {
seen: Mutex<Vec<ChatMessage>>,
}
impl Provider for RecordingProvider {
fn name(&self) -> &str {
"recording"
}
fn complete(&self, request: &ProviderRequest) -> Result<ModelTurn, CoreError> {
*self.seen.lock().unwrap() = request.messages.clone();
Ok(ModelTurn {
directives: vec![ModelDirective::FinalAnswer {
text: "done".to_string(),
}],
stop_reason: ModelStopReason::EndTurn,
usage: None,
telemetry: None,
})
}
}

let provider = Arc::new(RecordingProvider {
seen: Mutex::new(Vec::new()),
});

// Build a long history that exceeds the token threshold.
let mut messages = vec![ChatMessage::system("You are an agent.")];
messages.push(ChatMessage::user(
"IMPORTANT: the API key env var is ACME_TOKEN.",
));
messages.push(ChatMessage::assistant("Understood, will use ACME_TOKEN."));
for i in 0..30 {
messages.push(ChatMessage::user(format!("step {i} {}", "y".repeat(400))));
messages.push(ChatMessage::assistant(format!(
"log {i} {}",
"z".repeat(400)
)));
}
messages.push(ChatMessage::user("What env var holds the API key?"));
let input_len = messages.len();

let middleware = Arc::new(SummarizationMiddleware::new(SummarizationConfig {
token_threshold: 5_000,
recent_turns: 4,
result_char_threshold: 300,
}));

let orchestrator = Orchestrator::with_turn_middlewares(
provider.clone(),
ToolRegistry::default(),
vec![middleware],
OrchestratorConfig {
max_iterations: 2,
// Disable message-drop compaction so the summarization
// middleware is the only compressor under test.
context: None,
context_compiler: None,
},
);

let output = orchestrator.run(
RunInput {
run_id: "run-1".to_string(),
session_id: "s1".to_string(),
branch_id: "main".to_string(),
messages,
state: AppState::default(),
},
|_| {},
);

assert_eq!(output.reason, RunStopReason::Completed);

let seen = provider.seen.lock().unwrap().clone();
// The provider saw a compressed list...
assert!(
seen.len() < input_len,
"provider should have received a compressed message list ({} !< {input_len})",
seen.len()
);
// ...containing the summary marker...
assert!(
seen.iter()
.any(|m| m.content.contains(COMPRESSED_SUMMARY_HEADER)),
"compressed request must carry a summary block"
);
// ...the current request verbatim...
assert!(
seen.iter()
.any(|m| m.content == "What env var holds the API key?"),
"current request must survive compression"
);
// ...and the key early decision (in the summary).
let joined: String = seen.iter().map(|m| m.content.clone()).collect();
assert!(
joined.contains("ACME_TOKEN"),
"early decision must survive in the summary"
);

// The orchestrator's durable output history is NOT compressed — full
// fidelity is preserved upstream of the per-call request.
assert!(
output.messages.len() > seen.len(),
"durable history stays uncompressed ({} > {})",
output.messages.len(),
seen.len()
);
}

#[test]
fn budget_exceeded_when_iterations_exhausted() {
// Provider always returns ToolUse but no tool call directives → continues loop
Expand Down
Loading
Loading