refactor(auto-fin): one note per topic, merged by a separate digest step (#558)

* fix(auto-fin): parse topic IDs from fenced JSON replies

* fix(auto-fin): limit report agent tool calls in prompt

* refactor(auto-fin): research news by topic before market open

* fix(config): update default model version for claude_code backend

- Change model version from qwen3.8-max to qwen3.7-plus
- Use environment variable LLM_MODEL_NAME to allow override
- Ensure backend configuration reflects updated model setting

* refactor(auto-fin): write one note per topic before the daily digest

The merge step did two jobs at once: it researched every topic and
combined the results into a single report. Split it the way daily-paper
separates analysis from its brief, so each topic earns a durable note of
its own.

- auto_fin_research_step writes one note per topic that had relevant
  news, tagged `kind: auto-fin-topic` and `topic` in frontmatter
- auto_fin_digest_step merges those notes into the day's brief with no
  tools of its own and appends a `## 主题详解` section linking back to
  each note
- a same-day rerun finds a topic's note by its `topic` frontmatter and
  replaces it in place, deleting the old file when the title changed
- base.py now owns the shared Markdown layer: title sanitizing, report
  normalizing, wikilink validation, note lookup, atomic frontmatter
  writes, and change tracking, so both steps share one write path
- the DingTalk step maps `auto_fin_digest_path` to `markdown_path`
  explicitly instead of relying on whichever step ran last
- drop the unused AutoFinTopicOutput schema and read `job_tools` from
  the step config rather than hardcoding `search`

Co-Authored-By: Claude <noreply@anthropic.com>

* fix(dingtalk): surface rejection details when delivery fails

A failed group send only reported HTTPStatusError, so an operator had to
reproduce the request by hand to learn why DingTalk refused it. Include
the status code and the whitelisted error keys from the response body in
both the log line and the raised RuntimeError.

Only `code`, `message`, and `requestid` are reported: the request body
carries the message content and credentials, so an error response that
echoes it back must not reach the log. Detail is truncated to 200 chars.

Co-Authored-By: Claude <noreply@anthropic.com>

* fix(test): make the suite green on CI

- Point the cookbook claude_code model assertion at qwen3.7-plus, the
  default commit 9404e600 set, so the pre-existing red stops blocking
- Satisfy pylint on the auto-fin tests: prefer implicit booleaness for
  the recorded Agent calls and drop an unused tmp_path fixture

Co-Authored-By: Claude <noreply@anthropic.com>

* refactor(auto-fin): name notes after the topic, not the Agent title

The research Agent returned a whole paragraph as its title; that became a
filename and blew past the filesystem's 255-byte name limit, failing with
ENAMETOOLONG inside resolve_note_path. Topics are configured values, so
they are short and predictable - use them for file names and keep the
Agent title in frontmatter.

- Name topic notes after the topic and the digest after the run date
- Fold a byte budget into normalize_title as a safety net for long topics
- Take an AutoFinReportOutput in _write_report instead of loose fields
- Ask both prompts for a short title now that it is display-only

Co-Authored-By: Claude <noreply@anthropic.com>

* fix: isolate per-topic research failures and scope frontmatter reads

A single failing topic used to fail the whole cron job and discard the news
already gathered for the topics that had not run yet -- the 09-19 09:24 run
lost its robot notes that way. Research now logs the failure, continues with
the remaining topics, and only fails the run when no topic produced a note.

frontmatter_read was the only frontmatter step without the _allowed_paths
check that read, write, edit, and frontmatter_update already honour, so an
Agent scoped to one file could still read another file's metadata.

- Isolate per-topic research failures and report them as failed_topics
- Fail loudly when every topic fails so an empty brief is never sent
- Apply _check_path_permission in FrontmatterReadStep

Co-Authored-By: Claude <noreply@anthropic.com>

* fix(auto-fin): keep hand-edited notes and wikilink delimiters from breaking the run

Addresses three review findings on the topic-per-note rework.

- `find_note` and `read_note` now skip a note whose YAML frontmatter does
  not parse. A hand-edited note in the day directory raised
  `yaml.parser.ParserError`, which per-topic isolation surfaced as
  "Auto Fin research failed for every topic" and took the run down with it.
- `normalize_title` also strips `[`, `]` and `#`, which `WikilinkHandler`
  treats as target delimiters. `AI[算力]` used to emit a trailer link the
  parser could not read at all, and `C#` resolved to `.../C` plus an anchor.
- `_write_report` returns the body it actually wrote, and both callers
  propagate it, so the digest answer and the note handed to the digest Agent
  no longer carry links that validation had already downgraded on disk.

Co-Authored-By: Claude <noreply@anthropic.com>

---------

Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
jinliyl 2026-09-20 14:52:56 +08:00 • committed by GitHub
parent 5231f3970c
commit cae613d1c4
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
24 changed files with 1315 additions and 381 deletions

View file

@ -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/<topic>.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

View file

@ -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 并继续其余主题;没有相关当前新闻则成功跳过。
## 验证

View file

@ -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",
]

View file

@ -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<wikilink>\[\[(?P<inner>[^\[\]\n]+)\]\])\((?P<destination>[^()\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 `- [[<path>]]` 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

View file

@ -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 {

View file

@ -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

View file

@ -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。

View file

@ -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<wikilink>\[\[(?P<inner>[^\[\]\n]+)\]\])\((?P<destination>[^()\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

View file

@ -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。

View file

@ -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

View file

@ -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

View file

@ -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。

View file

@ -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

View file

@ -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

View file

@ -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}
```

View file

@ -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("<p>甲&amp;乙</p><style>隐藏</style><p>丙</p>") == "甲&乙 丙"
def test_plain_text_drops_hidden_and_unescapes_entities():
assert plain_text("<p>甲&amp;乙</p><style>隐藏</style><p>丙</p>") == "甲&乙 丙"
@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>) "
"[[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"}

View file

@ -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

View file

@ -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)

View file

@ -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:

View file

@ -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"

View file

@ -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 []

View file

@ -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 * * *"

View file

@ -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."""

View file

@ -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."""