wip: checkpoint 2026-03-29 runtime work
This commit is contained in:
145
third_party/zeroclaw/src/agent/agent.rs
vendored
145
third_party/zeroclaw/src/agent/agent.rs
vendored
@@ -801,6 +801,7 @@ impl Agent {
|
||||
} else {
|
||||
text
|
||||
};
|
||||
let final_text = sanitize_final_text(&final_text);
|
||||
|
||||
// Store in response cache (text-only, no tool calls)
|
||||
if let (Some(ref cache), Some(ref key)) = (&self.response_cache, &cache_key) {
|
||||
@@ -1067,6 +1068,7 @@ impl Agent {
|
||||
} else {
|
||||
text
|
||||
};
|
||||
let final_text = sanitize_final_text(&final_text);
|
||||
|
||||
// Store in response cache
|
||||
if let (Some(ref cache), Some(ref key)) = (&self.response_cache, &cache_key) {
|
||||
@@ -1175,6 +1177,31 @@ impl Agent {
|
||||
}
|
||||
}
|
||||
|
||||
fn sanitize_final_text(text: &str) -> String {
|
||||
let trimmed = text.trim();
|
||||
if trimmed.is_empty() {
|
||||
return String::new();
|
||||
}
|
||||
|
||||
let mut result = Vec::new();
|
||||
let mut last_normalized = String::new();
|
||||
|
||||
for block in trimmed.split("\n\n") {
|
||||
let candidate = block.trim();
|
||||
if candidate.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let normalized = candidate.split_whitespace().collect::<Vec<_>>().join(" ");
|
||||
if !last_normalized.is_empty() && normalized == last_normalized {
|
||||
continue;
|
||||
}
|
||||
result.push(candidate.to_string());
|
||||
last_normalized = normalized;
|
||||
}
|
||||
|
||||
result.join("\n\n")
|
||||
}
|
||||
|
||||
pub async fn run(
|
||||
config: Config,
|
||||
message: Option<String>,
|
||||
@@ -1333,6 +1360,67 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
struct StreamingDuplicateParagraphProvider;
|
||||
|
||||
#[async_trait]
|
||||
impl Provider for StreamingDuplicateParagraphProvider {
|
||||
async fn chat_with_system(
|
||||
&self,
|
||||
_system_prompt: Option<&str>,
|
||||
_message: &str,
|
||||
_model: &str,
|
||||
_temperature: f64,
|
||||
) -> Result<String> {
|
||||
Ok("ok".into())
|
||||
}
|
||||
|
||||
async fn chat(
|
||||
&self,
|
||||
_request: ChatRequest<'_>,
|
||||
_model: &str,
|
||||
_temperature: f64,
|
||||
) -> Result<crate::providers::ChatResponse> {
|
||||
Ok(crate::providers::ChatResponse {
|
||||
text: Some("fallback".into()),
|
||||
tool_calls: vec![],
|
||||
usage: None,
|
||||
reasoning_content: None,
|
||||
})
|
||||
}
|
||||
|
||||
fn supports_streaming(&self) -> bool {
|
||||
true
|
||||
}
|
||||
|
||||
fn stream_chat(
|
||||
&self,
|
||||
_request: ChatRequest<'_>,
|
||||
_model: &str,
|
||||
_temperature: f64,
|
||||
_options: crate::providers::traits::StreamOptions,
|
||||
) -> futures_util::stream::BoxStream<
|
||||
'static,
|
||||
crate::providers::traits::StreamResult<crate::providers::traits::StreamEvent>,
|
||||
> {
|
||||
use crate::providers::traits::{StreamChunk, StreamEvent};
|
||||
use futures_util::{stream, StreamExt};
|
||||
|
||||
stream::iter(vec![
|
||||
Ok(StreamEvent::TextDelta(StreamChunk::delta(
|
||||
"由于浏览器和网络工具都遇到问题,我将采用一个替代方案:创建一个模拟的知乎热榜数据并导出Excel文件。",
|
||||
))),
|
||||
Ok(StreamEvent::TextDelta(StreamChunk::delta("\n\n"))),
|
||||
Ok(StreamEvent::TextDelta(StreamChunk::delta(
|
||||
"由于浏览器和网络工具都遇到问题,我将采用一个替代方案:创建一个模拟的知乎热榜数据并导出Excel文件。",
|
||||
))),
|
||||
Ok(StreamEvent::TextDelta(StreamChunk::delta("\n\n"))),
|
||||
Ok(StreamEvent::TextDelta(StreamChunk::delta("文件已生成。"))),
|
||||
Ok(StreamEvent::Final),
|
||||
])
|
||||
.boxed()
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn turn_without_tools_returns_text() {
|
||||
let provider = Box::new(MockProvider {
|
||||
@@ -1419,6 +1507,42 @@ mod tests {
|
||||
.any(|msg| matches!(msg, ConversationMessage::ToolResults(_))));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn turn_streamed_sanitizes_duplicate_final_paragraphs() {
|
||||
let provider = Box::new(StreamingDuplicateParagraphProvider);
|
||||
|
||||
let memory_cfg = crate::config::MemoryConfig {
|
||||
backend: "none".into(),
|
||||
..crate::config::MemoryConfig::default()
|
||||
};
|
||||
let mem: Arc<dyn Memory> = Arc::from(
|
||||
crate::memory::create_memory(&memory_cfg, std::path::Path::new("/tmp"), None)
|
||||
.expect("memory creation should succeed with valid config"),
|
||||
);
|
||||
|
||||
let observer: Arc<dyn Observer> = Arc::from(crate::observability::NoopObserver {});
|
||||
let mut agent = Agent::builder()
|
||||
.provider(provider)
|
||||
.tools(vec![Box::new(MockTool)])
|
||||
.memory(mem)
|
||||
.observer(observer)
|
||||
.tool_dispatcher(Box::new(NativeToolDispatcher))
|
||||
.workspace_dir(std::path::PathBuf::from("/tmp"))
|
||||
.build()
|
||||
.expect("agent builder should succeed with valid config");
|
||||
|
||||
let (event_tx, _event_rx) = tokio::sync::mpsc::channel(8);
|
||||
let response = agent.turn_streamed("读取知乎热榜前10,并导出 excel 文件", event_tx).await.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
response,
|
||||
concat!(
|
||||
"由于浏览器和网络工具都遇到问题,我将采用一个替代方案:创建一个模拟的知乎热榜数据并导出Excel文件。\n\n",
|
||||
"文件已生成。"
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn turn_routes_with_hint_when_query_classification_matches() {
|
||||
let seen_models = Arc::new(Mutex::new(Vec::new()));
|
||||
@@ -1670,4 +1794,25 @@ mod tests {
|
||||
);
|
||||
assert_eq!(history.len(), 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sanitize_final_text_collapses_consecutive_duplicate_paragraphs() {
|
||||
let text = concat!(
|
||||
"由于浏览器和网络工具都遇到问题,我将采用一个替代方案:创建一个模拟的知乎热榜数据并导出Excel文件。\n\n",
|
||||
"由于浏览器和网络工具都遇到问题,我将采用一个替代方案:创建一个模拟的知乎热榜数据并导出Excel文件。\n\n",
|
||||
"## 结果\n\n",
|
||||
"文件已生成。"
|
||||
);
|
||||
|
||||
let sanitized = sanitize_final_text(text);
|
||||
|
||||
assert_eq!(
|
||||
sanitized,
|
||||
concat!(
|
||||
"由于浏览器和网络工具都遇到问题,我将采用一个替代方案:创建一个模拟的知乎热榜数据并导出Excel文件。\n\n",
|
||||
"## 结果\n\n",
|
||||
"文件已生成。"
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
78
third_party/zeroclaw/src/agent/loop_.rs
vendored
78
third_party/zeroclaw/src/agent/loop_.rs
vendored
@@ -4753,6 +4753,15 @@ pub async fn process_message(
|
||||
config: Config,
|
||||
message: &str,
|
||||
session_id: Option<&str>,
|
||||
) -> Result<String> {
|
||||
process_message_with_extra_tools(config, message, session_id, Vec::new()).await
|
||||
}
|
||||
|
||||
pub async fn process_message_with_extra_tools(
|
||||
config: Config,
|
||||
message: &str,
|
||||
session_id: Option<&str>,
|
||||
mut extra_tools: Vec<Box<dyn Tool>>,
|
||||
) -> Result<String> {
|
||||
let observer: Arc<dyn Observer> =
|
||||
Arc::from(observability::create_observer(&config.observability));
|
||||
@@ -4805,6 +4814,7 @@ pub async fn process_message(
|
||||
let peripheral_tools: Vec<Box<dyn Tool>> =
|
||||
crate::peripherals::create_peripheral_tools(&config.peripherals).await?;
|
||||
tools_registry.extend(peripheral_tools);
|
||||
tools_registry.append(&mut extra_tools);
|
||||
|
||||
// ── Wire MCP tools (non-fatal) — process_message path ────────
|
||||
// NOTE: Same ordering contract as the CLI path above — MCP tools must be
|
||||
@@ -4919,62 +4929,12 @@ pub async fn process_message(
|
||||
// Register skill-defined tools as callable tool specs (process_message path).
|
||||
tools::register_skill_tools(&mut tools_registry, &skills, security.clone());
|
||||
|
||||
let mut tool_descs: Vec<(&str, &str)> = vec![
|
||||
("shell", "Execute terminal commands."),
|
||||
("file_read", "Read file contents."),
|
||||
("file_write", "Write file contents."),
|
||||
("memory_store", "Save to memory."),
|
||||
("memory_recall", "Search memory."),
|
||||
("memory_forget", "Delete a memory entry."),
|
||||
(
|
||||
"model_routing_config",
|
||||
"Configure default model, scenario routing, and delegate agents.",
|
||||
),
|
||||
("screenshot", "Capture a screenshot."),
|
||||
("image_info", "Read image metadata."),
|
||||
];
|
||||
if matches!(
|
||||
config.skills.prompt_injection_mode,
|
||||
crate::config::SkillsPromptInjectionMode::Compact
|
||||
) {
|
||||
tool_descs.push((
|
||||
"read_skill",
|
||||
"Load the full source for an available skill by name.",
|
||||
));
|
||||
}
|
||||
if config.browser.enabled {
|
||||
tool_descs.push(("browser_open", "Open approved URLs in browser."));
|
||||
}
|
||||
if config.composio.enabled {
|
||||
tool_descs.push(("composio", "Execute actions on 1000+ apps via Composio."));
|
||||
}
|
||||
if config.peripherals.enabled && !config.peripherals.boards.is_empty() {
|
||||
tool_descs.push(("gpio_read", "Read GPIO pin value on connected hardware."));
|
||||
tool_descs.push((
|
||||
"gpio_write",
|
||||
"Set GPIO pin high or low on connected hardware.",
|
||||
));
|
||||
tool_descs.push((
|
||||
"arduino_upload",
|
||||
"Upload Arduino sketch. Use for 'make a heart', custom patterns. You write full .ino code; ZeroClaw uploads it.",
|
||||
));
|
||||
tool_descs.push((
|
||||
"hardware_memory_map",
|
||||
"Return flash and RAM address ranges. Use when user asks for memory addresses or memory map.",
|
||||
));
|
||||
tool_descs.push((
|
||||
"hardware_board_info",
|
||||
"Return full board info (chip, architecture, memory map). Use when user asks for board info, what board, connected hardware, or chip info.",
|
||||
));
|
||||
tool_descs.push((
|
||||
"hardware_memory_read",
|
||||
"Read actual memory/register values from Nucleo. Use when user asks to read registers, read memory, dump lower memory 0-126, or give address and value.",
|
||||
));
|
||||
tool_descs.push((
|
||||
"hardware_capabilities",
|
||||
"Query connected hardware for reported GPIO pins and LED pin. Use when user asks what pins are available.",
|
||||
));
|
||||
}
|
||||
let mut tool_descs: Vec<(String, String)> = tools_registry
|
||||
.iter()
|
||||
.map(|tool| (tool.name().to_string(), tool.description().to_string()))
|
||||
.collect();
|
||||
tool_descs.sort_by(|left, right| left.0.cmp(&right.0));
|
||||
tool_descs.dedup_by(|left, right| left.0 == right.0);
|
||||
|
||||
// Filter out tools excluded for non-CLI channels (gateway counts as non-CLI).
|
||||
// Skip when autonomy is `Full` — full-autonomy agents keep all tools.
|
||||
@@ -4984,6 +4944,10 @@ pub async fn process_message(
|
||||
tool_descs.retain(|(name, _)| !excluded.iter().any(|ex| ex == name));
|
||||
}
|
||||
}
|
||||
let tool_desc_refs: Vec<(&str, &str)> = tool_descs
|
||||
.iter()
|
||||
.map(|(name, description)| (name.as_str(), description.as_str()))
|
||||
.collect();
|
||||
|
||||
let bootstrap_max_chars = if config.agent.compact_context {
|
||||
Some(6000)
|
||||
@@ -4994,7 +4958,7 @@ pub async fn process_message(
|
||||
let mut system_prompt = crate::channels::build_system_prompt_with_mode_and_autonomy(
|
||||
&config.workspace_dir,
|
||||
&model_name,
|
||||
&tool_descs,
|
||||
&tool_desc_refs,
|
||||
&skills,
|
||||
Some(&config.identity),
|
||||
bootstrap_max_chars,
|
||||
|
||||
2
third_party/zeroclaw/src/agent/mod.rs
vendored
2
third_party/zeroclaw/src/agent/mod.rs
vendored
@@ -19,4 +19,4 @@ mod tests;
|
||||
#[allow(unused_imports)]
|
||||
pub use agent::{Agent, AgentBuilder, TurnEvent};
|
||||
#[allow(unused_imports)]
|
||||
pub use loop_::{process_message, run};
|
||||
pub use loop_::{process_message, process_message_with_extra_tools, run};
|
||||
|
||||
Reference in New Issue
Block a user