diff --git a/plugins/auto-fin/README.md b/plugins/auto-fin/README.md index 690be651..71a48169 100644 --- a/plugins/auto-fin/README.md +++ b/plugins/auto-fin/README.md @@ -2,12 +2,12 @@ [中文](README_ZH.md) -Auto Fin fetches a rolling window of CLS telegraph news (24 hours by default), selects items related to configured -topics, searches ReMe for useful historical context, and writes one Chinese Markdown report with validated wikilinks. -Current news and topic selection stay in runtime memory; only the final report becomes durable memory. This directory -is an independent Python distribution. Its single `reme.plugins` entry point exposes a `plugin.yaml` containing the -three Step backends and their Job configuration under `application_defaults`. Enable the installed plugin explicitly -through `plugins=["auto-fin"]`. +Auto Fin fetches a rolling window of CLS telegraph news (24 hours by default), groups items by configured +topics, researches each topic against ReMe history into its own note, and merges those notes into one dated brief +whose validated wikilinks point back at every note. Current news and topic selection stay in runtime memory; only the +notes and the brief become durable memory. This directory is an independent Python distribution. Its single +`reme.plugins` entry point exposes a `plugin.yaml` containing the four Step backends and their Job configuration under +`application_defaults`. Enable the installed plugin explicitly through `plugins=["auto-fin"]`. > Auto Fin has no reliable market-price feed. It does not calculate returns, targets, or entry points and is not > investment advice. @@ -60,7 +60,7 @@ reme start plugins='["auto-fin"]' \ ``` Custom application configs must provide `agent_wrapper.default`, a `file_store.default` with an enabled tag index, and -the `search`, `read`, `list_tags`, `frontmatter_read`, and `frontmatter_update` Jobs used by Auto Fin and automatic +the `search`, `list_tags`, `frontmatter_read`, and `frontmatter_update` Jobs used by Auto Fin and automatic tagging. ## Pipeline @@ -70,13 +70,12 @@ CLS public telegraph endpoint (rolling 24 hours) ↓ normalize and deduplicate in RuntimeContext ↓ -topic Agent selects real news IDs in bounded batches +topic Agent returns topic-to-news-ID mappings in prompt-sized batches ↓ -research Agent uses search + read on historical memory +one research Agent per topic examines its latest 20 articles, searches history up to three times, +and writes one note per topic with validated historical wikilinks ↓ -validate historical wikilinks in code - ↓ -daily/YYYY-MM-DD/auto_fin.md +digest Agent merges the notes into today's brief; code appends a link back to every note ↓ generate memory tags; the background file watcher refreshes indexes ``` @@ -85,20 +84,24 @@ generate memory tags; the background file watcher refreshes indexes stops only after covering the exact preceding 24 hours. Requests are rate-limited and retried; malformed records and 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 +`auto_fin_topic_step` batches current news under a 100,000-character full-prompt limit and returns related `news_id` +values for each topic. Code ignores unknown IDs and deduplicates repeated IDs, then preserves source-news order. One article may belong to multiple topics. If nothing is relevant, the job succeeds as a 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 -wikilinks to historical Markdown actually used by the Agent; the code-level boundary independently keeps only existing, -workspace-relative Markdown targets. Missing, absolute, escaping, backslash, and self-referential targets are degraded -to their readable aliases. +`auto_fin_research_step` researches each nonempty topic with its latest 20 articles and writes one note per topic. It +exposes only `search`, enforces a three-call search budget per topic, and keeps current CLS IDs, times, and titles as +plain evidence. The prompt limits wikilinks to historical Markdown actually used by the Agent; the code-level boundary +independently keeps only existing, workspace-relative Markdown targets. Missing, absolute, escaping, backslash, and +self-referential targets are degraded to their readable aliases. -Same-day reruns use the existing report as context and replace it with the revised result. The final write is atomic and -refreshes the daily index. The workflow then runs `auto_tag_step` to update the generated report's memory-tag -frontmatter; the normal background file watcher observes that source-file change and refreshes derived indexes. No -JSONL, intermediate Markdown, or structured Agent output is written. +`auto_fin_digest_step` merges those notes into today's brief with no tools of its own: the notes already carry the +research. Code appends a `## 主题详解` section linking back to every note, refreshes the daily index, and publishes the +brief path for downstream steps. Both steps stop early when the topic step found no relevant news. + +Same-day reruns reuse the note a topic already produced, feeding its body back to the research Agent and replacing it +in place; the brief is replaced the same way. Writes are atomic. The workflow then runs `auto_tag_step` to update the +generated files' memory-tag frontmatter; the normal background file watcher observes those source-file changes and +refreshes derived indexes. No JSONL, intermediate Markdown, or structured Agent output is written. ## Parameters @@ -111,17 +114,24 @@ JSONL, intermediate Markdown, or structured Agent output is written. | `request_interval` | `10` | Minimum delay in seconds after every CLS request attempt; may be zero | | `max_retries` | `3` | Maximum attempts for each CLS page request; must be at least one | -The plugin cron Job starts with the application and runs daily at 18:00 in the application timezone. +The plugin cron Job starts with the application and runs daily at 09:00 in the application timezone, which defaults to `Asia/Shanghai`. The rolling window uses timestamps and may cross calendar days. Report completion depends on news volume and model latency. ## Output ```text -.reme/daily/YYYY-MM-DD/auto_fin.md +.reme/daily/YYYY-MM-DD/.md # one per topic that had relevant news +.reme/daily/YYYY-MM-DD/主题新闻观察(YYYY-MM-DD).md # today's merged brief, linking back to every note ``` -The report includes a title, description, current CLS evidence, historical analysis, contextual wikilinks, and a fixed -non-investment disclaimer. Network errors and invalid Agent output fail explicitly; no relevant current news is a -successful skip. +File names come from the configured topics and the run date rather than the Agent's title: an Agent title is free text +and can outgrow the filesystem limit for one name component. Names are sanitized, byte-truncated, and disambiguated; +the Agent title is kept in the `title` frontmatter field. Every file carries `kind` frontmatter (`auto-fin-topic` or +`auto-fin-digest`), so a rerun finds and replaces the notes +it produced instead of duplicating them. Each file includes a title, description, current CLS evidence, historical +analysis, contextual wikilinks, and a fixed non-investment disclaimer; the brief ends with a `## 主题详解` list linking +to the topic notes. A fetch failure fails the run explicitly, and so does a run in which every topic fails; a single +failing topic only logs a warning and the remaining topics continue. No relevant current news is a successful +skip. ## Validation diff --git a/plugins/auto-fin/README_ZH.md b/plugins/auto-fin/README_ZH.md index ba0aa70d..02ce5f78 100644 --- a/plugins/auto-fin/README_ZH.md +++ b/plugins/auto-fin/README_ZH.md @@ -2,9 +2,9 @@ [English](README.md) -Auto Fin 自动拉取一个滚动时间窗口内的财联社电报(默认 24 小时),按配置 topics 筛选相关新闻,搜索 ReMe 中有回顾价值的历史材料,最后写入一份带校验 -wikilink 的中文 Markdown 报告。当前新闻和筛选结果只存在于本次运行内存中,只有最终报告成为持久记忆。本目录是一个独立 Python -distribution:单个 `reme.plugins` entry point 暴露 `plugin.yaml`,其中声明三个 Step backend,并在 +Auto Fin 自动拉取一个滚动时间窗口内的财联社电报(默认 24 小时),按配置 topics 归类相关新闻,逐主题搜索 ReMe 中有回顾价值的历史材料并为每个主题写一份笔记, +最后把所有主题笔记合并成一份带校验 wikilink、且回链到各主题笔记的当日总览。当前新闻和筛选结果只存在于本次运行内存中,只有主题笔记和总览成为持久记忆。本目录是一个独立 Python +distribution:单个 `reme.plugins` entry point 暴露 `plugin.yaml`,其中声明四个 Step backend,并在 `application_defaults` 下提供 Job 配置;通过 `plugins=["auto-fin"]` 显式启用这个已安装插件。 > Auto Fin 没有可靠行情数据,不计算收益、目标价或买卖点,也不提供投资建议。 @@ -55,7 +55,7 @@ reme start plugins='["auto-fin"]' \ ``` 自定义应用配置需要提供 `agent_wrapper.default`、启用 tag index 的 `file_store.default`,以及 -Auto Fin 和自动标签使用的 `search`、`read`、`list_tags`、`frontmatter_read` 和 +Auto Fin 和自动标签使用的 `search`、`list_tags`、`frontmatter_read` 和 `frontmatter_update` Jobs。 ## 流程 @@ -65,13 +65,12 @@ Auto Fin 和自动标签使用的 `search`、`read`、`list_tags`、`frontmatter ↓ 在 RuntimeContext 中规范化和去重 ↓ -Topic Agent 分批选择真实 news_id +Topic Agent 按提示词长度分批输出“主题 → news_id” ↓ -Research Agent 使用 search + read 检索历史记忆 +每个有新闻的主题由独立 Research Agent 研究最新 20 篇,最多搜索 3 次, +并写出一份带校验历史 wikilink 的主题笔记 ↓ -代码校验历史 wikilink - ↓ -daily/YYYY-MM-DD/auto_fin.md +Digest Agent 把各主题笔记合并成当日总览,代码在末尾追加回链到每一份笔记 ↓ 生成记忆标签,由后台文件 watcher 刷新索引 ``` @@ -79,15 +78,18 @@ daily/YYYY-MM-DD/auto_fin.md `auto_fin_data_step` 使用财联社网页同源接口的签名和分页方式,从分析时刻开始向前翻页,直到完整覆盖严格的最近 24 小时。请求带有限速和重试;损坏记录及窗口外记录会被丢弃。 -`auto_fin_topic_step` 分批接收当前新闻,只返回相关的 `news_id`。代码会忽略未知 ID、去除重复 ID,并保持源新闻顺序。如果没有相关新闻,Job -会成功跳过,不写报告。 +`auto_fin_topic_step` 按完整提示词的 10 万字符上限分批接收当前新闻,返回每个主题的相关 `news_id`。代码会忽略未知 ID、去除重复 ID,并保持源新闻顺序;一条新闻可属于多个主题。如果没有相关新闻,Job +会成功跳过,不写任何文件。 -`auto_fin_merge_step` 只接收筛选后的当前新闻,并向 Agent 开放 `search` 和 `read`。当前新闻以 CLS ID、时间和标题作为普通证据。 +`auto_fin_research_step` 为每个有新闻的主题研究最新 20 篇,并各写一份主题笔记;只向 Agent 开放 `search`,按主题限制最多 3 次搜索。当前新闻以 CLS ID、时间和标题作为普通证据。 Prompt 要求 Agent 只链接实际使用过的历史 Markdown;代码边界则独立保证只保留真实存在、相对 workspace 的 Markdown 目标。不存在、绝对路径、越界、带反斜杠和自引用的目标都会降级为可读 alias。 -同日重跑会参考当天已有报告并覆盖为修订结果。最终写入使用原子替换并刷新当天索引,随后通过 `auto_tag_step` -更新报告的记忆标签 frontmatter;常规后台文件 watcher 会观察该源文件变化并刷新派生索引。流程不会写入 JSONL、 +`auto_fin_digest_step` 把各主题笔记合并成当日总览,自身不开放任何工具——研究结论已经在笔记里。代码在正文末尾追加 `## 主题详解` +章节,回链到每一份主题笔记,然后刷新当天索引并把总览路径交给下游步骤。两步在工作流已被判定跳过时都会提前返回。 + +同日重跑会按 frontmatter 中的 `topic` 找回该主题已有的笔记,把它的正文作为上下文重新研究并原地覆盖;总览同理。所有写入都是原子的,随后通过 `auto_tag_step` +更新这些文件的记忆标签 frontmatter;常规后台文件 watcher 会观察源文件变化并刷新派生索引。流程不会写入 JSONL、 中间 Markdown 或 Agent 结构化输出。 ## 参数 @@ -101,16 +103,20 @@ workspace 的 Markdown 目标。不存在、绝对路径、越界、带反斜杠 | `request_interval` | `10` | 每次财联社请求尝试后的最小等待秒数,可设为 0 | | `max_retries` | `3` | 每页财联社请求的最大尝试次数,至少为 1 | -插件的 cron Job 随应用启动,并按应用配置的时区在每天 18:00 运行。 +插件的 cron Job 随应用启动,并按应用配置的时区在每天 09:00 运行,默认时区为 `Asia/Shanghai`。滚动窗口按时间戳计算,允许跨自然日;09:00 是启动时间,报告完成时间取决于新闻量与模型耗时。 ## 产物 ```text -.reme/daily/YYYY-MM-DD/auto_fin.md +.reme/daily/YYYY-MM-DD/<主题>.md # 每个有相关新闻的主题一份 +.reme/daily/YYYY-MM-DD/主题新闻观察(YYYY-MM-DD).md # 合并总览,回链到每一份主题笔记 ``` -报告包含标题、说明、当前 CLS 证据、历史分析、上下文 wikilink 和固定非投资建议声明。网络错误与无效 Agent 输出 -会明确失败;没有相关当前新闻则成功跳过。 +文件名取自配置的主题(`topics`)和运行日期,而不是 Agent 生成的标题——Agent 标题是自由文本,可能长到超出文件名长度上限。文件名经过 +非法字符净化、字节数截断和重名消歧;Agent 标题保留在 frontmatter 的 `title` 中。每份文件都带 `kind` frontmatter(`auto-fin-topic` +或 `auto-fin-digest`),因此同日重跑会找回并覆盖自己产出的笔记,而不是重复生成。文件包含标题、说明、当前 CLS 证据、历史分析、上下文 +wikilink 和固定非投资建议声明;总览末尾另有一个 `## 主题详解` 列表,链接到各主题笔记。抓取失败与全部主题研究失败 +会明确失败;单个主题研究失败只记 warning 并继续其余主题;没有相关当前新闻则成功跳过。 ## 验证 diff --git a/plugins/auto-fin/src/reme_auto_fin/__init__.py b/plugins/auto-fin/src/reme_auto_fin/__init__.py index 2c820d0d..63689d0c 100644 --- a/plugins/auto-fin/src/reme_auto_fin/__init__.py +++ b/plugins/auto-fin/src/reme_auto_fin/__init__.py @@ -1,14 +1,16 @@ """Auto Fin news research workflow.""" from .data import AutoFinDataStep -from .merge import AutoFinMergeStep -from .schema import AutoFinReportOutput, AutoFinTopicOutput +from .digest import AutoFinDigestStep +from .research import AutoFinResearchStep +from .schema import AutoFinNote, AutoFinReportOutput from .topic import AutoFinTopicStep __all__ = [ - "AutoFinReportOutput", "AutoFinDataStep", - "AutoFinMergeStep", - "AutoFinTopicOutput", + "AutoFinDigestStep", + "AutoFinNote", + "AutoFinReportOutput", + "AutoFinResearchStep", "AutoFinTopicStep", ] diff --git a/plugins/auto-fin/src/reme_auto_fin/base.py b/plugins/auto-fin/src/reme_auto_fin/base.py index 37388b00..6eac42fe 100644 --- a/plugins/auto-fin/src/reme_auto_fin/base.py +++ b/plugins/auto-fin/src/reme_auto_fin/base.py @@ -1,24 +1,62 @@ -"""Shared helpers for the Auto Fin workflow.""" +"""Shared Markdown, context, and Agent helpers for the Auto Fin workflow.""" from __future__ import annotations +from datetime import datetime, timezone from html.parser import HTMLParser import json import os from pathlib import Path +import re from time import perf_counter -from typing import Any +from typing import Any, NamedTuple from uuid import uuid4 +import aiofiles +import frontmatter +import yaml from pydantic import BaseModel from reme.steps import BaseStep +from reme.steps.file_io import get_path_lock, validate_filename_component + +from .schema import AutoFinReportOutput AGENT_INPUT_LOG_LIMIT = 2000 AGENT_OUTPUT_LOG_LIMIT = 4000 +NOTE_CHAR_LIMIT = 30_000 +TITLE_BYTE_LIMIT = 180 + +_WIKILINK = re.compile(r"\[\[([^\[\]\n]+)\]\]") +_HYBRID_WIKILINK = re.compile(r"(?P\[\[(?P[^\[\]\n]+)\]\])\((?P[^()\n]+)\)") +_HEADING = re.compile(r"^#+\s*") +# A note's stem has to survive two consumers: the filesystem, which rejects +# `<>:"/\|?*` and control characters, and `WikilinkHandler`, whose targets stop +# at `[`, `]` and `#`. Leaving the latter in place turns `- [[]]` into a +# link that resolves to some other note (a trailing `#` reads as an anchor) +# or to no link at all. +_UNSAFE_FILENAME = re.compile(r'[<>:"/\\|?*\[\]#\x00-\x1f]') + +FIRST_RUN_NOTICE = "今日暂无更早时段的推荐,本次为当日首次生成。" +DISCLAIMER = "> 未接入可靠行情数据;本文只提供新闻研究和回顾线索,不提供收益、目标价或买卖建议。" + + +class WrittenReport(NamedTuple): + """One persisted note, described by what actually reached the file.""" + + path: str + """Workspace-relative POSIX path of the note.""" + + sources: list[str] + """Workspace-relative paths of the links that survived validation.""" + + body: str + """The validated body as written, with dangling wikilinks downgraded to plain text.""" class _TextExtractor(HTMLParser): + """Flatten HTML into whitespace-separated text.""" + def __init__(self) -> None: super().__init__(convert_charrefs=True) self.parts: list[str] = [] @@ -41,25 +79,161 @@ class _TextExtractor(HTMLParser): self.parts.append(data) -def _plain_text(value: str) -> str: +def plain_text(value: str) -> str: + """Return one HTML fragment as collapsed plain text.""" parser = _TextExtractor() parser.feed(value) parser.close() return " ".join("".join(parser.parts).split()) -def _write(path: Path, text: str) -> None: - path.parent.mkdir(parents=True, exist_ok=True) - temporary = path.with_name(f".{path.name}.{uuid4().hex}.tmp") +def utc_now_iso() -> str: + """Return the current UTC time as ISO-8601 for note frontmatter.""" + return datetime.now(timezone.utc).isoformat() + + +def normalize_title(raw: str, fallback: str) -> str: + """Return a topic or Agent label as a safe, filesystem- and wikilink-safe filename stem.""" + title = _UNSAFE_FILENAME.sub("-", _HEADING.sub("", str(raw or "").strip())) + title = re.sub(r"\s+", " ", title).strip(" .-") + if title.lower().endswith(".md"): + title = title[:-3].strip(" .-") + title = _truncate(title or fallback).strip(" .-") or fallback + if error := validate_filename_component(title, kind="title"): + raise ValueError(f"Unable to produce a safe Auto Fin title from {raw!r}: {error}") + return title + + +def _truncate(title: str) -> str: + """Fit a title into the byte budget one filename component accepts, without splitting a character.""" + kept: list[str] = [] + size = 0 + for character in title: + size += len(character.encode()) + if size > TITLE_BYTE_LIMIT: + break + kept.append(character) + return "".join(kept) + + +def normalize_report(output: AutoFinReportOutput) -> AutoFinReportOutput: + """Strip a duplicated H1 and fill in the fallbacks for one Agent report.""" + body = output.body.strip() + if body.startswith("# "): + body = body.partition("\n")[2].strip() + return output.model_copy( + update={ + "title": _HEADING.sub("", output.title.strip()) or "主题新闻观察", + "description": output.description.strip() or "基于当前新闻与历史记忆的主题研究。", + "body": body or "## 结论\n\n暂无可用结论。", + }, + ) + + +def normalize_hybrid_wikilinks(body: str) -> str: + """Drop a redundant Markdown destination from an unambiguous wikilink hybrid.""" + + def replace(match: re.Match[str]) -> str: + inner = match.group("inner").strip() + raw_target = inner.partition("|")[0].strip() + destination = match.group("destination").strip().removeprefix("<").removesuffix(">").strip() + return ( + match.group("wikilink") + if destination in {raw_target, raw_target.partition("#")[0].strip()} + else match.group(0) + ) + + return _HYBRID_WIKILINK.sub(replace, body) + + +def _valid_source_path(path: str) -> bool: + """Return whether one wikilink target could be a workspace-relative Markdown path.""" + return bool( + path + and not path.startswith("/") + and "\\" not in path + and path.endswith(".md") + and "." not in Path(path).parts + and ".." not in Path(path).parts + and not any(character in path for character in "[]|"), + ) + + +def validate_wikilinks(body: str, workspace: Path, exclude: Path) -> tuple[str, list[str]]: + """Keep real in-workspace Markdown links and downgrade invalid links to plain text.""" + workspace = workspace.resolve() + exclude = exclude.resolve() + sources: list[str] = [] + + def replace(match: re.Match[str]) -> str: + raw_target, separator, raw_alias = match.group(1).strip().partition("|") + path = raw_target.strip().partition("#")[0].strip() + alias = (raw_alias.strip() if separator else "") or Path(path).stem.replace("_", " ") + if not _valid_source_path(path): + return alias + resolved = (workspace / path).resolve() + if not resolved.is_relative_to(workspace) or not resolved.is_file() or resolved == exclude: + return alias + if path not in sources: + sources.append(path) + return match.group(0) + + return _WIKILINK.sub(replace, body), sources + + +def resolve_note_path(day_dir: Path, title: str, *, existing: Path | None) -> tuple[str, Path]: + """Return a title and path that neither reuse nor overwrite an unrelated note.""" + path = day_dir / f"{title}.md" + index = 2 + while path != existing and path.exists(): + path = day_dir / f"{title}({index}).md" + index += 1 + return path.stem, path + + +def find_note(day_dir: Path, *, kind: str, **matches: Any) -> Path | None: + """Return this workflow's note in one day directory whose frontmatter matches. + + Every Markdown file in the directory is a candidate, including notes the user + is editing by hand, so an unreadable or malformed one is skipped rather than + allowed to abort the run. + """ + for path in sorted(day_dir.glob("*.md")) if day_dir.is_dir() else (): + try: + metadata = frontmatter.load(path).metadata + except (OSError, UnicodeError, ValueError, yaml.YAMLError): + continue + if metadata.get("kind") == kind and all(metadata.get(key) == value for key, value in matches.items()): + return path + return None + + +def read_note(path: Path | None) -> str: + """Return an earlier note's body for intra-day refinement, or the first-run notice.""" + if path is None: + return FIRST_RUN_NOTICE try: - temporary.write_text(text, encoding="utf-8") - os.replace(temporary, path) - finally: - temporary.unlink(missing_ok=True) + return frontmatter.load(path).content.strip()[:NOTE_CHAR_LIMIT] or FIRST_RUN_NOTICE + except (OSError, UnicodeError, ValueError, yaml.YAMLError): + return FIRST_RUN_NOTICE + + +async def write_markdown(path: Path, body: str, metadata: dict[str, Any]) -> None: + """Serialize one frontmatter Markdown document atomically under its path lock.""" + path.parent.mkdir(parents=True, exist_ok=True) + rendered = frontmatter.dumps(frontmatter.Post(body.strip(), **metadata)) + async with await get_path_lock(path): + temporary = path.with_name(f".{path.name}.{uuid4().hex}.tmp") + try: + async with aiofiles.open(temporary, "w", encoding="utf-8") as stream: + await stream.write(rendered if rendered.endswith("\n") else f"{rendered}\n") + os.replace(temporary, path) + finally: + temporary.unlink(missing_ok=True) class AutoFinStep(BaseStep): - """Shared Auto Fin helpers.""" + """Shared helpers for the steps in one Auto Fin RuntimeContext.""" def _value(self, key: str, default: Any = None) -> Any: assert self.context is not None @@ -71,14 +245,66 @@ class AutoFinStep(BaseStep): raise RuntimeError(f"Auto Fin data is missing: {key}") return value + @property + def day_dir(self) -> Path: + """Return the dated directory that holds this run's notes.""" + return self.workspace_path / str(self.config_value("daily_dir")) / str(self._required("auto_fin_date")) + + def _track(self, path: Path, existing: Path | None) -> None: + """Record this write, and any note it replaced, for the auto_tag step.""" + assert self.context is not None + changes = list(self.context.get("changes") or []) + if existing is not None and existing != path: + existing.unlink(missing_ok=True) + changes.append({"change": "deleted", "path": existing.relative_to(self.workspace_path).as_posix()}) + changes.append( + { + "change": "modified" if existing == path else "added", + "path": path.relative_to(self.workspace_path).as_posix(), + }, + ) + self.context["changes"] = changes + + async def _write_report( + self, + filename: str, + output: AutoFinReportOutput, + *, + kind: str, + existing: Path | None = None, + trailer: str = "", + **metadata: Any, + ) -> WrittenReport: + """Write one report note, returning its path, valid sources, and validated body.""" + title, path = resolve_note_path(self.day_dir, filename, existing=existing) + body, sources = validate_wikilinks(normalize_hybrid_wikilinks(output.body), self.workspace_path, path) + await write_markdown( + path, + "\n\n".join(section for section in (body, trailer, DISCLAIMER) if section), + { + "name": title, + "title": output.title, + "description": output.description, + "kind": kind, + "generated_at": utc_now_iso(), + **metadata, + }, + ) + self._track(path, existing) + relative = path.relative_to(self.workspace_path).as_posix() + self.logger.info(f"[{self.name}] wrote kind={kind} path={relative} chars={len(body)} sources={len(sources)}") + return WrittenReport(path=relative, sources=sources, body=body) + async def _reply( self, prompt_name: str, model: type[BaseModel], job_tools: list[str] | None = None, injected_job_kwargs: dict[str, Any] | None = None, - **values: str, + tool_context_id: str | None = None, + **values: Any, ) -> BaseModel: + """Ask the Agent for one structured report and validate it against ``model``.""" if self.agent_wrapper is None: raise RuntimeError("Auto Fin analysis requires an agent_wrapper") prompt = self.prompt_format(prompt_name, **values) @@ -87,20 +313,26 @@ class AutoFinStep(BaseStep): f"[{self.name}] agent input prompt={prompt_name} schema={model.__name__} " f"query={self._preview(prompt, AGENT_INPUT_LOG_LIMIT)}", ) - kwargs: dict[str, Any] = {"output_schema": model} - if job_tools: - kwargs["job_tools"] = job_tools + kwargs: dict[str, Any] = { + "output_schema": model, + "builtin_tools": [], + "use_builtin_tools": False, + "skills": [], + "job_tools": list(self.kwargs.get("job_tools") or []) if job_tools is None else job_tools, + } if injected_job_kwargs: kwargs["injected_job_kwargs"] = injected_job_kwargs + if tool_context_id: + kwargs["tool_context_id"] = tool_context_id result = await self.agent_wrapper.reply(prompt, **kwargs) if not isinstance(result, dict) or result.get("structured_output") is None: raise ValueError(f"Auto Fin Agent returned no structured output: {self._preview(result)}") value = result["structured_output"] output = value if isinstance(value, model) else model.model_validate(value) - output_preview = self._preview(output.model_dump(), AGENT_OUTPUT_LOG_LIMIT) + rendered = self._preview(output.model_dump(), AGENT_OUTPUT_LOG_LIMIT) self.logger.info( f"[{self.name}] agent output prompt={prompt_name} schema={model.__name__} " - f"elapsed={perf_counter() - started_at:.2f}s output={output_preview}", + f"elapsed={perf_counter() - started_at:.2f}s output={rendered}", ) return output diff --git a/plugins/auto-fin/src/reme_auto_fin/data.py b/plugins/auto-fin/src/reme_auto_fin/data.py index 79551b92..b061750a 100644 --- a/plugins/auto-fin/src/reme_auto_fin/data.py +++ b/plugins/auto-fin/src/reme_auto_fin/data.py @@ -9,7 +9,7 @@ from typing import Any import httpx -from .base import AutoFinStep, _plain_text +from .base import AutoFinStep, plain_text API_URL = "https://www.cls.cn/v1/roll/get_roll_list" HEADERS = { @@ -72,6 +72,11 @@ class AutoFinDataStep(AutoFinStep): interval = max(0.0, float(self._value("request_interval", 10))) max_retries = max(1, int(self._value("max_retries", 3))) records: dict[str, dict[str, str]] = {} + page_count = 0 + self.logger.info( + f"[{self.name}] fetching CLS news window_start={cutoff.isoformat()} " + f"window_end={decision_at.isoformat()} request_interval={interval}s", + ) async with httpx.AsyncClient(headers=HEADERS, timeout=httpx.Timeout(20, connect=5)) as client: while True: for attempt in range(max_retries): @@ -81,12 +86,17 @@ class AutoFinDataStep(AutoFinStep): except (httpx.HTTPError, ValueError, RuntimeError) as exc: if attempt + 1 == max_retries: raise RuntimeError(f"CLS request failed after {max_retries} attempts: {exc}") from exc + self.logger.warning( + f"[{self.name}] CLS page request failed cursor={cursor} " + f"attempt={attempt + 1}/{max_retries}: {exc}", + ) await asyncio.sleep(2**attempt) finally: if interval: await asyncio.sleep(interval) if not rows: raise RuntimeError("CLS API returned no news before the 24-hour window was covered") + page_count += 1 timestamps = [int(row["ctime"]) for row in rows if str(row.get("ctime", "")).isdigit()] if not timestamps: raise RuntimeError("CLS API page contained no valid timestamps") @@ -97,9 +107,14 @@ class AutoFinDataStep(AutoFinStep): normalized = self._normalize(row, cutoff, decision_at) if normalized is not None: records.setdefault(normalized["news_id"], normalized) + self.logger.info( + f"[{self.name}] CLS page={page_count} rows={len(rows)} oldest={oldest} " + f"unique_in_window={len(records)}", + ) if oldest <= int(cutoff.timestamp()): break cursor = oldest + self.logger.info(f"[{self.name}] completed CLS fetch pages={page_count} news={len(records)}") return sorted(records.values(), key=lambda row: (row["event_time"], row["news_id"])) @staticmethod @@ -111,8 +126,8 @@ class AutoFinDataStep(AutoFinStep): return None if not start <= published_at <= end: return None - content = _plain_text(str(row.get("content") or row.get("brief") or "")) - title = _plain_text(str(row.get("title") or row.get("brief") or content)) + content = plain_text(str(row.get("content") or row.get("brief") or "")) + title = plain_text(str(row.get("title") or row.get("brief") or content)) if not title and not content: return None return { diff --git a/plugins/auto-fin/src/reme_auto_fin/digest.py b/plugins/auto-fin/src/reme_auto_fin/digest.py new file mode 100644 index 00000000..1be75795 --- /dev/null +++ b/plugins/auto-fin/src/reme_auto_fin/digest.py @@ -0,0 +1,62 @@ +"""Merge the topic notes into one dated brief that links back to each of them.""" + +from __future__ import annotations + +import json +from types import SimpleNamespace + +from reme.steps.file_io import refresh_day_index + +from .base import AutoFinStep, find_note, normalize_report, normalize_title, read_note +from .schema import AutoFinNote, AutoFinReportOutput + + +class AutoFinDigestStep(AutoFinStep): + """Summarize the topic notes and index the day once every note is in place.""" + + async def execute(self): + assert self.context is not None + if self.context.get("auto_fin_skipped"): + return self.context.response + notes: list[AutoFinNote] = self._required("auto_fin_notes") + run_date = str(self._required("auto_fin_date")) + earlier = find_note(self.day_dir, kind="auto-fin-digest") + output = normalize_report( + await self._reply( + "digest_user", + AutoFinReportOutput, + job_tools=[], + decision_at=str(self._required("auto_fin_decision_at")), + window_start=str(self._required("auto_fin_window_start")), + notes=json.dumps([note.model_dump() for note in notes], ensure_ascii=False), + earlier_brief=read_note(earlier), + ), + ) + written = await self._write_report( + normalize_title(f"主题新闻观察({run_date})", "主题新闻观察"), + output, + kind="auto-fin-digest", + existing=earlier, + trailer="## 主题详解\n\n" + "\n".join(f"- [[{note.path}]]" for note in notes), + date=run_date, + source_notes=[note.path for note in notes], + ) + await refresh_day_index( + SimpleNamespace(workspace_path=self.workspace_path), + run_date, + str(self.config_value("daily_dir")), + ) + self.context["markdown_path"] = self.context["auto_fin_digest_path"] = written.path + # The caller gets the validated body, not the Agent's raw reply: a link the + # note dropped must not reappear in the API response. + self.context.response.answer = written.body + self.context.response.metadata.update( + { + "markdown_path": written.path, + "digest_path": written.path, + "source_paths": written.sources, + "note_paths": [note.path for note in notes], + "selected_news_count": len(self._required("auto_fin_selected_news")), + }, + ) + return self.context.response diff --git a/plugins/auto-fin/src/reme_auto_fin/digest.yaml b/plugins/auto-fin/src/reme_auto_fin/digest.yaml new file mode 100644 index 00000000..afec4f90 --- /dev/null +++ b/plugins/auto-fin/src/reme_auto_fin/digest.yaml @@ -0,0 +1,20 @@ +digest_user: | + ## 输入材料 + 研究窗口:{window_start} 至 {decision_at} + 各主题笔记(JSON 列表,含 topic 与完整 Markdown body): + {notes} + 今天早些时段的总览: + {earlier_brief} + + ## 任务指令 + 你是当日总览 Agent。各主题笔记已经完成,你只负责把它们合成一篇总览,不再重新研究新闻,也没有任何搜索工具。 + + 写作要求: + - 用 `## 标题` 分层:先给跨主题的共同线索,再逐主题提炼要点。 + - 同一事件被多个主题收录时只写一次,并说明它同时影响哪些主题。 + - 只使用上面笔记中的事实,不得虚构行情、收益、价格或未提供的数据;证据不足处标明待核实。 + - 复核今天早些时段总览中仍然成立的判断;与笔记冲突时以笔记为准。 + - 做提炼而不是摘抄,不要整段复制笔记原文,也不要给投资建议。 + - 不要输出 wikilink,代码会在正文末尾自动附上各主题笔记的链接。 + + title 用不超过 20 字的短标题,description 用一句话概括;按结构化输出契约返回 title、description 和完整中文 Markdown body。 \ No newline at end of file diff --git a/plugins/auto-fin/src/reme_auto_fin/merge.py b/plugins/auto-fin/src/reme_auto_fin/merge.py deleted file mode 100644 index 124bb934..00000000 --- a/plugins/auto-fin/src/reme_auto_fin/merge.py +++ /dev/null @@ -1,155 +0,0 @@ -"""Research current news with ReMe and save a wikilink-backed report.""" - -from __future__ import annotations - -import json -import re -from datetime import date, timedelta -from pathlib import Path -from types import SimpleNamespace - -from reme.steps.file_io import refresh_day_index - -from .base import AutoFinStep, _write -from .schema import AutoFinReportOutput - -_WIKILINK_RE = re.compile(r"\[\[([^\[\]\n]+)\]\]") -_HYBRID_WIKILINK_RE = re.compile( - r"(?P\[\[(?P[^\[\]\n]+)\]\])\((?P[^()\n]+)\)", -) - - -class AutoFinMergeStep(AutoFinStep): - """Give one Agent read-only ReMe tools, then validate links in its Markdown.""" - - def _report_path(self, run_date: date) -> Path: - return self.workspace_path / str(self.config_value("daily_dir")) / str(run_date) / "auto_fin.md" - - def _current_report(self, run_date: date) -> str: - """Return today's existing report so intra-day reruns refine it, not replace it.""" - path = self._report_path(run_date) - if path.is_file(): - return path.read_text(encoding="utf-8") - return "今日暂无更早时段的推荐,本次为当日首次生成。" - - def _normalize_hybrid_wikilinks(self, body: str) -> str: - """Remove a redundant Markdown destination from an unambiguous wikilink hybrid.""" - - def replace(match: re.Match[str]) -> str: - inner = match.group("inner").strip() - raw_target = inner.partition("|")[0].strip() - target_path = raw_target.partition("#")[0].strip() - destination = match.group("destination").strip() - if destination.startswith("<") and destination.endswith(">"): - destination = destination[1:-1].strip() - if destination in {raw_target, target_path}: - return match.group("wikilink") - return match.group(0) - - try: - return _HYBRID_WIKILINK_RE.sub(replace, body) - except Exception as exc: # Defensive boundary: report generation must not depend on cosmetic normalization. - self.logger.warning(f"[{self.name}] failed to normalize hybrid wikilinks; keeping original body: {exc}") - return body - - @staticmethod - def _normalize(output: AutoFinReportOutput) -> AutoFinReportOutput: - title = re.sub(r"^#+\s*", "", output.title.strip()) or "主题新闻观察" - description = output.description.strip() or "基于当前新闻与历史记忆的主题研究。" - body = output.body.strip() or "## 结论\n\n暂无可用结论。" - if body.startswith("# "): - body = body.partition("\n")[2].lstrip() or "## 结论\n\n暂无可用结论。" - return output.model_copy(update={"title": title, "description": description, "body": body}) - - def _validate_wikilinks(self, body: str, run_date: date) -> tuple[str, list[str]]: - """Keep real in-workspace Markdown links and downgrade invalid links to text.""" - source_paths: list[str] = [] - report = self._report_path(run_date).resolve() - workspace = self.workspace_path.resolve() - - def replace(match: re.Match[str]) -> str: - inner = match.group(1).strip() - raw_target, separator, raw_alias = inner.partition("|") - target = raw_target.strip() - path = target.partition("#")[0].strip() - alias = (raw_alias.strip() if separator else "") or Path(path).stem.replace("_", " ") - if not self._valid_wikilink_path(path): - return alias - resolved = (workspace / path).resolve() - try: - resolved.relative_to(workspace) - except ValueError: - return alias - if not resolved.is_file() or resolved == report: - return alias - if path not in source_paths: - source_paths.append(path) - return match.group(0) - - return _WIKILINK_RE.sub(replace, body), source_paths - - @staticmethod - def _valid_wikilink_path(path: str) -> bool: - parts = Path(path).parts - return bool( - path - and not path.startswith("/") - and "\\" not in path - and path.endswith(".md") - and "." not in parts - and ".." not in parts - and not any(character in path for character in "[]|"), - ) - - async def execute(self): - """Research the selected news and persist the validated report.""" - assert self.context is not None - self.context["changes"] = [] - if self.context.get("auto_fin_skipped"): - return self.context.response - run_date = date.fromisoformat(str(self._required("auto_fin_date"))) - historical_search = { - "limit": 5, - "min_score": 0.0, - "start_date": None, - "end_date": (run_date - timedelta(days=1)).isoformat(), - } - output = await self._reply( - "merge_user", - AutoFinReportOutput, - job_tools=list(self.kwargs.get("job_tools") or []), - injected_job_kwargs=historical_search, - decision_at=str(self._required("auto_fin_decision_at")), - window_start=str(self._required("auto_fin_window_start")), - topics=json.dumps(self._required("auto_fin_topics"), ensure_ascii=False), - news=json.dumps(self._required("auto_fin_selected_news"), ensure_ascii=False), - current_report=self._current_report(run_date), - ) - output = self._normalize(output) - output = output.model_copy(update={"body": self._normalize_hybrid_wikilinks(output.body)}) - body, source_paths = self._validate_wikilinks(output.body, run_date) - output = output.model_copy(update={"body": body}) - markdown = f"# {output.title}\n\n> {output.description}\n\n{output.body}\n\n" - markdown += "> 未接入可靠行情数据;本文只提供新闻研究和回顾线索,不提供收益、目标价或买卖建议。\n" - report = self._report_path(run_date) - change = "modified" if report.is_file() else "added" - _write(report, markdown) - await refresh_day_index( - SimpleNamespace(workspace_path=self.workspace_path), - str(run_date), - str(self.config_value("daily_dir")), - ) - relative = report.relative_to(self.workspace_path).as_posix() - self.context["changes"] = [{"change": change, "path": relative}] - self.context["markdown_path"] = relative - self.context["auto_fin_digest_path"] = relative - self.context.response.answer = output.body - self.context.response.metadata.update( - { - "markdown_path": relative, - "digest_path": relative, - "source_paths": source_paths, - "selected_news_count": len(self._required("auto_fin_selected_news")), - }, - ) - return self.context.response diff --git a/plugins/auto-fin/src/reme_auto_fin/merge.yaml b/plugins/auto-fin/src/reme_auto_fin/merge.yaml deleted file mode 100644 index 01d13fee..00000000 --- a/plugins/auto-fin/src/reme_auto_fin/merge.yaml +++ /dev/null @@ -1,19 +0,0 @@ -merge_user: | - 你是主题新闻研究 Agent。当前新闻已经按 topics 做过语义筛选。你可以使用 `search` 搜索历史记忆, - 并使用 `read` 阅读可能相关的完整 Markdown。不得使用外部搜索,不得虚构行情、收益、价格或未提供的数据。 - - 研究窗口:{window_start} 至 {decision_at} - topics:{topics} - 当前新闻:{news} - 今天早些时段的报告(如有,请保留仍成立的判断,只修订变化部分): - {current_report} - - 先围绕 topics 和当前重要事件多次调用 `search` 检索历史记忆。 - 只对明显相关的结果调用 `read`。说明历史事件与当前事件的相同点、关键差异,以及旧判断是否仍适用。 - 给出值得回顾的新闻、应继续观察的信息,以及哪些条件会强化或推翻判断,但不要给出投资建议。 - - 当前新闻不是 workspace 文件:在正文中保留其 CLS news_id、发布时间和标题,不要为它虚构 wikilink。 - 对实际搜索、阅读并支持判断的历史 Markdown,把完整 workspace-relative path 作为 contextual wikilink 直接写进 - 相关句子,只写一次。不得写裸链接行。代码会把不存在、越界或指向当前报告的链接降级为普通文本。 - - 按结构化输出契约返回 title、description 和完整中文 Markdown body。 diff --git a/plugins/auto-fin/src/reme_auto_fin/plugin.yaml b/plugins/auto-fin/src/reme_auto_fin/plugin.yaml index 67c29516..903b4d60 100644 --- a/plugins/auto-fin/src/reme_auto_fin/plugin.yaml +++ b/plugins/auto-fin/src/reme_auto_fin/plugin.yaml @@ -1,7 +1,8 @@ backends: auto_fin_data_step: reme_auto_fin.data:AutoFinDataStep auto_fin_topic_step: reme_auto_fin.topic:AutoFinTopicStep - auto_fin_merge_step: reme_auto_fin.merge:AutoFinMergeStep + auto_fin_research_step: reme_auto_fin.research:AutoFinResearchStep + auto_fin_digest_step: reme_auto_fin.digest:AutoFinDigestStep application_defaults: jobs: @@ -41,11 +42,12 @@ application_defaults: steps: &auto_fin_steps - backend: auto_fin_data_step - backend: auto_fin_topic_step - - backend: auto_fin_merge_step - job_tools: [search, read] + - backend: auto_fin_research_step + job_tools: [search] + - backend: auto_fin_digest_step - backend: auto_tag_step auto_fin_cron: backend: cron - cron: "0 18 * * *" + cron: "0 9 * * *" steps: *auto_fin_steps diff --git a/plugins/auto-fin/src/reme_auto_fin/research.py b/plugins/auto-fin/src/reme_auto_fin/research.py new file mode 100644 index 00000000..b3c7beae --- /dev/null +++ b/plugins/auto-fin/src/reme_auto_fin/research.py @@ -0,0 +1,106 @@ +"""Research each topic's current news and persist one note per topic.""" + +from __future__ import annotations + +import json +from datetime import date, timedelta +from typing import Any +from uuid import uuid4 + +from .base import AutoFinStep, find_note, normalize_report, normalize_title, read_note +from .schema import AutoFinNote, AutoFinReportOutput + +NEWS_LIMIT = 20 +SEARCH_LIMIT = 3 + + +class AutoFinResearchStep(AutoFinStep): + """Research every topic independently so each one earns its own note.""" + + @staticmethod + def _search_budget(run_date: str) -> dict[str, Any]: + """Return the read-only historical search budget granted to one topic.""" + return { + "limit": 5, + "min_score": 0.0, + "start_date": None, + "end_date": (date.fromisoformat(run_date) - timedelta(days=1)).isoformat(), + "max_search_calls": SEARCH_LIMIT, + } + + async def _research(self, topic: str, related: list[dict], run_date: str) -> AutoFinNote: + """Research one topic and write its validated note.""" + recent = sorted(related, key=lambda row: (row["event_time"], row["news_id"]), reverse=True)[:NEWS_LIMIT] + earlier = find_note(self.day_dir, kind="auto-fin-topic", topic=topic) + context_id = f"auto_fin:{uuid4().hex}" + self.logger.info( + f"[{self.name}] research topic={topic} related={len(related)} " + f"sent={len(recent)} omitted={len(related) - len(recent)}", + ) + try: + output = normalize_report( + await self._reply( + "research_user", + AutoFinReportOutput, + job_tools=["search"], + injected_job_kwargs=self._search_budget(run_date), + tool_context_id=context_id, + decision_at=str(self._required("auto_fin_decision_at")), + window_start=str(self._required("auto_fin_window_start")), + topic=topic, + news=json.dumps(recent, ensure_ascii=False), + omitted_news_count=str(len(related) - len(recent)), + earlier_note=read_note(earlier), + ), + ) + finally: + if self.app_context is not None: + self.app_context.metadata.get("__search_call_budgets", {}).pop(context_id, None) + written = await self._write_report( + normalize_title(topic, "主题观察"), + output, + kind="auto-fin-topic", + existing=earlier, + date=run_date, + topic=topic, + source_news_ids=[row["news_id"] for row in recent], + ) + return AutoFinNote( + topic=topic, + title=output.title, + description=output.description, + body=written.body, + path=written.path, + ) + + async def execute(self): + """Write one note per topic that had relevant news, isolating per-topic failures.""" + assert self.context is not None + self.context["changes"] = [] + if self.context.get("auto_fin_skipped"): + return self.context.response + run_date = str(self._required("auto_fin_date")) + by_topic = self._required("auto_fin_news_by_topic") + notes: list[AutoFinNote] = [] + failures: list[dict[str, str]] = [] + for topic in self._required("auto_fin_topics"): + if not by_topic[topic]: + self.logger.info(f"[{self.name}] skipping topic={topic} reason=no_related_news") + continue + try: + notes.append(await self._research(topic, by_topic[topic], run_date)) + except Exception as exc: + failures.append({"topic": topic, "error": str(exc)}) + self.logger.warning(f"[{self.name}] topic={topic} failed; continuing with the rest: {exc}") + if failures and not notes: + raise RuntimeError(f"Auto Fin research failed for every topic: {failures}") + self.context["auto_fin_notes"] = notes + self.context.response.answer = f"Researched {len(notes)} topic(s): {', '.join(note.path for note in notes)}" + self.context.response.metadata.update( + { + "note_paths": [note.path for note in notes], + "failed_topics": failures, + "selected_news_count": len(self._required("auto_fin_selected_news")), + }, + ) + return self.context.response diff --git a/plugins/auto-fin/src/reme_auto_fin/research.yaml b/plugins/auto-fin/src/reme_auto_fin/research.yaml new file mode 100644 index 00000000..c06394aa --- /dev/null +++ b/plugins/auto-fin/src/reme_auto_fin/research.yaml @@ -0,0 +1,25 @@ +research_user: | + ## 输入材料 + 研究窗口:{window_start} 至 {decision_at} + 当前主题:{topic} + 该主题最新的新闻(最多 20 篇,另有 {omitted_news_count} 篇未送入本次研究): + {news} + 今天早些时段该主题的笔记: + {earlier_note} + + ## 任务指令 + 你是单主题新闻研究 Agent。当前新闻已按主题筛选,你只研究「{topic}」这一个主题。 + 你只能使用 `search` 搜索历史记忆,最多调用 3 次。不得使用外部搜索,不得虚构行情、收益、价格或未提供的数据。 + 搜索片段不足以支持的细节应标明待核实。 + 复核今天早些时段笔记中与当前主题有关的判断;本次证据与之冲突时,说明冲突点和哪一方更可信。 + + 围绕当前主题和重要事件检索历史记忆;用完搜索额度后直接完成笔记。 + 说明历史事件与当前事件的相同点、关键差异,以及旧判断是否仍适用。 + 给出值得回顾的新闻、应继续观察的信息,以及哪些条件会强化或推翻判断,但不要给出投资建议。 + + 当前新闻不是 workspace 文件:在正文中保留其 CLS news_id、发布时间和标题,不要为它虚构 wikilink。 + 对搜索结果中实际支持判断的历史 Markdown,把完整 workspace-relative path 作为 contextual wikilink 写进 + 相关句子,例如:这与[[daily/2026-08-01/auto_fin.md|此前的供给判断]]相似,但本次政策范围不同。 + 不得写裸链接行或虚构路径。代码会把不存在、越界或指向本文的链接降级为普通文本。 + + title 用不超过 20 字的短标题,description 用一句话概括;按结构化输出契约返回 title、description 和完整中文 Markdown body。 \ No newline at end of file diff --git a/plugins/auto-fin/src/reme_auto_fin/schema.py b/plugins/auto-fin/src/reme_auto_fin/schema.py index 3db10c5b..3e5db0f8 100644 --- a/plugins/auto-fin/src/reme_auto_fin/schema.py +++ b/plugins/auto-fin/src/reme_auto_fin/schema.py @@ -18,14 +18,18 @@ class AutoFinAgentModel(AutoFinModel): class AutoFinReportOutput(AutoFinAgentModel): - """Final Markdown returned by the agentic news-research Agent.""" + """One Chinese Markdown report returned by an Auto Fin Agent.""" title: str description: str body: str -class AutoFinTopicOutput(AutoFinAgentModel): - """CLS news identifiers that are semantically related to configured topics.""" +class AutoFinNote(AutoFinModel): + """One topic note the research step hands to the digest step.""" - news_ids: list[str] + topic: str + title: str + description: str + body: str + path: str diff --git a/plugins/auto-fin/src/reme_auto_fin/topic.py b/plugins/auto-fin/src/reme_auto_fin/topic.py index d5c5e680..993ebaa3 100644 --- a/plugins/auto-fin/src/reme_auto_fin/topic.py +++ b/plugins/auto-fin/src/reme_auto_fin/topic.py @@ -2,38 +2,120 @@ from __future__ import annotations +from collections.abc import Iterator import json +import re +from time import perf_counter -from .base import AutoFinStep -from .schema import AutoFinTopicOutput +from .base import AGENT_INPUT_LOG_LIMIT, AGENT_OUTPUT_LOG_LIMIT, AutoFinStep class AutoFinTopicStep(AutoFinStep): """Filter current news in bounded Agent batches without writing files.""" + PROMPT_CHAR_LIMIT = 100_000 + + @staticmethod + def _parse_news_ids(value: object, topics: list[str]) -> dict[str, list[str]]: + if not isinstance(value, str): + raise ValueError("Auto Fin topic Agent returned no text") + match = re.search(r"```json\s*(.*?)```", value, re.IGNORECASE | re.DOTALL) + if match is None: + raise ValueError("Auto Fin topic Agent returned no JSON code block") + try: + mapping = json.loads(match.group(1).strip()) + except json.JSONDecodeError as exc: + raise ValueError("Auto Fin topic Agent returned invalid JSON") from exc + if not isinstance(mapping, dict) or set(mapping) != set(topics): + raise ValueError("Auto Fin topic Agent must return exactly the configured topic keys") + if any(not isinstance(ids, list) or any(not isinstance(item, str) for item in ids) for ids in mapping.values()): + raise ValueError("Auto Fin topic Agent must return string news ID arrays") + return mapping + + async def _select_news_ids(self, prompt: str, topics: list[str]) -> dict[str, list[str]]: + """Request and validate plain-text topic IDs, retrying one malformed reply.""" + if self.agent_wrapper is None: + raise RuntimeError("Auto Fin analysis requires an agent_wrapper") + self.logger.info( + f"[{self.name}] agent input prompt=topic_user query={self._preview(prompt, AGENT_INPUT_LOG_LIMIT)}", + ) + for attempt in range(2): + started_at = perf_counter() + result = await self.agent_wrapper.reply(prompt) + try: + ids = self._parse_news_ids(result.get("result") if isinstance(result, dict) else None, topics) + except ValueError as exc: + if attempt: + raise ValueError(f"Auto Fin topic Agent returned invalid news IDs: {exc}") from exc + self.logger.warning(f"[{self.name}] invalid topic JSON; retrying once: {exc}") + continue + self.logger.info( + f"[{self.name}] agent output prompt=topic_user elapsed={perf_counter() - started_at:.2f}s " + f"output={self._preview(ids, AGENT_OUTPUT_LOG_LIMIT)}", + ) + return ids + raise RuntimeError("Auto Fin topic Agent produced no response") + + def _prompt(self, news: list[dict], topics: list[str], window_hours: str) -> str: + return self.prompt_format( + "topic_user", + topics=json.dumps(topics, ensure_ascii=False), + news=json.dumps(news, ensure_ascii=False), + window_hours=window_hours, + output_example=json.dumps({topic: [] for topic in topics}, ensure_ascii=False), + ) + + def _batches(self, news: list[dict], topics: list[str], window_hours: str) -> Iterator[list[dict]]: + batch: list[dict] = [] + for row in news: + item = {**row, "title": str(row.get("title") or "")[:300], "content": str(row.get("content") or "")[:1000]} + if len(self._prompt([*batch, item], topics, window_hours)) > self.PROMPT_CHAR_LIMIT: + if not batch: + raise ValueError("Auto Fin news item exceeds the topic Agent prompt limit") + yield batch + batch = [item] + if len(self._prompt(batch, topics, window_hours)) > self.PROMPT_CHAR_LIMIT: + raise ValueError("Auto Fin news item exceeds the topic Agent prompt limit") + else: + batch.append(item) + if batch: + yield batch + async def execute(self): + """Select relevant news from each batch for the current invocation.""" assert self.context is not None news = list(self._required("auto_fin_news")) topics = list(self._required("auto_fin_topics")) window_hours = float(self._value("auto_fin_window_hours", 24)) formatted_hours = f"{window_hours:g}" - batch_size = max(1, int(self._value("topic_batch_size", 50))) - selected: set[str] = set() - for start in range(0, len(news), batch_size): - batch = [ - {**row, "content": str(row.get("content") or "")[:1000]} for row in news[start : start + batch_size] - ] - output = await self._reply( - "topic_user", - AutoFinTopicOutput, - topics=json.dumps(topics, ensure_ascii=False), - news=json.dumps(batch, ensure_ascii=False), - window_hours=formatted_hours, + selected: dict[str, set[str]] = {topic: set() for topic in topics} + batch_count = 0 + for batch in self._batches(news, topics, formatted_hours): + batch_count += 1 + valid_ids = {row["news_id"] for row in batch} + prompt = self._prompt(batch, topics, formatted_hours) + self.logger.info( + f"[{self.name}] filtering batch={batch_count} news={len(batch)} prompt_chars={len(prompt)} " + f"first_id={batch[0]['news_id']} last_id={batch[-1]['news_id']}", ) - selected.update(output.news_ids) - relevant = [row for row in news if row["news_id"] in selected] + for topic, ids in (await self._select_news_ids(prompt, topics)).items(): + unknown = len(set(ids) - valid_ids) + if unknown: + self.logger.warning( + f"[{self.name}] ignored unknown IDs batch={batch_count} topic={topic} count={unknown}", + ) + selected[topic].update(valid_ids.intersection(ids)) + by_topic = {topic: [row for row in news if row["news_id"] in selected[topic]] for topic in topics} + relevant_ids = set().union(*selected.values()) if selected else set() + relevant = [row for row in news if row["news_id"] in relevant_ids] + self.context["auto_fin_news_by_topic"] = by_topic self.context["auto_fin_selected_news"] = relevant self.context.response.metadata["relevant_news_count"] = len(relevant) + self.context.response.metadata["topic_batch_count"] = batch_count + self.logger.info( + f"[{self.name}] topic filtering complete batches={batch_count} selected_unique={len(relevant)} " + f"per_topic={{{', '.join(f'{topic!r}: {len(rows)}' for topic, rows in by_topic.items())}}}", + ) if not relevant: reason = f"最近{formatted_hours}小时没有与 {', '.join(topics)} 相关的财联社新闻。" self.context["auto_fin_skipped"] = True diff --git a/plugins/auto-fin/src/reme_auto_fin/topic.yaml b/plugins/auto-fin/src/reme_auto_fin/topic.yaml index 638e9a7d..dd16c9cf 100644 --- a/plugins/auto-fin/src/reme_auto_fin/topic.yaml +++ b/plugins/auto-fin/src/reme_auto_fin/topic.yaml @@ -1,6 +1,14 @@ topic_user: | - 从下面最近{window_hours}小时的财联社新闻中,选择与至少一个 topics 存在真实、可解释关系的新闻。 - 只返回输入中真实存在的 news_id;仅出现关键词但没有实质关系的新闻不要选择。 + ## 输入材料 + 时间范围:最近 {window_hours} 小时 + 主题:{topics} + 财联社新闻: + {news} - topics:{topics} - 新闻:{news} + ## 任务指令 + 按主题归类新闻。只选择与主题存在真实、可解释关系的新闻;仅出现关键词而无实质关系的不要选择。 + 同一新闻可以属于多个主题。只返回材料中真实存在的 news_id,每个输入主题都必须出现,没有相关新闻时使用空数组。 + 在 ```json 代码块中只返回 JSON 对象,不要补充解释。使用以下键结构,并填入真实 news_id: + ```json + {output_example} + ``` diff --git a/plugins/auto-fin/tests/test_auto_fin.py b/plugins/auto-fin/tests/test_auto_fin.py index a273d68d..d84d9f68 100644 --- a/plugins/auto-fin/tests/test_auto_fin.py +++ b/plugins/auto-fin/tests/test_auto_fin.py @@ -2,6 +2,8 @@ # pylint: disable=missing-function-docstring,protected-access +import json +import re from datetime import datetime from pathlib import Path from zoneinfo import ZoneInfo @@ -9,26 +11,69 @@ from zoneinfo import ZoneInfo import pytest import yaml -from reme_auto_fin.base import _plain_text, _write +from reme_auto_fin.base import ( + FIRST_RUN_NOTICE, + normalize_hybrid_wikilinks, + normalize_title, + plain_text, + read_note, + write_markdown, +) from reme_auto_fin.data import AutoFinDataStep -from reme_auto_fin.merge import AutoFinMergeStep -from reme_auto_fin.schema import AutoFinReportOutput, AutoFinTopicOutput +from reme_auto_fin.digest import AutoFinDigestStep +from reme_auto_fin.research import AutoFinResearchStep +from reme_auto_fin.schema import AutoFinNote, AutoFinReportOutput from reme_auto_fin.topic import AutoFinTopicStep from reme.components import ApplicationContext from reme.components.agent_wrapper.base_agent_wrapper import BaseAgentWrapper from reme.components.runtime_context import RuntimeContext +from reme.utils.wikilink_handler import WikilinkHandler SHANGHAI = ZoneInfo("Asia/Shanghai") PLUGIN_MANIFEST = yaml.safe_load( (Path(__file__).parents[1] / "src" / "reme_auto_fin" / "plugin.yaml").read_text(encoding="utf-8"), ) +LINKED_BODY = ( + "## 今日判断\n\n" + "CLS 1(09:00,黄金上涨)与 " + "[[daily/2026-08-01/auto_fin.md|历史黄金观察]]" + "(daily/2026-08-01/auto_fin.md) 背景相似。\n\n" + "无效引用 [[daily/missing.md|缺失文章]] 和 [[../../outside.md|越界文章]] 应降级。" +) def _row(news_id: int, value: datetime, title: str = "新闻", content: str = "正文") -> dict: - return {"id": news_id, "ctime": int(value.timestamp()), "title": title, "content": content} + return { + "id": news_id, + "ctime": int(value.timestamp()), + "title": title, + "content": content, + } -def test_atomic_write_preserves_existing_file_on_failure(tmp_path: Path, monkeypatch): +def _news(news_id: str, event_time: str, title: str = "新闻") -> dict: + return {"news_id": news_id, "event_time": event_time, "title": title, "content": "正文"} + + +def _context(**kwargs) -> RuntimeContext: + defaults = { + "auto_fin_date": "2026-08-10", + "auto_fin_decision_at": "2026-08-10T09:30:00+08:00", + "auto_fin_window_start": "2026-08-09T09:30:00+08:00", + "auto_fin_topics": ["黄金", "机器人", "半导体"], + "auto_fin_selected_news": [_news("1", "2026-08-10T09:00:00+08:00")], + } + return RuntimeContext(**{**defaults, **kwargs}) + + +def _history(tmp_path: Path) -> None: + historical = tmp_path / "daily" / "2026-08-01" / "auto_fin.md" + historical.parent.mkdir(parents=True, exist_ok=True) + historical.write_text("# 历史黄金观察\n", encoding="utf-8") + + +@pytest.mark.asyncio +async def test_write_markdown_preserves_existing_file_on_failure(tmp_path: Path, monkeypatch): path = tmp_path / "result.md" path.write_text("existing", encoding="utf-8") monkeypatch.setattr( @@ -37,11 +82,14 @@ def test_atomic_write_preserves_existing_file_on_failure(tmp_path: Path, monkeyp ) with pytest.raises(OSError): - _write(path, "replacement") + await write_markdown(path, "replacement", {"name": "标题"}) assert path.read_text(encoding="utf-8") == "existing" assert not list(tmp_path.glob(".*.tmp")) - assert _plain_text("

甲&乙

丙

") == "甲&乙 丙" + + +def test_plain_text_drops_hidden_and_unescapes_entities(): + assert plain_text("

甲&乙

丙

") == "甲&乙 丙" @pytest.mark.asyncio @@ -94,24 +142,34 @@ async def test_data_step_uses_configurable_window_hours(tmp_path: Path, monkeypa class _TopicAgent(BaseAgentWrapper): - def __init__(self, news_ids: list[str], **kwargs): + def __init__(self, selected: dict[str, list[str]], **kwargs): super().__init__(**kwargs) - self.news_ids = news_ids + self.selected = selected self.calls = [] async def reply(self, inputs, **kwargs): self.calls.append((str(inputs), kwargs)) - return {"structured_output": AutoFinTopicOutput(news_ids=self.news_ids)} + return {"result": f"筛选结果:\n```json\n{json.dumps(self.selected)}\n```\n以上是相关 ID。"} @pytest.mark.asyncio async def test_topic_step_keeps_real_ids_in_memory_only(tmp_path: Path): app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") - agent = _TopicAgent(["2", "missing", "2"], app_context=app_context) + agent = _TopicAgent({"黄金": ["2", "missing", "2"]}, app_context=app_context) context = RuntimeContext( auto_fin_news=[ - {"news_id": "1", "event_time": "2026-08-10T08:00:00+08:00", "title": "甲", "content": "甲"}, - {"news_id": "2", "event_time": "2026-08-10T09:00:00+08:00", "title": "乙", "content": "乙"}, + { + "news_id": "1", + "event_time": "2026-08-10T08:00:00+08:00", + "title": "甲", + "content": "甲", + }, + { + "news_id": "2", + "event_time": "2026-08-10T09:00:00+08:00", + "title": "乙", + "content": "乙", + }, ], auto_fin_topics=["黄金"], ) @@ -119,7 +177,10 @@ async def test_topic_step_keeps_real_ids_in_memory_only(tmp_path: Path): response = await AutoFinTopicStep(app_context=app_context, agent_wrapper=agent)(context) assert [row["news_id"] for row in context["auto_fin_selected_news"]] == ["2"] - assert agent.calls[0][1] == {"output_schema": AutoFinTopicOutput} + assert [row["news_id"] for row in context["auto_fin_news_by_topic"]["黄金"]] == ["2"] + assert agent.calls[0][1] == {} + assert '```json\n{"黄金": []}' in agent.calls[0][0] + assert agent.calls[0][0].index("## 输入材料") < agent.calls[0][0].index("## 任务指令") assert response.metadata["relevant_news_count"] == 1 assert not list(tmp_path.rglob("*.*")) @@ -127,10 +188,15 @@ async def test_topic_step_keeps_real_ids_in_memory_only(tmp_path: Path): @pytest.mark.asyncio async def test_topic_step_marks_empty_selection_as_successful_skip(tmp_path: Path): app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") - agent = _TopicAgent([], app_context=app_context) + agent = _TopicAgent({"黄金": []}, app_context=app_context) context = RuntimeContext( auto_fin_news=[ - {"news_id": "1", "event_time": "2026-08-10T08:00:00+08:00", "title": "甲", "content": "甲"}, + { + "news_id": "1", + "event_time": "2026-08-10T08:00:00+08:00", + "title": "甲", + "content": "甲", + }, ], auto_fin_topics=["黄金"], auto_fin_window_hours=12, @@ -144,114 +210,425 @@ async def test_topic_step_marks_empty_selection_as_successful_skip(tmp_path: Pat assert context["auto_fin_skipped"] is True assert response.metadata["skipped"] is True assert response.answer == "最近12小时没有与 黄金 相关的财联社新闻。" - assert "最近12小时" in agent.calls[0][0] + assert "最近 12 小时" in agent.calls[0][0] assert not list(tmp_path.rglob("*.md")) -class _ResearchAgent(BaseAgentWrapper): +@pytest.mark.asyncio +async def test_topic_step_retries_invalid_json_once(tmp_path: Path): + class RetryAgent(_TopicAgent): + """Return one malformed response before a fenced topic mapping.""" + + async def reply(self, inputs, **kwargs): + self.calls.append((str(inputs), kwargs)) + return {"result": ('{"new_ids": "[\\"1\\"]"}' if len(self.calls) == 1 else '```json\n{"黄金": ["1"]}\n```')} + + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + agent = RetryAgent({"黄金": []}, app_context=app_context) + context = RuntimeContext( + auto_fin_news=[ + { + "news_id": "1", + "event_time": "2026-08-10T08:00:00+08:00", + "title": "甲", + "content": "甲", + }, + ], + auto_fin_topics=["黄金"], + ) + + await AutoFinTopicStep(app_context=app_context, agent_wrapper=agent)(context) + + assert len(agent.calls) == 2 + assert [row["news_id"] for row in context["auto_fin_selected_news"]] == ["1"] + + +@pytest.mark.parametrize( + "value", + ['{"new_ids": ["1"]}', "```json\n[1]\n```", '```json\n{"黄金": [1]}\n```', "not json", ""], +) +def test_topic_step_rejects_non_array_or_non_string_ids(value: str): + with pytest.raises(ValueError): + AutoFinTopicStep._parse_news_ids(value, ["黄金"]) + + +@pytest.mark.asyncio +async def test_topic_step_batches_by_prompt_length_and_merges_topics(tmp_path: Path): + app_context = ApplicationContext(workspace_dir=str(tmp_path)) + agent = _TopicAgent({"黄金": ["1", "2", "2"], "机器人": ["2"]}, app_context=app_context) + news = [ + { + "news_id": str(index), + "event_time": f"2026-08-10T0{index}:00:00+08:00", + "title": "新闻", + "content": "正文" * 80, + } + for index in (1, 2, 3) + ] + step = AutoFinTopicStep(app_context=app_context, agent_wrapper=agent) + one_item_length = len(step._prompt([{**news[0], "content": news[0]["content"][:1000]}], ["黄金", "机器人"], "24")) + step.PROMPT_CHAR_LIMIT = one_item_length + 5 + context = RuntimeContext(auto_fin_news=news, auto_fin_topics=["黄金", "机器人"]) + + response = await step(context) + + assert len(agent.calls) == 3 + assert all(len(prompt) <= step.PROMPT_CHAR_LIMIT for prompt, _ in agent.calls) + assert [row["news_id"] for row in context["auto_fin_news_by_topic"]["黄金"]] == ["1", "2"] + assert [row["news_id"] for row in context["auto_fin_news_by_topic"]["机器人"]] == ["2"] + assert [row["news_id"] for row in context["auto_fin_selected_news"]] == ["1", "2"] + assert response.metadata["topic_batch_count"] == 3 + + +class _ReportAgent(BaseAgentWrapper): + """Return one titled Markdown report per call, keyed off the research topic.""" + def __init__(self, **kwargs): super().__init__(**kwargs) self.calls = [] async def reply(self, inputs, **kwargs): - self.calls.append((str(inputs), kwargs)) + prompt = str(inputs) + self.calls.append((prompt, kwargs)) + topic = re.search(r"当前主题:(\S+)", prompt) + title = f"{topic.group(1)}观察" if topic else "主题新闻观察" return { "structured_output": AutoFinReportOutput( - title="# 主题新闻观察", - description="关注黄金政策变化。", - body=( - "## 今日判断\n\n" - "CLS 1(09:00,黄金上涨)与 " - "[[daily/2026-08-01/auto_fin.md|历史黄金观察]]" - "(daily/2026-08-01/auto_fin.md) 背景相似。\n\n" - "无效引用 [[daily/missing.md|缺失文章]] 和 [[../../outside.md|越界文章]] 应降级。" - ), + title=f"# {title}", + description="关注政策变化。", + body=LINKED_BODY, ), } @pytest.mark.asyncio -async def test_merge_writes_only_final_report_and_validates_historical_links(tmp_path: Path): - historical = tmp_path / "daily" / "2026-08-01" / "auto_fin.md" - historical.parent.mkdir(parents=True) - historical.write_text("# 历史黄金观察\n", encoding="utf-8") +async def test_research_writes_one_note_per_topic_with_latest_twenty_news(tmp_path: Path): + _history(tmp_path) app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") - agent = _ResearchAgent(app_context=app_context) - context = RuntimeContext( - auto_fin_date="2026-08-10", - auto_fin_decision_at="2026-08-10T09:30:00+08:00", - auto_fin_window_start="2026-08-09T09:30:00+08:00", - auto_fin_topics=["黄金"], - auto_fin_selected_news=[ - { - "news_id": "1", - "event_time": "2026-08-10T09:00:00+08:00", - "title": "黄金上涨", - "content": "避险需求增强", - }, - ], + agent = _ReportAgent(app_context=app_context) + gold = [_news(str(index), f"2026-08-10T09:{index:02}:00+08:00", "黄金") for index in range(25)] + context = _context( + auto_fin_news_by_topic={"黄金": gold, "机器人": [_news("30", "2026-08-10T08:00:00+08:00")], "半导体": []}, + auto_fin_selected_news=[*gold, _news("30", "2026-08-10T08:00:00+08:00")], ) - response = await AutoFinMergeStep( - app_context=app_context, - agent_wrapper=agent, - job_tools=["search", "read"], - )(context) + response = await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(context) - prompt, kwargs = agent.calls[0] - assert "end_date" not in prompt - assert "调用 `search`" in prompt - assert "调用 `read`" in prompt - assert kwargs == { - "output_schema": AutoFinReportOutput, - "job_tools": ["search", "read"], - "injected_job_kwargs": { - "limit": 5, - "min_score": 0.0, - "start_date": None, - "end_date": "2026-08-09", - }, + assert len(agent.calls) == 2 + gold_prompt, gold_kwargs = agent.calls[0] + assert '"news_id": "24"' in gold_prompt + assert '"news_id": "4"' not in gold_prompt + assert "另有 5 篇" in gold_prompt + assert "当前主题:机器人" in agent.calls[1][0] + assert gold_kwargs["job_tools"] == ["search"] + assert gold_kwargs["output_schema"] == AutoFinReportOutput + assert gold_kwargs["tool_context_id"].startswith("auto_fin:") + assert gold_kwargs["tool_context_id"] != agent.calls[1][1]["tool_context_id"] + assert gold_kwargs["injected_job_kwargs"] == { + "limit": 5, + "min_score": 0.0, + "start_date": None, + "end_date": "2026-08-09", + "max_search_calls": 3, } - report = (tmp_path / "daily" / "2026-08-10" / "auto_fin.md").read_text(encoding="utf-8") - assert "[[daily/2026-08-01/auto_fin.md|历史黄金观察]]" in report - assert "](daily/2026-08-01/auto_fin.md)" not in report - assert "缺失文章" in report and "越界文章" in report - assert "missing.md" not in report and "outside.md" not in report - assert not (tmp_path / "daily" / "2026-08-10" / "auto_fin_news.md").exists() - assert not (tmp_path / "resource").exists() - assert context["changes"] == [{"change": "added", "path": "daily/2026-08-10/auto_fin.md"}] - assert response.metadata["source_paths"] == ["daily/2026-08-01/auto_fin.md"] + + day = tmp_path / "daily" / "2026-08-10" + assert sorted(path.name for path in day.glob("*.md")) == ["机器人.md", "黄金.md"] + note = (day / "黄金.md").read_text(encoding="utf-8") + assert "kind: auto-fin-topic" in note and "topic: 黄金" in note + assert "title: 黄金观察" in note + assert "[[daily/2026-08-01/auto_fin.md|历史黄金观察]]" in note + assert "](daily/2026-08-01/auto_fin.md)" not in note + assert "缺失文章" in note and "越界文章" in note + assert "missing.md" not in note and "outside.md" not in note + assert "不提供收益、目标价或买卖建议" in note + + assert [item.path for item in context["auto_fin_notes"]] == [ + "daily/2026-08-10/黄金.md", + "daily/2026-08-10/机器人.md", + ] + assert [item.title for item in context["auto_fin_notes"]] == ["黄金观察", "机器人观察"] + assert context["changes"] == [ + {"change": "added", "path": "daily/2026-08-10/黄金.md"}, + {"change": "added", "path": "daily/2026-08-10/机器人.md"}, + ] + assert response.metadata["selected_news_count"] == 26 + assert response.metadata["note_paths"] == [item.path for item in context["auto_fin_notes"]] -def test_hybrid_wikilink_normalization_is_conservative_and_failure_safe(tmp_path: Path, monkeypatch): - import reme_auto_fin.merge as merge_module +@pytest.mark.asyncio +async def test_research_names_the_note_after_the_topic_not_the_agent_title(tmp_path: Path): + """A whole-paragraph Agent title used to become a filename and fail with ENAMETOOLONG.""" - step = AutoFinMergeStep( - app_context=ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai"), + class _VerboseAgent(BaseAgentWrapper): + async def reply(self, _prompt, **_kwargs): + return { + "structured_output": AutoFinReportOutput( + title="机器人主题 9-18 研究:" + "量产与政策" * 60, + description="关注政策变化。", + body=LINKED_BODY, + ), + } + + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + context = _context( + auto_fin_news_by_topic={"黄金": [], "机器人": [_news("1", "2026-08-10T09:00:00+08:00")], "半导体": []}, ) + + await AutoFinResearchStep(app_context=app_context, agent_wrapper=_VerboseAgent(app_context=app_context))(context) + + day = tmp_path / "daily" / "2026-08-10" + assert [path.name for path in day.glob("*.md")] == ["机器人.md"] + assert "title: 机器人主题 9-18 研究:量产与政策" in (day / "机器人.md").read_text(encoding="utf-8") + + +@pytest.mark.asyncio +async def test_research_rerun_replaces_the_same_note_for_a_topic(tmp_path: Path): + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + agent = _ReportAgent(app_context=app_context) + news = {"黄金": [_news("1", "2026-08-10T09:00:00+08:00")], "机器人": [], "半导体": []} + step = AutoFinResearchStep(app_context=app_context, agent_wrapper=agent) + + first = _context(auto_fin_news_by_topic=news) + await step(first) + second = _context(auto_fin_news_by_topic=news) + await step(second) + + assert len(agent.calls) == 2 + assert "本次为当日首次生成。" in agent.calls[0][0] + assert "## 今日判断" in agent.calls[1][0] + assert second["changes"] == [{"change": "modified", "path": "daily/2026-08-10/黄金.md"}] + assert [path.name for path in (tmp_path / "daily" / "2026-08-10").glob("*.md")] == ["黄金.md"] + + +@pytest.mark.asyncio +async def test_research_skips_topics_without_news_and_honours_the_skip_flag(tmp_path: Path): + app_context = ApplicationContext(workspace_dir=str(tmp_path)) + agent = _ReportAgent(app_context=app_context) + context = _context(auto_fin_news_by_topic={"黄金": [], "机器人": [], "半导体": []}) + + response = await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(context) + + assert not agent.calls + assert context["auto_fin_notes"] == [] + assert response.metadata["note_paths"] == [] + assert not (tmp_path / "daily").exists() + + skipped = _context(auto_fin_skipped=True, auto_fin_news_by_topic={"黄金": []}) + await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(skipped) + assert len(agent.calls) == 0 + + +@pytest.mark.asyncio +async def test_research_keeps_the_other_topics_when_one_fails(tmp_path: Path): + """One broken topic must not discard the notes the other topics already produced.""" + + class _FlakyAgent(_ReportAgent): + async def reply(self, inputs, **kwargs): + if "当前主题:机器人" in str(inputs): + raise RuntimeError("agent exploded") + return await super().reply(inputs, **kwargs) + + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + agent = _FlakyAgent(app_context=app_context) + context = _context( + auto_fin_news_by_topic={ + "黄金": [_news("1", "2026-08-10T09:00:00+08:00")], + "机器人": [_news("2", "2026-08-10T09:05:00+08:00")], + "半导体": [], + }, + ) + + response = await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(context) + + assert [item.path for item in context["auto_fin_notes"]] == ["daily/2026-08-10/黄金.md"] + assert response.success is True + assert response.metadata["failed_topics"] == [{"topic": "机器人", "error": "agent exploded"}] + assert [path.name for path in (tmp_path / "daily" / "2026-08-10").glob("*.md")] == ["黄金.md"] + + +@pytest.mark.asyncio +async def test_research_ignores_notes_with_unparsable_frontmatter(tmp_path: Path): + """A note the user is editing by hand must not abort every topic in the run.""" + + day = tmp_path / "daily" / "2026-08-10" + day.mkdir(parents=True) + hand_edited = "---\ntags: [unfinished\n---\n\n正在编辑的笔记。\n" + (day / "手记.md").write_text(hand_edited, encoding="utf-8") + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + context = _context( + auto_fin_news_by_topic={"黄金": [_news("1", "2026-08-10T09:00:00+08:00")], "机器人": [], "半导体": []}, + ) + + response = await AutoFinResearchStep(app_context=app_context, agent_wrapper=_ReportAgent(app_context=app_context))( + context, + ) + + assert response.success is True + assert response.metadata["failed_topics"] == [] + assert [item.path for item in context["auto_fin_notes"]] == ["daily/2026-08-10/黄金.md"] + assert (day / "手记.md").read_text(encoding="utf-8") == hand_edited + + +def test_read_note_falls_back_when_the_frontmatter_is_unparsable(tmp_path: Path): + broken = tmp_path / "手记.md" + broken.write_text("---\ntags: [unfinished\n---\n\n正文\n", encoding="utf-8") + + assert read_note(broken) == FIRST_RUN_NOTICE + assert read_note(tmp_path / "missing.md") == FIRST_RUN_NOTICE + assert read_note(None) == FIRST_RUN_NOTICE + + +@pytest.mark.asyncio +async def test_research_fails_the_run_when_every_topic_fails(tmp_path: Path): + """A run that produced no note at all must fail instead of sending an empty brief.""" + + class _DeadAgent(BaseAgentWrapper): + async def reply(self, *_args, **_kwargs): + raise RuntimeError("agent exploded") + + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + context = _context( + auto_fin_news_by_topic={"黄金": [_news("1", "2026-08-10T09:00:00+08:00")], "机器人": [], "半导体": []}, + ) + + with pytest.raises(RuntimeError, match="failed for every topic"): + await AutoFinResearchStep(app_context=app_context, agent_wrapper=_DeadAgent(app_context=app_context))(context) + + +@pytest.mark.asyncio +async def test_digest_merges_notes_and_links_back_to_each_of_them(tmp_path: Path): + _history(tmp_path) + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + agent = _ReportAgent(app_context=app_context) + context = _context( + auto_fin_news_by_topic={"黄金": [_news("1", "2026-08-10T09:00:00+08:00")], "机器人": [], "半导体": []}, + ) + await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(context) + note_path = context["auto_fin_notes"][0].path + + response = await AutoFinDigestStep(app_context=app_context, agent_wrapper=agent)(context) + + prompt, kwargs = agent.calls[-1] + assert kwargs["job_tools"] == [] + assert '"topic": "黄金"' in prompt + assert "各主题笔记" in prompt + assert "当前主题:" not in prompt + assert "本次为当日首次生成。" in prompt + + digest_path = "daily/2026-08-10/主题新闻观察(2026-08-10).md" + digest = (tmp_path / digest_path).read_text(encoding="utf-8") + assert "kind: auto-fin-digest" in digest + assert "title: 主题新闻观察" in digest + assert "## 主题详解" in digest + assert f"- [[{note_path}]]" in digest + assert "[[daily/2026-08-01/auto_fin.md|历史黄金观察]]" in digest + assert "缺失文章" in digest and "missing.md" not in digest + assert digest.rstrip().endswith("不提供收益、目标价或买卖建议。") + + assert context["markdown_path"] == digest_path + assert context["changes"][-1] == {"change": "added", "path": digest_path} + assert response.metadata["digest_path"] == context["markdown_path"] + assert response.metadata["source_paths"] == ["daily/2026-08-01/auto_fin.md"] + assert response.metadata["note_paths"] == [note_path] + + +@pytest.mark.parametrize("topic", ["黄金", "AI[算力]", "C#", "新能源/储能", "#热点", "a|b"]) +@pytest.mark.asyncio +async def test_digest_trailer_resolves_to_each_topic_note(tmp_path: Path, topic: str): + """A topic name must not smuggle a wikilink delimiter into the trailer it lands in.""" + + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + agent = _ReportAgent(app_context=app_context) + context = _context( + auto_fin_topics=[topic], + auto_fin_news_by_topic={topic: [_news("1", "2026-08-10T09:00:00+08:00")]}, + ) + await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(context) + note_path = context["auto_fin_notes"][0].path + + response = await AutoFinDigestStep(app_context=app_context, agent_wrapper=agent)(context) + + digest = (tmp_path / response.metadata["digest_path"]).read_text(encoding="utf-8") + targets = [match.target for match in WikilinkHandler.iter_matches(digest)] + assert (tmp_path / note_path).is_file() + assert note_path in targets + # Every ``[[`` the digest emits opens a link the workspace parser can read. + assert digest.count("[[") == len(targets) + + +@pytest.mark.asyncio +async def test_digest_returns_the_validated_body_to_the_caller(tmp_path: Path): + """The API answer must not carry a link that the note itself downgraded to plain text.""" + + app_context = ApplicationContext(workspace_dir=str(tmp_path), timezone="Asia/Shanghai") + agent = _ReportAgent(app_context=app_context) + context = _context( + auto_fin_news_by_topic={"黄金": [_news("1", "2026-08-10T09:00:00+08:00")], "机器人": [], "半导体": []}, + ) + await AutoFinResearchStep(app_context=app_context, agent_wrapper=agent)(context) + + response = await AutoFinDigestStep(app_context=app_context, agent_wrapper=agent)(context) + + digest = (tmp_path / response.metadata["digest_path"]).read_text(encoding="utf-8") + assert "缺失文章" in response.answer + assert "missing.md" not in response.answer + assert response.answer in digest + + +@pytest.mark.asyncio +async def test_digest_skips_when_the_run_was_already_skipped(tmp_path: Path): + app_context = ApplicationContext(workspace_dir=str(tmp_path)) + agent = _ReportAgent(app_context=app_context) + context = _context(auto_fin_skipped=True) + + await AutoFinDigestStep(app_context=app_context, agent_wrapper=agent)(context) + + assert not agent.calls + assert not (tmp_path / "daily").exists() + + +def test_normalize_hybrid_wikilinks_is_conservative(): body = ( "[[digest/wiki/gold.md]](digest/wiki/gold.md) " "[[digest/wiki/gold.md|黄金]]() " "[[digest/wiki/gold.md#L2|黄金]](digest/wiki/gold.md) " "[[digest/wiki/gold.md]](digest/wiki/other.md)" ) - assert step._normalize_hybrid_wikilinks(body) == ( + + assert normalize_hybrid_wikilinks(body) == ( "[[digest/wiki/gold.md]] " "[[digest/wiki/gold.md|黄金]] " "[[digest/wiki/gold.md#L2|黄金]] " "[[digest/wiki/gold.md]](digest/wiki/other.md)" ) - class _BrokenPattern: - @staticmethod - def sub(_replace, _body): - raise RuntimeError("normalization failed") - monkeypatch.setattr(merge_module, "_HYBRID_WIKILINK_RE", _BrokenPattern()) - assert step._normalize_hybrid_wikilinks(body) == body +def test_normalize_title_sanitizes_agent_titles_and_keeps_notes_importable(): + assert normalize_title("# 黄金/政策:观察", "黄金观察") == "黄金-政策:观察" + assert normalize_title(" ", "黄金观察") == "黄金观察" + assert normalize_title("解读.md", "黄金观察") == "解读" + assert AutoFinNote(topic="黄金", title="标题", description="说明", body="正文", path="a.md").topic == "黄金" -def test_plugin_config_has_default_topics_and_no_intermediate_index_step(): +def test_normalize_title_keeps_wikilink_delimiters_out_of_the_stem(): + """`AI[算力]` used to produce a stem the wikilink parser could not read at all.""" + + assert normalize_title("AI[算力]", "主题观察") == "AI-算力" + assert normalize_title("C#", "主题观察") == "C" + assert normalize_title("# 热点", "主题观察") == "热点" + for raw in ("AI[算力]", "C#", "# 热点", "a|b"): + assert not set("[]#|") & set(normalize_title(raw, "主题观察")) + + +def test_normalize_title_fits_the_filename_component_byte_budget(): + """A whole-paragraph title used to reach os.stat and fail with ENAMETOOLONG.""" + assert normalize_title("长" * 200, "黄金观察") == "长" * 60 + assert len(normalize_title("long" * 100, "黄金观察").encode()) <= 180 + assert normalize_title(" / ", "黄金观察") == "黄金观察" + + +def test_plugin_config_has_default_topics_and_two_report_steps(): jobs = PLUGIN_MANIFEST["application_defaults"]["jobs"] job = jobs["auto_fin"] assert job["parameters"]["properties"]["topics"]["default"] == "黄金,机器人,半导体" @@ -259,15 +636,14 @@ def test_plugin_config_has_default_topics_and_no_intermediate_index_step(): assert job["parameters"]["properties"]["request_interval"]["default"] == 10 assert job["parameters"]["properties"]["max_retries"]["default"] == 3 assert "news_file" not in job["parameters"]["properties"] - assert [step["backend"] for step in job["steps"]] == [ - "auto_fin_data_step", - "auto_fin_topic_step", - "auto_fin_merge_step", - "auto_tag_step", + assert job["steps"] == [ + {"backend": "auto_fin_data_step"}, + {"backend": "auto_fin_topic_step"}, + {"backend": "auto_fin_research_step", "job_tools": ["search"]}, + {"backend": "auto_fin_digest_step"}, + {"backend": "auto_tag_step"}, ] - assert job["steps"][2]["job_tools"] == ["search", "read"] - assert job["steps"][3] == {"backend": "auto_tag_step"} - assert jobs["auto_fin_cron"]["cron"] == "0 18 * * *" + assert jobs["auto_fin_cron"]["cron"] == "0 9 * * *" assert jobs["auto_fin_cron"]["steps"] == job["steps"] assert ( not { @@ -279,11 +655,8 @@ def test_plugin_config_has_default_topics_and_no_intermediate_index_step(): ) -def test_agent_schemas_are_small_and_required(): - topic = AutoFinTopicOutput.model_json_schema() +def test_report_schema_is_small_and_required(): report = AutoFinReportOutput.model_json_schema() - assert topic["required"] == ["news_ids"] - assert set(topic["properties"]) == {"news_ids"} assert report["required"] == ["title", "description", "body"] assert set(report["properties"]) == {"title", "description", "body"} diff --git a/plugins/dingtalk/src/reme_dingtalk/send.py b/plugins/dingtalk/src/reme_dingtalk/send.py index a88870c7..e81635aa 100644 --- a/plugins/dingtalk/src/reme_dingtalk/send.py +++ b/plugins/dingtalk/src/reme_dingtalk/send.py @@ -12,11 +12,38 @@ from reme.steps.file_io._path import gate_md, resolve_path _GROUP_SEND_URL = "https://api.dingtalk.com/v1.0/robot/groupMessages/send" +_MAX_ERROR_DETAIL = 200 +_ERROR_BODY_KEYS = ("code", "message", "requestid") + def _conversation_ids(value: str) -> list[str]: return [item.strip() for item in value.split(",") if item.strip()] +def _failure_detail(exc: Exception) -> str: + """Summarize a delivery failure for logs and response metadata. + + DingTalk reports rejections (an oversized ``msgParam``, say) as an HTTP error whose + body carries a machine-readable ``code``; without it the operator only sees + ``HTTPStatusError`` and has to reproduce the request by hand. Only that whitelist of + keys is reported -- the request body holds the message content and credentials, so a + DingTalk error that echoes it back must not reach the log. + """ + if not isinstance(exc, httpx.HTTPStatusError): + return f"{type(exc).__name__}: {exc}"[:_MAX_ERROR_DETAIL] + response = exc.response + detail = f"HTTP {response.status_code}" + try: + body = response.json() + except ValueError: + body = None + if isinstance(body, dict): + for key in _ERROR_BODY_KEYS: + if body.get(key): + detail += f" {key}={body[key]}" + return detail[:_MAX_ERROR_DETAIL] + + class DingTalkMarkdownSendStep(BaseStep): """Send one Markdown document serially to configured DingTalk groups.""" @@ -105,10 +132,10 @@ class DingTalkMarkdownSendStep(BaseStep): if not isinstance(result, dict) or not result.get("processQueryKey"): raise ValueError("missing processQueryKey") except (httpx.HTTPError, ValueError) as exc: - failures.append(f"recipient {index}: {type(exc).__name__}") + detail = _failure_detail(exc) + failures.append(f"recipient {index}: {detail}") self.logger.warning( - f"[{self.name}] DingTalk delivery failed recipient={index}/{len(recipients)} " - f"error_type={type(exc).__name__}", + f"[{self.name}] DingTalk delivery failed recipient={index}/{len(recipients)} {detail}", ) continue self.context.response.metadata["dingtalk_sent_count"] += 1 @@ -117,6 +144,9 @@ class DingTalkMarkdownSendStep(BaseStep): sent_count = self.context.response.metadata["dingtalk_sent_count"] if failures: self.context.response.metadata["dingtalk_delivery_errors"] = failures - raise RuntimeError(f"DingTalk Markdown delivery failed for {len(failures)} of {len(recipients)} recipients") + raise RuntimeError( + f"DingTalk Markdown delivery failed for {len(failures)} of {len(recipients)} recipients: " + f"{'; '.join(failures)}", + ) self.logger.info(f"[{self.name}] DingTalk Markdown delivery complete sent={sent_count} total={len(recipients)}") return self.context.response diff --git a/plugins/dingtalk/tests/test_dingtalk.py b/plugins/dingtalk/tests/test_dingtalk.py index 8d7d0e3c..282ac25e 100644 --- a/plugins/dingtalk/tests/test_dingtalk.py +++ b/plugins/dingtalk/tests/test_dingtalk.py @@ -98,6 +98,57 @@ async def test_markdown_send_without_conversations_is_a_noop(tmp_path): } +@pytest.mark.asyncio +async def test_markdown_send_surfaces_dingtalk_rejection_detail(tmp_path, monkeypatch): + report = tmp_path / "daily" / "report.md" + report.parent.mkdir() + report.write_text(frontmatter.dumps(frontmatter.Post("# Report\n\nBody")), encoding="utf-8") + + async def handler(_request: httpx.Request) -> httpx.Response: + return httpx.Response( + 400, + json={ + "code": "InvalidmsgParam", + "requestid": "01A0B51A-96E1-735D-A1A9-F043A281EC5B", + "message": "Specified parameter msgParam is not valid.", + }, + ) + + dingtalk_stream = importlib.import_module("dingtalk_stream") + monkeypatch.setattr( + dingtalk_stream.DingTalkStreamClient, + "get_access_token", + lambda _client: "access-token", + ) + monkeypatch.setattr( + dingtalk_send.httpx, + "AsyncHTTPTransport", + lambda **kwargs: httpx.MockTransport(handler), + ) + step = DingTalkMarkdownSendStep( + app_context=ApplicationContext(workspace_dir=str(tmp_path)), + app_key="app-key", + app_secret="app-secret", + robot_code="robot-code", + conversation_ids="group-one", + ) + step.logger = MagicMock() + + detail = ( + "HTTP 400 code=InvalidmsgParam" + " message=Specified parameter msgParam is not valid." + " requestid=01A0B51A-96E1-735D-A1A9-F043A281EC5B" + ) + with pytest.raises(RuntimeError, match="code=InvalidmsgParam"): + await step(RuntimeContext(markdown_path="daily/report.md")) + + assert step.context.response.metadata["dingtalk_sent_count"] == 0 + assert step.context.response.metadata["dingtalk_delivery_errors"] == [f"recipient 1: {detail}"] + warnings = "\n".join(call.args[0] for call in step.logger.warning.call_args_list) + assert f"recipient=1/1 {detail}" in warnings + assert all(value not in warnings for value in ("app-key", "app-secret", "robot-code", "group-one")) + + class _AgentWrapper(BaseAgentWrapper): def __init__(self, **kwargs): super().__init__(**kwargs) diff --git a/reme/config/cookbook.yaml b/reme/config/cookbook.yaml index 65255785..201d66b0 100644 --- a/reme/config/cookbook.yaml +++ b/reme/config/cookbook.yaml @@ -86,7 +86,7 @@ components: agent_wrapper: claude_code: backend: claude_code - model: ${LLM_MODEL_NAME:-qwen3.8-max} + model: ${LLM_MODEL_NAME:-qwen3.7-plus} api_key: ${LLM_API_KEY:-} base_url: ${CLAUDE_CODE_BASE_URL:-https://dashscope.aliyuncs.com/apps/anthropic} permission_mode: bypassPermissions @@ -122,10 +122,13 @@ jobs: steps: &auto_fin_steps - backend: auto_fin_data_step - backend: auto_fin_topic_step - - backend: auto_fin_merge_step - job_tools: [search, read] + - backend: auto_fin_research_step + job_tools: [search] + - backend: auto_fin_digest_step - backend: auto_tag_step - <<: *dingtalk_send_defaults + input_mapping: + auto_fin_digest_path: markdown_path title: ReMe Auto Fin auto_fin_cron: diff --git a/reme/steps/file_io/frontmatter_read.py b/reme/steps/file_io/frontmatter_read.py index 0b127986..fea6312d 100644 --- a/reme/steps/file_io/frontmatter_read.py +++ b/reme/steps/file_io/frontmatter_read.py @@ -13,7 +13,7 @@ from pathlib import Path import frontmatter import yaml -from ._path import display_path, gate_md, resolve_path +from ._path import _check_path_permission, display_path, gate_md, resolve_path from ..base_step import BaseStep from ...components import R @@ -41,6 +41,14 @@ class FrontmatterReadStep(BaseStep): if target != original_target: resolved["resolved_path"] = display_path(workspace_dir, target) probed = resolved.get("resolved_path", path) + if not _check_path_permission(workspace_dir, target, self.context.get("_allowed_paths")): + self.context.response.success = False + self.context.response.answer = "Error: no permission to access this file" + self.context.response.metadata.update( + {"path": path, "error": "no permission to access this file", **resolved}, + ) + self.logger.info(f"[{self.name}] path={path} error=no_permission") + return if not target.is_file(): self.context.response.success = False self.context.response.answer = f"Error: {probed} not found" diff --git a/reme/steps/index/search.py b/reme/steps/index/search.py index a7bc42a9..b5aafd3b 100644 --- a/reme/steps/index/search.py +++ b/reme/steps/index/search.py @@ -225,6 +225,7 @@ class SearchStep(BaseStep): expand_links_enabled: bool = bool(self.kwargs.get("expand_links", True)) max_links_per_direction: int = int(self.kwargs.get("max_links_per_direction", 10)) tool_context_id: str = (self.context.get("tool_context_id", "") or "").strip() + max_search_calls = self.context.get("max_search_calls") strict_date_filter: bool = bool( self.context.get("strict_date_filter") or self.kwargs.get("strict_date_filter", False), ) @@ -235,6 +236,27 @@ class SearchStep(BaseStep): return self.context.response assert limit > 0, f"limit must be positive, got {limit}" + if max_search_calls is not None: + if not tool_context_id or self.app_context is None: + raise ValueError("max_search_calls requires an application and tool_context_id") + maximum = int(max_search_calls) + if maximum < 1: + raise ValueError("max_search_calls must be positive") + budgets = self.app_context.metadata.setdefault("__search_call_budgets", {}) + budget = budgets.setdefault(tool_context_id, {"count": 0, "lock": asyncio.Lock()}) + async with budget["lock"]: + if budget["count"] >= maximum: + self.logger.warning( + f"[{self.name}] search budget exhausted context={tool_context_id} limit={maximum}", + ) + self.context.response.success = False + self.context.response.answer = f"Error: search call limit of {maximum} reached" + return self.context.response + budget["count"] += 1 + self.logger.info( + f"[{self.name}] search budget context={tool_context_id} " f"call={budget['count']}/{maximum}", + ) + candidates = min(_MAX_CANDIDATES, max(1, int(limit * candidate_multiplier))) search_filter: dict = dict(self.context.get("search_filter", {}) or {}) raw_tags = self.context.get("tags", []) or [] diff --git a/tests/unit/test_cookbook_config.py b/tests/unit/test_cookbook_config.py index a70bb3cc..37b14b63 100644 --- a/tests/unit/test_cookbook_config.py +++ b/tests/unit/test_cookbook_config.py @@ -86,7 +86,7 @@ def test_cookbook_enables_embedding_and_separate_agent_backends(monkeypatch): assert components["agent_wrapper"]["default"]["backend"] == "agentscope" assert components["agent_wrapper"]["claude_code"] == { "backend": "claude_code", - "model": "qwen3.8-max", + "model": "qwen3.7-plus", "api_key": "llm-api-key", "base_url": "https://dashscope.aliyuncs.com/apps/anthropic", "permission_mode": "bypassPermissions", @@ -124,11 +124,15 @@ def test_cookbook_appends_dingtalk_to_business_pipelines(monkeypatch): assert [step["backend"] for step in auto_fin_steps] == [ "auto_fin_data_step", "auto_fin_topic_step", - "auto_fin_merge_step", + "auto_fin_research_step", + "auto_fin_digest_step", "auto_tag_step", "dingtalk_markdown_send_step", ] assert jobs["auto_fin_cron"]["steps"] == auto_fin_steps + assert auto_fin_steps[-1]["input_mapping"] == { + "auto_fin_digest_path": "markdown_path", + } assert auto_fin_steps[-1]["title"] == "ReMe Auto Fin" assert [step["backend"] for step in daily_paper_steps] == [ @@ -196,7 +200,7 @@ def test_cookbook_overrides_merge_with_pure_plugin_defaults(monkeypatch): assert config.jobs["auto_fin"].backend == "base" assert config.jobs["auto_fin"].parameters["properties"]["topics"]["default"] == "黄金,机器人,半导体" assert config.jobs["auto_fin_cron"].backend == "cron" - assert config.jobs["auto_fin_cron"].model_extra["cron"] == "0 18 * * *" + assert config.jobs["auto_fin_cron"].model_extra["cron"] == "0 9 * * *" assert config.jobs["daily_paper"].backend == "base" assert config.jobs["daily_paper_cron"].backend == "cron" assert config.jobs["daily_paper_cron"].model_extra["cron"] == "0 8 * * *" diff --git a/tests/unit/test_frontmatter_steps.py b/tests/unit/test_frontmatter_steps.py index 51a49f34..f95504d8 100644 --- a/tests/unit/test_frontmatter_steps.py +++ b/tests/unit/test_frontmatter_steps.py @@ -76,6 +76,23 @@ async def test_read_no_suffix_autoappends_md(): await store.close() +@pytest.mark.asyncio +async def test_read_honors_injected_path_scope(): + """An injected allowed-path scope denies frontmatter reads outside of it.""" + with tempfile.TemporaryDirectory() as tmp, temp_chdir(tmp): + _seed(Path(tmp), NOTE, BODY) + store = await _make_store() + + allowed = await _run(FrontmatterReadStep, store, path=NOTE, _allowed_paths=[NOTE]) + assert allowed.success is True + assert allowed.metadata["frontmatter"] == {"name": "n", "tags": ["a"]} + + denied = await _run(FrontmatterReadStep, store, path=NOTE, _allowed_paths=["notes/other.md"]) + assert denied.success is False + assert denied.metadata["error"] == "no permission to access this file" + await store.close() + + @pytest.mark.asyncio async def test_update_no_suffix_autoappends_md(): """frontmatter_update on a suffix-less path resolves to the ``.md`` file.""" diff --git a/tests/unit/test_search_step.py b/tests/unit/test_search_step.py index bf9165b0..5e218e4d 100644 --- a/tests/unit/test_search_step.py +++ b/tests/unit/test_search_step.py @@ -93,6 +93,32 @@ class FakeSearchStore(BaseFileStore): return self.keyword_results[:limit] +def test_search_call_budget_is_scoped_to_tool_context(tmp_path): + """An injected budget rejects a fourth call before touching the store.""" + store = FakeSearchStore() + app_context = ApplicationContext(workspace_dir=str(tmp_path)) + + async def run(): + for _ in range(3): + response = await SearchStep(app_context=app_context, file_store=store, expand_links=False)( + RuntimeContext(query="gold", tool_context_id="topic-a", max_search_calls=3), + ) + assert response.success + calls = len(store.calls) + rejected = await SearchStep(app_context=app_context, file_store=store, expand_links=False)( + RuntimeContext(query="gold", tool_context_id="topic-a", max_search_calls=3), + ) + assert not rejected.success + assert "limit of 3" in rejected.answer + assert len(store.calls) == calls + other = await SearchStep(app_context=app_context, file_store=store, expand_links=False)( + RuntimeContext(query="gold", tool_context_id="topic-b", max_search_calls=3), + ) + assert other.success + + asyncio.run(run()) + + class TaggedFakeSearchStore(FakeSearchStore): """Fake store with a tag index and ordinary file-store filtering."""