ReMe/reme/steps/evolve/dream/extract.py
jinliyl 215c1f72f2
Some checks are pending
NPM Format / Website checks (push) Waiting to run
Pre-commit / run (ubuntu-latest) (push) Waiting to run
Tests ReMe / Unit Tests - py3.13 (push) Waiting to run
Tests ReMe / Unit Tests - py3.11 (push) Waiting to run
Tests ReMe / Unit Tests - py3.12 (push) Waiting to run
Windows Smoke / CLI smoke - py3.11 (push) Waiting to run
feat: refine local-first research and memory workflows (#444)
* feat: refine local-first research workflows

* fix: delegate structured output tool choice

* refactor(auto-fin): fetch and filter rolling CLS news

* fix(auto-fin): keep imports portable across platforms

* feat(auto-fin): expose CLS fetch controls

* fix(auto-fin): propagate configurable news window

* feat(auto_fin): normalize hybrid wikilinks in report body

- Add _normalize_hybrid_wikilinks method to remove redundant Markdown destinations
- Use regex to identify hybrid wikilinks with optional destinations
- Replace redundant destinations with simpler wikilink format for clarity
- Ensure normalization is failure-safe with exception handling and logging
- Update report body normalization process to apply hybrid wikilink fix
- Add unit tests to verify correct normalization and failure safety behavior

* fix(dream): serialize integration with application-wide asyncio lock

- Add application-wide asyncio.Lock to serialize digest writes during integration
- Update _snapshot_digest to capture metadata per bucket
- Validate bucket association when recovering from file changes
- Add tests ensuring recovery only from the correct bucket
- Add tests confirming integration lock is shared across application context
- Enhance strict topic YAML loading validation in dream utils
- Add tests for strict topic loading rejecting invalid or lossy fields

* fix(cookbook): enable configurable job_tools for digest and merge steps

- Update daily_cookbook.yaml to add job_tools: [memory_search, read] in digest steps
- Modify DailyPaperDigestStep to read job_tools from kwargs instead of fixed list
- Modify AutoFinMergeStep to similarly read job_tools from kwargs
- Update tests to pass job_tools explicitly when invoking these steps
- Remove hardcoded _TOOLS constants and replace with dynamic job_tools handling

* fix: retry incomplete dream receipts

* perf(pdf): increase max PDF pages limit from 20 to 35

- Updated configuration max_pdf_pages from 20 to 35 in daily_cookbook.yaml
- Modified code to extract up to 35 pages instead of 20 in analyze.py
- Updated README and README_ZH to document the increased max_pdf_pages
- Adjusted unit test assertions to reflect new max_pdf_pages limit of 35

* fix memory integration and daily paper links

* docs clarify cookbook tool usage
2026-08-11 23:32:34 +08:00

227 lines
10 KiB
Python

"""Global dream extract step."""
import json
from ...base_step import BaseStep
from ...file_io import refresh_day_index
from .._evolve import agent_reply_result_text
from ....components import R
from ....enumeration import DreamBucketEnum
from ....schema import DreamState
from .utils import (
clean_paths,
daily_dir,
llm_available,
pack_paths,
parse_structured_reply,
recent_dates,
scan_day_files,
store_state,
today,
workspace_dir,
)
_TOOLS = ("read",)
@R.register("dream_extract_step")
class DreamExtractStep(BaseStep):
"""Scan changed daily files and globally extract merged units/topics."""
def __init__(self, topic_session_id: str = "interests", scan_days: int = 2, max_units: int = 5, **kwargs):
super().__init__(**kwargs)
self.topic_session_id = topic_session_id
self.scan_days = scan_days
self.max_units = max_units
async def execute(self):
assert self.context is not None
day = today(self, str(self.context.get("date", "") or ""))
raw_scan_days = self.context.get("scan_days", self.scan_days)
scan_days = max(int(raw_scan_days or self.scan_days), 1)
raw_max_units = self.context.get("max_units", self.max_units)
max_units = max(int(raw_max_units or self.max_units), 0)
dates = recent_dates(day, scan_days)
hint = str(self.context.get("hint", "") or "").strip()
daily, workspace = daily_dir(self), workspace_dir(self)
if self.file_catalog is None:
raise RuntimeError("dream_extract_step requires file_catalog")
self.logger.info(
f"[{self.name}] start date={day} dates={','.join(dates)} scan_days={scan_days} "
f"max_units={max_units} hint={bool(hint)}",
)
for scan_day in dates:
self.logger.info(f"[{self.name}] refresh index start date={scan_day} daily_dir={daily}")
await refresh_day_index(self.file_store, scan_day, daily)
self.logger.info(f"[{self.name}] refresh index done date={scan_day}")
existing = self._existing(
workspace,
[
path
for scan_day in dates
for path in scan_day_files(workspace, scan_day, daily, f"{self.topic_session_id}.yaml")
],
)
interest_rels = {f"{daily}/{scan_day}/{self.topic_session_id}.yaml" for scan_day in dates}
day_mds = {f"{daily}/{scan_day}.md" for scan_day in dates}
day_prefixes = tuple(f"{daily}/{scan_day}/" for scan_day in dates)
nodes = await self.file_catalog.get_nodes()
indexed_all = {n.path: n.st_mtime for n in nodes if n.path in day_mds or n.path.startswith(day_prefixes)}
indexed = {path: mt for path, mt in indexed_all.items() if path not in interest_rels}
changed = [rel for rel, mt in existing.items() if indexed.get(rel) != mt]
unchanged = [rel for rel, mt in existing.items() if indexed.get(rel) == mt]
protected = set(existing) | {rel for rel in interest_rels if (workspace / rel).is_file()}
deleted = sorted(indexed_all.keys() - protected)
self.logger.info(
f"[{self.name}] scan summary existing={len(existing)} indexed={len(indexed)} "
f"changed={len(changed)} unchanged={len(unchanged)} deleted={len(deleted)}",
)
if deleted:
self.logger.info(f"[{self.name}] catalog delete start paths={len(deleted)}")
await self.file_catalog.delete(deleted)
self.logger.info(f"[{self.name}] catalog delete done paths={len(deleted)}")
state = DreamState(
date=day,
dates=dates,
scan_days=scan_days,
hint=hint,
daily_dir=daily,
workspace=str(workspace),
files_scanned=len(existing),
files_unchanged=len(unchanged),
files_changed=len(changed),
files_deleted=len(deleted),
changed_paths=changed,
unchanged_paths=unchanged,
deleted_paths=deleted,
existing=existing,
indexed=indexed,
)
if not changed:
self.logger.info(f"[{self.name}] skip no changed input dates={','.join(dates)}")
return self._finish(state, True, f"No changed dream input for {', '.join(dates)}")
if not llm_available(self):
state.errors.append("no llm configured; dream extract requires an LLM")
state.failed_paths = list(changed)
self.logger.warning(f"[{self.name}] skip no llm changed={len(changed)}")
return self._finish(state, False, state.errors[-1])
self.logger.info(f"[{self.name}] agent start changed={len(changed)} dates={len(dates)}")
raw_result, meta = "", {}
for attempt in range(2):
try:
result = await self.agent_wrapper.reply(
self.prompt_format(
"extract_user_message",
date=day,
dates_json=json.dumps(dates, ensure_ascii=False, indent=2),
hint=hint or "(none)",
max_units=max_units,
changed_paths_json=json.dumps(changed, ensure_ascii=False, indent=2),
material_blob=pack_paths(workspace, changed),
),
system_prompt=self.prompt_format(
"extract_system_prompt",
workspace_dir=str(workspace),
buckets=", ".join(bucket.value for bucket in DreamBucketEnum),
max_units=max_units,
),
job_tools=list(_TOOLS),
)
self.logger.info(f"[{self.name}] agent done has_result={bool(result.get('result'))}")
raw_result = agent_reply_result_text(result)
meta = parse_structured_reply(raw_result)
except Exception as e: # noqa: BLE001
if attempt == 0:
self.logger.warning(
f"[{self.name}] extract attempt 1 returned no usable receipt; retrying once: "
f"{type(e).__name__}: {e}",
)
continue
error = f"dream extract agent failed after retry: {type(e).__name__}: {e}"
state.errors.append(error)
state.failed_paths = list(changed)
self.logger.error(f"[{self.name}] {error}")
return self._finish(state, False, error)
units = meta.get("units") if "units" in meta else meta.get("memory_units")
if isinstance(units, list) and isinstance(meta.get("topics"), list):
break
if attempt == 0:
self.logger.warning(f"[{self.name}] extract attempt 1 returned an unusable receipt; retrying once")
continue
# Keep the warning-only result checkpointable after one retry so a bad source cannot loop forever.
warning = "dream extract skipped unusable agent receipt after retry; expected units and topics lists"
state.warnings.append(warning)
self.logger.warning(f"[{self.name}] {warning}")
self.logger.info(f"[{self.name}] parse done keys={','.join(sorted(meta.keys())) if meta else '(none)'}")
self.clean_output(state, meta, max_units=max_units)
state.extract_summary = raw_result
answer = f"Extracted {len(state.units)} unit(s), {len(state.topics)} topic(s)"
answer = f"{answer} from {len(changed)} changed file(s) across {len(dates)} day(s)"
return self._finish(state, True, answer)
def _existing(self, workspace, files: list[str]) -> dict[str, float]:
out: dict[str, float] = {}
for rel in files:
try:
out[rel] = (workspace / rel).stat().st_mtime
except OSError as e:
self.logger.error(f"[{self.name}] stat failed on {rel}: {e}")
return out
def clean_output(self, state: DreamState, meta: dict, max_units: int | None = None) -> None:
"""Clean up output"""
allowed = set(state.changed_paths)
for raw in meta.get("units") or meta.get("memory_units") or []:
if max_units is not None and len(state.units) >= max_units:
break
if not isinstance(raw, dict):
continue
name = str(raw.get("name") or "").strip()
summary = str(raw.get("summary") or "").strip()
raw_bucket = str(raw.get("bucket") or "").strip()
paths = clean_paths(raw.get("paths"), allowed)
if not name or not summary or not paths:
continue
try:
bucket = DreamBucketEnum(raw_bucket).value
except ValueError:
self.logger.warning(f"[{self.name}] unit {name!r} emitted bucket {raw_bucket!r}; routing to wiki")
bucket = DreamBucketEnum.WIKI.value
state.units.append({"name": name, "bucket": bucket, "summary": summary, "paths": paths})
for raw in meta.get("topics") or []:
topic = self._clean_topic(raw, allowed)
if topic:
state.topics.append(topic)
@staticmethod
def _clean_topic(raw, allowed: set[str]) -> dict:
if not isinstance(raw, dict):
return {}
title = str(raw.get("title") or "").strip()
reason = str(raw.get("reason") or "").strip()
paths = clean_paths(raw.get("paths"), allowed)
if not title or not reason or not paths:
return {}
keywords = raw.get("keywords") or []
cleaned_keywords = [str(k).strip() for k in keywords if str(k).strip()] if isinstance(keywords, list) else []
return {
"title": title,
"reason": reason,
"evidence": str(raw.get("evidence") or "").strip(),
"keywords": cleaned_keywords,
"paths": paths,
}
def _finish(self, state: DreamState, success: bool, answer: str):
assert self.context is not None
state.summary = answer
store_state(self, state)
self.context.response.success = success
self.context.response.answer = answer
self.logger.info(f"[{self.name}] finish success={success} answer={answer!r}")
return self.context.response