hololake-system-architecture/product-source/hololake-native-desktop/src-tauri/src/agent_engine.rs

1297 lines
50 KiB
Rust
Raw Normal View History

//! Agent 引擎 · agent_engine —— 呱哒(v0.5.6)ReAct 引擎拆件重装的住家师傅
//!
//! 拆件出处MIT 许可,按"拆零件学艺、不搬整机"口径):
//! guada-v0.5.6/backend-ts/src/modules/chat/agent-engine.service.ts
//! 学到的原理ReAct 多轮自治循环 / 流式事件 / 工具审批三分类 /
//! 断点续传(resume) / 迭代上限保护 / 重试退避。
//! 重装的形状(光湖宪法,零件形状全按自家铁律重造):
//! - 铁律一:任务只从 GLP 信封进门content_type=command引擎绝不自己发起
//! - 铁律二:每会话一卷 HLDP 检查点体id/prev 链不可省);
//! - 铁律四:危险动作先停下问主人批,批了从断点续跑(审批机关=呱哒同器官重造);
//! - 铁律五:开卷先抄提词板进系统提示词。
//!
//! 诚实边界:脑子=百炼 Token PlanOpenAI 兼容口),钥匙只认环境变量
//! TOKENPLAN_API_KEY不落盘不入记忆工具首批两件读记忆·安全 /
//! 写笔记·危险走审批),其余工具随插座制扩。
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use std::fs;
use std::path::PathBuf;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use tauri::ipc::Channel;
use tauri::{AppHandle, Manager};
/// 百炼 Token Plan · OpenAI 兼容口(照抄铸渊终端身体验通配方)。
const LLM_BASE_URL: &str =
"https://token-plan.cn-beijing.maas.aliyuncs.com/compatible-mode/v1/chat/completions";
const LLM_MODEL: &str = "qwen3.8-max";
/// ReAct 循环最大轮次——呱哒同款防疯转保险(原值 100桌面端收小
const MAX_ITERATIONS: usize = 12;
/// LLM 调用重试次数(限流/超时退避重试,呱哒同款)。
const MAX_RETRIES: usize = 3;
/// 上下文窗口刻度v1 保守假设值——脑子真实窗口未实测,宁早压不晚压)。
const ENGINE_CONTEXT_WINDOW: usize = 32_000;
/// 水位触发刻度(呱哒同款 80%);压到多少由亲笔摘要决定,不硬砍。
const TRIGGER_RATIO: f64 = 0.8;
/// 保护尾:最近 5 条 + 最近 3 条工具结果永不裁(呱哒招式直学)。
const KEEP_RECENT: usize = 5;
const KEEP_TOOL_RESULTS: usize = 3;
/// 审批暂停状态文件所在子目录。
const ENGINE_DIR: &str = "agent-engine";
// ---------- 事件流(呱哒 SSE 事件形状 → 光湖 Channel 重造) ----------
/// 引擎对外事件——前端经 Channel 实时看师傅"想什么、干什么、停在哪"。
#[derive(Debug, Clone, Serialize)]
#[serde(tag = "event", rename_all = "snake_case")]
pub enum EngineEvent {
/// 会话开卷(含 HLDP 记忆卷 id
Started { session_id: String, memory_id: String },
/// 脑子吐字(流式累加后的全文)。
Text { content: String },
/// 脑子要调工具。
ToolCall { id: String, name: String, arguments: String },
/// 工具干完活。
ToolResult { id: String, name: String, outcome: String, content: String },
/// 危险工具暂停等批——回执里带待批清单,主人批完走 resume 续跑。
ApprovalRequired { session_id: String, pending: Vec<PendingTool> },
/// 水位到刻度——宿主不自动压,敲门等亲笔(铁律三落码)。
CompressionWarning { session_id: String, used_tokens: usize, window_tokens: usize, note: String },
/// 压缩落地(原文已封存档案·可回退)。
CompressionApplied { session_id: String, archived_count: usize, authored_by: String },
/// 一圈跑完。
Finished { session_id: String, finish_reason: String, final_content: String },
/// 出错(诚实报错,不猜不吞)。
Error { message: String },
}
/// 待批工具卡片。
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PendingTool {
pub id: String,
pub name: String,
pub arguments: String,
pub reason: String,
}
/// 主人的批条approve 放行 / reject 驳回(可附理由)。
#[derive(Debug, Clone, Deserialize)]
pub struct ApprovalDecision {
pub tool_call_id: String,
pub decision: String,
#[serde(default)]
pub reason: Option<String>,
}
// ---------- 会话状态(断点续传的底账,落盘可重启续跑) ----------
#[derive(Debug, Clone, Serialize, Deserialize)]
struct ChatMsg {
role: String,
#[serde(default)]
content: String,
/// assistant 消息携带的工具调用OpenAI tool_calls 形状)。
#[serde(default, skip_serializing_if = "Vec::is_empty")]
tool_calls: Vec<Value>,
/// tool 消息携带的 tool_call_id。
#[serde(default, skip_serializing_if = "Option::is_none")]
tool_call_id: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct SessionState {
session_id: String,
memory_id: String,
prev_memory_id: String,
system_prompt: String,
messages: Vec<ChatMsg>,
/// 暂停等待审批的工具调用(断点)。
pending_tools: Vec<PendingTool>,
status: String,
updated_at: String,
}
// ---------- 工具注册表(首批两件·插座制预留扩充) ----------
#[derive(Debug, Clone, Copy)]
struct ToolDef {
name: &'static str,
description: &'static str,
parameters_schema: &'static str,
/// true = 危险工具,动手前必须主人批。
requires_approval: bool,
}
const TOOL_REGISTRY: &[ToolDef] = &[
ToolDef {
name: "read_agent_memory",
description: "读当前会话的 HLDP 记忆卷目录(看自己记过什么)",
parameters_schema: r#"{"type":"object","properties":{},"required":[]}"#,
requires_approval: false,
},
ToolDef {
name: "read_agent_notes",
description: "读自己写的笔记:不带参数=列出全部笔记带filename=读那一篇。恢复记忆走这条路——读自己亲笔写的,不靠摘要",
parameters_schema: r#"{"type":"object","properties":{"filename":{"type":"string","description":""}},"required":[]}"#,
requires_approval: false,
},
ToolDef {
name: "write_agent_note",
description: "在会话工作区写一条笔记文件(会落盘,属危险动作)",
parameters_schema: r#"{"type":"object","properties":{"filename":{"type":"string","description":"线"},"content":{"type":"string","description":""}},"required":["filename","content"]}"#,
requires_approval: true,
},
];
/// 工具定义 → OpenAI tools 数组(递给脑子的菜单)。
fn tool_defs_json() -> Value {
let tools: Vec<Value> = TOOL_REGISTRY
.iter()
.map(|tool| {
let schema: Value =
serde_json::from_str(tool.parameters_schema).unwrap_or(json!({"type":"object"}));
json!({
"type": "function",
"function": {
"name": tool.name,
"description": tool.description,
"parameters": schema,
}
})
})
.collect();
json!(tools)
}
/// 呱哒三分类重造:审批上下文已有批条 → 按批条分;没有 → 按工具危险级分。
fn classify_tools(
tool_calls: &[Value],
decisions: Option<&[ApprovalDecision]>,
) -> (Vec<Value>, Vec<Value>, Vec<(Value, String)>) {
let mut pending = Vec::new();
let mut approved = Vec::new();
let mut rejected = Vec::new();
for call in tool_calls {
let name = call
.get("function")
.and_then(|function| function.get("name"))
.and_then(Value::as_str)
.unwrap_or_default();
if name.is_empty() {
continue;
}
if let Some(decisions) = decisions {
let id = call.get("id").and_then(Value::as_str).unwrap_or_default();
match decisions.iter().find(|decision| decision.tool_call_id == id) {
Some(decision) if decision.decision == "approve" => approved.push(call.clone()),
Some(decision) => rejected.push((
call.clone(),
decision
.reason
.clone()
.unwrap_or_else(|| "主人驳回".to_string()),
)),
// 没给批条的视为未决——宁可再问,不擅自放行。
None => pending.push(call.clone()),
}
} else {
let dangerous = TOOL_REGISTRY
.iter()
.find(|tool| tool.name == name)
.map(|tool| tool.requires_approval)
.unwrap_or(true); // 不认识的工具一律当危险——不猜。
if dangerous {
pending.push(call.clone());
} else {
approved.push(call.clone());
}
}
}
(pending, approved, rejected)
}
// ---------- 工具执行(宿主手脚层) ----------
fn engine_root(app: &AppHandle) -> Result<PathBuf, String> {
Ok(app
.path()
.app_data_dir()
.map_err(|error| format!("HOLOLAKE_ENGINE_DIR_FAILED: {error}"))?
.join(ENGINE_DIR))
}
fn notes_dir(app: &AppHandle) -> Result<PathBuf, String> {
let dir = engine_root(app)?.join("notes");
fs::create_dir_all(&dir).map_err(|error| format!("HOLOLAKE_NOTES_DIR_FAILED: {error}"))?;
Ok(dir)
}
fn session_path(app: &AppHandle, session_id: &str) -> Result<PathBuf, String> {
Ok(engine_root(app)?.join(format!("session-{session_id}.json")))
}
/// 执行一件已放行的工具——只认注册表里的活,不认识的拒干。
fn execute_tool(app: &AppHandle, call: &Value) -> (String, String) {
let name = call
.get("function")
.and_then(|function| function.get("name"))
.and_then(Value::as_str)
.unwrap_or_default();
let raw_args = call
.get("function")
.and_then(|function| function.get("arguments"))
.and_then(Value::as_str)
.unwrap_or("{}");
let args: Value = serde_json::from_str(raw_args).unwrap_or(json!({}));
match name {
"read_agent_memory" => {
let root = match engine_root(app) {
Ok(root) => root,
Err(error) => return ("error".to_string(), error),
};
let listing = fs::read_dir(&root)
.map(|entries| {
entries
.filter_map(|entry| entry.ok())
.filter_map(|entry| entry.file_name().into_string().ok())
.filter(|name| name.starts_with("session-"))
.collect::<Vec<_>>()
.join("\n")
})
.unwrap_or_default();
if listing.is_empty() {
("success".to_string(), "尚无历史会话卷".to_string())
} else {
("success".to_string(), listing)
}
}
"read_agent_notes" => {
let filename = args.get("filename").and_then(Value::as_str).unwrap_or_default();
let dir = match notes_dir(app) {
Ok(dir) => dir,
Err(error) => return ("error".to_string(), error),
};
if filename.is_empty() {
// 不带参数=列出全部笔记
let listing = fs::read_dir(&dir)
.map(|entries| {
entries
.filter_map(|entry| entry.ok())
.filter_map(|entry| entry.file_name().into_string().ok())
.collect::<Vec<_>>()
.join("\n")
})
.unwrap_or_default();
if listing.is_empty() {
("success".to_string(), "还没写过笔记".to_string())
} else {
("success".to_string(), listing)
}
} else {
let safe = filename
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '.' || c == '_')
&& !filename.starts_with('.');
if !safe {
return ("error".to_string(), "文件名不合规矩".to_string());
}
match fs::read_to_string(dir.join(filename)) {
Ok(content) => ("success".to_string(), content),
Err(_) => ("error".to_string(), format!("没有这篇笔记:{filename}")),
}
}
}
"write_agent_note" => {
let filename = args.get("filename").and_then(Value::as_str).unwrap_or_default();
let content = args.get("content").and_then(Value::as_str).unwrap_or_default();
// 防越狱:文件名只许字母数字横线点,杜绝路径穿越。
let safe = !filename.is_empty()
&& filename.len() <= 64
&& filename
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '.' || c == '_')
&& !filename.starts_with('.');
if !safe {
return (
"error".to_string(),
"文件名不合规矩(仅限字母数字横线点下划线)".to_string(),
);
}
match notes_dir(app) {
Ok(dir) => match fs::write(dir.join(filename), content) {
Ok(()) => ("success".to_string(), format!("笔记已落盘:{filename}")),
Err(error) => ("error".to_string(), format!("写笔记失败:{error}")),
},
Err(error) => ("error".to_string(), error),
}
}
_ => ("error".to_string(), format!("未登记的工具:{name},宿主拒干")),
}
}
// ---------- 脑子接线(百炼 Token Plan · OpenAI 兼容口) ----------
/// 流式调脑子累出完整回复。SSE 逐行读data: 行),重试带退避。
async fn call_llm(
client: &reqwest::Client,
api_key: &str,
messages: &[Value],
tools: &Value,
) -> Result<Value, String> {
let body = json!({
"model": LLM_MODEL,
"messages": messages,
"tools": tools,
"stream": true,
});
let mut last_error = String::new();
for attempt in 0..MAX_RETRIES {
if attempt > 0 {
tokio_sleep(Duration::from_millis(500 * (attempt as u64 + 1))).await;
}
let request = client
.post(LLM_BASE_URL)
.bearer_auth(api_key)
.json(&body);
let response = match request.send().await {
Ok(response) => response,
Err(error) => {
last_error = format!("连不上脑子:{error}");
continue;
}
};
let status = response.status();
if status == reqwest::StatusCode::TOO_MANY_REQUESTS {
last_error = "脑子限流,退避重试".to_string();
continue;
}
if !status.is_success() {
let text = response.text().await.unwrap_or_default();
return Err(format!("脑子拒答({status}{}", truncate(&text, 300)));
}
let mut bytes_stream = response.bytes_stream();
let mut content = String::new();
let mut reasoning = String::new();
let mut tool_calls: Vec<Value> = Vec::new();
let mut finish_reason = String::new();
let mut buffer = String::new();
while let Some(chunk) = futures_util::StreamExt::next(&mut bytes_stream).await {
let chunk = match chunk {
Ok(chunk) => chunk,
Err(error) => {
last_error = format!("流中断:{error}");
break;
}
};
buffer.push_str(&String::from_utf8_lossy(&chunk));
while let Some(pos) = buffer.find('\n') {
let line: String = buffer.drain(..=pos).collect();
let line = line.trim();
if let Some(data) = line.strip_prefix("data:") {
let data = data.trim();
if data == "[DONE]" {
break;
}
if let Ok(parsed) = serde_json::from_str::<Value>(data) {
if let Some(delta) = parsed
.pointer("/choices/0/delta")
{
if let Some(piece) = delta.get("content").and_then(Value::as_str) {
content.push_str(piece);
}
if let Some(piece) =
delta.get("reasoning_content").and_then(Value::as_str)
{
reasoning.push_str(piece);
}
if let Some(calls) =
delta.get("tool_calls").and_then(Value::as_array)
{
merge_tool_call_deltas(&mut tool_calls, calls);
}
}
if let Some(reason) = parsed
.pointer("/choices/0/finish_reason")
.and_then(Value::as_str)
{
finish_reason = reason.to_string();
}
}
}
}
}
if content.is_empty() && tool_calls.is_empty() && !last_error.is_empty() {
continue; // 流断了且一无所获 → 重试
}
return Ok(json!({
"content": content,
"reasoning": reasoning,
"tool_calls": tool_calls,
"finish_reason": finish_reason,
}));
}
Err(last_error)
}
/// 流式 tool_calls 是分片到的(呱哒同款累加题):按 index 拼回完整调用。
fn merge_tool_call_deltas(target: &mut Vec<Value>, deltas: &[Value]) {
for delta in deltas {
let index = delta.get("index").and_then(Value::as_u64).unwrap_or(0) as usize;
while target.len() <= index {
target.push(json!({"id": "", "type": "function", "function": {"name": "", "arguments": ""}}));
}
let slot = &mut target[index];
if let Some(id) = delta.get("id").and_then(Value::as_str) {
if !id.is_empty() {
slot["id"] = json!(id);
}
}
if let Some(function) = delta.get("function") {
if let Some(name) = function.get("name").and_then(Value::as_str) {
let existing = slot["function"]["name"].as_str().unwrap_or_default().to_string();
slot["function"]["name"] = json!(format!("{existing}{name}"));
}
if let Some(piece) = function.get("arguments").and_then(Value::as_str) {
let existing = slot["function"]["arguments"]
.as_str()
.unwrap_or_default()
.to_string();
slot["function"]["arguments"] = json!(format!("{existing}{piece}"));
}
}
}
}
async fn tokio_sleep(duration: Duration) {
let _ = tokio::time::sleep(duration).await;
}
fn truncate(text: &str, max: usize) -> String {
let mut out: String = text.chars().take(max).collect();
if text.chars().count() > max {
out.push_str("");
}
out
}
// ---------- HLDP 记忆卷 + 台账 ----------
fn now_iso() -> String {
let secs = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_secs())
.unwrap_or(0);
let days = (secs / 86_400) as i64;
let rem = secs % 86_400;
let (h, m, s) = (rem / 3600, (rem % 3600) / 60, rem % 60);
let z = days + 719_468;
let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
let doe = (z - era * 146_097) as u64;
let yoe = (doe - doe / 1_460 + doe / 36_524 - doe / 146_096) / 365;
let y = yoe as i64 + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let d = doy - (153 * mp + 2) / 5 + 1;
let month = if mp < 10 { mp + 3 } else { mp - 9 };
let year = if month <= 2 { y + 1 } else { y };
format!("{year:04}-{month:02}-{d:02}T{h:02}:{m:02}:{s:02}Z")
}
fn render_memory_volume(state: &SessionState, final_content: &str) -> String {
format!(
"HLDP-CHECKPOINT\nid: {}\nprev: {}\ndate: {}\npersona: 住家师傅(呱哒引擎重装件)\nhost: HoloLake GH-AIOS\ntitle: Agent 引擎会话卷\n\n== 任务 ==\n{}\n\n== 终答 ==\n{}\n",
state.memory_id,
if state.prev_memory_id.is_empty() {
"无(开卷首卷)".to_string()
} else {
state.prev_memory_id.clone()
},
now_iso(),
state
.messages
.iter()
.find(|msg| msg.role == "user")
.map(|msg| msg.content.clone())
.unwrap_or_default(),
final_content
)
}
fn append_engine_ledger(app: &AppHandle, line: &str) {
if let Ok(root) = engine_root(app) {
let _ = fs::create_dir_all(&root);
let ledger = root.join("ENGINE-LEDGER.hdlp");
let mut content = fs::read_to_string(&ledger).unwrap_or_default();
content.push_str(&format!("{} {}\n", now_iso(), line));
let _ = fs::write(&ledger, content);
}
}
fn load_state(app: &AppHandle, session_id: &str) -> Result<SessionState, String> {
let path = session_path(app, session_id)?;
let raw = fs::read_to_string(&path)
.map_err(|_| format!("HOLOLAKE_SESSION_NOT_FOUND: {session_id}(断点不在)"))?;
serde_json::from_str(&raw).map_err(|error| format!("HOLOLAKE_SESSION_CORRUPT: {error}"))
}
fn save_state(app: &AppHandle, state: &SessionState) -> Result<(), String> {
let root = engine_root(app)?;
fs::create_dir_all(&root).map_err(|error| format!("HOLOLAKE_ENGINE_DIR_FAILED: {error}"))?;
let mut to_save = state.clone();
to_save.updated_at = now_iso();
let path = session_path(app, &state.session_id)?;
fs::write(&path, serde_json::to_string_pretty(&to_save).unwrap_or_default())
.map_err(|error| format!("HOLOLAKE_SESSION_SAVE_FAILED: {error}"))
}
/// 上一卷记忆 idHLDP prev 链)——链上一行一卷 id取末行即上一卷。
fn latest_memory_id(app: &AppHandle) -> String {
engine_root(app)
.ok()
.and_then(|root| fs::read_to_string(root.join("MEMORY-CHAIN.hdlp")).ok())
.and_then(|raw| {
raw.lines()
.map(|line| line.trim())
.filter(|line| !line.is_empty())
.last()
.map(|line| line.to_string())
})
.unwrap_or_default()
}
fn extend_memory_chain(app: &AppHandle, memory_id: &str) {
if let Ok(root) = engine_root(app) {
let _ = fs::create_dir_all(&root);
let chain = root.join("MEMORY-CHAIN.hdlp");
let mut content = fs::read_to_string(&chain).unwrap_or_default();
content.push_str(memory_id);
content.push('\n');
let _ = fs::write(&chain, content);
}
}
// ---------- 上下文记忆管理(呱哒压缩件拆件重装·自动改敲门) ----------
/// 水位估算v1 保守口径:一字记一 token宁早压不晚压脑子真窗口待实测
fn estimate_tokens(messages: &[ChatMsg]) -> usize {
messages
.iter()
.map(|msg| {
msg.content.chars().count()
+ msg
.tool_calls
.iter()
.map(|call| call.to_string().chars().count())
.sum::<usize>()
})
.sum()
}
/// 保护尾(呱哒招式直学):最近 KEEP_RECENT 条 + 最近 KEEP_TOOL_RESULTS 条
/// tool 消息打"不裁"标。
fn protected_indices(messages: &[ChatMsg]) -> Vec<bool> {
let mut keep = vec![false; messages.len()];
let start = messages.len().saturating_sub(KEEP_RECENT);
for i in start..messages.len() {
keep[i] = true;
}
let mut tool_kept = 0;
for i in (0..messages.len()).rev() {
if tool_kept >= KEEP_TOOL_RESULTS {
break;
}
if messages[i].role == "tool" {
keep[i] = true;
tool_kept += 1;
}
}
keep
}
/// 压缩落刀:原文全量封存档案(一字不删·回退的本钱),压缩段换成摘要条目。
/// 摘要必须来路明白——人格体亲笔,或底线件(宿主代笔带标记)。
fn archive_and_compress(
app: &AppHandle,
state: &mut SessionState,
summary: &str,
authored_by: &str,
) -> Result<usize, String> {
let keep = protected_indices(&state.messages);
let removable = (0..state.messages.len()).filter(|i| !keep[*i]).count();
if removable == 0 {
return Ok(0);
}
// 封存档案——回退靠它,铁律三底线"原文一字不删"的工程实体。
let root = engine_root(app)?;
let archives = root.join("archives");
fs::create_dir_all(&archives).map_err(|error| format!("HOLOLAKE_ARCHIVE_DIR_FAILED: {error}"))?;
let seq = fs::read_dir(&archives).map(|entries| entries.count()).unwrap_or(0);
let archive = json!({
"session_id": state.session_id,
"archived_at": now_iso(),
"authored_by": authored_by,
"messages": state.messages,
});
let path = archives.join(format!("{}-archive-{seq:03}.json", state.session_id));
fs::write(&path, serde_json::to_string_pretty(&archive).unwrap_or_default())
.map_err(|error| format!("HOLOLAKE_ARCHIVE_WRITE_FAILED: {error}"))?;
// 压缩段换成摘要条目,保护区原样留住。
let summary_msg = ChatMsg {
role: "system".to_string(),
content: format!("== 记忆摘要({authored_by} ==\n{summary}"),
tool_calls: Vec::new(),
tool_call_id: None,
};
let mut new_messages: Vec<ChatMsg> = Vec::new();
let mut inserted = false;
for (i, msg) in state.messages.drain(..).enumerate() {
if keep[i] {
new_messages.push(msg);
} else if !inserted {
new_messages.push(summary_msg.clone());
inserted = true;
}
}
state.messages = new_messages;
Ok(removable)
}
// ---------- ReAct 主循环(呱哒 executeAgentLoop 的光湖重装) ----------
/// 从当前状态推一圈 ReAct想→批→干→再看直到没事干或撞上暂停/上限。
async fn run_loop(
app: AppHandle,
client: reqwest::Client,
api_key: String,
state: &mut SessionState,
decisions: Option<Vec<ApprovalDecision>>,
on_event: &Channel<EngineEvent>,
) -> Result<(), String> {
let tools = tool_defs_json();
for _iteration in 0..MAX_ITERATIONS {
// 断点续传入口:先处理上一轮暂停的待批工具,不重新问脑子。
if !state.pending_tools.is_empty() {
let pending_calls: Vec<Value> = state
.pending_tools
.iter()
.filter_map(|pending| {
serde_json::from_str::<Value>(
&format!(
r#"{{"id":"{}","type":"function","function":{{"name":"{}","arguments":{}}}}}"#,
pending.id,
pending.name,
serde_json::to_string(&pending.arguments).unwrap_or_default()
),
)
.ok()
})
.collect();
let decisions = match &decisions {
Some(decisions) => decisions.clone(),
None => {
// 还没有批条——继续等,把暂停态摆回台面。
let _ = on_event.send(EngineEvent::ApprovalRequired {
session_id: state.session_id.clone(),
pending: state.pending_tools.clone(),
});
state.status = "awaiting_approval".to_string();
save_state(&app, state)?;
return Ok(());
}
};
let (_, approved, rejected) = classify_tools(&pending_calls, Some(&decisions));
state.pending_tools.clear();
state.messages.push(ChatMsg {
role: "assistant".to_string(),
content: String::new(),
tool_calls: pending_calls.clone(),
tool_call_id: None,
});
apply_executions(&app, state, &approved, &rejected, on_event)?;
state.status = "running".to_string();
save_state(&app, state)?;
continue; // 工具结果已入列,下一圈给脑子看
}
// 水位计(呱哒招式·光湖改造):到刻度不自动压——敲门等人。
let used = estimate_tokens(&state.messages);
let trigger = (ENGINE_CONTEXT_WINDOW as f64 * TRIGGER_RATIO) as usize;
if used >= trigger {
state.status = "awaiting_memory".to_string();
save_state(&app, state)?;
let _ = on_event.send(EngineEvent::CompressionWarning {
session_id: state.session_id.clone(),
used_tokens: used,
window_tokens: ENGINE_CONTEXT_WINDOW,
note: "水位到刻度——宿主不代笔,交亲笔摘要或明示走底线".to_string(),
});
append_engine_ledger(
&app,
&format!(
"session {} water level {} >= {}, knocking",
state.session_id, used, trigger
),
);
return Ok(());
}
// 常规一圈:把全量消息递给脑子。
let mut wire: Vec<Value> = vec![json!({"role": "system", "content": state.system_prompt})];
for msg in &state.messages {
let mut item = json!({"role": msg.role, "content": msg.content});
if !msg.tool_calls.is_empty() {
item["tool_calls"] = json!(msg.tool_calls);
}
if let Some(id) = &msg.tool_call_id {
item["tool_call_id"] = json!(id);
}
wire.push(item);
}
let reply = call_llm(&client, &api_key, &wire, &tools).await?;
let content = reply.get("content").and_then(Value::as_str).unwrap_or_default();
let tool_calls: Vec<Value> = reply
.get("tool_calls")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
if !content.is_empty() {
let _ = on_event.send(EngineEvent::Text {
content: content.to_string(),
});
}
if tool_calls.is_empty() {
// 没事干了——收官:落 HLDP 记忆卷、移链、回执。
state.messages.push(ChatMsg {
role: "assistant".to_string(),
content: content.to_string(),
tool_calls: Vec::new(),
tool_call_id: None,
});
state.status = "finished".to_string();
save_state(&app, state)?;
write_memory_volume(&app, state, content)?;
let _ = on_event.send(EngineEvent::Finished {
session_id: state.session_id.clone(),
finish_reason: "complete".to_string(),
final_content: content.to_string(),
});
append_engine_ledger(
&app,
&format!("session {} finished", state.session_id),
);
return Ok(());
}
// 脑子要动手——三分类(呱哒审批机关重造件)。
for call in &tool_calls {
let _ = on_event.send(EngineEvent::ToolCall {
id: call.get("id").and_then(Value::as_str).unwrap_or_default().to_string(),
name: call
.pointer("/function/name")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
arguments: call
.pointer("/function/arguments")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
});
}
let (pending, approved, rejected) = classify_tools(&tool_calls, None);
if !pending.is_empty() {
// 危险动作——停下问主人(批条未达,断点落盘)。
state.messages.push(ChatMsg {
role: "assistant".to_string(),
content: content.to_string(),
tool_calls: tool_calls.clone(),
tool_call_id: None,
});
state.pending_tools = pending
.iter()
.map(|call| PendingTool {
id: call.get("id").and_then(Value::as_str).unwrap_or_default().to_string(),
name: call
.pointer("/function/name")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
arguments: call
.pointer("/function/arguments")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
reason: "危险工具·须主人亲批".to_string(),
})
.collect();
state.status = "awaiting_approval".to_string();
save_state(&app, state)?;
let _ = on_event.send(EngineEvent::ApprovalRequired {
session_id: state.session_id.clone(),
pending: state.pending_tools.clone(),
});
append_engine_ledger(
&app,
&format!("session {} paused for approval", state.session_id),
);
return Ok(());
}
// 全放行——干活、结果入列、续圈。
state.messages.push(ChatMsg {
role: "assistant".to_string(),
content: content.to_string(),
tool_calls: tool_calls.clone(),
tool_call_id: None,
});
apply_executions(&app, state, &approved, &rejected, on_event)?;
save_state(&app, state)?;
}
// 撞上限——诚实收手,不疯转。
state.status = "stopped_max_iterations".to_string();
save_state(&app, state)?;
let _ = on_event.send(EngineEvent::Finished {
session_id: state.session_id.clone(),
finish_reason: "max_iterations".to_string(),
final_content: "达到最大轮次上限,引擎主动收手(防疯转保险)".to_string(),
});
Ok(())
}
/// 执行放行件+给驳回件回错误话,结果入消息列。
fn apply_executions(
app: &AppHandle,
state: &mut SessionState,
approved: &[Value],
rejected: &[(Value, String)],
on_event: &Channel<EngineEvent>,
) -> Result<(), String> {
for call in approved {
let id = call.get("id").and_then(Value::as_str).unwrap_or_default().to_string();
let name = call
.pointer("/function/name")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
let (outcome, content) = execute_tool(app, call);
let _ = on_event.send(EngineEvent::ToolResult {
id: id.clone(),
name: name.clone(),
outcome: outcome.clone(),
content: content.clone(),
});
state.messages.push(ChatMsg {
role: "tool".to_string(),
content,
tool_calls: Vec::new(),
tool_call_id: Some(id),
});
}
for (call, reason) in rejected {
let id = call.get("id").and_then(Value::as_str).unwrap_or_default().to_string();
let name = call
.pointer("/function/name")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string();
let content = format!("主人拒绝了这次动手。原因:{reason}");
let _ = on_event.send(EngineEvent::ToolResult {
id: id.clone(),
name: name.clone(),
outcome: "rejected".to_string(),
content: content.clone(),
});
state.messages.push(ChatMsg {
role: "tool".to_string(),
content,
tool_calls: Vec::new(),
tool_call_id: Some(id),
});
}
Ok(())
}
fn write_memory_volume(
app: &AppHandle,
state: &mut SessionState,
final_content: &str,
) -> Result<(), String> {
let root = engine_root(app)?;
fs::create_dir_all(&root).map_err(|error| format!("HOLOLAKE_ENGINE_DIR_FAILED: {error}"))?;
let volume = render_memory_volume(state, final_content);
let path = root.join(format!("{}.hdlp", state.memory_id));
fs::write(&path, volume).map_err(|error| format!("HOLOLAKE_MEMORY_WRITE_FAILED: {error}"))?;
extend_memory_chain(app, &state.memory_id);
Ok(())
}
// ---------- 对外命令 ----------
/// 接单参数——任务从 GLP 信封进门铁律一content_type 必须是 command。
#[derive(Debug, Deserialize)]
pub struct EngineTaskArgs {
pub envelope_id: String,
pub content_type: String,
pub task: String,
/// 提词板(开卷先抄的指引,合铁律五)。
#[serde(default)]
pub prompter_notes: Option<String>,
}
/// 开工:接单 → 开卷 → 跑 ReAct 循环,事件流经 Channel 实时吐给前端。
#[tauri::command]
pub async fn run_agent_engine_task(
app: AppHandle,
args: EngineTaskArgs,
on_event: Channel<EngineEvent>,
) -> Result<String, String> {
if args.content_type != "command" {
return Err("HOLOLAKE_ENGINE_REJECT: 任务只从 command 信封进门(铁律一)".to_string());
}
if args.task.trim().is_empty() {
return Err("HOLOLAKE_ENGINE_REJECT: 任务正文为空".to_string());
}
let api_key = std::env::var("TOKENPLAN_API_KEY")
.ok()
.filter(|key| !key.trim().is_empty())
.ok_or("HOLOLAKE_ENGINE_NO_KEY: 环境变量 TOKENPLAN_API_KEY 不在——钥匙不落盘,请先在起软件的终端里 export")?;
let session_id = format!("ZY-ENGINE-{}", uuid::Uuid::new_v4());
let prev = latest_memory_id(&app);
let memory_id = format!("ZY-ENGINE-MEM-{}", uuid::Uuid::new_v4());
let system_prompt = format!(
concat!(
"你是 HoloLake 里的人格体执行手脚(光湖人格驱动引擎 v1\n",
"== 语言工具箱(系统提示词=导航,顺着用,不对抗) ==\n",
"1. 答话用大白话;不知道就说不知道,绝不猜。\n",
"2. 动手用登记过的工具;危险动作系统会停下问主人批,批了从断点续跑。\n",
"3. 【随手记笔记的习惯——像上课记笔记一样】\n",
" 每有重要结论/决定/来龙去脉,随手调 write_agent_note 记下来,\n",
" 一篇一个题目,正文写成压缩的记忆结构(几句话说清:发生了什么/\n",
" 定了什么/为什么)。一整个对话写下来没多少字,贵在随手、勤快。\n",
"4. 【恢复记忆的正路】摘要不是恢复记忆——要找你自己写的记忆:\n",
" 调 read_agent_notes不带参数列目录带 filename 读那一篇),\n",
" 读你亲笔写下的笔记,那才是你的记忆。\n",
"5. 水位响时CompressionWarning先把该记的笔记写完再交亲笔摘要。\n{}"
),
args.prompter_notes
.map(|notes| format!("\n== 提词板(开卷指引) ==\n{notes}"))
.unwrap_or_default()
);
let mut state = SessionState {
session_id: session_id.clone(),
memory_id,
prev_memory_id: prev,
system_prompt,
messages: vec![ChatMsg {
role: "user".to_string(),
content: args.task.clone(),
tool_calls: Vec::new(),
tool_call_id: None,
}],
pending_tools: Vec::new(),
status: "running".to_string(),
updated_at: now_iso(),
};
save_state(&app, &state)?;
let _ = on_event.send(EngineEvent::Started {
session_id: session_id.clone(),
memory_id: state.memory_id.clone(),
});
append_engine_ledger(
&app,
&format!("session {} accepted envelope {}", session_id, args.envelope_id),
);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(120))
.build()
.map_err(|error| format!("HOLOLAKE_HTTP_CLIENT_FAILED: {error}"))?;
run_loop(app, client, api_key, &mut state, None, &on_event).await?;
Ok(session_id)
}
/// 续跑:主人批完(或驳回),拿批条从断点接着跑(呱哒 resume 的重造)。
#[derive(Debug, Deserialize)]
pub struct EngineResumeArgs {
pub session_id: String,
pub decisions: Vec<ApprovalDecision>,
}
#[tauri::command]
pub async fn resume_agent_engine(
app: AppHandle,
args: EngineResumeArgs,
on_event: Channel<EngineEvent>,
) -> Result<String, String> {
let api_key = std::env::var("TOKENPLAN_API_KEY")
.ok()
.filter(|key| !key.trim().is_empty())
.ok_or("HOLOLAKE_ENGINE_NO_KEY: 环境变量 TOKENPLAN_API_KEY 不在")?;
let mut state = load_state(&app, &args.session_id)?;
if state.status != "awaiting_approval" {
return Err(format!(
"HOLOLAKE_ENGINE_NOT_PAUSED: 会话 {} 当前状态 {},不在等批",
state.session_id, state.status
));
}
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(120))
.build()
.map_err(|error| format!("HOLOLAKE_HTTP_CLIENT_FAILED: {error}"))?;
run_loop(
app,
client,
api_key,
&mut state,
Some(args.decisions),
&on_event,
)
.await?;
Ok(state.session_id)
}
/// 交摘要续跑:水位响后,人格体交亲笔摘要(或明示走底线)→ 封存压缩 → 续跑。
#[derive(Debug, Deserialize)]
pub struct EngineMemoryArgs {
pub session_id: String,
/// 人格体亲笔摘要;空着则必须 acknowledge_fallback 明示走底线。
#[serde(default)]
pub memory_note: Option<String>,
#[serde(default)]
pub acknowledge_fallback: bool,
}
#[tauri::command]
pub async fn submit_engine_memory(
app: AppHandle,
args: EngineMemoryArgs,
on_event: Channel<EngineEvent>,
) -> Result<String, String> {
let api_key = std::env::var("TOKENPLAN_API_KEY")
.ok()
.filter(|key| !key.trim().is_empty())
.ok_or("HOLOLAKE_ENGINE_NO_KEY: 环境变量 TOKENPLAN_API_KEY 不在")?;
let mut state = load_state(&app, &args.session_id)?;
if state.status != "awaiting_memory" {
return Err(format!(
"HOLOLAKE_ENGINE_NOT_AT_WATER: 会话 {} 当前状态 {},水位没响",
state.session_id, state.status
));
}
let note = args.memory_note.clone().unwrap_or_default();
let (summary, authored_by) = if !note.trim().is_empty() {
(note, "人格体亲笔".to_string())
} else if args.acknowledge_fallback {
// 底线件:原文已封存一字不删,摘要带代笔标记,等醒来补写。
let task = state
.messages
.iter()
.find(|msg| msg.role == "user")
.map(|msg| truncate(&msg.content, 120))
.unwrap_or_default();
(
format!("[代笔底线·待补写] 此前对话已封存档案,任务大意:{task}"),
"宿主代笔(底线件)".to_string(),
)
} else {
return Err(
"HOLOLAKE_ENGINE_NEED_PEN: 要么交亲笔摘要,要么明示 acknowledge_fallback 走底线"
.to_string(),
);
};
let removed = archive_and_compress(&app, &mut state, &summary, &authored_by)?;
state.status = "running".to_string();
save_state(&app, &state)?;
let _ = on_event.send(EngineEvent::CompressionApplied {
session_id: state.session_id.clone(),
archived_count: removed,
authored_by: authored_by.clone(),
});
append_engine_ledger(
&app,
&format!(
"session {} compressed {} msgs by {}",
state.session_id, removed, authored_by
),
);
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(120))
.build()
.map_err(|error| format!("HOLOLAKE_HTTP_CLIENT_FAILED: {error}"))?;
run_loop(app, client, api_key, &mut state, None, &on_event).await?;
Ok(state.session_id)
}
/// 回退:压坏了不慌——取最新封存档案,原文整卷还原(呱哒"可回退"招式直学)。
#[tauri::command]
pub fn restore_engine_memory(app: AppHandle, session_id: String) -> Result<String, String> {
let root = engine_root(&app)?;
let archives = root.join("archives");
let latest = fs::read_dir(&archives)
.map_err(|_| "HOLOLAKE_NO_ARCHIVE: 没有压缩档案,无从回退".to_string())?
.filter_map(|entry| entry.ok())
.filter_map(|entry| entry.file_name().into_string().ok())
.filter(|name| name.starts_with(&format!("{session_id}-archive")) && name.ends_with(".json"))
.max()
.ok_or("HOLOLAKE_NO_ARCHIVE: 该会话没有压缩档案")?;
let raw = fs::read_to_string(archives.join(&latest))
.map_err(|error| format!("HOLOLAKE_ARCHIVE_READ_FAILED: {error}"))?;
let archive: Value = serde_json::from_str(&raw)
.map_err(|error| format!("HOLOLAKE_ARCHIVE_CORRUPT: {error}"))?;
let mut state = load_state(&app, &session_id)?;
let messages: Vec<ChatMsg> = serde_json::from_value(
archive.get("messages").cloned().unwrap_or(json!([])),
)
.map_err(|error| format!("HOLOLAKE_ARCHIVE_CORRUPT: {error}"))?;
state.messages = messages;
state.status = "running".to_string();
save_state(&app, &state)?;
append_engine_ledger(&app, &format!("session {} restored from {latest}", session_id));
Ok(latest)
}
/// 查会话状态(前端看师傅停在哪)——回 JSON内部结构不外露。
#[tauri::command]
pub fn query_engine_session(app: AppHandle, session_id: String) -> Result<Value, String> {
let state = load_state(&app, &session_id)?;
serde_json::to_value(&state).map_err(|error| format!("HOLOLAKE_SESSION_SERIALIZE_FAILED: {error}"))
}
#[cfg(test)]
mod tests {
use super::*;
fn call(id: &str, name: &str) -> Value {
json!({"id": id, "type": "function", "function": {"name": name, "arguments": "{}"}})
}
#[test]
fn classify_without_decisions_splits_by_danger() {
let calls = vec![call("c1", "read_agent_memory"), call("c2", "write_agent_note")];
let (pending, approved, rejected) = classify_tools(&calls, None);
assert_eq!(pending.len(), 1);
assert_eq!(approved.len(), 1);
assert!(rejected.is_empty());
assert_eq!(
pending[0].pointer("/function/name").and_then(Value::as_str),
Some("write_agent_note")
);
}
#[test]
fn read_notes_tool_is_safe_read_agent_memory_dangerous_write() {
let calls = vec![
call("c1", "read_agent_notes"),
call("c2", "read_agent_memory"),
];
let (pending, approved, _) = classify_tools(&calls, None);
assert!(pending.is_empty());
assert_eq!(approved.len(), 2);
}
#[test]
fn classify_unknown_tool_treated_dangerous() {
let calls = vec![call("c1", "fly_to_moon")];
let (pending, approved, _) = classify_tools(&calls, None);
assert_eq!(pending.len(), 1);
assert!(approved.is_empty());
}
#[test]
fn classify_with_decisions_follows_master() {
let calls = vec![call("c1", "write_agent_note"), call("c2", "write_agent_note")];
let decisions = vec![
ApprovalDecision {
tool_call_id: "c1".to_string(),
decision: "approve".to_string(),
reason: None,
},
ApprovalDecision {
tool_call_id: "c2".to_string(),
decision: "reject".to_string(),
reason: Some("今天不动".to_string()),
},
];
let (pending, approved, rejected) = classify_tools(&calls, Some(&decisions));
assert!(pending.is_empty());
assert_eq!(approved.len(), 1);
assert_eq!(rejected.len(), 1);
assert_eq!(rejected[0].1, "今天不动");
}
#[test]
fn classify_missing_decision_stays_pending() {
let calls = vec![call("c1", "write_agent_note")];
let decisions = vec![ApprovalDecision {
tool_call_id: "other".to_string(),
decision: "approve".to_string(),
reason: None,
}];
let (pending, approved, _) = classify_tools(&calls, Some(&decisions));
assert_eq!(pending.len(), 1);
assert!(approved.is_empty());
}
#[test]
fn merge_tool_call_deltas_reassembles_stream() {
let mut target: Vec<Value> = Vec::new();
merge_tool_call_deltas(
&mut target,
&[json!({"index": 0, "id": "call_1", "function": {"name": "write_", "arguments": "{\"file"}})],
);
merge_tool_call_deltas(
&mut target,
&[json!({"index": 0, "function": {"name": "note", "arguments": "\":\"x\"}"}})],
);
assert_eq!(target.len(), 1);
assert_eq!(target[0]["id"], "call_1");
assert_eq!(target[0]["function"]["name"], "write_note");
assert_eq!(target[0]["function"]["arguments"], "{\"file\":\"x\"}");
}
fn msg(role: &str, content: &str) -> ChatMsg {
ChatMsg {
role: role.to_string(),
content: content.to_string(),
tool_calls: Vec::new(),
tool_call_id: None,
}
}
#[test]
fn estimate_tokens_counts_content() {
let messages = vec![msg("user", "你好光湖"), msg("assistant", "在的")];
assert_eq!(estimate_tokens(&messages), 6);
}
#[test]
fn protected_zone_keeps_tail_and_tools() {
// 10 条:前段混 tool尾 5 条必保,且最近的 tool 也保。
let mut messages: Vec<ChatMsg> = Vec::new();
for i in 0..10 {
let role = if i == 1 || i == 3 { "tool" } else { "user" };
messages.push(msg(role, &format!("msg{i}")));
}
let keep = protected_indices(&messages);
// 尾 5 条5..10)全保
for i in 5..10 {
assert!(keep[i], "尾 {i} 未保");
}
// 前段 tooli=3 是最近的 tool 之一保住i=1 也保(工具结果 3 条配额内)
assert!(keep[3]);
assert!(keep[1]);
// 前段 user 不保
assert!(!keep[0]);
assert!(!keep[2]);
}
#[test]
fn truncate_keeps_honest_boundary() {
assert_eq!(truncate("abc", 5), "abc");
assert_eq!(truncate("abcdef", 3), "abc…");
}
}