diff --git a/crates/astra-turn-core/src/prompt_facing.rs b/crates/astra-turn-core/src/prompt_facing.rs index 7b077b7154..8765d71723 100644 --- a/crates/astra-turn-core/src/prompt_facing.rs +++ b/crates/astra-turn-core/src/prompt_facing.rs @@ -158,6 +158,26 @@ pub fn sanitize_prompt_facing_messages_with_state( /// orphaned tool frames are still removed at this trust boundary. pub fn sanitize_canonical_continuation_messages_with_turn_semantics( messages: Vec, +) -> Result, astra_turn_types::UserTurnSemanticsError> { + sanitize_canonical_continuation_messages_impl(messages, false) +} + +/// Project an already-compacted current-turn delta for resumable canonical +/// commit. +/// +/// Tiered compaction has already removed the compacted middle from this +/// vector. Only an explicit typed replacement in the retained tail supersedes +/// the protected head. Refinements, corrections, continuations, and unjudged +/// user messages retain the objective they depend on. +pub fn sanitize_compacted_canonical_continuation_messages_with_turn_semantics( + messages: Vec, +) -> Result, astra_turn_types::UserTurnSemanticsError> { + sanitize_canonical_continuation_messages_impl(messages, true) +} + +fn sanitize_canonical_continuation_messages_impl( + messages: Vec, + preserve_compacted_head: bool, ) -> Result, astra_turn_types::UserTurnSemanticsError> { for message in &messages { if message @@ -168,7 +188,11 @@ pub fn sanitize_canonical_continuation_messages_with_turn_semantics( } } - let start = latest_compaction_boundary_start(&messages).unwrap_or(0); + let start = if preserve_compacted_head { + compacted_canonical_turn_start(&messages)? + } else { + latest_compaction_boundary_start(&messages).unwrap_or(0) + }; let messages = messages .into_iter() .skip(start) @@ -225,8 +249,8 @@ pub fn sanitize_canonical_continuation_messages_with_turn_semantics( /// Completed tool call/result frames are execution evidence owned by the run /// transcript and recovery checkpoint, not by every future model request. This /// projection intentionally has no generic recent-message cap: the caller -/// passes one canonical turn delta, and dropping its opening user message would -/// make the resulting suffix structurally invalid. +/// passes one canonical turn delta. Its opening objective is retained unless +/// a typed replacement after compaction explicitly supersedes it. pub fn sanitize_completed_canonical_turn_messages_with_turn_semantics( messages: Vec, ) -> Result, astra_turn_types::UserTurnSemanticsError> { @@ -239,7 +263,7 @@ pub fn sanitize_completed_canonical_turn_messages_with_turn_semantics( } } - let start = latest_compaction_boundary_start(&messages).unwrap_or(0); + let start = compacted_canonical_turn_start(&messages)?; let mut out = Vec::new(); let mut has_user_context = false; for message in messages.into_iter().skip(start) { @@ -542,6 +566,26 @@ fn latest_compaction_boundary_start(messages: &[Value]) -> Option { }) } +fn compacted_canonical_turn_start( + messages: &[Value], +) -> Result { + let Some(boundary_index) = latest_compaction_boundary_start(messages) else { + return Ok(0); + }; + for (index, message) in messages.iter().enumerate().skip(boundary_index + 1).rev() { + if !is_runtime_owned_message(message) + && message.get("role").and_then(Value::as_str) == Some("user") + && canonical_text_message(message, "user", false).is_some() + && astra_turn_types::user_turn_semantics(message)?.is_some_and(|semantics| { + semantics.objective_relation == astra_turn_types::ObjectiveRelation::Replace + }) + { + return Ok(index); + } + } + Ok(0) +} + fn prompt_facing_content_for_role(role: &str, content: &str) -> Option { let _ = role; let content = content.trim().to_string(); @@ -703,6 +747,125 @@ mod tests { ); } + fn compacted_turn_projections(messages: Vec) -> Vec> { + vec![ + super::sanitize_compacted_canonical_continuation_messages_with_turn_semantics( + messages.clone(), + ) + .expect("valid compacted turn"), + super::sanitize_completed_canonical_turn_messages_with_turn_semantics(messages) + .expect("valid completed turn"), + ] + } + + fn typed_user(content: &str, relation: astra_turn_types::ObjectiveRelation) -> Value { + let mut message = json!({"role": "user", "content": content}); + astra_turn_types::mark_user_turn_semantics( + &mut message, + astra_turn_types::UserTurnSemantics::new(relation, None), + ); + message + } + + #[test] + fn compacted_current_turn_preserves_objective_until_typed_replacement() { + use astra_turn_types::ObjectiveRelation; + for relation in [ + None, + Some(ObjectiveRelation::Unknown), + Some(ObjectiveRelation::Acknowledge), + Some(ObjectiveRelation::Continue), + Some(ObjectiveRelation::Refine), + Some(ObjectiveRelation::Correct), + Some(ObjectiveRelation::Replace), + ] { + let head = typed_user("initial objective", ObjectiveRelation::Replace); + // Identical text deliberately cannot tell the projection what to do. + let tail = relation.map_or_else( + || json!({"role": "user", "content": "follow-up input"}), + |relation| typed_user("follow-up input", relation), + ); + let answer = json!({"role": "assistant", "content": "done"}); + let messages = vec![ + head.clone(), + json!({"role": "system", "content": "boundary", "_compact_boundary": true}), + tail.clone(), + answer.clone(), + ]; + let mut expected = vec![tail, answer]; + if relation != Some(ObjectiveRelation::Replace) { + expected.insert(0, head); + } + for projected in compacted_turn_projections(messages) { + assert_eq!(projected, expected, "relation: {relation:?}"); + } + } + } + + #[test] + fn compacted_current_turn_starts_at_last_explicit_replacement() { + use astra_turn_types::ObjectiveRelation; + let replacement = typed_user("replacement objective", ObjectiveRelation::Replace); + let refinement = typed_user("additional constraint", ObjectiveRelation::Refine); + let messages = vec![ + typed_user("obsolete head", ObjectiveRelation::Replace), + json!({"role": "system", "content": "boundary", "_compact_boundary": true}), + typed_user("obsolete refinement", ObjectiveRelation::Refine), + typed_user("superseded replacement", ObjectiveRelation::Replace), + json!({"role": "assistant", "content": "obsolete answer"}), + replacement.clone(), + refinement.clone(), + ]; + for projected in compacted_turn_projections(messages) { + assert_eq!(projected, vec![replacement.clone(), refinement.clone()]); + } + } + + #[test] + fn compacted_current_turn_ignores_runtime_owned_replacement() { + use astra_turn_types::ObjectiveRelation; + let head = typed_user("initial objective", ObjectiveRelation::Replace); + let mut control = + runtime_owned_message("user", "control", RuntimeMessageDelivery::EphemeralControl); + astra_turn_types::mark_user_turn_semantics( + &mut control, + astra_turn_types::UserTurnSemantics::new(ObjectiveRelation::Replace, None), + ); + let messages = vec![ + head.clone(), + json!({"role": "system", "content": "boundary", "_compact_boundary": true}), + control, + typed_user(" ", ObjectiveRelation::Replace), + ]; + for projected in compacted_turn_projections(messages) { + assert_eq!(projected, vec![head.clone()]); + } + } + + #[test] + fn compacted_current_turn_rejects_corrupt_semantics() { + let messages = vec![ + json!({"role": "user", "content": "objective"}), + json!({"role": "system", "content": "boundary", "_compact_boundary": true}), + json!({ + "role": "user", "content": "input", + (astra_turn_types::USER_TURN_SEMANTICS_FIELD): { + "schema_version": 1, "objective_relation": "invalid", + }, + }), + ]; + assert!( + super::sanitize_compacted_canonical_continuation_messages_with_turn_semantics( + messages.clone() + ) + .is_err() + ); + assert!( + super::sanitize_completed_canonical_turn_messages_with_turn_semantics(messages) + .is_err() + ); + } + #[test] fn canonical_continuation_normalizes_structured_tool_results_to_provider_neutral_text() { let messages = vec![ @@ -826,7 +989,7 @@ mod tests { } #[test] - fn recovery_drops_only_corrupt_semantics_then_runs_the_full_sanitizer() { + fn recovery_drops_corrupt_semantics_and_pre_boundary_history() { let messages = vec![ json!({"role": "user", "content": "stale objective"}), json!({"role": "system", "content": "boundary", "_compact_boundary": true}), diff --git a/crates/runtime/src/server/run/lifecycle/tests.rs b/crates/runtime/src/server/run/lifecycle/tests.rs index c8568ab54d..371c1ce121 100644 --- a/crates/runtime/src/server/run/lifecycle/tests.rs +++ b/crates/runtime/src/server/run/lifecycle/tests.rs @@ -674,6 +674,223 @@ fn compaction_commits_the_complete_replacement_projection() { assert_eq!(packs.concat(), compacted); } +#[test] +fn single_user_tool_turn_remains_committable_after_real_tiered_compaction() { + let mut messages = vec![json!({"role": "user", "content": "translate the document"})]; + for index in 0..6 { + messages.push(json!({ + "role": "assistant", + "content": null, + "tool_calls": [{ + "id": format!("call-{index}"), + "type": "function", + "function": { + "name": "bash", + "arguments": format!(r#"{{"chunk":{index}}}"#), + }, + }], + })); + messages.push(json!({ + "role": "tool", + "tool_call_id": format!("call-{index}"), + "content": format!("translated chunk {index}"), + })); + } + messages.push(json!({"role": "assistant", "content": "translation complete"})); + + let mut engine = crate::turn::CompactionEngine::new(); + engine.add_layer(Box::new(crate::turn::cloud::TieredCompaction::new(2, 0.0))); + let outcome = engine.compress_if_needed( + &mut messages, + &crate::turn::TokenBudget { + max_prompt_tokens: 64_000, + last_measured_tokens: 100_000, + current_round_index: None, + now_secs: 10_000_000, + }, + ); + + assert!( + outcome.total_tokens_freed > 0, + "the real pipeline must compact" + ); + let boundary_index = messages + .iter() + .position(|message| message["_compact_boundary"].as_bool() == Some(true)) + .expect("tiered compaction must leave its canonical boundary marker"); + assert_eq!( + boundary_index, 1, + "the protected first user precedes the boundary" + ); + assert!( + messages[boundary_index + 1..] + .iter() + .all(|message| message["role"] != "user"), + "one agentic turn has no second user anchor after compaction" + ); + + let (mode, packs) = canonical_commit_delta(&[], false, &messages, None, false) + .expect("a compacted single-user turn must remain committable") + .expect("a completed turn must produce a canonical delta"); + let committed = packs.concat(); + + assert_eq!(mode, astra_turn_types::CanonicalDeltaModeV1::Append); + assert_eq!( + committed.first(), + Some(&json!({ + "role": "user", + "content": "translate the document", + })) + ); + assert_eq!( + committed.last(), + Some(&json!({ + "role": "assistant", + "content": "translation complete", + })) + ); + assert!( + committed + .iter() + .all(|message| message.get("_compact_boundary").is_none()), + "the marker is runtime state, not canonical conversation content" + ); + + let (_, resumable_packs) = canonical_commit_delta(&[], false, &messages, None, true) + .expect("a compacted interrupted turn must remain committable") + .expect("an interrupted turn must preserve its resumable canonical delta"); + let resumable = resumable_packs.concat(); + assert_eq!( + resumable.first(), + Some(&json!({ + "role": "user", + "content": "translate the document", + })) + ); + assert!( + resumable.iter().any(|message| message["role"] == "tool"), + "an interrupted turn retains complete tool evidence for resume" + ); +} + +#[test] +fn typed_objective_relations_survive_real_tiered_compaction() { + use astra_turn_core::resume_hydration::objective_context_from_messages; + use astra_turn_types::{ObjectiveRelation, UserTurnSemantics, mark_user_turn_semantics}; + + for relation in [ + None, + Some(ObjectiveRelation::Unknown), + Some(ObjectiveRelation::Acknowledge), + Some(ObjectiveRelation::Continue), + Some(ObjectiveRelation::Refine), + Some(ObjectiveRelation::Correct), + Some(ObjectiveRelation::Replace), + ] { + let mut head = json!({"role": "user", "content": "initial objective"}); + mark_user_turn_semantics( + &mut head, + UserTurnSemantics::new(ObjectiveRelation::Replace, None), + ); + let mut messages = vec![head.clone()]; + for index in 0..6 { + messages.push(json!({ + "role": "assistant", "content": null, + "tool_calls": [{ + "id": format!("call-{index}"), "type": "function", + "function": {"name": "bash", "arguments": "{}"}, + }], + })); + messages.push(json!({ + "role": "tool", "tool_call_id": format!("call-{index}"), + "content": format!("result {index}"), + })); + } + let mut tail = json!({"role": "user", "content": "follow-up input"}); + if let Some(relation) = relation { + mark_user_turn_semantics(&mut tail, UserTurnSemantics::new(relation, None)); + } + messages.push(tail.clone()); + messages.push(json!({ + "role": "assistant", "content": null, + "tool_calls": [{ + "id": "tail-call", "type": "function", + "function": {"name": "bash", "arguments": "{}"}, + }], + })); + messages + .push(json!({"role": "tool", "tool_call_id": "tail-call", "content": "tail result"})); + messages.push(json!({"role": "assistant", "content": "done"})); + let expected_objective = objective_context_from_messages(&messages).unwrap(); + + // Admit an actual prefix, then prove the real compaction rewrites it. + let prior = messages[..7].to_vec(); + let root = astra_turn_types::canonical_conversation_root(&prior); + let mut proof = CanonicalRewriteProof::new(&prior, &root, 0); + let permit = proof.begin(&messages); + let mut engine = crate::turn::CompactionEngine::new(); + engine.add_layer(Box::new(crate::turn::cloud::TieredCompaction::new(2, 0.0))); + let outcome = engine.compress_if_needed( + &mut messages, + &crate::turn::TokenBudget { + max_prompt_tokens: 64_000, + last_measured_tokens: 100_000, + current_round_index: None, + now_secs: 10_000_000, + }, + ); + proof.finish(permit, &messages); + + assert!(outcome.total_tokens_freed > 0, "real compaction must run"); + assert_eq!( + messages[0], head, + "the compactor retains the protected head" + ); + assert_eq!(messages[1]["_compact_boundary"], true); + assert_eq!(messages[2], tail, "typed tail survives real compaction"); + + for preserve_execution_scratch in [false, true] { + let (mode, packs) = canonical_commit_delta( + &prior, + true, + &messages, + Some(&proof), + preserve_execution_scratch, + ) + .expect("verified compacted turn") + .expect("nonempty canonical delta"); + assert_eq!(mode, astra_turn_types::CanonicalDeltaModeV1::Replace); + let committed = packs.concat(); + let users = committed + .iter() + .filter(|message| message["role"] == "user") + .cloned() + .collect::>(); + let expected_users = if relation == Some(ObjectiveRelation::Replace) { + vec![tail.clone()] + } else { + vec![head.clone(), tail.clone()] + }; + assert_eq!(users, expected_users, "relation: {relation:?}"); + assert_eq!(committed.first(), expected_users.first()); + assert_eq!(committed.last().unwrap()["content"], "done"); + assert_eq!( + committed.iter().any(|message| message["role"] == "tool"), + preserve_execution_scratch, + "retained execution evidence follows the commit contract", + ); + + let serialized = serde_json::to_string(&committed).unwrap(); + let restored: Vec = serde_json::from_str(&serialized).unwrap(); + assert_eq!( + objective_context_from_messages(&restored).unwrap(), + expected_objective, + "commit/restore must preserve typed objective semantics: {relation:?}", + ); + } + } +} + #[test] fn unexplained_canonical_prefix_shrink_remains_rejected() { let prior = vec![json!({"role": "user", "content": "committed"})]; diff --git a/crates/runtime/src/turn/canonical_commit.rs b/crates/runtime/src/turn/canonical_commit.rs index 4469a331b8..ccf9552aa6 100644 --- a/crates/runtime/src/turn/canonical_commit.rs +++ b/crates/runtime/src/turn/canonical_commit.rs @@ -114,7 +114,7 @@ pub(crate) fn canonical_commit_delta( let canonical_changed = if preserve_execution_scratch { // An interrupted turn may resume from its tool boundary, so retain // complete call/result groups until recovery has settled it. - astra_turn_core::prompt_facing::sanitize_canonical_continuation_messages_with_turn_semantics( + astra_turn_core::prompt_facing::sanitize_compacted_canonical_continuation_messages_with_turn_semantics( changed_messages.to_vec(), ) } else {