refactor(plugins): decouple DingTalk notifications (#550)

This commit is contained in:
jinliyl 2026-09-15 18:29:06 +08:00 • committed by GitHub
parent b9caae1e50
commit 6f4bdfd416
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 8 additions and 243 deletions

View file

@ -87,7 +87,7 @@ records outside the window are discarded.
`auto_fin_topic_step` receives batches of current news and returns only related `news_id` values. Code ignores unknown
IDs and deduplicates repeated IDs, then preserves the source-news order. If nothing is relevant, the job succeeds as a
skip without writing or sending a report.
skip without writing a report.
`auto_fin_merge_step` receives only selected current news. It exposes `search` and `read`, and keeps current CLS IDs,
times, and titles as plain evidence. The prompt limits

View file

@ -80,7 +80,7 @@ daily/YYYY-MM-DD/auto_fin.md
小时。请求带有限速和重试;损坏记录及窗口外记录会被丢弃。
`auto_fin_topic_step` 分批接收当前新闻,只返回相关的 `news_id`。代码会忽略未知 ID、去除重复 ID,并保持源新闻顺序。如果没有相关新闻,Job
会成功跳过,不写报告也不发送通知。
会成功跳过,不写报告。
`auto_fin_merge_step` 只接收筛选后的当前新闻,并向 Agent 开放 `search` 和 `read`。当前新闻以 CLS ID、时间和标题作为普通证据。
Prompt 要求 Agent 只链接实际使用过的历史 Markdown;代码边界则独立保证只保留真实存在、相对

View file

@ -69,8 +69,6 @@ download and parse arXiv PDFs, then write three Chinese analyses
use search + read to connect prior memory and generate a brief
↓
generate memory tags; the background file watcher refreshes indexes
↓
optionally send the brief to DingTalk
```
`daily_paper_collect_step` concurrently reads the weekly and monthly rankings for the run date plus the strictly
@ -89,8 +87,7 @@ PDFs and files without a text layer fail explicitly.
`search` and `read` tools for linking earlier memory. Code validates historical wikilinks, appends links to all
three source notes, and rebuilds the daily index. The workflow then runs `auto_tag_step` to update the memory-tag
frontmatter of all three analyses and the final brief. The normal background file watcher observes those source-file
changes and refreshes derived indexes before the optional `dingtalk_markdown_send_step` sends the brief. DingTalk
delivery skips without side effects when conversation IDs are not configured.
changes and refreshes derived indexes.
## Parameters
@ -108,15 +105,11 @@ Step-level defaults are `candidate_limit=20`, `rrf_k=60`, `hf_timeout=600`, `hf_
The data clients automatically honor `HTTP_PROXY`, `HTTPS_PROXY`, and `NO_PROXY`. Manual runs enable the Hugging Face
mirror with `use_hf_mirror=true`; the cron Job enables it by default and can use the official service with
`DAILY_PAPER_USE_HF_MIRROR=false`. These environment variables override data sources and DingTalk settings:
`DAILY_PAPER_USE_HF_MIRROR=false`. These environment variables override the data sources:
```dotenv
HF_MIRROR_URL=https://hf-mirror.com
ARXIV_MIRROR_URL=https://export.arxiv.org
DINGTALK_APP_KEY=your-app-key
DINGTALK_APP_SECRET=your-app-secret
DINGTALK_ROBOT_CODE=your-robot-code
DINGTALK_CONVERSATION_IDS=cid-group-one,cid-group-two
```
## Output
@ -142,4 +135,4 @@ Network errors, too few candidates, invalid Agent output, and unparseable PDFs f
python -m pytest plugins/daily_paper -v
```
Unit tests mock the Hugging Face, arXiv, AgentScope, and DingTalk boundaries and do not contact external services.
Unit tests mock the Hugging Face, arXiv, and AgentScope boundaries and do not contact external services.

View file

@ -65,8 +65,6 @@ RRF 排序后由 Agent 精选三篇
使用 search + read 关联历史记忆并生成简报
↓
生成记忆标签,由后台文件 watcher 刷新索引
↓
按需发送到钉钉
```
`daily_paper_collect_step` 并发读取运行日期所在周和所在月的榜单,以及严格前一日的 Daily Papers。候选按 arXiv ID
@ -80,8 +78,7 @@ RRF 排序后由 Agent 精选三篇
`daily_paper_digest_step` 以本次生成的三篇解读为事实来源,只开放只读的 `search` 和 `read` 来关联较早记忆。
代码会校验历史 wikilink、追加三篇源笔记链接,并重建当日索引。随后 `auto_tag_step` 会更新三篇解读及最终简报的
记忆标签 frontmatter,常规后台文件 watcher 会观察这些源文件变化并刷新派生索引,再由可选的
`dingtalk_markdown_send_step` 发送最终简报;未配置群会话时无副作用跳过。
记忆标签 frontmatter,常规后台文件 watcher 会观察这些源文件变化并刷新派生索引。
## 参数
@ -98,15 +95,11 @@ RRF 排序后由 Agent 精选三篇
`pdf_timeout=600`、`max_pdf_bytes=52428800`、`max_pdf_pages=35` 和 `max_pdf_chars=300000`。
数据客户端自动使用 `HTTP_PROXY`、`HTTPS_PROXY` 和 `NO_PROXY`。手动任务通过 `use_hf_mirror=true` 启用 Hugging Face
镜像;定时任务默认启用,可设置 `DAILY_PAPER_USE_HF_MIRROR=false` 改用官方服务。以下环境变量可覆盖数据源和钉钉配置:
镜像;定时任务默认启用,可设置 `DAILY_PAPER_USE_HF_MIRROR=false` 改用官方服务。以下环境变量可覆盖数据源配置:
```dotenv
HF_MIRROR_URL=https://hf-mirror.com
ARXIV_MIRROR_URL=https://export.arxiv.org
DINGTALK_APP_KEY=your-app-key
DINGTALK_APP_SECRET=your-app-secret
DINGTALK_ROBOT_CODE=your-robot-code
DINGTALK_CONVERSATION_IDS=cid-group-one,cid-group-two
```
## 产物
@ -131,4 +124,4 @@ Markdown 和 PDF 都通过同目录临时文件原子写入。`force=true` 会
python -m pytest plugins/daily_paper -v
```
单元测试 mock Hugging Face、arXiv、AgentScope 和钉钉边界,不访问外部服务。
单元测试 mock Hugging Face、arXiv 和 AgentScope 边界,不访问外部服务。

View file

@ -55,15 +55,6 @@ application_defaults:
- backend: daily_paper_digest_step
job_tools: [search, read]
- backend: auto_tag_step
- backend: dingtalk_markdown_send_step
input_mapping:
daily_paper_digest_path: markdown_path
app_key: ${DINGTALK_APP_KEY:-}
app_secret: ${DINGTALK_APP_SECRET:-}
robot_code: ${DINGTALK_ROBOT_CODE:-}
conversation_ids: ${DINGTALK_CONVERSATION_IDS:-}
title: ReMe Daily Paper
timeout: 15
daily_paper_cron:
backend: cron

View file

@ -1,11 +1,7 @@
"""Focused tests for the Daily Paper plugin."""
import datetime as dt
import importlib
import json
from pathlib import Path
import subprocess
import sys
from unittest.mock import AsyncMock, MagicMock
import frontmatter
@ -41,8 +37,6 @@ from reme.components import ApplicationContext
from reme.components.agent_wrapper.base_agent_wrapper import BaseAgentWrapper
from reme.components.runtime_context import RuntimeContext
from reme.config import expand_env_vars
from reme.steps.cookbook.dingtalk import DingTalkMarkdownSendStep
from reme.steps.cookbook.dingtalk import send as dingtalk_send
PLUGIN_MANIFEST = yaml.safe_load(
(Path(__file__).parents[1] / "src" / "reme_daily_paper" / "plugin.yaml").read_text(encoding="utf-8"),
@ -72,7 +66,6 @@ def test_plugin_manifest_declares_complete_runtime_surface():
"daily_paper_analyze_step",
"daily_paper_digest_step",
"auto_tag_step",
"dingtalk_markdown_send_step",
]
assert jobs["daily_paper"]["steps"][5] == {"backend": "auto_tag_step"}
assert jobs["daily_paper_cron"]["steps"] == jobs["daily_paper"]["steps"]
@ -576,27 +569,6 @@ async def test_daily_paper_steps_construct_source_clients_without_proxy(
assert "proxy_url" not in arxiv_kwargs[0]
def test_daily_paper_config_passes_dingtalk_environment(monkeypatch):
"""The notifier receives all proactive-message settings from the environment."""
values = {
"DINGTALK_APP_KEY": "app-key",
"DINGTALK_APP_SECRET": "app-secret",
"DINGTALK_ROBOT_CODE": "robot-code",
"DINGTALK_CONVERSATION_IDS": "group-one,group-two",
}
for name, value in values.items():
monkeypatch.setenv(name, value)
step = _plugin_config()["jobs"]["daily_paper"]["steps"][-1]
assert {key: step[key] for key in ("app_key", "app_secret", "robot_code", "conversation_ids")} == {
"app_key": "app-key",
"app_secret": "app-secret",
"robot_code": "robot-code",
"conversation_ids": "group-one,group-two",
}
def test_daily_paper_uses_agentscope_without_tools():
"""Daily Paper uses the shared tool-free agent."""
config = _plugin_config()
@ -730,33 +702,6 @@ async def test_selection_retries_invalid_id_and_keeps_three_candidates(tmp_path:
assert "仅将这些 topics 作为主题偏好" in agent.calls[1]["inputs"]
def test_reme_import_does_not_require_optional_dingtalk_stream():
"""Importing ReMe must not eagerly load the core-only DingTalk dependency."""
script = """
import builtins
original_import = builtins.__import__
def guarded_import(name, *args, **kwargs):
if name == "dingtalk_stream":
raise ModuleNotFoundError("blocked optional dependency")
return original_import(name, *args, **kwargs)
builtins.__import__ = guarded_import
import reme
"""
result = subprocess.run(
[sys.executable, "-c", script],
cwd=Path(__file__).parents[3],
capture_output=True,
text=True,
check=False,
)
assert result.returncode == 0, result.stderr
@pytest.mark.asyncio
async def test_pipeline_filters_strict_yesterday_and_writes_outputs(
tmp_path: Path,
@ -1071,160 +1016,3 @@ async def test_digest_validates_model_generated_historical_wikilinks(tmp_path: P
assert "非 daily 节点" in rendered and "[[digest/wiki/相关概念.md" not in rendered
assert "越界路径" in rendered and "../outside.md" not in rendered
assert "[[daily/2026-07-21/论文解读1.md]]" in rendered
@pytest.mark.asyncio
async def test_dingtalk_markdown_sends_groups_serially_in_configured_order(
tmp_path: Path,
monkeypatch,
):
"""The notifier gets one app token and posts once per group in list order."""
digest_path = tmp_path / "daily" / "2026-07-21" / "daily-paper-brief.md"
digest_path.parent.mkdir(parents=True)
digest_path.write_text(
frontmatter.dumps(
frontmatter.Post("# 今日论文\n\n测试内容", name="daily-paper-brief"),
),
encoding="utf-8",
)
token_calls = 0
seen_payloads: list[dict] = []
def get_access_token(client):
nonlocal token_calls
token_calls += 1
assert client.credential.client_id == "app-key"
assert client.credential.client_secret == "app-secret"
return "app-access-token"
async def handler(request: httpx.Request) -> httpx.Response:
assert request.url.path == "/v1.0/robot/groupMessages/send"
assert request.headers["x-acs-dingtalk-access-token"] == "app-access-token"
seen_payloads.append(json.loads(request.content))
return httpx.Response(
200,
json={"processQueryKey": f"query-{len(seen_payloads)}"},
)
transport = httpx.MockTransport(handler)
transport_kwargs: dict = {}
def ipv4_transport(**kwargs):
transport_kwargs.update(kwargs)
return transport
dingtalk_stream = importlib.import_module("dingtalk_stream")
monkeypatch.setattr(
dingtalk_stream.DingTalkStreamClient,
"get_access_token",
get_access_token,
)
monkeypatch.setattr(dingtalk_send.httpx, "AsyncHTTPTransport", ipv4_transport)
app_context = ApplicationContext(workspace_dir=str(tmp_path))
context = RuntimeContext(markdown_path="daily/2026-07-21/daily-paper-brief.md")
step = DingTalkMarkdownSendStep(
app_context=app_context,
app_key="app-key",
app_secret="app-secret",
robot_code="robot-code",
conversation_ids=" group-one,group-two ",
title="ReMe Daily Paper",
)
step.logger = MagicMock()
response = await step(context)
assert token_calls == 1
assert transport_kwargs == {"local_address": "0.0.0.0"}
assert [payload["openConversationId"] for payload in seen_payloads] == [
"group-one",
"group-two",
]
assert all(payload["robotCode"] == "robot-code" for payload in seen_payloads)
assert all(payload["msgKey"] == "sampleMarkdown" for payload in seen_payloads)
assert [json.loads(payload["msgParam"]) for payload in seen_payloads] == [
{"title": "ReMe Daily Paper", "text": "# 今日论文\n\n测试内容"},
] * 2
assert response.metadata["dingtalk_configured_count"] == 2
assert response.metadata["dingtalk_sent_count"] == 2
logs = "\n".join(call.args[0] for call in step.logger.info.call_args_list)
assert "sending DingTalk Markdown" in logs
assert "delivery complete sent=2 total=2" in logs
assert all(value not in logs for value in ("app-key", "app-secret", "robot-code", "group-one", "group-two"))
@pytest.mark.asyncio
async def test_dingtalk_markdown_without_conversations_is_a_noop(tmp_path: Path):
"""An empty conversation list keeps daily-paper generation usable without DingTalk."""
context = RuntimeContext(markdown_path="missing.md")
response = await DingTalkMarkdownSendStep(
app_context=ApplicationContext(workspace_dir=str(tmp_path)),
)(context)
assert response.success is True
assert response.metadata["dingtalk_configured_count"] == 0
assert response.metadata["dingtalk_sent_count"] == 0
@pytest.mark.asyncio
async def test_existing_daily_paper_is_reused_and_sent_to_dingtalk(
tmp_path: Path,
monkeypatch,
):
"""An idempotent daily-paper run skips generation but still notifies DingTalk."""
digest_path = tmp_path / "daily" / "2026-07-22" / "daily-paper-brief.md"
digest_path.parent.mkdir(parents=True)
digest_path.write_text(
frontmatter.dumps(
frontmatter.Post("# 已有日报\n\n复用正文", name="daily-paper-brief"),
),
encoding="utf-8",
)
seen_payloads: list[dict] = []
dingtalk_stream = importlib.import_module("dingtalk_stream")
monkeypatch.setattr(
dingtalk_stream.DingTalkStreamClient,
"get_access_token",
lambda _client: "app-access-token",
)
async def handler(request: httpx.Request) -> httpx.Response:
seen_payloads.append(json.loads(request.content))
return httpx.Response(200, json={"processQueryKey": "query-1"})
transport = httpx.MockTransport(handler)
monkeypatch.setattr(
dingtalk_send.httpx,
"AsyncHTTPTransport",
lambda **_kwargs: transport,
)
app_context = ApplicationContext(workspace_dir=str(tmp_path))
context = RuntimeContext(date="2026-07-22")
await DailyPaperCollectStep(app_context=app_context)(context)
response = await DingTalkMarkdownSendStep(
app_context=app_context,
input_mapping={"daily_paper_digest_path": "markdown_path"},
app_key="app-key",
app_secret="app-secret",
robot_code="robot-code",
conversation_ids="existing-group",
title="ReMe Daily Paper",
)(context)
assert response.metadata["skipped"] is True
assert response.metadata["dingtalk_sent_count"] == 1
assert seen_payloads == [
{
"robotCode": "robot-code",
"openConversationId": "existing-group",
"msgKey": "sampleMarkdown",
"msgParam": json.dumps(
{"title": "ReMe Daily Paper", "text": "# 已有日报\n\n复用正文"},
ensure_ascii=False,
),
},
]