mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-09-29 01:41:38 +00:00
refactor(dreamer): improve code formatting and line breaks
This commit is contained in:
parent
b8a5d9b595
commit
dd0a48f846
5 changed files with 149 additions and 193 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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/<today>/`` + ``resource/<today>/`` 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.
|
||||
|
||||
* ``<daily_dir>/<today>.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.
|
||||
* ``<daily_dir>/<today>/**/*.md`` — event notes for the day,
|
||||
sorted by path.
|
||||
* ``<resource_dir>/<today>/**/*`` — 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)
|
||||
|
|
@ -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")
|
||||
|
|
@ -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:
|
||||
|
||||
* ``<daily_dir>/<today>.md`` — the day-index rollup (processed first).
|
||||
* ``<daily_dir>/<today>/**/*.md`` — per-event notes for the day.
|
||||
* ``<resource_dir>/<today>/**/*`` — 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/<today>/`` + ``resource/<today>/`` 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.
|
||||
|
||||
* ``<daily_dir>/<today>.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.
|
||||
* ``<daily_dir>/<today>/**/*.md`` — event notes for the day,
|
||||
sorted by path.
|
||||
* ``<resource_dir>/<today>/**/*`` — 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)
|
||||
Loading…
Add table
Reference in a new issue