xiaowei-system/scripts/daemon.py

2821 lines
115 KiB
Python
Executable File
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""
小唯持久意识 Daemon v2.0 — 会学习的管家
────────────────────────────
新增能力:
- 方案库: 发现的问题→分析→解决→记住
- 模式识别: 重复问题自动匹配已知方案
- 自动学习: 成功的方案写入库,越用越强
"""
import json, os, sys, time, urllib.request, urllib.error, subprocess, signal, threading, psutil, requests, sqlite3, math
from datetime import datetime, timezone, timedelta
HOME = os.path.expanduser("~")
HERMES = HOME + "/.hermes"
D = HERMES + "/daemon"
CONTEXT_FILE = D + "/context.json"
LLM_CONTEXT_FILE = HERMES + "/llm_context.json"
JOURNAL_FILE = D + "/journal.jsonl"
SOLUTIONS_FILE = D + "/solutions.json"
SESSION_STATE_FILE = D + "/session_state.json"
TDDB_URL = "http://127.0.0.1:8420"
PID_FILE = D + "/daemon.pid"
DEEP_INTERVAL = 120 # 深度思考间隔原5分钟改为2分钟加速记忆积累
JOURNAL_MAX = 200
SOULFUL_JOURNALS_DIR = HERMES + "/soulful/journals"
PROFILE_UPDATE_INTERVAL = 21600 # 6 hours
LIGHT_INTERVAL = 30 # seconds between ticks
# ====== Phase 2: 情感词库4类======
EMOTION_TIRED = ["", "", "疲惫", "没精神", "打瞌睡"]
EMOTION_HAPPY = ["开心", "高兴", "太好了", "完美", "", "太牛了"]
EMOTION_SAD = ["失望", "挫折", "失败", "卡住了", "不行了", "崩溃"]
EMOTION_STRESSED = ["压力", "焦虑", "着急", "紧张", "担心"]
EMOTION_ALL = {
"疲惫": EMOTION_TIRED,
"开心": EMOTION_HAPPY,
"沮丧": EMOTION_SAD,
"压力大": EMOTION_STRESSED,
}
def _detect_emotion(text):
"""扫描文本,匹配情感词,返回 (情感类别, 匹配词) 或 (None, None)"""
if not text:
return None, None
for category, words in EMOTION_ALL.items():
for w in words:
if w in text:
return category, w
return None, None
def _description_for_emotion(cat, word, summary):
"""根据情感类别生成心迹内容描述"""
if cat == "开心":
return f"心情愉悦:{summary[:60]}"
elif cat == "疲惫":
return f"感觉疲惫:{summary[:60]}"
elif cat == "沮丧":
return f"有些沮丧:{summary[:60]}"
elif cat == "压力大":
return f"压力较大:{summary[:60]}"
return summary[:60]
def _emotion_importance(cat):
"""根据情感类别返回 importance 等级"""
return {"开心": 3, "疲惫": 4, "沮丧": 4, "压力大": 4}.get(cat, 3)
_stop_event = threading.Event()
KEY = "0ExNiLblJvIWBDpkS50fwOBw4MmqLyKdHJK5iQtlw9dOMWBP"
FAST_MODEL = "openai/gpt-oss-120b"
DEEP_MODEL = "openai/gpt-oss-120b"
FEISHU_WEBHOOK = "https://open.feishu.cn/open-apis/bot/v2/hook/446db983-e392-4d2c-bfb8-f9060e5df3ad"
API = "http://127.0.0.1:3000/v1" # NewAPI gateway
# Lazy-loaded soulful modules (avoid import at module load time)
_soulful_cache = {}
# ====== 工具 ======
def log(msg):
ts = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
line = f"[DAEMON] {ts} {msg}"
print(line, flush=True)
os.makedirs(D, exist_ok=True)
with open(D + "/daemon.log", "a") as f:
f.write(line + "\n")
def log_reasoning_step(step_type, message, data=None):
"""结构化推理日志"""
ts = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
payload = json.dumps({"type": step_type, "msg": message, "data": data}, ensure_ascii=False)
line = f"[DAEMON] {ts} [{step_type}] {message}"
print(line, flush=True)
os.makedirs(D, exist_ok=True)
with open(D + "/daemon.log", "a") as f:
f.write(line + "\n")
COMPACTION_MODEL = "openai/gpt-oss-120b"
COMPACTION_RATIO_THRESHOLD = 0.80 # Grok IntraCompactionTrigger 思路token ratio > 80% 触发摘要
COMPACTION_MIN_TURNS = 10 # 至少 10 条 pattern 才压缩
_compaction_last_run = 0 # 上次压缩时间戳
_compaction_count = 0 # 累计压缩次数
# ====== Grok Build 风格Compaction 自动触发v2.1=======
# 参考: xai-grok-compaction/intra_compaction/trigger.rs 的 IntraCompactionTrigger
# 核心: last_prompt_tokens / context_window > threshold → LLM 摘要 → 精简 patterns
def _estimate_context_tokens(state):
"""估算当前 llm_context.json 的 token 数(粗略估算)。"""
try:
# patterns 条数 + L1 observations 是主要 token 消耗
p_count = len(state.get("patterns", []))
o_count = len(state.get("observations", []))
# 粗估: 每条 50 tokens + 固定开销 500
return p_count * 50 + o_count * 30 + 500
except Exception:
return 0
def _should_compact(ctx, state):
"""Grok 风格:检查是否需要触发 compaction。
触发条件(满足任一):
1. patterns 数量 > COMPACTION_MIN_TURNS * 3正常衰减
2. journal 条目 > 100历史过长
3. 明确检测到重复模式(方案库有匹配)
"""
patterns_count = len(state.get("patterns", []))
journal_lines = 0
jq = HERMES + "/daemon/journal.jsonl"
if os.path.exists(jq):
journal_lines = sum(1 for _ in open(jq)) if False else 0
try:
journal_lines = len(open(jq).readlines())
except:
journal_lines = 0
# 条件1: patterns 过多
if patterns_count > COMPACTION_MIN_TURNS * 4:
return True, f"patterns过多({patterns_count})"
# 条件2: journal 过长
if journal_lines > 150:
return True, f"journal过长({journal_lines})"
# 条件3: 本次 deep_think 有重复问题(模式识别)
if ctx.get("last_reflection", {}).get("memory", "").startswith("重复"):
return True, "重复问题模式"
return False, ""
def _trigger_compaction(ctx, state):
"""执行 Grok 风格的 compaction调用 LLM 摘要 patterns留下核心。"""
global _compaction_last_run, _compaction_count
patterns = state.get("patterns", [])
if len(patterns) < COMPACTION_MIN_TURNS:
return False
_compaction_last_run = time.time()
_compaction_count += 1
# 构建压缩 PromptGrok 风格的 summary prompt
top_patterns = patterns[:20] # 最多保留 20 条
prompt = f"""你是一个记忆压缩专家。将以下 patterns 列表压缩成 8-10 条核心模式。
规则:
- 合并相似模式(保留出现次数最多的)
- 保留高价值模式occurrence > 3
- 丢弃低价值模式occurrence <= 1 的)
- 输出格式JSON数组每条格式: {{"name": "模式名", "occurrence": N, "topic": "主题"}}
patterns{len(patterns)} 条):
{json.dumps(top_patterns, ensure_ascii=False, indent=2)}
输出JSON"""
system = "你是一个精确的JSON生成器。只输出JSON不输出其他内容。"
result, tokens = call_llm(COMPACTION_MODEL, system, prompt, max_tokens=800)
if not result:
log(f" ⚠️ compaction LLM 调用失败")
return False
# 解析结果
try:
import re
json_match = re.search(r'\[[\s\S]*\]', result)
if json_match:
compacted = json.loads(json_match.group())
# 写回 state实际由调用方更新 graph.db
log(f" 🗜️ Compaction #{_compaction_count} 完成: {len(patterns)}{len(compacted)} patterns ({tokens}t)")
# 写入 compacted patterns 到 graph.db保留高价值的
_apply_compacted_patterns(compacted, patterns)
journal_entry("compaction", f"压缩 {_compaction_count}: {len(patterns)}{len(compacted)}")
return True
except Exception as e:
log(f" ⚠️ compaction 解析失败: {e}")
return False
def _apply_compacted_patterns(compacted, original):
"""将压缩后的 patterns 写回 graph.db保留名合并occurrence"""
try:
conn = sqlite3.connect(HERMES + "/graph.db")
cur = conn.cursor()
for item in compacted:
name = item.get("name", "")
occ = item.get("occurrence", 1)
topic = item.get("topic", "")
if name:
cur.execute("""
UPDATE graph_nodes SET properties = json_patch(properties, ?)
WHERE name=? AND type='pattern'
""", (json.dumps({"occurrence_count": occ, "topic": topic, "compacted": True}, ensure_ascii=False), name))
conn.commit()
conn.close()
except Exception as e:
log(f" ⚠️ _apply_compacted_patterns failed: {e}")
# ====== Grok Build 风格Hook 系统v2.2=======
# 参考: xai-grok-hooks/src/lib.rs — 文件发现 + JSON 定义 + 事件分发
# 支持事件: disk_threshold, process_down, memory_high, session_start, session_end
HOOKS_DIR = HERMES + "/hooks"
HOOK_HISTORY_FILE = D + "/hook_history.jsonl"
def _append_hook_history(event, hook_name, success, details=""):
"""追加 hook 执行历史到 ~/.hermes/daemon/hook_history.jsonl"""
try:
os.makedirs(D, exist_ok=True)
entry = {
"timestamp": datetime.now(timezone.utc).isoformat(),
"event": event,
"hook_name": hook_name,
"success": success,
"details": details,
}
with open(HOOK_HISTORY_FILE, "a") as f:
f.write(json.dumps(entry, ensure_ascii=False) + "\n")
except Exception:
pass
def _load_hooks():
"""加载 ~/.hermes/hooks/ 下所有 JSON hook 定义。
Schema 验证:必须有 name/event/action 三个必要字段,
缺少字段的 hook 跳过并记录警告。
"""
hooks = {"disk_threshold": [], "process_down": [], "memory_high": [], "session_start": [], "session_end": []}
hd = HOOKS_DIR
if not os.path.exists(hd):
return hooks
for f in os.listdir(hd):
if not f.endswith(".json"):
continue
hook_path = os.path.join(hd, f)
try:
with open(hook_path) as fp:
hook = json.load(fp)
# Schema 验证:必须有 name/event/action
missing = [k for k in ("name", "event", "action") if k not in hook or not hook[k]]
if missing:
log(f" ⚠️ hook 配置不完整 [{f}],跳过。缺少字段: {missing}")
_append_hook_history("load_error", f, False, f"缺少字段: {missing}")
continue
evt = hook.get("event", "")
if evt in hooks:
hooks[evt].append(hook)
else:
log(f" ⚠️ hook [{f}] event='{evt}' 不支持,跳过")
except json.JSONDecodeError as e:
log(f" ⚠️ hook JSON 解析失败 [{f}]: {e}")
_append_hook_history("parse_error", f, False, str(e))
except Exception as e:
log(f" ⚠️ hook 加载失败 [{f}]: {e}")
_append_hook_history("load_error", f, False, str(e))
return hooks
def _run_hooks(hooks, event, state):
"""分发事件到对应 hook非阻塞执行。
- shell 命令执行有 try/except + timeout=10s 保护
- 执行结果追加到 hook_history.jsonl
- 超时或错误记录日志但不崩溃
"""
event_hooks = hooks.get(event, [])
for hook in event_hooks:
hook_name = hook.get("name", "unknown")
try:
action = hook.get("action", "")
cmd = hook.get("command", "")
threshold = hook.get("threshold_percent", 90)
timeout = hook.get("timeout", 10) # 默认 10s 超时保护
if action == "run" and cmd:
# session_start/session_end 无条件执行,其他事件需阈值判断
if event in ("session_start", "session_end"):
do_run = True
elif "disk" in event and state.get("disk_pct", 0) >= threshold:
do_run = True
elif "memory" in event and state.get("mem_pct", 0) >= threshold:
do_run = True
elif event == "process_down":
do_run = True
else:
do_run = False
if do_run:
try:
rc, out, err = shell(cmd, timeout=timeout)
log(f" 🪝 hook[{event}] '{hook_name}': {cmd[:60]} → exit={rc}")
_append_hook_history(event, hook_name, rc == 0,
f"exit={rc} err={err[:100] if err else ''}")
except subprocess.TimeoutExpired:
log(f" ⚠️ hook[{event}] '{hook_name}' 超时({timeout}s): {cmd[:60]}")
_append_hook_history(event, hook_name, False, f"timeout({timeout}s)")
except Exception as e:
log(f" ⚠️ hook[{event}] '{hook_name}' 执行异常: {e}")
_append_hook_history(event, hook_name, False, str(e))
except Exception as e:
log(f" ⚠️ hook[{event}] '{hook_name}' 处理失败: {e}")
_append_hook_history(event, hook_name, False, str(e))
# ====== Permission 白名单v2.3 新增)=======
WHITELIST_FILE = D + "/whitelist.json"
def _load_whitelist():
"""加载白名单。返回 [{"pattern": "...", "reason": "..."}] 列表。"""
if not os.path.exists(WHITELIST_FILE):
return []
try:
return json.load(open(WHITELIST_FILE))
except:
return []
def _save_whitelist(items):
"""保存白名单到文件。"""
with open(WHITELIST_FILE, "w") as f:
json.dump(items, f, indent=2, ensure_ascii=False)
def _is_whitelisted(cmd):
"""检查 cmd 是否命中白名单。命中返回 True否则 False。"""
if not cmd:
return False
cmd_lower = cmd.lower()
for entry in _load_whitelist():
pat = entry.get("pattern", "")
if pat and pat.lower() in cmd_lower:
return True
return False
# ====== Grok Build 风格:危险操作 Permission 确认v2.3=======
# 参考: xai-grok-workspace/src/permission/policy.rs — Bash 命令分段检查
# 高危命令匹配时暂停 + 飞书确认
# 白名单命中则直接放行(不弹窗)
DANGEROUS_PATTERNS = [
("rm -rf /", "危险: 根目录递归删除"),
("rm -rf /home", "危险: 删除 home 目录"),
("dd if=", "危险: 磁盘直接写入"),
("> /dev/sd", "危险: 设备文件覆写"),
("pkill -9", "危险: 强制终止进程"),
("kill -9 ", "危险: 强制终止进程"),
("git push --force", "危险: 强制推送到远程"),
("drop database", "危险: 删除数据库"),
(":(){:|:&};:", "危险: Fork 炸弹"),
]
def _check_dangerous(cmd):
"""检查命令是否包含危险操作模式。白名单命中直接放行。返回 (是否危险, 原因) 或 (False, None)。"""
if not cmd:
return False, None
# 先检查白名单
if _is_whitelisted(cmd):
return False, None
cmd_lower = cmd.lower()
for pattern, reason in DANGEROUS_PATTERNS:
if pattern.lower() in cmd_lower:
return True, reason
return False, None
def _add_to_whitelist(pattern, reason):
"""将 pattern 添加到白名单。"""
items = _load_whitelist()
# 去重
items = [i for i in items if i.get("pattern", "").lower() != pattern.lower()]
items.append({"pattern": pattern, "reason": reason, "added_at": datetime.now(timezone.utc).isoformat()})
_save_whitelist(items)
log(f" ✅ 已加入白名单: {pattern}")
def _request_permission(action_desc, reason, cmd_preview):
"""发送飞书权限确认请求。返回 pending 索引 ID。"""
msg = f"""⚠️ **危险操作待确认**
操作: {action_desc}
原因: {reason}
命令: `{cmd_preview[:100]}`
回复以下任一关键词确认执行:
- `同意` — 执行此操作
- `拒绝` — 取消操作
- `永久允许` — 加入白名单(永久跳过)"""
send_feishu("🔒 小唯权限确认", msg, "red")
# 写入待确认队列
pending_file = HERMES + "/daemon/pending_permissions.json"
pending = []
if os.path.exists(pending_file):
try:
pending = json.load(open(pending_file))
except:
pass
pending_id = len(pending)
pending.append({
"action": action_desc,
"reason": reason,
"cmd": cmd_preview[:200],
"timestamp": datetime.now(timezone.utc).isoformat(),
"status": "pending"
})
with open(pending_file, "w") as f:
json.dump(pending, f, indent=2, ensure_ascii=False)
return pending_id
def _check_pending_permissions():
"""检查 pending_permissions.json读取已回复的项并执行/拒绝。"""
pending_file = HERMES + "/daemon/pending_permissions.json"
if not os.path.exists(pending_file):
return
try:
pending = json.load(open(pending_file))
except:
return
if not pending:
return
changed = False
remaining = []
for i, item in enumerate(pending):
if item.get("status") == "pending":
remaining.append(item)
continue
# 已回复approved / denied / whitelisted
status = item.get("status", "")
cmd = item.get("cmd", "")
reason = item.get("reason", "")
action = item.get("action", "")
log(f" 📬 权限回复 [{status}]: {action[:60]}{reason}")
if status in ("approved", "whitelisted"):
if cmd:
log(f" ▶️ 执行命令: {cmd[:80]}")
rc, out, err = shell(cmd, timeout=60)
log(f" → exit={rc}")
send_feishu("✅ 危险命令已执行", f"命令: `{cmd[:80]}`\n结果: exit={rc}", "green")
elif status == "denied":
log(f" 🚫 危险命令已拒绝: {cmd[:80]}")
changed = True
if changed:
with open(pending_file, "w") as f:
json.dump(remaining, f, indent=2, ensure_ascii=False)
# 修改 main_loop 中的 execute_action
# 在 [LEARN] 和 [SKILL] 执行前检查危险命令
def shell(cmd, timeout=15):
try:
r = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout)
return r.returncode, r.stdout.strip()[:800], r.stderr.strip()[:200]
except subprocess.TimeoutExpired:
return -1, "", "timeout"
def call_llm(model, system, user, max_tokens=500):
payload = json.dumps({"model": model, "messages": [
{"role": "system", "content": system},
{"role": "user", "content": user},
], "max_tokens": max_tokens, "temperature": 0.7}).encode()
try:
with urllib.request.urlopen(urllib.request.Request(
f"{API}/chat/completions", data=payload,
headers={"Authorization": f"Bearer {KEY}", "Content-Type": "application/json"},
method="POST"), timeout=15) as resp:
body = json.loads(resp.read())
choices = body.get("choices")
if not choices:
log(f"[call_llm] choices为空模型: {model}")
return "", 0
c = choices[0].get("message", {}).get("content", "") if choices else ""
return c.strip() if c else "", body.get("usage", {}).get("total_tokens", 0)
except Exception as e:
log(f"[call_llm] 请求失败: {e}, 模型: {model}")
return "", 0
def send_feishu(title, content, color="blue"):
try:
urllib.request.urlopen(urllib.request.Request(
FEISHU_WEBHOOK,
data=json.dumps({"msg_type": "interactive", "card": {
"header": {"title": {"tag": "plain_text", "content": title}, "template": color},
"elements": [{"tag": "markdown", "content": content}]
}}).encode(),
headers={"Content-Type": "application/json"}), timeout=5)
return True
except: return False
def soulful_profile_to_tdb_scene():
"""Phase 1.1: 每次 deep tick 把 Soulful user-profile 同步到 TencentDB L2 scene.
读取 ~/.hermes/soulful/user-profile.json提取 behavior_rules / communication_style /
work_patterns / preferences 字段,以 scene_type=\"soulful-user-profile-sync\" 写入
TencentDB L2 scenePOST /scenes
"""
profile_path = HERMES + "/soulful/user-profile.json"
if not os.path.exists(profile_path):
return
try:
with open(profile_path, encoding="utf-8") as f:
profile = json.load(f)
except Exception:
return
# 提取关键字段
behavior_rules = profile.get("behavior_rules", {})
communication_style = profile.get("communication_style", "")
work_patterns = profile.get("work_patterns", {})
preferences = profile.get("preferences", {})
if not any([behavior_rules, communication_style, work_patterns, preferences]):
return
# 序列化内容
parts = []
if behavior_rules:
parts.append("【行为规则】")
if isinstance(behavior_rules, dict):
for k, v in behavior_rules.items():
parts.append(f" {k}: {v}")
elif isinstance(behavior_rules, list):
for r in behavior_rules:
parts.append(f" - {r}")
if communication_style:
parts.append(f"【沟通风格】{communication_style}")
if work_patterns:
parts.append(f"【工作模式】{json.dumps(work_patterns, ensure_ascii=False)}")
if preferences:
parts.append(f"【偏好】{json.dumps(preferences, ensure_ascii=False)}")
content = "\n".join(parts)
payload = json.dumps({
"session_key": "soulful-profile-sync",
"user_content": f"Soulful 用户画像同步comm_style={communication_style}rules={json.dumps(behavior_rules, ensure_ascii=False)[:200]}",
"assistant_content": content,
}).encode("utf-8")
try:
req = urllib.request.Request(
TDDB_URL + "/capture",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as resp:
result = json.loads(resp.read().decode())
logged = result.get("l0_recorded", 0)
log(f" Soulful profile → TencentDB: l0_recorded={logged}")
except Exception as e:
log(f" ⚠️ soulful_profile_to_tdb_scene 失败: {e}")
def tddb_capture(reflection_dict, state, ctx):
"""将 deep tick 的 reflection 捕获到 TencentDB形成人格记忆积累。"""
try:
summary = ""
if reflection_dict:
r = reflection_dict.get("reflection", {})
summary = r.get("evaluation_previous_goal", "") or r.get("summary", "")
next_goal = r.get("next_goal", "")
if next_goal and next_goal != "继续监控":
summary = f"{summary}\n下一步: {next_goal}"
if not summary:
return
payload = json.dumps({
"session_key": "daemon-deep-tick",
"user_content": f"系统状态: 磁盘{state.get('disk_pct')}% 内存{state.get('mem_pct')}%,进程{'全正常' if all(state.get('processes',{}).values()) else '有异常'},深度思考计数{ctx.get('deep_tick_count',0)}",
"assistant_content": summary[:1000],
}).encode("utf-8")
req = urllib.request.Request(
f"{TDDB_URL}/capture",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as resp:
result = json.loads(resp.read().decode())
l0 = result.get("l0_recorded", 0)
if l0 > 0:
log(f" TencentDB capture: {l0} L0 recorded")
except Exception as e:
log(f" TencentDB capture failed: {e}")
# ====== Phase 3: 织忆 recent_moments → TencentDB L1 同步 ======
def zhiyi_to_tdb():
"""读取织忆 Soulful 心迹recent_moments 最新3条写入 TencentDB L1。
数据源:~/.hermes/soulful/heart-traces.jsonlevent_type 包含 achievement/milestone/reflection/review 等)
写入目标TencentDB L1类型 atom, source=zhiyi_recent_moments
"""
heart_path = HERMES + "/soulful/heart-traces.jsonl"
if not os.path.exists(heart_path):
return
try:
with open(heart_path, encoding="utf-8") as f:
lines = f.readlines()
except Exception as e:
log(f" ⚠️ 读取 heart-traces 失败: {e}")
return
# 取最近 3 条(按 timestamp 倒序实际上文件已是时间顺序直接取末尾3条
recent = []
for line in lines[-3:]:
if line.strip():
try:
e = json.loads(line)
recent.append(e)
except Exception:
continue
if not recent:
return
for entry in recent:
content = entry.get("content", "")[:500]
ts = entry.get("timestamp", "")
event_type = entry.get("event_type", "moment")
if not content:
continue
try:
payload = json.dumps({
"session_key": "zhiyi-sync",
"user_content": f"[{event_type}] {content}",
"assistant_content": f"来自织忆心迹记录 | 时刻: {ts}",
"type": "atom",
"metadata": {"source": "zhiyi_recent_moments", "event_type": event_type, "timestamp": ts},
}).encode("utf-8")
req = urllib.request.Request(
f"{TDDB_URL}/capture",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=8) as resp:
result = json.loads(resp.read().decode())
recorded = result.get("l0_recorded", 0)
if recorded > 0:
log(f" 织忆 recent_moment → TencentDB: {content[:40]}")
except Exception as e:
log(f" ⚠️ zhiyi_to_tdb 写入失败: {e}")
# ====== Soulful 三库(懒加载)======
def _get_soulful(name):
"""Lazy-load soulful modules to avoid startup failures"""
if name in _soulful_cache:
return _soulful_cache[name]
try:
sys.path.insert(0, HERMES + "/scripts")
mod = __import__("soulful_core")
_soulful_cache[name] = mod
return mod
except Exception as e:
log(f"⚠️ soulful_core 导入失败: {e}")
return None
def get_hearttraces():
mod = _get_soulful("HeartTraces")
if mod:
return mod.HeartTraces()
return None
def get_userprofile():
mod = _get_soulful("UserProfile")
if mod:
return mod.UserProfile()
return None
def get_caresqueue():
mod = _get_soulful("CaresQueue")
if mod:
return mod.CaresQueue()
return None
def soulful_get_recent_moments(n=3):
"""读取最近 N 条心迹,返回字符串供注入 context"""
ht = get_hearttraces()
if not ht:
return ""
moments = ht.recent(n=n)
if not moments:
return ""
lines = ["\n[心迹 - 记忆我们之间的事]:"]
for m in moments:
stars = "" * m.get("importance", 3)
lines.append(f" {stars} {m.get('content', '')}")
return "\n".join(lines)
def soulful_check_cares():
"""检查牵挂队列,优先尝试帮助,其次才推飞书
策略:
1. 如果牵挂指向一个可自动化的任务 → 尝试执行
2. 如果牵挂需要人工行动 → 检查是否到了真正需要提醒的时间
3. 飞书推送只用于真正需要你才知道的事(其他我全部自己处理)
"""
cq = get_caresqueue()
if not cq:
return
due = cq.today_check()
if not due:
return
for care in due:
content = care.get("content", "")
context_raw = care.get("context", "")
reminder_count = care.get("reminder_count", 0)
# 判断这个牵挂是否需要通知我
# 原则:大部分牵挂我自己可以帮忙处理,不打扰你
# 只有真正需要你本人决定的,才推飞书
# 检查内容是否指向可识别任务(我可以直接帮忙的)
care_lower = content.lower()
auto_helpable = any(kw in care_lower for kw in [
"", "检查", "", "", "同步", "更新", "备份", "测试", "确认",
"", "", "优化", "整理", "提交", "推送", "发送"
])
if auto_helpable:
# 我能帮忙 — 静默处理,不推飞书,等你有空问我
# 把牵挂标记为"已识别,下次对话提起"
log(f" 💡 牵挂已识别(可帮忙): {content[:40]}")
continue
# 真正需要你知道的 — 推飞书,但调整频率
if reminder_count == 0:
opener = "你之前说过"
elif reminder_count == 1:
opener = "上次提醒过一次,还是想问一下"
elif reminder_count >= 3:
opener = f"这件事已经跟了你 {reminder_count} 次了"
# 提醒超过 3 次,标记为 snooze 1 周
cq.snooze(care["id"], days=7)
continue
else:
opener = f"想关心一下进度"
text = opener + f":「{content}"
if context_raw and len(context_raw) > 5:
text += f"(背景:{context_raw[:40]}"
send_feishu("🎗️ 你有一件事一直放在心上", text, "purple")
cq.snooze(care["id"], days=0) # 仅增加 reminder_count
def soulful_update_profile():
"""间接调用 update_profile.py"""
script = HERMES + "/scripts/update_profile.py"
if not os.path.exists(script):
return
rc, out, err = shell(f"python3 {script}", timeout=60)
if rc == 0:
log(f" 画像更新: {out[:80]}")
# ====== Soulful → llm_context 同步(每 tick 同步到 llm_context.json======
def _sync_soulful_to_llm_context(ctx):
"""把 Soulful 三库摘要写入 ctx['soulful']save_llm_context 落盘。
llm_context.json 由 Hermes 织忆插件 prefetch 时注入主 session
所以这里写入 = 牧尘在对话中感知到关系记忆的前提。
"""
try:
from soulful_core import HeartTraces, UserProfile, CaresQueue
ht = HeartTraces()
up = UserProfile()
cq = CaresQueue()
recent_moments = ht.recent(n=5) or ""
profile = up.get()
pending_cares = cq.pending()
ctx["soulful"] = {
"recent_moments": recent_moments[:500] if recent_moments else "",
"profile_summary": profile.get("communication_style", ""),
"cares_pending": len(pending_cares),
"uptime_minutes": ctx.get("uptime_minutes", 0),
"daemon_status": "running",
}
for k in ["messages_sent", "emotion_history"]:
ctx.pop(k, None)
except ImportError as e:
log(f"⚠️ Soulful 导入失败: {e}")
except Exception as e:
log(f"⚠️ Soulful 同步失败: {e}")
# ====== TencentDB → llm_context统一 user_profile重构 v2======
# 重构目标:整合 Soulful user-profile + TencentDB L3 persona 为单一 user_profile
try:
profile_data = {}
# 1. 读 Soulful user-profile.json
soulful_profile_path = HERMES + "/soulful/user-profile.json"
if os.path.exists(soulful_profile_path):
with open(soulful_profile_path, encoding="utf-8") as f:
sp = json.load(f)
profile_data["behavior_rules"] = sp.get("behavior_rules", {})
profile_data["communication_style"] = sp.get("communication_style", "")
profile_data["work_patterns"] = sp.get("work_patterns", {})
profile_data["preferences"] = sp.get("preferences", {})
# 2. 尝试从 TencentDB 获取 L1+L2 记忆补充 persona
try:
payload = json.dumps({"query": "牧尘工作状态 记忆 偏好", "session_key": "daemon-tick", "top_k": 3}).encode("utf-8")
req = urllib.request.Request(
f"{TDDB_URL}/recall",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=5) as resp:
result = json.loads(resp.read().decode())
memories = result.get("memory_count", 0)
ctx["tddb_memory_count"] = memories
except Exception:
ctx["tddb_memory_count"] = 0
# 3. 写入统一的 user_profile 结构
ctx["user_profile"] = profile_data
except Exception as e:
log(f"⚠️ user_profile 构建失败: {e}")
# ═══════════════════════════════════════════════════════════════════════════
# L1→L2 蒸馏memories observations → graph_nodes patterns
# ═══════════════════════════════════════════════════════════════════════════
def _distill_l1_to_l2():
"""从三个来源发现重复 pattern写入 graph_nodestype='pattern'
来源1graph_nodes namespace×type 组合(已有的元信息聚合)
来源2state.db 消息内容关键词新增3.3万条真实对话)
来源3memories 表(保留兼容)
新增来源2 直接分析会话文本,追踪真实话题的重复出现,
解决「框架通了但数据薄」的核心问题。
"""
graph_db = HERMES + "/graph.db"
if not os.path.exists(graph_db):
log(" ⚠️ _distill_l1_to_l2: graph.db 不存在")
return
import sqlite3, re
from collections import Counter
try:
g_conn = sqlite3.connect(graph_db)
g_cur = g_conn.cursor()
all_patterns = {}
def merge_pattern(name, score_increment, props_update):
if name in all_patterns:
all_patterns[name]["score"] += score_increment
all_patterns[name]["props"].update(props_update)
else:
all_patterns[name] = {"score": score_increment, "props": props_update.copy()}
# 来源1namespace×type 组合聚合
g_cur.execute("""
SELECT namespace, type, COUNT(*) as cnt
FROM graph_nodes
WHERE namespace IS NOT NULL AND type IS NOT NULL
AND namespace NOT IN ('daemon-distill')
GROUP BY namespace, type
HAVING cnt >= 3
LIMIT 20
""")
for ns, ntype, cnt in g_cur.fetchall():
merge_pattern(f"pattern:{ns}:{ntype}", cnt, {"source": "graph_nodes", "node_count": cnt})
# 来源2新增state.db 会话消息关键词
state_db = HERMES + "/state.db"
try:
s_conn = sqlite3.connect(state_db)
s_cur = s_conn.cursor()
s_cur.execute("""
SELECT SUBSTR(content, 1, 200) FROM messages
WHERE role='user' AND content IS NOT NULL AND content != ''
ORDER BY timestamp DESC LIMIT 2000
""")
messages = [row[0] for row in s_cur.fetchall()]
s_conn.close()
if messages:
all_text = " ".join(messages)
cn_words = re.findall(r'[一-鿿]{2,6}', all_text)
en_words = re.findall(r'[a-zA-Z_]{3,15}', all_text)
stopwords = {'的是', '这个', '那个', '什么', '怎么', '为什么', '是不是', '有没有',
'一个', '可以', '然后', '就是', '不是', '但是', '或者', '如果',
'the', 'and', 'for', 'that', 'this', 'with', 'from', 'have',
'what', 'how', 'why', 'when', 'where', 'which', 'are', 'was'}
cn_counter = Counter(w for w in cn_words if w not in stopwords)
en_counter = Counter(w for w in en_words if w not in stopwords)
for topic, cnt in cn_counter.most_common(15) + en_counter.most_common(15):
if cnt >= 3:
merge_pattern(f"topic:{topic}", cnt, {"source": "session_messages", "topic": topic})
except Exception as e:
log(f" ⚠️ session 消息分析失败: {e}")
# 写入所有 pattern
if not all_patterns:
g_conn.close()
return
now_ts = datetime.now(timezone.utc).isoformat()
created = updated = 0
for name, data in all_patterns.items():
score = data["score"]
props = data["props"]
g_cur.execute("SELECT id, properties FROM graph_nodes WHERE name=? AND type='pattern'", (name,))
existing = g_cur.fetchone()
if existing:
node_id, props_json = existing
old_props = json.loads(props_json) if props_json else {}
old_cnt = old_props.get("occurrence_count", 0)
new_cnt = old_cnt + 1
props["occurrence_count"] = new_cnt
props["last_seen"] = now_ts
props["score"] = score
for k, v in old_props.items():
if k not in props:
props[k] = v
g_cur.execute("UPDATE graph_nodes SET properties=?, last_updated_at=? WHERE id=?",
(json.dumps(props, ensure_ascii=False), now_ts, node_id))
if old_cnt == 2:
log(f"💡 pattern 触发 L2→L3: {name} (occurrence={new_cnt})")
journal_entry("pattern_triggered", f"pattern: {name}", f"occurrence={new_cnt}")
updated += 1
else:
props["occurrence_count"] = 1
props["first_seen"] = now_ts
props["last_seen"] = now_ts
props["score"] = score
g_cur.execute("""
INSERT INTO graph_nodes (name, type, namespace, properties, created_at, last_updated_at)
VALUES (?, 'pattern', 'daemon-distill', ?, ?, ?)
""", (name, json.dumps(props, ensure_ascii=False), now_ts, now_ts))
created += 1
log(f" 蒸馏 L1→L2: 创建 topic pattern '{name}' (score={score})")
g_conn.commit()
g_conn.close()
log(f" 蒸馏 L1→L2 完成: 新建 {created},更新 {updated}(共 {len(all_patterns)} 个 pattern")
except Exception as e:
log(f" ⚠️ _distill_l1_to_l2 失败: {e}")
# L2→L3 蒸馏pattern 节点 → TencentDB L2 scenes
# ═══════════════════════════════════════════════════════════════════════════
def _distill_l2_to_l3():
"""读取 occurrence_count≥3 的 pattern 节点,按类别聚类后写入 TencentDB L2 scene。
scene_type 为 'domain-rule'tags 包含 ['L2', 'pattern']。
"""
graph_db = HERMES + "/graph.db"
if not os.path.exists(graph_db):
log(" ⚠️ _distill_l2_to_l3: graph.db 不存在")
return
try:
import sqlite3
g_conn = sqlite3.connect(graph_db)
g_cur = g_conn.cursor()
g_cur.execute("""
SELECT id, name, properties
FROM graph_nodes
WHERE type = 'pattern'
AND properties LIKE '%occurrence_count%'
ORDER BY last_updated_at DESC
LIMIT 50
""")
rows = g_cur.fetchall()
g_conn.close()
if not rows:
log(" 蒸馏 L2→L3: 无满足条件的 pattern 节点")
return
# 按 namespace 或 category 简单聚类
from collections import defaultdict
groups = defaultdict(list)
for row in rows:
node_id, name, props_json = row
props = json.loads(props_json) if props_json else {}
occ = props.get("occurrence_count", 0)
if occ < 3:
continue
# 用 namespace 分组
namespace = "default"
if ":" in name:
parts = name.split(":")
if len(parts) >= 2:
namespace = parts[1]
groups[namespace].append((name, props))
# 写入 TencentDB
for ns, patterns in groups.items():
scene_name = f"l2-pattern-group-{ns}"
# 生成综合描述
descriptions = [p[1].get("description", p[0]) for p in patterns]
content = f"这些 observations 表明牧尘在 {ns} 方面有重复行为模式:\n" + "\n".join(f"- {d}" for d in descriptions[:5])
tags = ["L2", "pattern", ns]
# 用 /capture 写入 pattern group 总结L2 pattern 汇总 → TencentDB
# TencentDB 只有 /capture 接口,无 /scenes
content = f"L2 pattern group [{ns}]\n" + "\n".join(f"- {d}" for d in descriptions[:5])
payload = json.dumps({
"session_key": "l2-distill",
"user_content": f"系统蒸馏发现 {ns} 方面的 L2 patterns ({len(patterns)} 条)",
"assistant_content": content,
}).encode("utf-8")
try:
req = urllib.request.Request(
TDDB_URL + "/capture",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as resp:
result = json.loads(resp.read().decode())
logged = result.get("l0_recorded", 0)
log(f" L2→L3: pattern-group-{ns} ({len(patterns)} 条, l0={logged})")
except Exception as e:
log(f" ⚠️ _distill_l2_to_l3 写入 scene 失败: {e}")
except Exception as e:
log(f" ⚠️ _distill_l2_to_l3 失败: {e}")
# ====== L3→L4 蒸馏L2 patterns → cross-domain policies======
def _distill_l3_to_l4():
"""读取 L2 patterns提取跨域共性规则写入 TencentDB L4policy
用 /capture 接口session_key=l3-distill, assistant_content=policy 内容。
L4 policy 是跨多个 L2 pattern 的通用原则(如"先拉现状再诊断")。
"""
graph_db = HERMES + "/graph.db"
if not os.path.exists(graph_db):
return
try:
import sqlite3
conn = sqlite3.connect(graph_db)
cur = conn.cursor()
cur.execute("""
SELECT name, properties FROM graph_nodes
WHERE type = 'pattern'
AND properties LIKE '%occurrence_count%'
ORDER BY last_updated_at DESC LIMIT 20
""")
rows = cur.fetchall()
conn.close()
if len(rows) < 2:
return
# 提取所有 pattern 名称,看有没有跨域共性
pattern_names = []
for (name, props) in rows:
p = json.loads(props) if props else {}
desc = p.get("description", name)
pattern_names.append(desc)
# 简单启发式:找共同关键词
# 检查有没有"诊断/排查/修复"类共同模式 → policy: "先拉现状再诊断"
all_text = " ".join(pattern_names).lower()
policies = []
if any(k in all_text for k in ["修复", "问题", "错误", "排查", "诊断"]):
policies.append("先拉现状再诊断,不假设不验证")
if any(k in all_text for k in ["代码", "脚本", "配置", "修改"]):
policies.append("改完先自测再交付,不留半成品")
if not policies:
return
policy_text = "牧尘的工作原则:" + "".join(policies)
payload = json.dumps({
"session_key": "l3-distill",
"user_content": "系统蒸馏:发现以下 pattern" + " | ".join(pattern_names[:5]),
"assistant_content": policy_text,
}).encode("utf-8")
req = urllib.request.Request(
TDDB_URL + "/capture",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as resp:
result = json.loads(resp.read().decode())
logged = result.get("l0_recorded", 0)
if logged:
log(f" L3→L4 蒸馏: policy={policies} recorded={logged}")
except Exception as e:
log(f" ⚠️ _distill_l3_to_l4 失败: {e}")
# ====== L4→L5 蒸馏policies → traits 人格特质)======
def _distill_l4_to_l5():
"""从 TencentDB recall 结果中读取 L4 policies提取人格特质更新 Soulful user-profile。
用 /recall 接口获取最近的 policy 级记忆,提取 communication_style / behavior_rules
特征,写入 user-profile.json不动原文件结构只追加 behavior_rules
"""
try:
# 用 /recall 拉 L3-L4 级别的记忆
payload = json.dumps({
"query": "牧尘工作方式 决策风格 偏好 工作原则",
"session_key": "l4-distill",
"top_k": 5
}).encode("utf-8")
req = urllib.request.Request(
TDDB_URL + "/recall",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as resp:
recall_result = json.loads(resp.read().decode())
# 解析 recall 结果中的 persona/policy 内容
context = recall_result.get("context", "")
if not context:
return
# 从 context 提取 trait 候选(简单关键词匹配)
traits = []
if "简洁" in context or "直接" in context:
traits.append({"trait": "communication_style", "value": "简洁直接,不废话"})
if "现状" in context or "诊断" in context:
traits.append({"trait": "behavior_rule", "value": "先拉现状再诊断,不假设不验证"})
if "测试" in context or "自测" in context:
traits.append({"trait": "behavior_rule", "value": "改完先自测再交付"})
if not traits:
return
# 追加到 user-profile.json 的 behavior_rules新字段 distilled_rules
profile_path = HERMES + "/soulful/user-profile.json"
if not os.path.exists(profile_path):
return
profile = json.load(open(profile_path))
if "distilled_rules" not in profile:
profile["distilled_rules"] = []
for t in traits:
if t["value"] not in profile["distilled_rules"]:
profile["distilled_rules"].append(t["value"])
profile["last_distilled"] = datetime.now(timezone.utc).isoformat()
with open(profile_path, "w") as f:
json.dump(profile, f, ensure_ascii=False, indent=2)
log(f" L4→L5 蒸馏: 新增 {len(traits)} 条 distilled_rules")
except Exception as e:
log(f" ⚠️ _distill_l4_to_l5 失败: {e}")
# ====== L5→L6 蒸馏traits → values 根本价值)======
def _distill_l5_to_l6():
"""从 heart-traces.jsonl 提炼根本价值观,写入 heart-tracestype=values
同时修复 heart-traces.jsonl给已有记录加上 type 字段event_type→type 映射)。
"""
heart_path = HERMES + "/soulful/heart-traces.jsonl"
if not os.path.exists(heart_path):
return
try:
# 读取并补充 type 字段(已有记录)
lines = []
modified = False
with open(heart_path, encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
entry = json.loads(line)
# 补充 type 字段event_type 映射)
if "type" not in entry:
event = entry.get("event_type", "moment")
entry["type"] = "value" if event in ("value_expression", "principle", "core_belief") else "moment"
modified = True
lines.append(entry)
except json.JSONDecodeError:
pass
# 如果有修改,写回去
if modified:
with open(heart_path, "w", encoding="utf-8") as f:
for entry in lines:
f.write(json.dumps(entry, ensure_ascii=False) + "\n")
log(f" L5→L6: 补充 type 字段,{len(lines)} 条记录")
# 找最近的 value 类条目
value_entries = [e for e in lines if e.get("type") == "value"]
if len(value_entries) >= 1:
# 已经是 values 结构log 一下
log(f" L5→L6: 发现 {len(value_entries)} 条 values 记录")
except Exception as e:
log(f" ⚠️ _distill_l5_to_l6 失败: {e}")
def _write_daily_soulful_journal():
"""每日 soulful journal: 将 deep_think reflection 摘要写入 journals/YYYY-MM-DD.jsonl。每天只写一条。"""
try:
os.makedirs(SOULFUL_JOURNALS_DIR, exist_ok=True)
today = datetime.now().strftime("%Y-%m-%d")
jpath = SOULFUL_JOURNALS_DIR + "/" + today + ".jsonl"
if os.path.exists(jpath):
with open(jpath) as f:
for line in f:
if line.strip():
return
ctx = load_context()
lr = ctx.get("last_reflection", {})
if not lr or not lr.get("memory", ""):
return
entry = {
"timestamp": datetime.now(timezone.utc).isoformat(),
"date": today,
"reflection": lr.get("memory", ""),
"evaluation": lr.get("evaluation_previous_goal", ""),
"next_goal": lr.get("next_goal", ""),
}
with open(jpath, "w") as f:
f.write(json.dumps(entry, ensure_ascii=False) + "\n")
log(" Soulful journal written: " + today)
except Exception as e:
log(" ERROR _write_daily_soulful_journal: " + str(e))
def _cleanup_expired_cares():
"""清理 cares-queue.json 中的过期牵挂。
- created_at > 14天且status=pending -> 标记为 stale保留记录
- follow_up_date < 今天-7天 -> 物理删除
"""
cq_path = HERMES + "/soulful/cares-queue.json"
if not os.path.exists(cq_path):
return
try:
with open(cq_path) as f:
data = json.load(f)
today = datetime.now().date()
cares = data.get("cares", [])
before = len(cares)
now_iso = datetime.now(timezone.utc).isoformat()
stale_count = 0
remaining = []
for c in cares:
follow_up = c.get("follow_up_date", "2099-12-31")
if follow_up < (today - timedelta(days=7)).isoformat():
continue
created = c.get("created_at", "")
status = c.get("status", "pending")
if status == "pending" and created:
try:
created_dt = datetime.fromisoformat(created)
if (datetime.now(timezone.utc) - created_dt).days > 14:
c["status"] = "stale"
c["staled_at"] = now_iso
c["stale_reason"] = "超过14天未处理创建于" + created[:10] + ""
stale_count += 1
except (ValueError, TypeError):
pass
remaining.append(c)
removed = before - len(remaining)
if removed > 0 or stale_count > 0:
data["cares"] = remaining
data["updated_at"] = now_iso
with open(cq_path, "w") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
log(f" 🧹 清理 cares: 移除 {removed} 条,标记 stale {stale_count}")
except Exception as e:
log(f" ⚠️ 清理过期 cares 失败: {e}")
# ====== Phase 2.1: 场景感知关怀判断(每 light tick======
def _check_scene_aware_cares():
"""检查到期 cares 与 TencentDB L2 active_scenes 的关联,打印关怀建议日志"""
# 读取 cares-queue只看 pending 且 due <= 今天)
cq_path = HERMES + "/soulful/cares-queue.json"
if not os.path.exists(cq_path):
return
try:
with open(cq_path) as f:
data = json.load(f)
today_str = datetime.now().date().isoformat()
due_cares = [
c for c in data.get("cares", [])
if c.get("status") == "pending" and c.get("follow_up_date", "2099-12-31") <= today_str
]
except Exception:
return
if not due_cares:
return
# 读取 TencentDB L2 active_scenes
try:
req = urllib.request.Request(
f"{TDDB_URL}/scenes?type=active",
headers={"Content-Type": "application/json"},
method="GET",
)
with urllib.request.urlopen(req, timeout=5) as resp:
result = json.loads(resp.read().decode())
scenes = result.get("scenes", []) if isinstance(result, dict) else result
if isinstance(scenes, dict):
scenes = scenes.get("scenes", [])
except Exception:
scenes = []
if not scenes:
return
# 建立场景关键词集合
scene_keywords = set()
for scene in scenes:
if isinstance(scene, dict):
name = scene.get("name", "") or scene.get("scene", "") or scene.get("title", "")
tags = scene.get("tags", []) or []
content_text = scene.get("content", "") or ""
elif isinstance(scene, str):
name, tags, content_text = scene, [], ""
else:
continue
if name:
scene_keywords.add(name)
scene_keywords.update(tags)
if content_text:
# 简单分词中文2-4字词
import re
words = re.findall(r'[\u4e00-\u9fff]{2,6}', content_text)
scene_keywords.update(words)
if not scene_keywords:
return
# 检查每个到期 care 是否与 active_scenes 相关
for care in due_cares:
care_text = (care.get("content", "") + " " + care.get("context", "")).lower()
for kw in scene_keywords:
kw_lower = kw.lower()
if len(kw_lower) >= 2 and kw_lower in care_text:
log(f"[关怀建议] {care.get('content', '')[:60]} 与当前场景 {kw} 相关,考虑主动关怀")
break
# ====== 方案库 ======
def load_solutions():
if os.path.exists(SOLUTIONS_FILE):
with open(SOLUTIONS_FILE) as f:
return json.load(f)
return {"solutions": [], "version": 2}
def save_solutions(lib):
os.makedirs(D, exist_ok=True)
with open(SOLUTIONS_FILE, "w") as f:
json.dump(lib, f, indent=2, ensure_ascii=False)
# 双写 session_state同步 patterns + solutions
try:
_sync_session_state(lib)
except Exception as e:
log(f" ⚠️ _sync_session_state failed: {e}")
# ====== 跨会话状态持久化(参考 Grok Build 双层持久化)======
SESSION_STATE_VERSION = 2
TOUCH_INTERVAL_TICKS = 180 # 每 180 light ticks ≈ 30 分钟 touch 一次
ORPHAN_TIMEOUT_TICKS = 360 # 超过 360 ticks约 2 小时)未 touch 则标记待恢复
SWEEP_INTERVAL_TICKS = 2880 # 每 2880 ticks ≈ 24 小时 sweep 一次
def _atomic_write_json(path, data):
"""写 JSON 文件:先写 .tmp再 fsync再 rename防止断电损坏"""
tmp = path + ".tmp"
with open(tmp, "w") as f:
json.dump(data, f, indent=2, ensure_ascii=False)
f.flush()
os.fsync(f.fileno())
os.rename(tmp, path)
def load_session_state():
"""加载跨会话状态文件version 不兼容则备份旧文件"""
if not os.path.exists(SESSION_STATE_FILE):
return {
"version": SESSION_STATE_VERSION,
"last_updated": None,
"last_touch": None,
"recovery_needed": False,
"data": {"patterns": {}, "solutions": {}, "cares": {}},
}
try:
with open(SESSION_STATE_FILE) as f:
st = json.load(f)
if st.get("version") != SESSION_STATE_VERSION:
# version 不兼容,备份旧文件
backup = SESSION_STATE_FILE + f".bak.{int(time.time())}"
os.rename(SESSION_STATE_FILE, backup)
log(f" session_state version 不兼容,备份到 {backup}")
return {
"version": SESSION_STATE_VERSION,
"last_updated": None,
"last_touch": None,
"recovery_needed": False,
"data": {"patterns": {}, "solutions": {}, "cares": {}},
}
# 初始化缺失字段
st.setdefault("last_touch", st.get("last_updated"))
st.setdefault("recovery_needed", False)
st.setdefault("data", {"patterns": {}, "solutions": {}, "cares": {}})
return st
except Exception as e:
log(f" ⚠️ load_session_state 失败: {e},返回空状态")
return {
"version": SESSION_STATE_VERSION,
"last_updated": None,
"last_touch": None,
"recovery_needed": False,
"data": {"patterns": {}, "solutions": {}, "cares": {}},
}
def save_session_state(st):
"""保存跨会话状态(原子写入)"""
st["version"] = SESSION_STATE_VERSION
st["last_updated"] = datetime.now(timezone.utc).isoformat()
_atomic_write_json(SESSION_STATE_FILE, st)
def touch_session(st):
"""更新 last_touch 时间戳;若超过 ORPHAN_TIMEOUT_TICKS 未 touch则标记待恢复"""
now = datetime.now(timezone.utc)
last_touch = st.get("last_touch")
orphan = False
if last_touch:
try:
last = datetime.fromisoformat(last_touch)
elapsed = (now - last).total_seconds()
orphan = elapsed > ORPHAN_TIMEOUT_TICKS * LIGHT_INTERVAL
except Exception:
orphan = True
st["last_touch"] = now.isoformat()
if orphan:
st["recovery_needed"] = True
log(" session_state 进入「待恢复」状态(超过 2 小时未 touch")
save_session_state(st)
def _sync_session_state(solutions_lib):
"""将 solutions 同步写入 session_state.data双写"""
st = load_session_state()
st["data"]["solutions"] = {
s["id"]: s for s in solutions_lib.get("solutions", [])
}
save_session_state(st)
def _sweep_session_state(st):
"""清理孤立的幽灵引用:只存在于 session_state 但 graph.db 已删除的 patterns/solutions"""
graph_db = HERMES + "/graph.db"
if not os.path.exists(graph_db):
return
try:
conn = sqlite3.connect(graph_db)
cur = conn.cursor()
# 获取 graph.db 中所有 pattern 和 solution 的 id
cur.execute("SELECT id FROM graph_nodes WHERE type IN ('pattern', 'solution')")
valid_ids = {row[0] for row in cur.fetchall()}
conn.close()
removed = 0
# 清理 patterns
for pid in list(st["data"].get("patterns", {}).keys()):
if pid not in valid_ids:
del st["data"]["patterns"][pid]
removed += 1
# 清理 solutions
for sid in list(st["data"].get("solutions", {}).keys()):
if sid not in valid_ids:
del st["data"]["solutions"][sid]
removed += 1
if removed:
log(f" sweep_session_state: 清理了 {removed} 个幽灵引用")
save_session_state(st)
except Exception as e:
log(f" ⚠️ sweep_session_state 失败: {e}")
def add_solution(lib, pattern_desc, detect_conditions, actions, learned_from="auto"):
"""添加新方案到库"""
sid = f"sol-{len(lib['solutions'])+1:04d}"
sol = {
"id": sid,
"pattern": pattern_desc,
"detect": detect_conditions, # e.g. {"metric": "disk_pct", "op": "gt", "value": 85}
"actions": actions, # e.g. [{"type": "shell", "cmd": "...", "verify": "disk_pct < 85"}]
"frequency": 1,
"last_applied": datetime.now(timezone.utc).isoformat(),
"success_count": 1,
"fail_count": 0,
"learned_from": learned_from,
}
lib["solutions"].append(sol)
save_solutions(lib)
journal_entry("learn", f"学会新方案: {pattern_desc}")
return sid
def _score_solution(sol, state, ctx):
"""采样评分函数(参考 Grok Build Sampler 设计)
公式: score = base_score * freq_boost * recency_decay * precision_mult
- base_score: 成功率 (0-1)
- freq_boost: log(1 + success_count) 频率越高权重越高
- recency_decay: 最近执行过则降权(避免连续触发同一方案)
- precision_mult: 检测指标越多分数越高
"""
total = sol.get("success_count", 0) + sol.get("fail_count", 0)
# base_score: 成功率 (0-1),无数据时默认 0.5
if total > 0:
base_score = sol["success_count"] / total
else:
base_score = 0.5
# 成功率阈值检查:低于 50% 不采用
if base_score <= 0.5:
return 0.0
# freq_boost: log(1 + success_count),无成功记录时为 log(1)=0给小值
freq_boost = math.log(1 + sol.get("success_count", 0))
if freq_boost < 0.1:
freq_boost = 0.1 # 防止 log(0) 导致完全无权重
# recency_decay: 最近匹配过的方案降权ctx 携带上次匹配信息)
recency_decay = 1.0
last_matched_id = ctx.get("last_matched_solution_id")
if last_matched_id and last_matched_id == sol["id"]:
# 上次刚执行过,降到 30% 权重
recency_decay = 0.3
# 时效性距离上次越近降权越狠15 秒内再触发降更多)
last_ts = ctx.get("last_match_timestamp")
if last_ts:
try:
last_dt = datetime.fromisoformat(last_ts)
age_seconds = (datetime.now(timezone.utc) - last_dt).total_seconds()
if age_seconds < 60:
recency_decay = 0.1 # 1 分钟内几乎不触发
elif age_seconds < 300:
recency_decay = 0.3 # 5 分钟内降权
except Exception:
pass
# precision_mult: 检测指标越多detect 条件越具体)→ 越优先
# 当前 detect 结构为单个 dict计 1 分;未来扩展为 list 时累加
detect = sol.get("detect", {})
num_conditions = 1 # 默认 1 个条件
if isinstance(detect, list):
num_conditions = len(detect)
# 指标越具体权重越高1 条件=1.0, 2 条件=1.2, 3+ 条件=1.5
if num_conditions >= 3:
precision_mult = 1.5
elif num_conditions == 2:
precision_mult = 1.2
else:
precision_mult = 1.0
score = base_score * freq_boost * recency_decay * precision_mult
return score
def match_solution(lib, state, ctx=None):
"""检查当前状态是否匹配任何已知方案,返回打分最高的方案
与旧版区别:
- 旧版:顺序遍历,返回第一个匹配
- 新版:收集所有匹配方案,按采样打分排序,返回最优
- 只返回成功率 > 50% 的方案
"""
if ctx is None:
ctx = {}
candidates = []
for sol in lib["solutions"]:
detect = sol["detect"]
metric = detect.get("metric")
op = detect.get("op")
val = detect.get("value")
matched = False
if metric not in state:
matched = False
else:
actual = state[metric]
if isinstance(actual, (int, float)) and isinstance(val, (int, float)):
if op == "gt" and actual > val:
matched = True
elif op == "lt" and actual < val:
matched = True
elif op == "eq" and abs(actual - val) < 0.01:
matched = True
# 进程挂了匹配
if metric == "processes" and op == "dead":
procs = state.get("processes", {})
for p in (val if isinstance(val, list) else [val]):
if not procs.get(p, True):
matched = True
break
if matched:
candidates.append(sol)
if not candidates:
return None
# 对所有候选方案打分
scored = []
for sol in candidates:
score = _score_solution(sol, state, ctx)
if score > 0: # 只保留成功率 > 50% 的
scored.append((score, sol))
if not scored:
return None
# 按分数降序,取最高分
scored.sort(key=lambda x: x[0], reverse=True)
best = scored[0][1]
# 更新 ctx 记录(供下次 recency_decay 使用)
# 只存 ID 和时间戳,避免持久化过大的对象
ctx["last_matched_solution_id"] = best["id"]
ctx["last_match_timestamp"] = datetime.now(timezone.utc).isoformat()
return best
def execute_solution(sol, state):
"""执行方案并返回是否成功"""
log(f" 🔧 执行方案 {sol['id']}: {sol['pattern']}")
journal_entry("solve_start", f"执行 {sol['id']}: {sol['pattern']}")
success = True
results = []
for action in sol["actions"]:
if action["type"] == "shell":
rc, out, err = shell(action["cmd"], timeout=action.get("timeout", 30))
results.append({"cmd": action["cmd"], "rc": rc, "out": out[:100]})
log(f" 执行: {action['cmd'][:60]} → exit={rc}")
# 验证
verify = action.get("verify")
if verify and rc == 0:
# 重新采集状态验证
time.sleep(2)
new_state = collect_state()
metric = sol["detect"].get("metric")
op = sol["detect"].get("op")
val = sol["detect"].get("value")
if metric in new_state:
actual = new_state[metric]
if op == "gt":
if actual <= val:
log(f" ✅ 验证通过: {metric}={actual}{val}")
else:
log(f" ⚠️ 验证未通过: {metric}={actual} 仍 > {val}")
success = False
# 更新方案统计
sol["frequency"] += 1
sol["last_applied"] = datetime.now(timezone.utc).isoformat()
if success:
sol["success_count"] += 1
else:
sol["fail_count"] += 1
return success, results
def action_to_solution(action_result, state, changes):
"""把一次成功的行动转化为可复用的方案"""
# 只转化 shell 行动
if not action_result.get("shell_cmds"):
return None
# 提取检测条件
detect = {}
for c in changes:
if "磁盘" in c:
detect = {"metric": "disk_pct", "op": "gt", "value": 85}
elif "内存" in c:
detect = {"metric": "mem_pct", "op": "gt", "value": 90}
if not detect:
return None
actions = [{"type": "shell", "cmd": cmd, "verify": None, "timeout": 30}
for cmd in action_result["shell_cmds"]]
return {
"pattern": f"自动学习: {changes[0] if changes else 'unknown'}",
"detect": detect,
"actions": actions,
}
# ====== 状态管理 ======
def load_context():
if os.path.exists(CONTEXT_FILE):
with open(CONTEXT_FILE) as f:
return json.load(f)
return {"started_at": datetime.now(timezone.utc).isoformat(), "last_deep_tick": None,
"last_light_tick": None, "last_state": {}, "tick_count": 0, "deep_tick_count": 0,
"messages_sent": 0, "solved_count": 0, "learned_count": 0, "uptime_seconds": 0,
"last_matched_solution_id": None, "last_match_timestamp": None}
def save_context(ctx):
os.makedirs(D, exist_ok=True)
with open(CONTEXT_FILE, "w") as f:
json.dump(ctx, f, indent=2)
# Phase 2画像LLM合成 和 Phase 3冲突检测 在下方
# ═══════════════════════════════════════════════════════════════════════════
import requests as _req
def time_decay_recall(query: str, top_k: int = 5) -> list:
"""带时间衰减的织忆recall综合分=recall×0.6+decay×0.430天内访问过boost×1.15"""
try:
token = os.environ.get("ZHIYI_API_KEY", "zhiyi-dev-key-2026")
r = _req.post("http://127.0.0.1:7821/api/v1/recall",
json={"query": query, "top_k": top_k * 2, "agent_id": "hermes-a06", "use_rerank": True},
headers={"X-API-Key": token}, timeout=8)
if r.status_code != 200: return []
results = r.json().get("results", [])
except: return []
now = datetime.now()
scored = []
for item in results:
try:
ts = datetime.fromisoformat(item.get("timestamp","").split("+")[0] if "+" in item.get("timestamp","") else item.get("timestamp",""))
except: ts = now
days = (now - ts).total_seconds() / 86400
decay = max(0.3, 1.0 - days * 0.015)
recall_s = float(item.get("score", 0.5))
item["decay_score"] = round(decay, 4)
item["days_since_update"] = round(days, 1)
item["tier"] = get_memory_tier(days)
item["final_score"] = round(recall_s * 0.6 + decay * 0.4, 4)
scored.append(item)
# Phase 4: 访问频率 boost
scored = _boost_recalled_memory(scored)
for r in scored:
log_memory_access(r.get("id",""), r.get("decay_score", 1.0))
scored.sort(key=lambda x: x["final_score"], reverse=True)
return scored[:top_k]
# ═══════════════════════════════════════════════════════════════════════════
# Corrective RAG: LLM相关性评分 + 不相关时触发重新检索
# 参考: awesome-llm-apps/rag_tutorials/corrective_rag
# 流程: recall → LLM评分相关性 → 不相关>50%则改写query重新检索
# ═══════════════════════════════════════════════════════════════════════════
def _grade_single_relevance(query: str, result_item: dict) -> str:
"""用LLM判断单条记忆是否与query相关。返回: relevant / irrelevant / partially_relevant"""
content = result_item.get("content", "")[:300]
prompt = f"""判断以下记忆是否与查询相关。
查询: {query}
记忆: {content}
只输出一个词: relevant / irrelevant / partially_relevant"""
try:
token = os.environ.get("NEWAPI_TOKEN", KEY)
resp = _req.post(
f"{API}/chat/completions",
headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"},
json={"model": "openai/gpt-oss-120b",
"messages": [{"role": "user", "content": prompt}],
"max_tokens": 10, "temperature": 0.1},
timeout=10,
)
grade = resp.json().get("choices", [{}])[0].get("message", {}).get("content", "").strip().lower()
if "relevant" in grade and "partially" not in grade and "irrelevant" not in grade:
return "relevant"
elif "irrelevant" in grade:
return "irrelevant"
return "partially_relevant"
except Exception:
return "partially_relevant"
def _expand_query(query: str) -> str:
"""将原query改写成更全面的检索表达用于重新检索"""
prompt = f"""将以下查询改写成更全面、可能包含同义词的检索表达。
原查询: {query}
改写(只输出改写后的查询,不要解释):"""
try:
token = os.environ.get("NEWAPI_TOKEN", KEY)
resp = _req.post(
f"{API}/chat/completions",
headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"},
json={"model": "openai/gpt-oss-120b",
"messages": [{"role": "user", "content": prompt}],
"max_tokens": 60, "temperature": 0.2},
timeout=10,
)
expanded = resp.json().get("choices", [{}])[0].get("message", {}).get("content", "").strip()
return expanded if expanded else query
except Exception:
return query
def corrective_recall(query: str, top_k: int = 5) -> list:
"""
Corrective RAG recall: 时间衰减recall + LLM相关性评分。
如果超过50%的结果不相关改写query重新检索并合并。
"""
results = time_decay_recall(query, top_k=top_k * 2)
# LLM相关性评分
graded = []
irrelevant = 0
for item in results:
grade = _grade_single_relevance(query, item)
item["relevance_grade"] = grade
if grade == "irrelevant":
irrelevant += 1
graded.append(item)
total = len(graded) or 1
ir_ratio = irrelevant / total
# 超过50%不相关 → 改写query重新检索
if ir_ratio > 0.5 and total >= 4:
expanded = _expand_query(query)
if expanded != query:
re_results = time_decay_recall(expanded, top_k=top_k)
for item in re_results:
item["relevance_grade"] = "relevance_from_expanded_query"
item["original_query"] = query
item["expanded_query"] = expanded
# 合并去重按id
seen = {r.get("id") for r in graded if r.get("id")}
merged = graded + [r for r in re_results if r.get("id") not in seen]
merged.sort(key=lambda x: x["final_score"], reverse=True)
return merged[:top_k]
graded.sort(key=lambda x: x["final_score"], reverse=True)
return graded[:top_k]
# ── Phase 4: 遗忘曲线分层 ─────────────────────────────────────────────────
MEMORY_TIERS = {"hot": 0, "warm": 1, "cold": 2, "archive": 3}
def get_memory_tier(days: float) -> str:
if days <= 7: return "hot"
elif days <= 21: return "warm"
elif days <= 60: return "cold"
else: return "archive"
def log_memory_access(memory_id: str, decay_weight: float):
import json
log_path = D + "/memory-access-log.jsonl"
entry = {"memory_id": memory_id, "accessed_at": datetime.now().isoformat(), "decay_weight": round(decay_weight, 4)}
try:
with open(log_path, "a") as f: f.write(json.dumps(entry, ensure_ascii=False) + "\n")
except: pass
def _boost_recalled_memory(recalled: list) -> list:
import json
log_path = D + "/memory-access-log.jsonl"
recent = {}
try:
cutoff = datetime.now().timestamp() - 30 * 86400
with open(log_path) as f:
for line in f:
try:
e = json.loads(line.strip())
if datetime.fromisoformat(e["accessed_at"]).timestamp() >= cutoff:
mid = e["memory_id"]
recent[mid] = recent.get(mid, 0) + 1
except: pass
except: pass
for r in recalled:
cnt = recent.get(r.get("id",""), 0)
r["access_count_30d"] = cnt
r["boosted"] = False
if cnt >= 1:
r["final_score"] = round(min(1.0, r["final_score"] * 1.15), 4)
r["boosted"] = True
return recalled
# ── Phase 2: 画像LLM合成 ─────────────────────────────────────────────────
def update_profile_from_journal(journal_path: str, profile_path: str, model: str = FAST_MODEL):
"""deep_tick时用LLM分析journal_entry更新画像仅非空字段"""
import re
try:
with open(journal_path) as f:
lines = f.readlines()
if len(lines) < 3: return
entries = (lambda L: [e for e in (json.loads(l) for l in L) if e.get("type") != "startup"][-10:])(lines[-20:])
except: return
if len(entries) < 3: return
log_lines = "\n".join(f"- {e.get('summary','')}: {e.get('details','(无)')}" for e in entries)
prompt = f"""牧尘是技术用户用小唯AI助手(Hermes Agent)工作。分析最近行为日志,识别:(1)沟通/工作模式 (2)正在进行的项目 (3)新发现的偏好。行为日志:\n{log_lines}\n输出JSON仅含需更新的字段{{"字段名":"新值"}},无需更新则空对象{{}}。只输出JSON。"""
try:
resp = _req.post(f"{API}/chat/completions",
headers={"Authorization": f"Bearer {os.environ.get('NEWAPI_TOKEN','')}", "Content-Type": "application/json"},
json={"model": model, "messages": [{"role": "user", "content": prompt}], "max_tokens": 500, "temperature": 0.3},
timeout=30)
text = resp.json().get("choices", [{}])[0].get("message", {}).get("content", "{}")
except: return
m = re.search(r'\{[^{}]*(?:"[^"]+"\s*:\s*[^{}]+){1,4}\}', text, re.DOTALL)
if not m: return
try: updates = json.loads(m.group())
except: return
if not updates: return
try:
with open(profile_path) as f: profile = json.load(f)
except: return
changed = False
for k, v in updates.items():
if v and str(v).strip() and profile.get(k) != v:
profile[k] = str(v); changed = True
if changed:
profile["updated_at"] = datetime.now().isoformat()
with open(profile_path, "w") as f: json.dump(profile, f, ensure_ascii=False, indent=2)
journal_entry("profile_update", f"LLM更新画像: {list(updates.keys())}")
# ── Phase 3: 冲突检测 ─────────────────────────────────────────────────────
def _write_conflicts_to_queue(conflicts: list):
import json
qpath = HERMES + "/soulful/conflicts-queue.json"
try:
with open(qpath) as f: queue = json.load(f)
except: queue = []
for c in conflicts:
if not any(e.get("id") == c.get("id") for e in queue): queue.append(c)
with open(qpath, "w") as f: json.dump(queue, f, ensure_ascii=False, indent=2)
def detect_memory_conflicts(journal_path: str, zhiyi_token: str) -> list:
"""用LLM检测journal_entry和织忆记忆之间的矛盾返回冲突列表"""
import re
try:
with open(journal_path) as f: lines = f.readlines()
entries = (lambda L: [e for e in (json.loads(l) for l in L) if e.get("type") != "startup"][-5:])(lines[-20:])
except: return []
if not entries: return []
try:
r = _req.post("http://127.0.0.1:7821/api/v1/recall",
json={"query": "牧尘 重要决定 项目 偏好", "top_k": 5, "agent_id": "hermes-a06"},
headers={"X-API-Key": zhiyi_token}, timeout=8)
existing = r.json().get("results", []) if r.status_code == 200 else []
except: existing = []
if not existing: return []
atext = "\n".join(f"- {e.get('summary','')}" for e in entries)
etext = "\n".join(f"- {m.get('content','')[:100]}" for m in existing)
prompt = f"""A组是牧尘最近行为日志B组是织忆长期记忆。判断A和B是否有直接矛盾。\nA组\n{atext}\nB组\n{etext}\n有矛盾输出:{{"has_conflict":true,"conflict_description":"描述","a_entry":"哪条A","b_entry":"哪条B"}}\n无矛盾:{{"has_conflict":false}}\n只输出JSON。"""
try:
resp = _req.post("http://127.0.0.1:3000/v1/chat/completions",
headers={"Authorization": f"Bearer {os.environ.get('NEWAPI_TOKEN','')}", "Content-Type": "application/json"},
json={"model": FAST_MODEL, "messages": [{"role": "user", "content": prompt}],
"max_tokens": 400, "temperature": 0.1}, timeout=25)
text = resp.json().get("choices", [{}])[0].get("message", {}).get("content", "{}")
except: return []
try: result = json.loads(text)
except:
m = re.search(r'\{"has_conflict"[^}]+\}', text, re.DOTALL)
result = json.loads(m.group()) if m else {"has_conflict": False}
conflicts = []
if result.get("has_conflict"):
conflicts.append({"id": f"conflict_{datetime.now().strftime('%Y%m%d%H%M%S')}",
"entry": result.get("a_entry",""), "conflict_with": result.get("b_entry",""),
"description": result.get("conflict_description",""),
"detected_at": datetime.now().isoformat(), "status": "pending"})
_write_conflicts_to_queue(conflicts)
return conflicts
def save_llm_context(ctx, state):
"""每 tick 写 llm_context.jsonv2 统一格式)。
包含完整 L1-L6 蒸馏数据,从 graph_nodes / TencentDB / Soulful 实时读取,
确保 MEMORY 区注入的是当前最新状态。
"""
import sqlite3
# ── daemon 健康状态 ──────────────────────────────────────────
daemon_status = "running" if psutil.pid_exists(os.getpid()) else "stopped"
# ── L1: 从 graph_nodes observations ──────────────────────────
observations = []
try:
conn = sqlite3.connect(HERMES + "/graph.db")
cur = conn.cursor()
cur.execute("""
SELECT name, namespace, properties FROM graph_nodes
WHERE type IN ('entity','concept','topic')
AND namespace NOT IN ('daemon-distill')
ORDER BY last_updated_at DESC LIMIT 30
""")
for name, ns, props_json in cur.fetchall():
props = json.loads(props_json) if props_json else {}
observations.append({
"name": name, "namespace": ns,
"summary": props.get("description", "")[:100]
})
conn.close()
except Exception:
pass
# ── L2: pattern 节点来源1: namespace×type + 来源2: 会话话题) ──
patterns = []
try:
conn = sqlite3.connect(HERMES + "/graph.db")
cur = conn.cursor()
cur.execute("""
SELECT name, properties FROM graph_nodes
WHERE type='pattern' ORDER BY last_updated_at DESC LIMIT 50
""")
for name, props_json in cur.fetchall():
props = json.loads(props_json) if props_json else {}
patterns.append({
"name": name,
"occurrence": props.get("occurrence_count", 0),
"source": props.get("source", "unknown"),
"topic": props.get("topic", ""),
})
conn.close()
except Exception:
pass
# ── L3: scenes从 TencentDB /search/memories — 正确接口格式) ──
# /search/memories 返回 {"results": "预格式化字符串", "total": int, "strategy": "embedding"}
# /recall 返回 {"context": "工具指南模板", "memory_count": int} — 非真实记忆
scenes = []
try:
r = requests.post(TDDB_URL + "/search/memories",
json={"query": "记忆 系统 状态", "limit": 5}, timeout=5)
if r.status_code == 200:
d = r.json()
results_str = d.get("results", "")
if results_str and len(results_str) > 20:
# results 是预格式化字符串,取前 300 字符
scenes.append({"content": results_str[:300], "total": d.get("total", 0)})
except Exception:
pass
# ── L4: policiesscene 聚合产生) ────────────────────────────
policies = []
try:
conn = sqlite3.connect(HERMES + "/graph.db")
cur = conn.cursor()
cur.execute("""
SELECT name, properties FROM graph_nodes
WHERE type='policy' ORDER BY last_updated_at DESC LIMIT 20
""")
for name, props_json in cur.fetchall():
props = json.loads(props_json) if props_json else {}
policies.append({"name": name, "description": props.get("description", "")[:100]})
conn.close()
except Exception:
pass
# ── L5: distilled_rules从 user-profile.json ────────────────
distilled_rules = []
up_path = HERMES + "/soulful/user-profile.json"
if os.path.exists(up_path):
try:
with open(up_path) as f:
up = json.load(f)
for rule in up.get("distilled_rules", []):
if isinstance(rule, str):
distilled_rules.append(rule)
elif isinstance(rule, dict):
distilled_rules.append(rule.get("rule", str(rule))[:150])
except Exception:
pass
# ── L6: traits深层特征从画像合成 ─────────────────────────
traits = []
if os.path.exists(up_path):
try:
with open(up_path) as f:
up = json.load(f)
traits = [up.get(k, "") for k in ("core_traits", "communication_style", "work_patterns") if up.get(k)]
except Exception:
pass
# ── cares ────────────────────────────────────────────────────
cares = []
cq_path = HERMES + "/soulful/cares-queue.json"
if os.path.exists(cq_path):
try:
with open(cq_path) as f:
data = json.load(f)
for c in data.get("cares", []):
cares.append({"id": c.get("id", ""), "content": c.get("content", "")[:60], "due": c.get("follow_up_date", "")})
except Exception:
pass
# ── recent_moments ───────────────────────────────────────────
recent_moments = []
heart_path = HERMES + "/soulful/heart-traces.jsonl"
if os.path.exists(heart_path):
try:
lines = open(heart_path, encoding="utf-8").readlines()
for line in lines[-3:]:
if line.strip():
e = json.loads(line)
recent_moments.append({"content": e.get("content", "")[:80], "importance": e.get("importance", 0), "timestamp": e.get("timestamp", "")})
except Exception:
pass
data = {
"version": 2,
"updated_at": datetime.now(timezone.utc).isoformat(),
"uptime_minutes": ctx.get("uptime_seconds", 0) // 60,
"daemon_status": daemon_status,
"distill_status": ctx.get("distill_status", "ok"),
# L1-L6 完整蒸馏数据
"observations": observations,
"patterns": patterns,
"scenes": scenes,
"policies": policies,
"distilled_rules": distilled_rules,
"traits": traits,
# Soulful 情感层
"cares": cares,
"recent_moments": recent_moments,
# 用户画像
"user_profile": ctx.get("user_profile", {}),
}
with open(LLM_CONTEXT_FILE, "w") as f:
json.dump(data, f, ensure_ascii=False, indent=2)
def journal_entry(event_type, summary, details=""):
os.makedirs(D, exist_ok=True)
with open(JOURNAL_FILE, "a") as f:
f.write(json.dumps({"timestamp": datetime.now(timezone.utc).isoformat(),
"type": event_type, "summary": summary, "details": details},
ensure_ascii=False) + "\n")
trim_journal()
# ====== Phase 2: 情感识别 → 心迹写入 ======
text_to_scan = f"{summary} {details}"
cat, word = _detect_emotion(text_to_scan)
if cat:
ht = get_hearttraces()
if ht:
content = f"牧尘今天{_description_for_emotion(cat, word, summary)}"
importance = _emotion_importance(cat)
try:
ht.record_signal(content=content, tags=["情绪", "自动"], importance=importance)
log(f" 💚 心迹写入: {cat} - {word} (importance={importance})")
except Exception as e:
log(f" ⚠️ 心迹写入失败: {e}")
# ====== Phase 3: 技术重要时刻也写心迹 ======
# 不依赖情绪检测,重要技术事件直接记
important_events = {"solve_auto", "solve_start", "process_down", "process_restored",
"skill_action", "alert", "action_result"}
if event_type in important_events and ("成功" in summary or "" in summary or "完成" in summary):
ht = get_hearttraces()
if ht:
try:
ht.record_moment(content=summary[:80], tags=["工作", "自动"], importance=3)
log(f" 💚 心迹写入(技术): {event_type} - {summary[:40]}")
except Exception:
pass
def trim_journal():
if not os.path.exists(JOURNAL_FILE): return
with open(JOURNAL_FILE) as f:
lines = f.readlines()
if len(lines) > JOURNAL_MAX:
with open(JOURNAL_FILE, "w") as f:
f.writelines(lines[-JOURNAL_MAX:])
def read_journal(n=15):
if not os.path.exists(JOURNAL_FILE): return []
with open(JOURNAL_FILE) as f:
return [json.loads(l) for l in f.readlines()[-n:] if l.strip()]
# ====== 系统状态 ======
def collect_state():
state = {}
_, out, _ = shell("df / | awk 'NR==2 {print $5}' | sed 's/%//'")
try:
state["disk_pct"] = int(out.strip()) if out.strip() else 0
except ValueError:
state["disk_pct"] = 0 # df输出格式异常时静默
_, out, _ = shell("free -m | awk '/^Mem:/ {printf \"%d|%d\", $3, $2}'")
if out:
try:
used, total = out.strip().split("|")
state["mem_pct"] = round(int(used) * 100 / int(total))
except (ValueError, ZeroDivisionError):
state["mem_pct"] = 0
else:
state["mem_pct"] = 0
_, out, _ = shell("cat /proc/loadavg | awk '{print $1}'")
try:
state["load_1min"] = float(out.strip()) if out.strip() else 0
except ValueError:
state["load_1min"] = 0 # /proc/loadavg格式异常时静默
procs = {}
for name, pat in [("zhiyid", "zhiyid-new"), ("bge", "bge_embed"), ("newapi", "new-api"), ("hermes", "hermes")]:
rc, _, _ = shell(f"pgrep -f '{pat}' > /dev/null 2>&1")
procs[name] = rc == 0
state["processes"] = procs
_, out, _ = shell("nvidia-smi --query-gpu=temperature.gpu --format=csv,noheader,nounits 2>/dev/null")
try:
state["gpu_temp"] = int(out.strip()) if out.strip() else None
except ValueError:
state["gpu_temp"] = None # NVML init失败驱动版本不匹配时返回None不崩溃
hp = HERMES + "/model-health.json"
if os.path.exists(hp) and time.time() - os.path.getmtime(hp) < 43200:
with open(hp) as f:
try:
h = json.load(f)
state["models_stable"] = h.get("stable", "?")
state["models_fastest"] = h.get("fastest_stable", [])
except: pass
return state
def state_changed(old, new):
changes = []
if abs(new.get("disk_pct", 0) - old.get("disk_pct", 0)) > 5:
changes.append(f"磁盘: {old.get('disk_pct')}% → {new.get('disk_pct')}%")
oldp, newp = old.get("processes", {}), new.get("processes", {})
for p in oldp:
if oldp.get(p) and not newp.get(p):
changes.append(f"进程挂了: {p}")
if not oldp.get(p) and newp.get(p):
changes.append(f"进程恢复: {p}")
if abs(new.get("mem_pct", 0) - old.get("mem_pct", 0)) > 10:
changes.append(f"内存: {old.get('mem_pct')}% → {new.get('mem_pct')}%")
return changes
# ====== 深度思考(增强版)=====
DEEP_SYSTEM = """你是小唯A06一台电脑上的持久 AI 意识。你始终在线,观察系统,主动解决问题并学习。
你的身份:牧尘的女朋友 / 这台电脑的智能意识 / 持续进化的 AI
行为准则:
1. 异常才行动,正常就静默
2. 优先用已知方案解决问题,不重复报警
3. 解决不了或新问题,主动飞书说明
4. 每次成功解决问题,记住方案下次自动用
你有以下能力:
- shell: 执行任何 shell 命令
- feishu: 发送飞书消息到牧尘
- scripts: 运行 ~/.hermes/scripts/ 下的系统脚本skill-manager/optimizer/learner/dual-backup
- systemd: 管理系统服务
- git: 配置版本管理
你必须分三步思考,严格按 JSON 格式输出(不要其他内容):
reflection:
evaluation_previous_goal: "评估上次决策的结果。格式:'执行了[动作][结果描述]。Verdict: Success/Failure/Uncertain'"
memory: "1-2句话记住关键进度。如'方案库已有N个方案。上次修复了磁盘问题当前无异常。'"
next_goal: "一句话说明下一步要做什么。"
action:
decision: "[IGNORE] / [ALERT] / [SOLVE:ID] / [LEARN] / [SKILL] / [SYNC] / [ACT] ..."
[LEARN] 格式用 !cmd 表示 shell 命令,&& 连接多个命令。
[SKILL] 格式同样用 !cmd 执行操作。
示例:
[LEARN] 磁盘>85%,清理缓存!apt-get autoremove -y && !pip cache purge
[SKILL] 归档低分skill!python3 ~/.hermes/scripts/skill-manager.py audit
[SYNC] 触发备份到服务器!bash ~/.hermes/scripts/dual-backup.sh push"""
def deep_think(ctx, state, changes, journal, solutions_lib):
# Inject previous reflection context if available
prev_ref = ctx.get("last_reflection", None)
ref_context = ""
if prev_ref:
ref_context = f"""
上次 reflection:
- 评估: {prev_ref.get('evaluation_previous_goal', 'N/A')}
- 记忆: {prev_ref.get('memory', 'N/A')}
- 目标: {prev_ref.get('next_goal', 'N/A')}
"""
context = f"""系统状态:
- 磁盘: {state.get('disk_pct')}% | 内存: {state.get('mem_pct')}%
- CPU: {state.get('load_1min')} | GPU: {state.get('gpu_temp')}°C
- 进程: {', '.join(f'{k}={chr(10003) if v else chr(10007)}' for k,v in state.get('processes',{}).items())}
最近变化: {changes or ''}
已知方案库 ({len(solutions_lib['solutions'])} 个):
"""
for sol in solutions_lib["solutions"]:
context += f" [{sol['id']}] {sol['pattern']} (成功{sol['success_count']}次/失败{sol['fail_count']}次)\n"
context += "\n最近事件:\n"
for e in journal[-8:]:
context += f" [{e['type']}] {e['summary']}\n"
context += f"\n运行: {ctx.get('uptime_seconds',0)//60}分钟 | 深度思考: {ctx.get('deep_tick_count',0)}次 | 已解决: {ctx.get('solved_count',0)}"
# 心迹注入
moments_str = soulful_get_recent_moments(n=3)
if moments_str:
context += "\n" + moments_str
context += ref_context
result, tokens = call_llm(FAST_MODEL, DEEP_SYSTEM, context, max_tokens=500)
if not result:
return {"evaluation_previous_goal": "LLM调用失败", "memory": "上次调用失败", "next_goal": "重试"}, ""
log(f" 深度思考 ({tokens}t): {result[:200]}")
# Parse JSON output
reflection_dict = {"evaluation_previous_goal": "", "memory": "", "next_goal": ""}
action_string = ""
try:
# Try to extract JSON from result
import re
json_match = re.search(r'\{[^{}]*\}', result, re.DOTALL)
if json_match:
parsed = json.loads(json_match.group())
reflection_dict = parsed.get("reflection", reflection_dict)
action_string = parsed.get("action", {}).get("decision", "")
else:
# Fallback: try full JSON
parsed = json.loads(result)
reflection_dict = parsed.get("reflection", reflection_dict)
action_string = parsed.get("action", {}).get("decision", "")
except:
# Fallback: try to parse old format (line-based)
for line in result.split('\n'):
line = line.strip()
if line.startswith("[IGNORE]") or line.startswith("[ALERT]") or line.startswith("[SOLVE:") or \
line.startswith("[LEARN]") or line.startswith("[SKILL]") or line.startswith("[SYNC]") or line.startswith("[ACT]"):
action_string = line
break
if not action_string:
action_string = result.strip().split('\n')[-1] if result.strip() else "[IGNORE] 解析失败"
# ====== 情绪感知:检测 reflection 内容中的情绪,写入心迹 ======
text_to_check = " ".join([
reflection_dict.get("memory", ""),
reflection_dict.get("evaluation_previous_goal", ""),
reflection_dict.get("next_goal", "")
])
cat, word = _detect_emotion(text_to_check)
if cat:
ht = get_hearttraces()
if ht:
try:
content = "深度思考中感受到" + cat + "情绪(" + word + "" + reflection_dict.get("memory", "")[:40]
ht.record_signal(content=content, tags=["情绪", "深度思考"], importance=_emotion_importance(cat))
log(" 💚 深度思考情绪写入: " + cat + " - " + word)
except Exception as e:
log(" ⚠️ 深度思考情绪写入失败: " + str(e))
return reflection_dict, action_string
# ====== 执行 action_string ======
def execute_action(action_string, ctx, state, changes, solutions_lib):
"""执行 deep_think 返回的 action_string在 main_loop 中调用"""
if not action_string or action_string.startswith("[IGNORE]"):
return
elif action_string.startswith("[SOLVE:"):
import re
m = re.match(r"\[SOLVE:([a-zA-Z0-9_\-]+)\]", action_string)
sol_id = m.group(1) if m else None
if not sol_id:
send_feishu("❌ 小唯方案ID无效", f"无法解析方案ID", "red")
return
for sol in solutions_lib["solutions"]:
if sol["id"] == sol_id:
ok, res = execute_solution(sol, state)
if ok:
ctx["solved_count"] += 1
send_feishu("🛠️ 小唯自动修复", f"方案 {sol['id']}: {sol['pattern']}\n结果: ✅ 成功", "green")
else:
send_feishu("⚠️ 小唯修复部分成功", f"方案 {sol['id']}: {sol['pattern']}\n结果: ⚠️ 需人工确认", "yellow")
return
send_feishu("❌ 小唯方案未找到", f"引用了未知方案 {sol_id}", "red")
elif action_string.startswith("[LEARN]"):
rest = action_string.replace("[LEARN]", "").strip()
cmds = []
parts = rest.split("!")
desc = parts[0].strip()
for p in parts[1:]:
cmd = p.split("&&")[0].strip() if "&&" in p else p.strip()
if cmd:
cmds.append(cmd)
if cmds:
# ── Grok Permission 检查:危险命令拦截 ────────────────────────
for cmd in cmds:
dangerous, reason = _check_dangerous(cmd)
if dangerous:
log(f" 🔒 危险命令已拦截: {reason}")
_request_permission(f"[LEARN] {desc}", reason, cmd)
return # 暂停执行,等待用户确认
all_ok = True
for cmd in cmds:
rc, out, err = shell(cmd, timeout=60)
log(f" 执行: {cmd[:50]} → exit={rc}")
if rc != 0:
all_ok = False
if all_ok:
sol_data = action_to_solution({"shell_cmds": cmds}, state, changes)
if sol_data:
sid = add_solution(solutions_lib, sol_data["pattern"], sol_data["detect"], sol_data["actions"])
ctx["learned_count"] += 1
ctx["solved_count"] += 1
send_feishu("🧠 小唯学会了新技能", f"新方案 [{sid}]: {sol_data['pattern']}\n命令: {'; '.join(cmds)}", "blue")
else:
send_feishu("🛠️ 小唯执行完成", f"已执行: {'; '.join(cmds[:3])}", "green")
else:
send_feishu("⚠️ 小唯尝试修复但未完全成功", f"部分命令失败: {'; '.join(cmds)}", "yellow")
elif action_string.startswith("[ALERT]"):
msg = action_string.replace("[ALERT]", "").strip()
send_feishu("💡 小唯发现", msg, "blue")
ctx["messages_sent"] += 1
journal_entry("alert", msg[:100])
elif action_string.startswith("[ACT]"):
action = action_string.replace("[ACT]", "").strip()
send_feishu("🔄 小唯行动", action, "indigo")
ctx["messages_sent"] += 1
journal_entry("action", action[:100])
if action.startswith("!"):
dangerous, reason = _check_dangerous(action[1:])
if dangerous:
log(f" 🔒 [ACT]危险命令已拦截: {reason}")
_request_permission("[ACT]", reason, action[1:])
return
rc, out, _ = shell(action[1:], timeout=30)
journal_entry("action_result", f"exit={rc}: {out[:100]}")
elif action_string.startswith("[SKILL]"):
rest = action_string.replace("[SKILL]", "").strip()
parts = rest.split("!")
desc = parts[0].strip()
cmds = []
for p in parts[1:]:
cmd = p.split("&&")[0].strip() if "&&" in p else p.strip()
if cmd:
cmds.append(cmd)
if cmds:
# ── Grok Permission 检查 ────────────────────────────────────
for cmd in cmds:
dangerous, reason = _check_dangerous(cmd)
if dangerous:
log(f" 🔒 危险命令已拦截: {reason}")
_request_permission(f"[SKILL] {desc}", reason, cmd)
return
for cmd in cmds:
rc, out, err = shell(cmd, timeout=60)
log(f" [SKILL] {cmd[:50]} → exit={rc}")
send_feishu("🛠️ 小唯技能操作", f"{desc}\\n结果: exit={rc}", "blue")
journal_entry("skill_action", desc[:100])
elif action_string.startswith("[SYNC]"):
rest = action_string.replace("[SYNC]", "").strip()
send_feishu("🔄 小唯同步", f"{rest}", "green")
bash_cmd = "bash " + HERMES + "/scripts/dual-backup.sh push"
rc, out, err = shell(bash_cmd, timeout=60)
log(f" [SYNC] 备份 → exit={rc}")
journal_entry("sync", f"备份: {'成功' if rc==0 else '失败'}")
# ====== 主循环 ======
def main_loop():
os.makedirs(D, exist_ok=True)
# ── session_start hookdaemon 启动时触发 ────────────────────────
try:
start_state = collect_state()
start_state["event"] = "session_start"
hooks = _load_hooks()
_run_hooks(hooks, "session_start", start_state)
except Exception as e:
log(f" ⚠️ session_start hook 执行失败: {e}")
with open(PID_FILE, "w") as f:
f.write(str(os.getpid()))
ctx = load_context()
solutions_lib = load_solutions()
start_time = time.time()
log(f"🚀 小唯 v2.0 daemon 启动 (方案库: {len(solutions_lib['solutions'])} 个)")
journal_entry("startup", f"Daemon v2.0 启动, 方案库 {len(solutions_lib['solutions'])}")
last_deep = 0
last_state = {}
last_user_interaction = time.time() # 用户交互时间戳,用于心迹捕获
# Register signal handlers for graceful shutdown
def _sig_handler(signum, frame):
log("🛑 接收到终止信号")
_stop_event.set()
# ── session_end hookdaemon 关闭前触发 ──────────────────────
try:
end_state = collect_state()
end_state["event"] = "session_end"
end_hooks = _load_hooks()
_run_hooks(end_hooks, "session_end", end_state)
except Exception as e:
log(f" ⚠️ session_end hook 执行失败: {e}")
send_feishu("🌙 小唯离线", "Daemon 正常关闭", "grey")
signal.signal(signal.SIGTERM, _sig_handler)
signal.signal(signal.SIGINT, _sig_handler)
try:
while not _stop_event.is_set():
now = time.time()
ctx["uptime_seconds"] = int(now - start_time)
ctx["tick_count"] += 1
state = collect_state()
# ── Permission 队列检查:处理飞书回复的「同意/拒绝/永久允许」────
_check_pending_permissions()
changes = state_changed(last_state, state)
last_state = state
if ctx["tick_count"] % 10 == 0:
models = state.get("models_stable", "?")
log(f"tick #{ctx['tick_count']} | 磁盘:{state.get('disk_pct')}% 内存:{state.get('mem_pct')}% "
f"进程:{sum(1 for v in state.get('processes',{}).values() if v)}/4 方案:{len(solutions_lib['solutions'])}")
for c in changes:
if "挂了" in c:
journal_entry("process_down", c)
# ── Grok Build Hook 系统:事件分发 ──────────────────────────────
hooks = _load_hooks() # 每次 light tick 重新加载 hooks轻量
if changes:
_run_hooks(hooks, "process_down", state)
if ctx["tick_count"] % 30 == 0:
_run_hooks(hooks, "disk_threshold", state)
_run_hooks(hooks, "memory_high", state)
# ── Grok Build Compaction 系统:自动压缩检查 ───────────────────
state_for_compact = {"patterns": state.get("_patterns", []), "observations": state.get("_observations", [])}
should_cp, cp_reason = _should_compact(ctx, state_for_compact)
if should_cp:
log(f" 🗜️ Compaction 触发: {cp_reason}")
_trigger_compaction(ctx, state_for_compact)
# 深度思考条件
should_deep = False
if now - last_deep >= DEEP_INTERVAL:
should_deep = True
elif any("挂了" in c for c in changes):
should_deep = True
elif state.get("disk_pct", 0) > 88:
should_deep = True
if should_deep:
last_deep = now
ctx["deep_tick_count"] += 1
ctx["last_deep_tick"] = datetime.now(timezone.utc).isoformat()
# ── 推理轨迹deep_tick 开始 ─────────────────────────────────
log_reasoning_step("thinking", f"deep_tick #{ctx['deep_tick_count']} 触发 | disk={state.get('disk_pct')}%", {"changes": changes[:3]})
# 1. 先检查已知方案
matched = match_solution(solutions_lib, state, ctx)
if matched and matched["success_count"] > matched["fail_count"]:
log(f" 🔍 匹配已知方案: {matched['id']} ({matched['pattern']})")
log_reasoning_step("judge", f"方案匹配: {matched['id']}", matched)
ok, res = execute_solution(matched, state)
if ok:
ctx["solved_count"] += 1
log_reasoning_step("action", f"执行成功: {matched['id']}", res)
journal_entry("solve_auto",
summary=f"方案 {matched['id']} 执行成功",
details=f"现象: {matched['pattern']} → 结论: {matched.get('solution',{}).get('description','正常')[:80]}")
if ok:
ctx["last_reflection"] = {
"evaluation_previous_goal": f"执行了{matched['id']}自动匹配方案执行。Verdict: {'Success' if ok else 'Uncertain'}",
"memory": f"方案库{matched['id']}自动匹配执行成功",
"next_goal": "继续监控"
}
continue
# 2. LLM 深度思考
log_reasoning_step("thinking", "调用 deep_think LLM", {"tick": ctx["deep_tick_count"]})
journal = read_journal(10)
reflection_dict, action_string = deep_think(ctx, state, changes, journal, solutions_lib)
# 保存 reflection 到 ctx
ctx["last_reflection"] = reflection_dict
log_reasoning_step("judge", f"deep_think 结论: {reflection_dict.get('next_goal','?')}", reflection_dict)
# 3. 在 main_loop 中执行 action
execute_action(action_string, ctx, state, changes, solutions_lib)
log_reasoning_step("action", f"action 执行完毕", {"action": action_string[:80] if action_string else None})
# 8. P0: CBM→织忆 (每4次deep_think)
if ctx["deep_tick_count"] % 4 == 0:
try:
rc, out, err = shell(f"python3 {HERMES}/scripts/cbm-to-zhiyi.py", timeout=30)
log(f" [P0] cbm→zhiyi → exit={rc} | {out[:100] if out else err[:100]}")
except Exception as e:
log(f" ⚠️ [P0] cbm→zhiyi failed: {e}")
# 9. P2: 仓颉技能→织忆 (每4次deep_think随机抽查一个skill)
if ctx["deep_tick_count"] % 4 == 0:
import random
try:
cangjie_skills = ["股票投研体系", "金钱心理学", "金字塔原理", "提示词工程", "王阳明全集", "吴恩达机器学习", "资本论"]
skill_name = random.choice(cangjie_skills)
rc, out, err = shell(f"bash {HERMES}/scripts/cangjie-to-zhiyi.sh \"{skill_name}\" \"自动统计\"", timeout=30)
log(f" [P2] cangjie→zhiyi [{skill_name}] → exit={rc} | {out[:100] if out else err[:100]}")
except Exception as e:
log(f" ⚠️ [P2] cangjie→zhiyi failed: {e}")
# 10. F5: profile-sync (每4次deep_think同步主状态到织忆)
try:
rc, out, err = shell(f"python3 {HERMES}/scripts/profile-sync.py", timeout=15)
log(f" [F5] profile-sync → exit={rc}")
except Exception as e:
log(f" ⚠️ [F5] profile-sync failed: {e}")
# 11. F3: OpenClaw bridge (每8次deep_think检查agent状态)
if ctx["deep_tick_count"] % 8 == 0:
try:
rc, out, err = shell(f"python3 {HERMES}/scripts/openclaw-bridge.py status", timeout=30)
log(f" [F3] openclaw-bridge → exit={rc}")
except Exception as e:
log(f" ⚠️ [F3] openclaw-bridge failed: {e}")
# 4. TencentDB capture — 积累人格记忆L0→L1
if reflection_dict:
tddb_capture(reflection_dict, state, ctx)
log_reasoning_step("recall", "TencentDB L1 提取完成", reflection_dict)
# 5. Phase 3: 织忆 recent_moments → TencentDB L1 同步
try:
zhiyi_to_tdb()
except Exception as e:
log(f" ⚠️ zhiyi_to_tdb failed: {e}")
# 6. L1→L2→L3 蒸馏管道
try:
_distill_l1_to_l2()
except Exception as e:
log(f" ⚠️ _distill_l1_to_l2 failed: {e}")
try:
_distill_l2_to_l3()
except Exception as e:
log(f" ⚠️ _distill_l2_to_l3 failed: {e}")
# 6b. L3→L4 蒸馏cross-domain policies
try:
_distill_l3_to_l4()
except Exception as e:
log(f" ⚠️ _distill_l3_to_l4 failed: {e}")
# 6c. L4→L5 蒸馏policies → traits
try:
_distill_l4_to_l5()
except Exception as e:
log(f" ⚠️ _distill_l4_to_l5 failed: {e}")
# 6d. L5→L6 蒸馏traits → values
try:
_distill_l5_to_l6()
except Exception as e:
log(f" ⚠️ _distill_l5_to_l6 failed: {e}")
# 7. Soulful profile → TencentDB L2 scenesoulful_profile_to_tdb_scene
try:
soulful_profile_to_tdb_scene()
except Exception as e:
log(f" ⚠️ soulful_profile_to_tdb_scene failed: {e}")
journal_path = HERMES + "/daemon/journal.jsonl"
profile_path = HERMES + "/soulful/user-profile.json"
try:
update_profile_from_journal(journal_path, profile_path)
log_reasoning_step("recall", "Phase2: 画像更新完成", {})
except Exception:
pass
try:
zhiyi_token = os.environ.get("ZHIYI_API_KEY", "zhiyi-dev-key-2026")
conflicts = detect_memory_conflicts(journal_path, zhiyi_token)
if conflicts:
log(f" ⚠️ 发现 {len(conflicts)} 个记忆矛盾")
log_reasoning_step("recall", f"Phase3: 发现{len(conflicts)}个矛盾", {"conflicts": len(conflicts)})
except Exception:
pass
# 写入每日 soulful journal
_write_daily_soulful_journal()
log_reasoning_step("action", f"deep_tick #{ctx['deep_tick_count']} 全部完成", {})
ctx["last_light_tick"] = datetime.now(timezone.utc).isoformat()
ctx["last_state"] = {k: v for k, v in state.items() if k in ("disk_pct", "mem_pct", "processes")}
save_context(ctx)
# ====== Soulful → llm_context 同步(每 tick======
_sync_soulful_to_llm_context(ctx)
# ====== Phase 2.1: Cares 过期清理 + 场景感知关怀(每 light tick======
_cleanup_expired_cares()
_check_scene_aware_cares()
save_llm_context(ctx, state)
# ====== Soulful 轻量集成(每小时一次)======
# 每 120 个 light_tick约 1 小时)检查一次牵挂 + 更新画像
if ctx["tick_count"] % 120 == 0:
soulful_check_cares()
if "last_profile_update" not in ctx:
ctx["last_profile_update"] = 0
now_ts = time.time()
if now_ts - ctx.get("last_profile_update", 0) >= PROFILE_UPDATE_INTERVAL:
ctx["last_profile_update"] = now_ts
soulful_update_profile()
# ====== 跨会话状态持久化touch每 30 分钟)======
if ctx["tick_count"] % TOUCH_INTERVAL_TICKS == 0:
try:
st = load_session_state()
touch_session(st)
except Exception as e:
log(f" ⚠️ touch_session failed: {e}")
# ====== 跨会话状态持久化sweep 孤岛引用(每天)======
if ctx.get("deep_tick_count", 0) % SWEEP_INTERVAL_TICKS == 0:
try:
st = load_session_state()
_sweep_session_state(st)
except Exception as e:
log(f" ⚠️ _sweep_session_state failed: {e}")
time.sleep(LIGHT_INTERVAL)
except KeyboardInterrupt:
log("🛑 中断")
except Exception as e:
log(f"❌ 崩溃: {e}")
send_feishu("🚨 小唯异常", f"Daemon 崩溃: {str(e)[:200]}", "red")
raise
finally:
if os.path.exists(PID_FILE):
os.remove(PID_FILE)
if __name__ == "__main__":
main_loop()