From dd0a48f846e80fd4c1447db60deae32aae9a2cdc Mon Sep 17 00:00:00 2001 From: huangsen Date: Mon, 1 Jun 2026 17:00:30 +0800 Subject: [PATCH] refactor(dreamer): improve code formatting and line breaks --- reme4/steps/__init__.py | 5 +- .../{dream/dreamer.py => auto_dream.py} | 163 ++++++++++++++++-- .../{dream/dreamer.yaml => auto_dream.yaml} | 0 reme4/steps/evolve/dream/__init__.py | 17 -- reme4/steps/evolve/dream/cron_dreamer.py | 157 ----------------- 5 files changed, 149 insertions(+), 193 deletions(-) rename reme4/steps/evolve/{dream/dreamer.py => auto_dream.py} (79%) rename reme4/steps/evolve/{dream/dreamer.yaml => auto_dream.yaml} (100%) delete mode 100644 reme4/steps/evolve/dream/__init__.py delete mode 100644 reme4/steps/evolve/dream/cron_dreamer.py diff --git a/reme4/steps/__init__.py b/reme4/steps/__init__.py index 428bd470..82228673 100644 --- a/reme4/steps/__init__.py +++ b/reme4/steps/__init__.py @@ -8,8 +8,7 @@ from .common.llm_demo import LLMDemoStep from .common.stream_demo import StreamDemoStep1, StreamDemoStep2 from .common.version import VersionStep from .evolve.auto_memory import AutoMemoryStep -from .evolve.dream.cron_dreamer import CronDreamer -from .evolve.dream.dreamer import Dreamer +from .evolve.auto_dream import CronDreamer, Dreamer from .file_io.daily_create import DailyCreateStep from .file_io.daily_list import DailyListStep from .file_io.daily_reindex import DailyReindexStep @@ -73,7 +72,7 @@ __all__ = [ "UpdateCatalogStep", "UpdateIndexStep", "WatchChangesStep", - # evolve.dream + # evolve (dream) "CronDreamer", "Dreamer", # transfer diff --git a/reme4/steps/evolve/dream/dreamer.py b/reme4/steps/evolve/auto_dream.py similarity index 79% rename from reme4/steps/evolve/dream/dreamer.py rename to reme4/steps/evolve/auto_dream.py index c3e37566..abbd949d 100644 --- a/reme4/steps/evolve/dream/dreamer.py +++ b/reme4/steps/evolve/auto_dream.py @@ -56,9 +56,9 @@ from agentscope.message import Msg, TextBlock from agentscope.tool import Toolkit, ToolResponse from pydantic import BaseModel, Field -from .._evolve import FlexReActAgent -from ...base_step import BaseStep -from ....components import R +from ._evolve import FlexReActAgent +from ..base_step import BaseStep +from ...components import R # Hard-coded bucket vocabulary. Phase 1 classifies each sub-unit into @@ -178,10 +178,7 @@ class IntegrateOutcome(BaseModel): ), ) target_path: str = Field( - description=( - "The digest path you wrote to — must match what your `write` / " - "`edit` call(s) targeted." - ), + description=("The digest path you wrote to — must match what your `write` / " "`edit` call(s) targeted."), ) note: str = Field( default="", @@ -509,10 +506,8 @@ class Dreamer(BaseStep): skipped=True, ) - self.logger.info( - f"[{self.name}] integrate phase: {len(self._units)} sub-unit(s): " - + ", ".join(f"{u['name']}/{u['bucket']}" for u in self._units), - ) + unit_handles = ", ".join(f"{u['name']}/{u['bucket']}" for u in self._units) + self.logger.info(f"[{self.name}] integrate phase: {len(self._units)} sub-unit(s): {unit_handles}") # Phase 2 — integrate, one fresh ReAct per sub-unit, dispatched to # the bucket-specific system prompt. Python-level loop, not agent @@ -534,13 +529,11 @@ class Dreamer(BaseStep): continue per_unit_lines.append(_render_outcome_line(name, bucket, outcome)) + per_unit_block = "\n".join(per_unit_lines) summary = ( - f"Declared {len(self._units)} sub-unit(s) " - + "(" - + ", ".join(f"{u['name']}/{u['bucket']}" for u in self._units) - + "); " + f"Declared {len(self._units)} sub-unit(s) ({unit_handles}); " f"created {len(self._created)}, updated {len(self._updated)}.\n" - + "\n".join(per_unit_lines) + f"{per_unit_block}" ) return DreamResult( @@ -573,3 +566,141 @@ class Dreamer(BaseStep): self.context.response.success = True self.context.response.answer = result.summary self.context.response.metadata.update(result.model_dump()) + + +# ============================================================ +# CronDreamer — daily-tick wrapper around Dreamer. +# +# Inherits the per-file pipeline from Dreamer.dream_one and adds the +# outer loop over today's daily/ + resource/ files. Cron scheduling +# itself is out of scope; this step is just the unit of work. +# +# Inputs (RuntimeContext): +# date (str, optional): YYYY-MM-DD to scan. Defaults to today +# in the dreamer's timezone. +# hint (str, optional): passed through to each per-file dream. +# ============================================================ + + +class CronDreamResult(BaseModel): + """Aggregated outcome of one cron tick.""" + + date: str = "" + files_scanned: int = 0 + files_dreamed: int = 0 + files_skipped: int = 0 + files_failed: int = 0 + per_file: list[DreamResult] = Field(default_factory=list) + summary: str = "" + + +@R.register("cron_dreamer_step") +class CronDreamer(Dreamer): + """Loop ``daily//`` + ``resource//`` and dream each file. + + Inherits :data:`auto_dream.yaml` from :class:`Dreamer` (no separate + yaml — there's no extra prompt for the outer loop). + """ + + async def execute(self): + assert self.context is not None + date_input: str = (self.context.get("date", "") or "").strip() + hint: str = (self.context.get("hint", "") or "").strip() + + # daily_dir / resource_dir come from app config — NOT tool params. + # Same convention as daily_create / daily_list / daily_reindex. + # resource_dir may be empty (default) — that just skips the resource scan. + cfg = self.app_context.app_config if self.app_context is not None else None + daily_dir = (cfg.daily_dir if cfg else "") or "daily" + resource_dir = cfg.resource_dir if cfg else "" + + today = date_input or self._now().strftime("%Y-%m-%d") + vault = self._vault_dir() + files = _scan_today_files(vault, today, daily_dir, resource_dir) + + result = CronDreamResult(date=today, files_scanned=len(files)) + self.logger.info( + f"[{self.name}] cron tick date={today} scanned={len(files)} file(s) under " + f"{daily_dir}/{today}/ + {resource_dir}/{today}/", + ) + + for rel_path in files: + try: + dr = await self.dream_one(rel_path, hint) + except Exception as e: # pylint: disable=broad-except + self.logger.error( + f"[{self.name}] dream_one failed on {rel_path}: {type(e).__name__}: {e}", + ) + dr = DreamResult( + path=rel_path, + error=f"{type(e).__name__}: {e}", + ) + result.per_file.append(dr) + if dr.error: + result.files_failed += 1 + elif dr.skipped: + result.files_skipped += 1 + else: + result.files_dreamed += 1 + + result.summary = _render_cron_summary(result) + self.context.response.success = result.files_failed == 0 + self.context.response.answer = result.summary + self.context.response.metadata.update(result.model_dump()) + + +def _scan_today_files( + vault: Path, + today: str, + daily_dir: str, + resource_dir: str, +) -> list[str]: + """Return vault-relative paths of today's daily notes + resource files. + + * ``/.md`` — the day-index file (auto-rebuilt + rollup of all of today's notes). Included first so its day-level + abstractions land before the per-event details. + * ``//**/*.md`` — event notes for the day, + sorted by path. + * ``//**/*`` — any file type ingested under + today's resource folder. Skipped when ``resource_dir`` is empty. + + Results are sorted for deterministic processing order within each + group; the day-index file leads. + """ + out: list[str] = [] + + if daily_dir: + day_index = vault / daily_dir / f"{today}.md" + if day_index.is_file(): + out.append(str(day_index.relative_to(vault))) + daily_root = vault / daily_dir / today + if daily_root.is_dir(): + for md in sorted(daily_root.rglob("*.md")): + if md.is_file(): + out.append(str(md.relative_to(vault))) + + if resource_dir: + resource_root = vault / resource_dir / today + if resource_root.is_dir(): + for f in sorted(p for p in resource_root.rglob("*") if p.is_file()): + out.append(str(f.relative_to(vault))) + + return out + + +def _render_cron_summary(r: CronDreamResult) -> str: + """One-line header + one line per file with its outcome.""" + lines = [ + f"[CronDreamer] date={r.date} scanned={r.files_scanned} " + f"dreamed={r.files_dreamed} skipped={r.files_skipped} failed={r.files_failed}", + ] + for dr in r.per_file: + if dr.error: + status = f"ERROR ({dr.error})" + elif dr.skipped: + status = "SKIP" + else: + status = f"OK (+{len(dr.nodes_created)} created, ~{len(dr.nodes_updated)} updated)" + lines.append(f" - {dr.path}: {status}") + return "\n".join(lines) diff --git a/reme4/steps/evolve/dream/dreamer.yaml b/reme4/steps/evolve/auto_dream.yaml similarity index 100% rename from reme4/steps/evolve/dream/dreamer.yaml rename to reme4/steps/evolve/auto_dream.yaml diff --git a/reme4/steps/evolve/dream/__init__.py b/reme4/steps/evolve/dream/__init__.py deleted file mode 100644 index 2691474e..00000000 --- a/reme4/steps/evolve/dream/__init__.py +++ /dev/null @@ -1,17 +0,0 @@ -"""dream — auto-dream pipeline: extract abstractions, then integrate per -sub-unit using bucket-specific Phase 2 prompts. - -Two steps: - - dreamer — 2-phase ReAct workflow (extract memory sub-units - tagged with bucket, then integrate per sub-unit - via the canonical write/edit tools). - cron_dreamer — daily wrapper around dreamer; scans today's - daily/ + resource/ files and runs dream_one on each. - -Phase 2 uses the canonical ``write`` / ``edit`` jobs directly — no -constrained variants. Bucket placement is prompt-level discipline. -""" - -from . import cron_dreamer # noqa: F401 -- @R.register("cron_dreamer_step") -from . import dreamer # noqa: F401 -- @R.register("dreamer_step") diff --git a/reme4/steps/evolve/dream/cron_dreamer.py b/reme4/steps/evolve/dream/cron_dreamer.py deleted file mode 100644 index 6f9ef3ec..00000000 --- a/reme4/steps/evolve/dream/cron_dreamer.py +++ /dev/null @@ -1,157 +0,0 @@ -"""``cron_dreamer_step`` — daily-tick wrapper around :class:`Dreamer`. - -Scans today's materials and runs the per-file dream pipeline on each one: - -* ``/.md`` — the day-index rollup (processed first). -* ``//**/*.md`` — per-event notes for the day. -* ``//**/*`` — resources ingested today. - -The per-file logic is inherited from :class:`Dreamer` (via -:meth:`Dreamer.dream_one`); this step just adds the outer loop. - -Cron scheduling itself is out of scope here — this step is the unit of -work executed when a cron fires (or when the operator manually invokes -the ``dream_today`` job). External schedulers (background-job watchers, -cron daemons, etc.) drive when it runs. - -The ``daily_dir`` and ``resource_dir`` subroots are NOT tool params — -they come from ``app_config.daily_dir`` / ``app_config.resource_dir`` -(same convention as the ``daily_*`` steps). - -Inputs (RuntimeContext): - date (str, optional): YYYY-MM-DD to scan. Defaults to today - in the dreamer's timezone. - hint (str, optional): passed through to each per-file dream. - -Output (Response.metadata): :class:`CronDreamResult` JSON. -""" - -from pathlib import Path - -from pydantic import BaseModel, Field - -from .dreamer import Dreamer, DreamResult -from ....components import R - - -class CronDreamResult(BaseModel): - """Aggregated outcome of one cron tick.""" - - date: str = "" - files_scanned: int = 0 - files_dreamed: int = 0 - files_skipped: int = 0 - files_failed: int = 0 - per_file: list[DreamResult] = Field(default_factory=list) - summary: str = "" - - -@R.register("cron_dreamer_step") -class CronDreamer(Dreamer): - """Loop ``daily//`` + ``resource//`` and dream each file.""" - - async def execute(self): - assert self.context is not None - date_input: str = (self.context.get("date", "") or "").strip() - hint: str = (self.context.get("hint", "") or "").strip() - - # daily_dir / resource_dir come from app config — NOT tool params. - # Same convention as daily_create / daily_list / daily_reindex. - # resource_dir may be empty (default) — that just skips the resource scan. - cfg = self.app_context.app_config if self.app_context is not None else None - daily_dir = (cfg.daily_dir if cfg else "") or "daily" - resource_dir = cfg.resource_dir if cfg else "" - - today = date_input or self._now().strftime("%Y-%m-%d") - vault = self._vault_dir() - files = _scan_today_files(vault, today, daily_dir, resource_dir) - - result = CronDreamResult(date=today, files_scanned=len(files)) - self.logger.info( - f"[{self.name}] cron tick date={today} scanned={len(files)} file(s) under " - f"{daily_dir}/{today}/ + {resource_dir}/{today}/", - ) - - for rel_path in files: - try: - dr = await self.dream_one(rel_path, hint) - except Exception as e: # pylint: disable=broad-except - self.logger.error( - f"[{self.name}] dream_one failed on {rel_path}: {type(e).__name__}: {e}", - ) - dr = DreamResult( - path=rel_path, - error=f"{type(e).__name__}: {e}", - ) - result.per_file.append(dr) - if dr.error: - result.files_failed += 1 - elif dr.skipped: - result.files_skipped += 1 - else: - result.files_dreamed += 1 - - result.summary = _render_summary(result) - self.context.response.success = result.files_failed == 0 - self.context.response.answer = result.summary - self.context.response.metadata.update(result.model_dump()) - - -def _scan_today_files( - vault: Path, - today: str, - daily_dir: str, - resource_dir: str, -) -> list[str]: - """Return vault-relative paths of today's daily notes + resource files. - - * ``/.md`` — the day-index file (auto-rebuilt - rollup of all of today's notes). Included first so its day-level - abstractions land before the per-event details. - * ``//**/*.md`` — event notes for the day, - sorted by path. - * ``//**/*`` — any file type ingested under - today's resource folder. Skipped when ``resource_dir`` is empty. - - Results are sorted for deterministic processing order within each - group; the day-index file leads. - """ - out: list[str] = [] - - if daily_dir: - day_index = vault / daily_dir / f"{today}.md" - if day_index.is_file(): - out.append(str(day_index.relative_to(vault))) - daily_root = vault / daily_dir / today - if daily_root.is_dir(): - for md in sorted(daily_root.rglob("*.md")): - if md.is_file(): - out.append(str(md.relative_to(vault))) - - if resource_dir: - resource_root = vault / resource_dir / today - if resource_root.is_dir(): - for f in sorted(p for p in resource_root.rglob("*") if p.is_file()): - out.append(str(f.relative_to(vault))) - - return out - - -def _render_summary(r: CronDreamResult) -> str: - """One-line header + one line per file with its outcome.""" - lines = [ - f"[CronDreamer] date={r.date} scanned={r.files_scanned} " - f"dreamed={r.files_dreamed} skipped={r.files_skipped} failed={r.files_failed}", - ] - for dr in r.per_file: - if dr.error: - status = f"ERROR ({dr.error})" - elif dr.skipped: - status = "SKIP" - else: - status = ( - f"OK (+{len(dr.nodes_created)} created, ~{len(dr.nodes_updated)} updated, " - f"!{len(dr.conservation_violations)} conservation)" - ) - lines.append(f" - {dr.path}: {status}") - return "\n".join(lines)