session-logリファクタのレビュー・修正
This commit is contained in:
@@ -324,11 +324,9 @@ impl PodController {
|
||||
// Broadcast the accepted user message so every
|
||||
// subscriber (including the submitter) can
|
||||
// render the turn header + user line from a
|
||||
// single source of truth. Mirror the segments
|
||||
// into shared_state so subsequent History fetches
|
||||
// can re-attach them to the corresponding worker
|
||||
// user_message item.
|
||||
shared_state.push_user_segments(input.clone());
|
||||
// single source of truth. shared_state's
|
||||
// `user_segments` is re-synced from `pod` after
|
||||
// the run completes, so we don't push here.
|
||||
let _ = event_tx.send(Event::UserMessage {
|
||||
segments: input.clone(),
|
||||
});
|
||||
@@ -377,6 +375,7 @@ impl PodController {
|
||||
|
||||
let items = pod.worker().history().to_vec();
|
||||
shared_state.update_history(items);
|
||||
shared_state.set_user_segments(pod.user_segments().to_vec());
|
||||
shared_state.set_status(new_status);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
let _ = runtime_dir.write_history(&shared_state).await;
|
||||
@@ -435,6 +434,7 @@ impl PodController {
|
||||
|
||||
let items = pod.worker().history().to_vec();
|
||||
shared_state.update_history(items);
|
||||
shared_state.set_user_segments(pod.user_segments().to_vec());
|
||||
shared_state.set_status(new_status);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
let _ = runtime_dir.write_history(&shared_state).await;
|
||||
@@ -490,6 +490,7 @@ impl PodController {
|
||||
|
||||
let items = pod.worker().history().to_vec();
|
||||
shared_state.update_history(items);
|
||||
shared_state.set_user_segments(pod.user_segments().to_vec());
|
||||
shared_state.set_status(new_status);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
let _ = runtime_dir.write_history(&shared_state).await;
|
||||
@@ -587,6 +588,7 @@ impl PodController {
|
||||
|
||||
let items = pod.worker().history().to_vec();
|
||||
shared_state.update_history(items);
|
||||
shared_state.set_user_segments(pod.user_segments().to_vec());
|
||||
shared_state.set_status(new_status);
|
||||
let _ = runtime_dir.write_status(&shared_state).await;
|
||||
let _ = runtime_dir.write_history(&shared_state).await;
|
||||
|
||||
@@ -1147,6 +1147,13 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
))
|
||||
});
|
||||
|
||||
// Count surviving user_messages before consuming `retained_items`
|
||||
// — needed to align `self.user_segments` after the swap below.
|
||||
let retained_user_msgs = retained_items
|
||||
.iter()
|
||||
.filter(|i| i.is_user_message())
|
||||
.count();
|
||||
|
||||
// Build new history: [summary, ...auto-read, references, ...retained].
|
||||
let mut new_history = Vec::with_capacity(
|
||||
1 + auto_read_messages.len()
|
||||
@@ -1197,6 +1204,19 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
if self.scope_allocation.is_some() {
|
||||
pod_registry::update_session(&self.manifest.pod.name, new_session_id)?;
|
||||
}
|
||||
// Align user_segments with the post-compaction history. Items
|
||||
// before `retain_from` (now folded into the summary) lose their
|
||||
// segments; only the user_messages surviving in retained_items
|
||||
// keep them. They are always the trailing K entries of
|
||||
// `self.user_segments` because submissions are appended in order.
|
||||
let drop_n = self
|
||||
.user_segments
|
||||
.len()
|
||||
.saturating_sub(retained_user_msgs);
|
||||
if drop_n > 0 {
|
||||
self.user_segments.drain(..drop_n);
|
||||
}
|
||||
|
||||
let worker = self.worker.as_mut().unwrap();
|
||||
worker.set_history(new_history);
|
||||
// Anchor the prompt cache at the summary item so that Anthropic
|
||||
|
||||
@@ -275,6 +275,37 @@ async fn agents_md_not_reread_after_compact() {
|
||||
assert!(!after_third.contains("mutated"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn compact_aligns_user_segments_with_retained_history() {
|
||||
// retained_tokens=0 folds the entire conversation into the summary,
|
||||
// so retained_items has zero user_messages and self.user_segments
|
||||
// must be drained to match. A subsequent run() then appends fresh
|
||||
// segments cleanly without ghost entries from the pre-compaction era.
|
||||
let client = MockClient::new(vec![
|
||||
single_text_events("a"),
|
||||
single_text_events("b"),
|
||||
write_summary_tool_use_events("call-1", "compacted summary"),
|
||||
single_text_events("done"),
|
||||
single_text_events("c"),
|
||||
]);
|
||||
let (mut pod, _pwd) = make_pod_with_body("BODY", client).await.unwrap();
|
||||
|
||||
pod.run_text("first").await.unwrap();
|
||||
pod.run_text("second").await.unwrap();
|
||||
assert_eq!(pod.user_segments().len(), 2);
|
||||
|
||||
pod.compact(0).await.unwrap();
|
||||
assert_eq!(
|
||||
pod.user_segments().len(),
|
||||
0,
|
||||
"compact(0) folds every user_message into the summary, so segments \
|
||||
must be drained to match retained_items"
|
||||
);
|
||||
|
||||
pod.run_text("third").await.unwrap();
|
||||
assert_eq!(pod.user_segments().len(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn compact_preserves_system_prompt() {
|
||||
let client = MockClient::new(vec![
|
||||
|
||||
Reference in New Issue
Block a user