Merge branch 'tui-system-message-render' into develop
This commit is contained in:
@@ -3,6 +3,7 @@ use std::sync::Arc;
|
||||
|
||||
use llm_worker::WorkerError;
|
||||
use llm_worker::llm_client::client::LlmClient;
|
||||
use llm_worker::llm_client::types::{Item, Role};
|
||||
use session_store::Store;
|
||||
use tokio::sync::{broadcast, mpsc, oneshot};
|
||||
|
||||
@@ -19,6 +20,16 @@ use crate::spawn::registry::SpawnedPodRegistry;
|
||||
use crate::spawn::tool::spawn_pod_tool;
|
||||
use protocol::{AlertLevel, AlertSource, ErrorCode, Event, Method, RunResult, TurnResult};
|
||||
|
||||
fn is_system_message_item(item: &Item) -> bool {
|
||||
matches!(
|
||||
item,
|
||||
Item::Message {
|
||||
role: Role::System,
|
||||
..
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// PodHandle — client-facing, Clone-able
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -246,6 +257,14 @@ impl PodController {
|
||||
alerter_for_worker.alert(AlertLevel::Warn, AlertSource::Worker, message.to_owned());
|
||||
});
|
||||
|
||||
let tx = event_tx.clone();
|
||||
worker.on_history_append(move |item| {
|
||||
if is_system_message_item(item) {
|
||||
let value = serde_json::to_value(item).expect("Item is Serialize");
|
||||
let _ = tx.send(Event::SystemMessage { item: value });
|
||||
}
|
||||
});
|
||||
|
||||
// Register the builtin file-manipulation tools (Read / Write /
|
||||
// Edit / Glob / Grep / Bash). `ScopedFs` carries the pod-
|
||||
// lifetime scope/pwd; `Tracker` is session-scoped — a fresh
|
||||
|
||||
+37
-10
@@ -6,7 +6,7 @@ use llm_worker::Item;
|
||||
use llm_worker::llm_client::RequestConfig;
|
||||
use llm_worker::llm_client::client::LlmClient;
|
||||
use llm_worker::state::Mutable;
|
||||
use llm_worker::{ToolOutputLimits, UsageRecord, Worker, WorkerError, WorkerResult};
|
||||
use llm_worker::{Role, ToolOutputLimits, UsageRecord, Worker, WorkerError, WorkerResult};
|
||||
use session_store::{EntryHash, PodScopeSnapshot, SessionId, SessionStartState, Store, StoreError};
|
||||
use tracing::{info, warn};
|
||||
|
||||
@@ -557,6 +557,20 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
}
|
||||
}
|
||||
|
||||
fn broadcast_system_message_item(&self, item: &Item) {
|
||||
if !matches!(
|
||||
item,
|
||||
Item::Message {
|
||||
role: Role::System,
|
||||
..
|
||||
}
|
||||
) {
|
||||
return;
|
||||
}
|
||||
let value = serde_json::to_value(item).expect("Item is Serialize");
|
||||
self.send_event(Event::SystemMessage { item: value });
|
||||
}
|
||||
|
||||
/// Push a `Method::Notify` (or rendered `Method::PodEvent`) entry
|
||||
/// onto the pending buffer.
|
||||
///
|
||||
@@ -1466,19 +1480,29 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
+ reference_message.is_some() as usize
|
||||
+ retained_items.len(),
|
||||
);
|
||||
new_history.push(Item::system_message(format!(
|
||||
"[Compacted context summary]\n\n{summary_text}"
|
||||
)));
|
||||
let mut compact_introduced_system_messages =
|
||||
Vec::with_capacity(2 + auto_read_messages.len() + reference_message.is_some() as usize);
|
||||
let summary_message =
|
||||
Item::system_message(format!("[Compacted context summary]\n\n{summary_text}"));
|
||||
compact_introduced_system_messages.push(summary_message.clone());
|
||||
compact_introduced_system_messages.extend(auto_read_messages.iter().cloned());
|
||||
if let Some(msg) = reference_message.as_ref() {
|
||||
compact_introduced_system_messages.push(msg.clone());
|
||||
}
|
||||
let task_snapshot_message = Item::system_message(format!(
|
||||
"[Session TaskStore snapshot]\n\n{task_snapshot_text}\n\n\
|
||||
This is the complete session task list preserved across compaction. \
|
||||
The following TaskList tool result presents the same state through the tool lane."
|
||||
));
|
||||
compact_introduced_system_messages.push(task_snapshot_message.clone());
|
||||
|
||||
new_history.push(summary_message);
|
||||
new_history.extend(auto_read_messages);
|
||||
if let Some(msg) = reference_message {
|
||||
new_history.push(msg);
|
||||
}
|
||||
new_history.extend(retained_items);
|
||||
new_history.push(Item::system_message(format!(
|
||||
"[Session TaskStore snapshot]\n\n{task_snapshot_text}\n\n\
|
||||
This is the complete session task list preserved across compaction. \
|
||||
The following TaskList tool result presents the same state through the tool lane."
|
||||
)));
|
||||
new_history.push(task_snapshot_message);
|
||||
new_history.push(Item::tool_call("compact-tasklist", "TaskList", "{}"));
|
||||
new_history.push(Item::tool_result_with_content(
|
||||
"compact-tasklist",
|
||||
@@ -1531,8 +1555,11 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
self.user_segments.drain(..drop_n);
|
||||
}
|
||||
|
||||
self.worker.as_mut().unwrap().set_history(new_history);
|
||||
for item in &compact_introduced_system_messages {
|
||||
self.broadcast_system_message_item(item);
|
||||
}
|
||||
let worker = self.worker.as_mut().unwrap();
|
||||
worker.set_history(new_history);
|
||||
// Anchor the prompt cache at the summary item so that Anthropic
|
||||
// can place a durable `cache_control` breakpoint there — our
|
||||
// compact layout guarantees history[0] is the summary.
|
||||
|
||||
@@ -14,6 +14,7 @@ use async_trait::async_trait;
|
||||
use futures::Stream;
|
||||
use llm_worker::Worker;
|
||||
use llm_worker::llm_client::event::{Event as LlmEvent, ResponseStatus, StatusEvent};
|
||||
use llm_worker::llm_client::types::Item;
|
||||
use llm_worker::llm_client::{ClientError, LlmClient, Request};
|
||||
use protocol::Event;
|
||||
use session_store::FsStore;
|
||||
@@ -176,6 +177,56 @@ fn drain(rx: &mut broadcast::Receiver<Event>) -> Vec<Event> {
|
||||
out
|
||||
}
|
||||
|
||||
fn system_event_text(event: &Event) -> Option<&str> {
|
||||
match event {
|
||||
Event::SystemMessage { item } => item["content"]
|
||||
.as_array()
|
||||
.and_then(|parts| parts.iter().filter_map(|p| p["text"].as_str()).next()),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn compact_broadcasts_only_new_system_messages_not_retained_ones() {
|
||||
let client = MockClient::new(vec![
|
||||
single_text_events("hi"),
|
||||
write_summary_tool_use_events("call-1", "summary"),
|
||||
single_text_events("done"),
|
||||
]);
|
||||
let mut pod = make_pod(client).await;
|
||||
|
||||
let (tx, mut rx) = broadcast::channel::<Event>(64);
|
||||
pod.attach_event_tx(tx);
|
||||
|
||||
pod.run_text("first").await.unwrap();
|
||||
let retained_message = Item::system_message("[Retained system]\nold");
|
||||
pod.worker_mut().push_item(retained_message);
|
||||
let _ = drain(&mut rx);
|
||||
|
||||
pod.compact(10_000).await.unwrap();
|
||||
|
||||
let events = drain(&mut rx);
|
||||
let system_texts: Vec<&str> = events.iter().filter_map(system_event_text).collect();
|
||||
assert!(
|
||||
system_texts
|
||||
.iter()
|
||||
.any(|text| text.starts_with("[Compacted context summary]")),
|
||||
"summary system message missing from {system_texts:?}"
|
||||
);
|
||||
assert!(
|
||||
system_texts
|
||||
.iter()
|
||||
.any(|text| text.starts_with("[Session TaskStore snapshot]")),
|
||||
"task snapshot system message missing from {system_texts:?}"
|
||||
);
|
||||
assert!(
|
||||
!system_texts
|
||||
.iter()
|
||||
.any(|text| text.starts_with("[Retained system]")),
|
||||
"retained system message should not be rebroadcast: {system_texts:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn post_run_compact_success_broadcasts_start_and_done() {
|
||||
// Responses: (1) first run returns short text, (2) compact worker
|
||||
|
||||
Reference in New Issue
Block a user