deerflow-code/batch-distill/distill_people.py
2026-09-07 18:24:55 +08:00

564 lines
22 KiB
Python
Raw Permalink 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
# -*- coding: utf-8 -*-
"""
批量人物蒸馏 → llmwiki markdown(完全独立的单文件客户端 + 内置定时器)
特点
----
- 串行(一个一个来)调用正在运行的 DeerFlow 服务问答接口
- 复用线上「女娲 / huashu-nuwa」技能做提炼;**不联网**,只用已配置的
知识库查询技能(knowledge-search-v2)检索资料
- agent 只负责产出内容,**markdown 由本脚本抓取最终回复后自己保存**
- 内置时间窗:每天 23:00 ~ 次日 08:00 才执行,到点自动暂停、次日继续
- 断点续跑:已完成的人自动跳过;中途加名单也会被带上
你只需要改一处:下面的 BASE_URL(线上服务地址)。其余都不用填。
名单
----
脚本同目录放 people.xlsx 或 people.csv,**第一列 = 人名**,其余列=可选提示。
运行
----
python distill_people.py # 常驻,按 23:00~08:00 窗口自动执行
python distill_people.py --now # 忽略时间窗,立刻把名单跑完(测试用)
python distill_people.py --redo # 忽略已完成,全部重做
"""
from __future__ import annotations
import argparse
import csv
import json
import re
import sys
import time
from datetime import datetime, time as dtime
from pathlib import Path
try:
import requests
except ImportError:
sys.exit("缺少依赖 requests:pip install requests")
# ════════════════════════════════════════════════════════════════════════
# 只改这一行:你的线上服务地址(不带结尾斜杠)
# ════════════════════════════════════════════════════════════════════════
BASE_URL = "http://47.94.209.59:2026/"
# ════════════════════════════════════════════════════════════════════════
# ---- 以下都已内置,正常无需改动 ----
USERNAME = "distiller" # 免密自动注册
PASSWORD = "123ewq" # 共享口令(你们开了口令门)
NUWA_SKILL = "huashu-nuwa" # 女娲技能
KB_SKILL = "knowledge-search-v2" # 知识库查询技能(不联网,用它检索)
DEPTH = "full" # full 完整 / fast 精简
HERE = Path(__file__).resolve().parent
OUTPUT_DIR = HERE / "output"
PEOPLE_FILE_DEFAULT = HERE / "people.example.csv"
# 执行时间窗:23:00 ~ 次日 08:00(跨午夜)
WIN_START = dtime(23, 0)
WIN_END = dtime(8, 0)
PER_PERSON_TIMEOUT = 1800 # 单人最长等待(秒)
POLL_INTERVAL = 5 # 轮询间隔(秒)
HTTP_TIMEOUT = 60
TERMINAL_STATUS = {"success", "error", "timeout", "interrupted"}
# ================================ llmwiki 模板 ================================
LLMWIKI_TEMPLATE = """\
---
name: <中文名/常用名>
aliases: [<别名/英文名...>]
type: person
era: <在世 | 历史人物>
domains: [<领域1>, <领域2>]
distilled_at: <YYYY-MM-DD>
sources_count: <检索到的来源条数>
primary_source_ratio: <一手来源占比, 0~1 的小数>
confidence: <high | medium | low>
---
# <人物名>
> **一句话定位**:<此人看世界的独特方式,一句话>
## 身份卡
- **活跃年代**:
- **核心身份**:
- **自述(50字, 以其语气)**:
## 核心观点(真信念 · 反复出现的主张)
1. ...
## 心智模型(3–7 个;每个含:一句话 / 证据场景≥2 / 应用 / 局限 / 来源)
### 1. <模型名>
- **一句话**:
- **证据场景**:
- **应用方式**:
- **局限**:
- **来源**:
## 决策启发式(5–10 条 · 如果 X 则 Y,带案例)
- 若 <X> → <Y>
## 表达 DNA
| 维度 | 特征 |
|------|------|
| 句式 | |
| 高频词 | |
| 幽默 | |
| 确定性 | |
| 引用习惯 | |
## 价值观与反模式
- **价值观排序**:
- **明确反对**:
## 内在张力(≥2 对矛盾)
- <A> ↔ <B>
## 经历时间线
| 年份 | 事件 | 思想意义 |
|------|------|---------|
| | | |
## 代表语录
- “...” —(出处/年份,标一手/二手)
## 智识谱系
- **受影响于**:
- **影响了**:
## 外部评价与争议
-
## 诚实边界
- <至少 3 条具体局限,含信息截止日期>
## 调研来源
**一手(N)**:
**二手(N)**:
"""
def build_prompt(name: str, hints: str, depth: str) -> str:
depth_note = (
"本次走【完整模式】:尽量覆盖著作/对话/表达/他者评价/决策/时间线 6 个维度。"
if depth == "full"
else "本次走【精简模式】:重点覆盖核心观点、心智模型、时间线三块即可。"
)
hint_block = f"\n补充提示(来自名单其它列,可作为检索关键词):{hints}\n" if hints.strip() else ""
return f"""\
你的任务:把人物【{name}】蒸馏成一份 llmwiki 风格的人物档案 markdown。
请复用本服务中已启用的「女娲 / {NUWA_SKILL}」技能的提炼方法论
(心智模型三重验证、表达DNA分析、诚实边界)。
⚠️ 本环境**无法联网,严禁使用 WebSearch 等任何联网工具**。
所有资料**只能通过已配置的知识库查询技能「{KB_SKILL}」检索获取**。
请充分调用该知识库技能,围绕该人物多角度检索(生平、观点、著作、语录、决策、评价等),
基于检索到的内容做提炼。知识库里查不到的部分,如实在「诚实边界」标注,不要编造。
{depth_note}
{hint_block}
输出要求:
1. **不要写任何文件,不要调用 write_file / present_files**。
2. 直接在你的回复正文里输出**完整中文 markdown 全文**,严格按下面模板结构
(保留各级标题,逐项填实,不要保留尖括号占位符),**不要任何额外解说、前言或结语**:
------------------ 模板开始 ------------------
{LLMWIKI_TEMPLATE}
------------------ 模板结束 ------------------
现在开始,独立完成,不要向我提问。"""
# ================================ HTTP 客户端 ================================
class DeerFlowClient:
def __init__(self, base_url: str):
self.base = base_url.rstrip("/")
self.s = requests.Session()
self.token: str | None = None
def _headers(self) -> dict:
h = {"Content-Type": "application/json"}
if self.token:
h["Authorization"] = f"Bearer {self.token}" # 只用 Bearer,不带 cookie → 跳过 CSRF
return h
def login(self) -> bool:
url = f"{self.base}/api/v1/auth/login/username"
params = {"password": PASSWORD} if PASSWORD else None
try:
r = self.s.post(url, params=params, json={"username": USERNAME}, timeout=HTTP_TIMEOUT)
except requests.RequestException as e:
print(f" [登录] 连接失败:{e}")
return False
if r.status_code == 200:
self.token = r.json().get("access_token")
self.s.cookies.clear()
print(f" [登录] 成功,用户={USERNAME}")
return True
print(f" [登录] 返回 {r.status_code}:{r.text[:200]}")
return False
def ensure_skill_enabled(self, name: str) -> None:
try:
r = self.s.get(f"{self.base}/api/skills", headers=self._headers(), timeout=HTTP_TIMEOUT)
if r.status_code == 200:
skills = r.json().get("skills", [])
hit = next((x for x in skills if x.get("name") == name), None)
if hit is None:
print(f" [技能] 警告:服务里没找到 '{name}'")
return
if hit.get("enabled"):
print(f" [技能] '{name}' 已启用")
return
pr = self.s.put(f"{self.base}/api/skills/{name}", headers=self._headers(),
json={"enabled": True}, timeout=HTTP_TIMEOUT)
print(f" [技能] 启用 '{name}' → {pr.status_code}")
except requests.RequestException as e:
print(f" [技能] 检查 '{name}' 失败(忽略):{e}")
def create_thread(self) -> str:
r = self.s.post(f"{self.base}/api/threads", headers=self._headers(), json={}, timeout=HTTP_TIMEOUT)
r.raise_for_status()
return r.json()["thread_id"]
def start_run(self, thread_id: str, prompt: str) -> str:
body = {
"input": {"messages": [{"role": "user", "content": prompt}]},
"excluded_tools": ["ask_clarification"], # 禁澄清,无人值守不卡住
}
r = self.s.post(f"{self.base}/api/threads/{thread_id}/runs",
headers=self._headers(), json=body, timeout=HTTP_TIMEOUT)
r.raise_for_status()
return r.json()["run_id"]
def poll_run(self, thread_id: str, run_id: str) -> str:
deadline = time.monotonic() + PER_PERSON_TIMEOUT
last = ""
while time.monotonic() < deadline:
try:
r = self.s.get(f"{self.base}/api/threads/{thread_id}/runs/{run_id}",
headers=self._headers(), timeout=HTTP_TIMEOUT)
if r.status_code == 200:
st = r.json().get("status", "")
if st != last:
print(f" ...状态:{st}")
last = st
if st in TERMINAL_STATUS:
return st
except requests.RequestException as e:
print(f" 轮询出错(重试):{e}")
time.sleep(POLL_INTERVAL)
return "timeout"
def fetch_result_text(self, thread_id: str, run_id: str) -> str | None:
"""抓取该 thread 的最终 AI 回复文本(由 py 保存为 md)。
注意:本部署的 `/runs/{rid}/messages`、`/messages`、`/events` 端点都返回空
(事件库按登录用户过滤,后台 worker 写入时未带 user_id → 读不到)。
可靠来源是 LangGraph 检查点状态 `/api/threads/{tid}/state` → values.messages,
其中每条消息形如 {"type": "human"|"ai"|"tool", "content": <str | 内容块列表>, ...}。
取最后一条有正文的 ai 消息即可。
"""
try:
r = self.s.get(f"{self.base}/api/threads/{thread_id}/state",
headers=self._headers(), timeout=HTTP_TIMEOUT)
if r.status_code != 200:
print(f" 取回复返回 {r.status_code}:{r.text[:200]}")
return None
messages = (r.json().get("values") or {}).get("messages") or []
# 从后往前找最后一条「有正文」的 AI 回复(跳过只含工具调用的中间消息)
for msg in reversed(messages):
if not isinstance(msg, dict) or msg.get("type") != "ai":
continue
txt = _extract_text(msg)
if txt and txt.strip():
return txt
except requests.RequestException as e:
print(f" 取回复出错:{e}")
return None
def _extract_text(payload) -> str:
"""从一条消息的 content 里提取纯文本。
payload 可能是:
- LangChain 消息的 model_dump 字典(正文在 payload["content"])
- 直接的字符串
- 内容块列表([{"type":"text","text":"..."}, ...])
"""
if payload is None:
return ""
# 若是整条消息的 model_dump,正文在其 content 字段
if isinstance(payload, dict):
c = payload.get("content", payload.get("text", ""))
else:
c = payload
if isinstance(c, str):
return c
if isinstance(c, list):
parts = []
for b in c:
if isinstance(b, dict):
# 只取文本块,跳过 tool_use / thinking 等非文本块
if b.get("type") in (None, "text", "text_block"):
parts.append(b.get("text") or b.get("content") or "")
elif b.get("text"):
parts.append(b["text"])
elif isinstance(b, str):
parts.append(b)
return "\n".join(p for p in parts if p)
return str(c) if c else ""
def _clean_markdown(text: str) -> str:
"""去掉 agent 可能套上的 ```markdown ... ``` 围栏,并规整 YAML frontmatter。"""
t = text.strip()
m = re.match(r"^```(?:markdown|md)?\s*\n(.*)\n```$", t, re.S)
if m:
t = m.group(1).strip()
return _normalize_frontmatter(t)
def _normalize_frontmatter(text: str) -> str:
"""修复模型常见的 frontmatter 跑偏,让 YAML 能正常解析。
模型经常把字段名写成 markdown 加粗(``**name:**``)、加空行、留行尾双空格,
这些都会让 frontmatter 解析失败。这里把开头 ``---`` 到下一个 ``---`` 之间的内容
规整为干净的 ``key: value`` 形式。正文(第二个 ``---`` 之后)保持原样。
"""
if not text.startswith("---"):
return text
# 拆出 frontmatter 块:开头 --- 与下一行单独的 --- 之间
m = re.match(r"^---[ \t]*\n(.*?)\n---[ \t]*\n?(.*)$", text, re.S)
if not m:
return text
body_fm, rest = m.group(1), m.group(2)
out_lines = []
for line in body_fm.splitlines():
s = line.rstrip() # 去行尾双空格
s = s.replace("**", "") # 去 markdown 加粗标记
if not s.strip():
continue # 丢掉空行
out_lines.append(s)
return "---\n" + "\n".join(out_lines) + "\n---\n\n" + rest.lstrip("\n")
# ================================ 名单读取 ================================
def slugify(name: str) -> str:
s = re.sub(r'[\\/:*?"<>|]+', "", name).strip()
s = re.sub(r"\s+", "_", s)
return s or "unnamed"
def read_people(path: Path) -> list[dict]:
if not path.exists():
return []
rows: list[list[str]] = []
if path.suffix.lower() in (".xlsx", ".xlsm"):
try:
import openpyxl
except ImportError:
sys.exit("读取 .xlsx 需要 openpyxl:pip install openpyxl(或把名单存成 .csv)")
wb = openpyxl.load_workbook(path, read_only=True, data_only=True)
for row in wb.active.iter_rows(values_only=True):
rows.append(["" if c is None else str(c).strip() for c in row])
elif path.suffix.lower() == ".csv":
with open(path, encoding="utf-8-sig", newline="") as f:
rows = [[c.strip() for c in r] for r in csv.reader(f)]
else:
sys.exit(f"不支持的名单格式:{path.suffix}(用 .xlsx 或 .csv)")
if rows and re.search(r"(姓名|名字|人物|name)", rows[0][0] if rows[0] else "", re.I):
rows = rows[1:]
people, seen = [], set()
for r in rows:
if not r or not r[0].strip():
continue
name = r[0].strip()
if name in seen:
continue
seen.add(name)
hints = " | ".join(x for x in r[1:] if x.strip())
people.append({"name": name, "hints": hints})
return people
# ================================ 时间窗 ================================
def in_window(now: datetime | None = None) -> bool:
now = now or datetime.now()
t = now.time()
return t >= WIN_START or t < WIN_END # 跨午夜
def seconds_until_window(now: datetime | None = None) -> int:
now = now or datetime.now()
if in_window(now):
return 0
target = now.replace(hour=WIN_START.hour, minute=WIN_START.minute, second=0, microsecond=0)
return max(1, int((target - now).total_seconds()))
def sleep_interruptible(total: int) -> None:
end = time.monotonic() + total
while time.monotonic() < end:
time.sleep(min(30, max(1, int(end - time.monotonic()))))
# ================================ 蒸馏单人 ================================
def distill_one(client: DeerFlowClient, person: dict, out_dir: Path,
manifest: dict, save_manifest, depth: str) -> None:
name, hints = person["name"], person["hints"]
slug = slugify(name)
md_path = out_dir / f"{slug}.md"
t0 = time.time()
rec = {"name": name, "slug": slug, "status": "running",
"started_at": _now(), "file": f"{slug}.md"}
manifest[name] = rec
save_manifest()
try:
tid = client.create_thread()
run_id = client.start_run(tid, build_prompt(name, hints, depth))
status = client.poll_run(tid, run_id)
rec.update(thread_id=tid, run_id=run_id, run_status=status)
text = client.fetch_result_text(tid, run_id) if status == "success" else None
if status == "success" and text:
md = _clean_markdown(text)
md_path.write_text(md, encoding="utf-8")
rec["status"] = "succeeded"
rec["chars"] = len(md)
print(f" ✅ 完成 → {md_path.name}({len(md)} 字,耗时 {int(time.time()-t0)}s)")
else:
rec["status"] = "failed"
rec["error"] = f"run={status}, text={'有' if text else '无'}"
if text:
(out_dir / f"{slug}.partial.md").write_text(text, encoding="utf-8")
print(f" ❌ 失败(run={status})")
except Exception as e:
rec["status"] = "failed"
rec["error"] = str(e)
print(f" ❌ 异常:{e}")
rec["finished_at"] = _now()
rec["seconds"] = int(time.time() - t0)
manifest[name] = rec
save_manifest()
# ================================ 主流程 ================================
def main():
ap = argparse.ArgumentParser(description="批量人物蒸馏 → llmwiki(定时窗 23:00~08:00)")
ap.add_argument("excel", nargs="?", default=str(PEOPLE_FILE_DEFAULT),
help="名单文件(.xlsx 或 .csv),默认同目录 people.xlsx")
ap.add_argument("--now", action="store_true", help="忽略时间窗,立刻跑完整批(测试用)")
ap.add_argument("--redo", action="store_true", help="忽略已完成,全部重做")
args = ap.parse_args()
excel = Path(args.excel)
OUTPUT_DIR.mkdir(parents=True, exist_ok=True)
manifest_path = OUTPUT_DIR / "manifest.json"
manifest: dict = {}
if manifest_path.exists():
try:
manifest = json.loads(manifest_path.read_text(encoding="utf-8"))
except Exception:
manifest = {}
# 校验名单存在
if not read_people(excel):
sys.exit(f"名单为空或不存在:{excel}\n请放 people.xlsx(第一列写人名)。")
# --redo:一次性清掉名单里这些人的旧记录,让他们重做一遍;
# 之后本次会话内跑成功的会正常计为「已完成」,循环才会往后推进
# (否则每轮 pending() 都把已完成的人当未完成,永远卡在第 1 个人)
if args.redo:
for p in read_people(excel):
manifest.pop(p["name"], None)
client = DeerFlowClient(BASE_URL)
if not client.login():
sys.exit("登录失败。请检查 BASE_URL 是否指向正在运行的服务、口令是否正确。")
client.ensure_skill_enabled(NUWA_SKILL)
client.ensure_skill_enabled(KB_SKILL)
def save_manifest():
manifest_path.write_text(json.dumps(manifest, ensure_ascii=False, indent=2), encoding="utf-8")
_write_index(OUTPUT_DIR, read_people(excel), manifest)
def pending() -> list[dict]:
people = read_people(excel)
out = []
for p in people:
r = manifest.get(p["name"], {})
done = r.get("status") == "succeeded" and (OUTPUT_DIR / f"{slugify(p['name'])}.md").exists()
if not done:
out.append(p)
return out
print(f"启动 | 服务={BASE_URL} | 输出={OUTPUT_DIR}")
print(f"模式={'立即跑完' if args.now else '定时窗 23:00~08:00'} | 名单={excel}")
try:
while True:
todo = pending()
if not todo:
print(f"[{_now()}] 全部已完成,待命中(每小时复查名单)...")
if args.now:
break
sleep_interruptible(3600)
continue
if not args.now and not in_window():
wait = seconds_until_window()
nxt = datetime.now().replace(hour=WIN_START.hour, minute=0, second=0, microsecond=0)
print(f"[{_now()}] 不在执行窗口,剩 {len(todo)} 人待蒸;"
f"休眠到 {nxt.strftime('%H:%M')}(约 {wait//60} 分钟)...")
sleep_interruptible(wait)
continue
# 在窗口内(或 --now):处理一个人,然后回到循环重新判断窗口/名单
p = todo[0]
idx_total = len(read_people(excel))
done_n = sum(1 for v in manifest.values() if v.get("status") == "succeeded")
print(f"[{_now()}] ({done_n}/{idx_total} 已完成) 蒸馏:{p['name']}")
distill_one(client, p, OUTPUT_DIR, manifest, save_manifest, DEPTH)
except KeyboardInterrupt:
print("\n已手动停止。重跑脚本会从未完成处继续。")
ok = sum(1 for r in manifest.values() if r.get("status") == "succeeded")
fail = sum(1 for r in manifest.values() if r.get("status") == "failed")
print(f"\n结束:成功 {ok} / 失败 {fail}")
print(f"结果目录:{OUTPUT_DIR} | 清单:{manifest_path} | 索引:{OUTPUT_DIR/'index.md'}")
def _now() -> str:
return time.strftime("%Y-%m-%d %H:%M:%S")
def _write_index(out_dir: Path, people: list[dict], manifest: dict) -> None:
lines = ["# 蒸馏执行清单", "", f"更新时间:{_now()}", "",
"| # | 人物 | 状态 | 字数 | 耗时 | 文件 |",
"|---|------|------|------|------|------|"]
badge = {"succeeded": "✅完成", "failed": "❌失败", "running": "⏳进行中"}
for i, p in enumerate(people, 1):
r = manifest.get(p["name"], {})
st = badge.get(r.get("status", ""), "⌛排队")
f = r.get("file", "")
link = f"[{f}]({f})" if r.get("status") == "succeeded" and f else "-"
lines.append(f"| {i} | {p['name']} | {st} | {r.get('chars','-')} | {r.get('seconds','-')} | {link} |")
(out_dir / "index.md").write_text("\n".join(lines) + "\n", encoding="utf-8")
if __name__ == "__main__":
main()