Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 38 additions & 9 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -66,11 +66,6 @@
# SONA_WEIBO_AISEARCH_REFER=weibo_aisearch
# SONA_REFERENCE_FETCH_TIMEOUT_SEC=12
# SONA_WEIBO_AISEARCH_PLAYWRIGHT_FALLBACK=true
# 是否输出 structured(叙事分槽 + report_bridge,供 report_html 与 wiki RAG 对齐模板)
# SONA_WEIBO_AISEARCH_STRUCTURE=true
# 每槽最多条数、写入 structured 的单条最大字符
# SONA_WEIBO_AISEARCH_SLOT_EACH=8
# SONA_WEIBO_AISEARCH_SNIPPET_IN_STRUCTURE=240

# -----------------------------------------------------------------------------
# 五、/wiki 知识库与 RAG(workflow/wiki_cli、workflow/wiki_rag)
Expand All @@ -91,9 +86,6 @@
# RAG 合成温度、注入微博文本最大字符数
# SONA_WIKI_LLM_TEMPERATURE=0.25
# SONA_WIKI_WEIBO_MAX_CHARS=3500
# 是否在 wiki 智搜附录中附带「叙事分槽」索引块(与 HTML 模板占位符对齐)
# SONA_WIKI_WEIBO_STRUCTURE_IN_AUX=true
# SONA_WIKI_WEIBO_STRUCTURE_MAX_CHARS=1600

# 半自动回流:高价值回答自动写入 wiki/output/_candidates(不自动入正式 output)
# SONA_WIKI_OUTPUT_CANDIDATE_ENABLED=true
Expand All @@ -120,27 +112,64 @@
# SONA_DEBUG_LOG_PATH=

# Graph RAG / Neo4j(可选)
# SONA_NEO4J_URI=bolt://127.0.0.1:7687
# SONA_NEO4J_URI=neo4j://127.0.0.1:7687
# SONA_NEO4J_USER=neo4j
# SONA_NEO4J_PASSWORD=
# SONA_NEO4J_DATABASE=
# SONA_ENABLE_GRAPH_RAG=auto

# Wiki 熊猫/大熊猫图谱问答(可选;与上方 Graph RAG 的 SONA_NEO4J_* 独立)
# SONA_WIKI_NEO4J_PRIORITY=1
# NEO4J_URI=bolt://127.0.0.1:7687
# NEO4J_USERNAME=neo4j
# NEO4J_PASSWORD=
# NEO4J_DATABASE=neo4j

# 专题监测本地 SQLite(默认启用;设 0 时使用内存 demo/test store)
# SONA_TOPIC_MONITOR_SQLITE=1

# -----------------------------------------------------------------------------
# 八、热点 / 多引擎(tools/hottopics、utils/hot_topics_env 等,按需)
# -----------------------------------------------------------------------------

# 热点事件簇领域标注:默认允许使用 LLM 批量标注;设 0 则只用关键词规则
# HOT_DOMAIN_USE_LLM=1

# INSIGHT_ENGINE_API_KEY=
# INSIGHT_ENGINE_BASE_URL=
# INSIGHT_ENGINE_MODEL_NAME=
# 热点流程通义默认:https://dashscope.aliyuncs.com/compatible-mode/v1(与 model.yaml 一致)
# 若 .env 中仍有旧版 baseurl=https://coding.dashscope.aliyuncs.com/v1 会被自动改写,建议删除该行

# 情感分析默认不复用 CSV 情感列;设 1 才允许列兜底/质量保护回退
# SONA_SENTIMENT_ALLOW_COLUMN_FALLBACK=0

# QUERY_ENGINE_API_KEY=
# QUERY_ENGINE_BASE_URL=
# QUERY_ENGINE_MODEL_NAME=
# 热点流程默认将 QUERY_* 与 INSIGHT_* 对齐;勿单独留错误的 QUERY key(否则 Forum 401 而 Insight 正常)。
# 若需论坛与舆情分析使用不同模型,可设 HOT_TOPICS_QUERY_SEPARATE=1 并自行配置 QUERY_*。
# 热点「事件聚类」领域:默认用 INSIGHT_ENGINE_* 对簇代表标题批量调用 qwen-plus 分领域;仅关键词规则可设 HOT_DOMAIN_USE_LLM=0
# HOT_DOMAIN_USE_LLM=0
# 聚类簇数量上限(25–100),减轻打满后「其他」大兜底簇:HOT_CLUSTER_MAX=55
# HOT_CLUSTER_MAX=50

# KIMI_BASE_URL=
# KIMI_MODEL_NAME=

# -----------------------------------------------------------------------------
# 八点五、情感分析(事件分析 Step7,默认仅 LLM,不用 CSV 情感列兜底)
# -----------------------------------------------------------------------------

# 流水线调用 analysis_sentiment 的超时(秒),默认 600
# SONA_SENTIMENT_TIMEOUT_SEC=600
# 工具内 LLM 批处理总墙钟(秒),默认 600(应 ≤ 流水线超时)
# SONA_SENTIMENT_MAX_WALLTIME_SEC=600
# 默认 1:强制走大模型;设为 0 且下方 ALLOW_COLUMN_FALLBACK=1 时才可回退 CSV 列
# SONA_SENTIMENT_FORCE_LLM=1
# 设为 1 才允许复用 CSV「情感」列或在 LLM 失败/质量保护时回退列统计(默认 0)
# SONA_SENTIMENT_ALLOW_COLUMN_FALLBACK=0

# -----------------------------------------------------------------------------
# 九、代理(数据采集走代理时)
# -----------------------------------------------------------------------------
Expand Down
18 changes: 17 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ share/python-wheels/
.installed.cfg
*.egg
MANIFEST
# 专题监测本地库
data/topic_monitor_local.db

.pytest_cache/
.mypy_cache/
.ruff_cache/
Expand Down Expand Up @@ -76,6 +79,14 @@ logs/
*.log
*.log.*

# Frontend / Node
node_modules/
.next/
out/
frontend/node_modules/
frontend/.next/
frontend/out/

# Eval artifacts
eval_results/

Expand All @@ -89,6 +100,10 @@ sanbox/
# Memory - 会话记忆
memory/

# Local SQLite app data
data/*.db
data/*.db-*

# 热点流程产物(独立流程)
data_langgraph/
data_langgraph_hourly/
Expand All @@ -110,10 +125,11 @@ Desktop.ini
*.bak
*.swp
*~
cache/

# 本地课程作业与知识库压缩包(不纳入版本库)
期末报告-智能体修改-作业/
opinion_analysis_kb.zip

# 舆情知识库目录(体积大、本地维护;已从 Git 索引移除,不再推送到远端)
opinion_analysis_kb/
opinion_analysis_kb/
70 changes: 66 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@
- **🧠 垂类知识库增强**:支持控烟、健康、交通、大熊猫等领域包,并通过 `workflow/domain_routing.json` 自动注入优先证据
- **🕸️ Graph RAG 可选增强**:通过 Neo4j Aura/本地 Neo4j 召回相似案例、理论框架和处置经验;连接失败时自动降级
- **📚 案例库与专题监测**:完整报告可自动沉淀到 `opinion_analysis_kb/references/wiki/cases/`,并提供 `/case` 相似案例检索与 `/monitor` 专题快照/日报周报演示
- **🧩 HTTP API 与轻量 GUI**:`sona serve` 提供 FastAPI 接口,`streamlit_app.py` 提供多页 Streamlit 查看器
- **🧩 HTTP API 与轻量 GUI**:`sona serve` 提供 FastAPI;`frontend/` 提供 Next.js 分析工作台;`streamlit_app.py` 保留多页控制台
- **📊 报告质量增强**:时间线证据/影响标签、情绪结构、四阶段行动清单、热点风险分级与案例候选输出

### 支持的模型提供商
Expand Down Expand Up @@ -366,7 +366,7 @@ sona serve --host 127.0.0.1 --port 8765
- `/hot` - 运行独立的热点抓取与态势感知流程(可选参数:配置路径)
- `/case` - 检索本地案例库,输出相似案例列表与横向对照
- `/monitor` - 运行专题监测命令,支持创建专题、查看状态、生成日报/周报;未配置外部库时可运行内存演示
- `/wiki` - 基于本地知识库(`opinion_analysis_kb/references/wiki/`)的问答,输出摘要与引用来源
- `/wiki` - 基于本地知识库(`opinion_analysis_kb/references/wiki/`)的问答,输出摘要与引用来源;可通过环境变量 `SONA_WIKI_AUTO_COMPILE`、`SONA_WIKI_AUTO_COMPILE_ON_QUERY` 控制是否在写入专家研判或每次检索前自动把 `expert_notes/**/*.md` 增量编译进 `wiki/sources`(默认开启;`tests/conftest.py` 对 CI 关闭「检索前编译」)
- `/wiki-approve` - 审核并回流高价值候选
- `/clear` - 清除 memory 和 sandbox
- `/exit` - 退出程序
Expand All @@ -382,18 +382,80 @@ sona serve --host 127.0.0.1 --port 8765
# 探活
curl http://127.0.0.1:8765/health

# 终端 2:启动 Streamlit 多页 GUI(需安装 streamlit)
# 终端 2A:启动 Next.js 前端(推荐)
npm --prefix frontend install
npm --prefix frontend run dev

# 终端 2B:或启动 Streamlit 多页 GUI(保留入口)
streamlit run streamlit_app.py
```

Next.js 前端访问地址:`http://127.0.0.1:3000`。

注意:`server.py` / `api.server:app` 不会由前端自动启动;必须先运行 `sona serve`。Next.js 只通过 BFF 代理访问后端 API。若后端端口不是 `8765`,在 `frontend/.env.local` 设置:

```env
SONA_API_BASE=http://127.0.0.1:8765
```

#### Next.js 分析工作台是什么

`frontend/` 是 React + Next.js App Router 前端,页面只访问 `/api/sona/...`,由 Next.js Route Handler 代理到 `sona serve` 的 FastAPI。它支持:

| 能力 | 对应 API |
|------|----------|
| 会话创建 / 恢复 | `POST /v1/chat/sessions`、`GET /v1/chat/sessions` |
| 普通对话流式输出 | `POST /v1/chat/sessions/{task_id}/messages:stream` |
| `/event` 事件分析 | `POST /v1/analyze-event` |
| `/wiki` / `/wiki-approve` | `POST /v1/wiki/query`、`POST /v1/wiki/approve` |
| `/case` 案例检索 | `POST /v1/cases/search` |
| `/hot` 热点态势 | `POST /v1/hot/run` |
| `/monitor` 专题监测 | `GET/POST /v1/monitor/...` |
| `/models` / `/tools` | `GET /v1/models`、`GET /v1/tools` |

`/set`、`/compress`、`/clear`、`/exit` 属于 CLI 管理或生命周期指令,前端 v1 不在聊天输入中执行。

#### 分析员控制台(Streamlit)是什么

多页 Streamlit **不是**「第二个舆情监测大屏」,而是给分析员用的 **控制台**:

| 区域 | 作用 |
|------|------|
| **仪表盘**(`streamlit_app.py`) | 探活 `sona serve`、展示本 API 进程内最近任务、工作流速查与环境变量说明 |
| **新建 / 报告 / 任务** | 走 `POST /v1/analyze-event` 与任务轮询(重操作,可能长时间同步) |
| **案例 / 专题** | 演示案例检索、编辑 `config/topics.yaml` |
| **经典会话** | 与 CLI 一致的对话界面,支持 `/hot`、`/wiki`、事件路由等 |

UI 采用与 BettaFish「微舆」类似的 **高对比、硬边框** 风格;空状态与 API 离线有统一提示块。

**仅跑控制台时的环境变量(常用)**

| 变量 | 说明 |
|------|------|
| `API_BASE` | Streamlit 访问的 API 根地址,默认 `http://127.0.0.1:8765` |
| `SONA_API_BASE` | Next.js BFF 访问的 FastAPI 根地址,默认 `http://127.0.0.1:8765` |
| `SONA_API_CORS_ORIGINS` | API 的 CORS 白名单(见 `docs/api_design.md`) |
| `.env` 内模型与采集 Key | 与 CLI 相同;缺省则事件分析或采集会失败 |

知识库目录 `opinion_analysis_kb/` 若被 `.gitignore` 排除,克隆仓库后需按团队约定自行放入或从网盘恢复。

主要 API:

- `GET /health`:服务探活
- `POST /v1/chat/sessions`:创建前端会话
- `POST /v1/chat/sessions/{task_id}/messages:stream`:流式对话与自动路由
- `POST /v1/analyze-event`:同步执行一次事件分析,返回 `task_id` 与报告路径
- `POST /v1/wiki/query`、`POST /v1/cases/search`:知识库问答与案例检索
- `POST /v1/hot/run`、`GET/POST /v1/monitor/...`:热点与专题监测
- `GET /v1/tasks`:查看当前 API 进程内存中的任务
- `GET /v1/tasks/{task_id}/report`:返回 HTML 报告

更多约定见 `docs/api_design.md` 和 `docs/gui_decision.md`。
更多约定见 `docs/api_design.md`、`docs/gui_decision.md` 和 `docs/frontend_next.md`。

**控制台 / 新建任务常见问题**

- **`WORKFLOW_ERROR` / `data_collect` / 微博 / `BrowserType.launch: Executable doesn't exist`**:未安装 Playwright 浏览器。在项目根执行 **`playwright install chromium`**(或 `playwright install`),见上文「数据采集」一节。另需配置 **NetInsight** 账号(`NETINSIGHT_USER` / `NETINSIGHT_PASS`),否则会出现「登录失败」。
- **`ModuleNotFoundError: utils.hot_time_parser`**:请拉取包含 `utils/hot_time_parser.py` 的版本;经典会话页已增加项目根 `sys.path` 兜底。

**`/hot` 热点流程说明**:
- 从公网聚合接口拉取各平台热搜(需本机可访问外网),再在本地用 **OpenAI 兼容 API** 做归纳与报告。
Expand Down
123 changes: 104 additions & 19 deletions agent/reactagent.py
Original file line number Diff line number Diff line change
Expand Up @@ -233,32 +233,105 @@ def _stream_mode_flow(
existing_data_path = opts.get("existing_data_path")
skip_data_collect = bool(opts.get("skip_data_collect", False))
force_fresh_start = opts.get("force_fresh_start")
yield {"type": "tool_call", "tool_name": "full_report_mode_node", "args": {"query": user_input}, "run_id": f"mode_full_{task_id or 'na'}"}
report_length = str(opts.get("report_length") or "").strip() or None
file_url_or_path = run_full_report_mode(
user_query=user_input,
task_id=task_id or "",
session_manager=session_manager,
debug=True,
existing_data_path=existing_data_path,
skip_data_collect=skip_data_collect,
force_fresh_start=force_fresh_start,
report_length=report_length,
)
progress_queue: queue.Queue = queue.Queue()
result_holder: Dict[str, Any] = {"value": None, "error": None}

def on_workflow_progress(event: Dict[str, Any]) -> Any:
hook = opts.get("_web_progress_hook")
hook_result = None
if callable(hook):
hook_result = hook(event)
progress_queue.put(("step", dict(event)))
return hook_result

def run_pipeline() -> None:
try:
report_length = str(opts.get("report_length") or "").strip() or None
result_holder["value"] = run_full_report_mode(
user_query=user_input,
task_id=task_id or "",
session_manager=session_manager,
debug=True,
existing_data_path=existing_data_path,
skip_data_collect=skip_data_collect,
force_fresh_start=force_fresh_start,
report_length=report_length,
progress_callback=on_workflow_progress,
skip_session_user_message=bool(opts.get("_skip_session_user_message")),
)
except Exception as exc: # noqa: BLE001
result_holder["error"] = exc
finally:
progress_queue.put(("done", None))

yield {
"type": "tool_call",
"tool_name": "full_report_mode_node",
"args": {"query": user_input},
"run_id": f"mode_full_{task_id or 'na'}",
}
worker = threading.Thread(target=run_pipeline, daemon=True)
worker.start()
while True:
try:
kind, payload = progress_queue.get(timeout=0.3)
except queue.Empty:
if not worker.is_alive():
break
continue
if kind == "step":
yield {"type": "workflow_step", **payload}
elif kind == "done":
break
worker.join()
if result_holder["error"] is not None:
raise result_holder["error"]
file_url_or_path = result_holder["value"]
result_text = str(file_url_or_path or "").strip()
if result_text.startswith("已完成舆情事件分析工作流。报告:") or result_text.startswith("这次在 "):
final_text = result_text
else:
final_text = f"已完成舆情事件分析工作流。报告:{result_text}"
yield {
"type": "tool_result",
"tool_name": "full_report_mode_node",
"result": str(file_url_or_path or ""),
"result": result_text,
"run_id": f"mode_full_{task_id or 'na'}",
}
final_text = (
"完整舆情报告流程已完成。\n"
f"- 报告地址:{str(file_url_or_path or '未返回')}\n"
"- 已复用事件分析工作流节点(采集/分析/报告生成)。"
)
yield {"type": "message", "message": AIMessage(content=final_text), "message_id": f"mode_full_msg_{task_id or 'na'}"}
yield {
"type": "message",
"message": AIMessage(content=final_text),
"message_id": f"mode_full_msg_{task_id or 'na'}",
"persist": False,
}
return


def _extract_reasoning_content(chunk: Any) -> str:
"""Best-effort extraction for OpenAI-compatible reasoning/thinking chunks."""
candidates = []
for attr in ("additional_kwargs", "response_metadata"):
value = getattr(chunk, attr, None)
if isinstance(value, dict):
candidates.append(value)
if isinstance(chunk, dict):
candidates.append(chunk)

keys = ("reasoning_content", "reasoning", "thinking_content", "thinking")
for data in candidates:
for key in keys:
value = data.get(key)
if isinstance(value, str) and value:
return value
delta = data.get("delta")
if isinstance(delta, dict):
for key in keys:
value = delta.get(key)
if isinstance(value, str) and value:
return value
return ""


# 创建带消息历史的 Agent
# 注意:此函数目前未使用,保留作为预留功能,用于未来可能需要直接使用带历史管理的 Agent 的场景
def _create_agent_with_history():
Expand Down Expand Up @@ -433,6 +506,7 @@ async def _stream_events():

current_message_id = None
current_content = ""
current_reasoning = ""
# 追踪当前正在执行的工具 run_id,用于过滤工具内部 LLM 的流式输出
active_tool_run_ids: set[str] = set()

Expand All @@ -457,6 +531,16 @@ async def _stream_events():
if active_tool_run_ids:
continue
chunk = data.get("chunk")
if chunk:
reasoning_delta = _extract_reasoning_content(chunk)
if reasoning_delta:
current_reasoning += reasoning_delta
result_queue.put({
"type": "thinking",
"content": reasoning_delta,
"message_id": event.get("run_id", ""),
"accumulated": current_reasoning
})
if chunk and hasattr(chunk, "content") and chunk.content:
if current_message_id is None:
current_message_id = event.get("run_id", "")
Expand All @@ -483,6 +567,7 @@ async def _stream_events():
})
current_message_id = None
current_content = ""
current_reasoning = ""

# 工具调用
elif event_type == "on_tool_start":
Expand Down
Loading