mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-10-10 03:30:56 +00:00
Refactor logging levels and add dream schema definitions (#291)
* chore(logging): change info logs to debug level for data loading operations - Changed stopwords loading log from info to debug level - Changed file catalog nodes loading log from info to debug level - Changed file graph nodes loading log from info to debug level * feat(dream): add dream schema definitions and enum for auto-dream functionality - Add DreamBucketEnum with procedure, personal, and wiki values - Create comprehensive dream-related Pydantic models including DreamUnit, DreamTopic, DreamExtractOutput, IntegrateOutcome, TopicSelectionOutput, ProactiveResult, and DreamState - Move schema definitions from local step module to shared schema package - Update dream extraction and integration steps to use new enum-based bucket validation - Initialize digest directories for each dream bucket type - Enhance embedding store health check with workspace directory logging * refactor(tests): update DreamState import path in test_auto_dream.py - Move DreamState import from reme.steps.evolve.dream.schema to reme.schema - Maintain same functionality with updated module reference - Align import with new schema location in project structure
This commit is contained in:
parent
a3bd81bde2
commit
7d86658f33
14 changed files with 59 additions and 23 deletions
|
|
@ -55,8 +55,8 @@ class LocalEmbeddingStore(BaseEmbeddingStore):
|
|||
async def _close(self) -> None:
|
||||
await self.dump()
|
||||
|
||||
async def health_check(self, timeout: float = 2.0) -> bool:
|
||||
tag = f"[EMBEDDING HEALTH CHECK] name={self.name}"
|
||||
async def health_check(self, timeout: float = 5.0) -> bool:
|
||||
tag = f"[EMBEDDING HEALTH CHECK] name={self.name} workspace_dir={self.workspace_path}"
|
||||
try:
|
||||
result = await asyncio.wait_for(self.as_embedding(["ping"]), timeout=timeout)
|
||||
if not result or result[0] is None:
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@ class LocalFileCatalog(BaseFileCatalog):
|
|||
if not self._catalog_file.exists():
|
||||
return
|
||||
await self._read_jsonl()
|
||||
self.logger.info(f"Loaded {len(self._nodes)} nodes from {self._catalog_file}")
|
||||
self.logger.debug(f"Loaded {len(self._nodes)} nodes from {self._catalog_file}")
|
||||
|
||||
async def dump(self) -> None:
|
||||
async with self._io_lock:
|
||||
|
|
|
|||
|
|
@ -35,7 +35,7 @@ class LocalFileGraph(BaseFileGraph):
|
|||
if line.strip():
|
||||
node = FileNode.model_validate_json(line)
|
||||
self._nodes[node.path] = node
|
||||
self.logger.info(f"Loaded {len(self._nodes)} nodes from {self._graph_file}")
|
||||
self.logger.debug(f"Loaded {len(self._nodes)} nodes from {self._graph_file}")
|
||||
except Exception as e:
|
||||
self.logger.exception(f"Failed to load {self._graph_file}: {e}")
|
||||
|
||||
|
|
|
|||
|
|
@ -38,7 +38,7 @@ class BaseTokenizer(BaseComponent):
|
|||
async with aiofiles.open(self.stopwords_path, encoding="utf-8") as f:
|
||||
content = await f.read()
|
||||
self._stopwords = {line.strip().lower() for line in content.splitlines() if line.strip()}
|
||||
self.logger.info(f"Loaded {len(self._stopwords)} stopwords from {self.stopwords_path}")
|
||||
self.logger.debug(f"Loaded {len(self._stopwords)} stopwords from {self.stopwords_path}")
|
||||
|
||||
async def _close(self) -> None:
|
||||
self._stopwords.clear()
|
||||
|
|
|
|||
|
|
@ -2,10 +2,12 @@
|
|||
|
||||
from .chunk_enum import ChunkEnum
|
||||
from .component_enum import ComponentEnum
|
||||
from .dream_bucket_enum import DreamBucketEnum
|
||||
from .link_scope_enum import LinkScopeEnum
|
||||
|
||||
__all__ = [
|
||||
"ChunkEnum",
|
||||
"ComponentEnum",
|
||||
"DreamBucketEnum",
|
||||
"LinkScopeEnum",
|
||||
]
|
||||
|
|
|
|||
11
reme/enumeration/dream_bucket_enum.py
Normal file
11
reme/enumeration/dream_bucket_enum.py
Normal file
|
|
@ -0,0 +1,11 @@
|
|||
"""Dream bucket enumeration module."""
|
||||
|
||||
from enum import Enum
|
||||
|
||||
|
||||
class DreamBucketEnum(str, Enum):
|
||||
"""Enumeration of digest memory buckets used by dream integration."""
|
||||
|
||||
PROCEDURE = "procedure"
|
||||
PERSONAL = "personal"
|
||||
WIKI = "wiki"
|
||||
|
|
@ -1,6 +1,15 @@
|
|||
"""Schema"""
|
||||
|
||||
from .application_config import ApplicationConfig, ComponentConfig, JobConfig
|
||||
from .dream import (
|
||||
DreamExtractOutput,
|
||||
DreamState,
|
||||
DreamTopic,
|
||||
DreamUnit,
|
||||
IntegrateOutcome,
|
||||
ProactiveResult,
|
||||
TopicSelectionOutput,
|
||||
)
|
||||
from .emb_node import EmbNode
|
||||
from .file_chunk import FileChunk
|
||||
from .file_front_matter import FileFrontMatter
|
||||
|
|
@ -14,12 +23,19 @@ __all__ = [
|
|||
"ApplicationConfig",
|
||||
"ComponentConfig",
|
||||
"JobConfig",
|
||||
"DreamExtractOutput",
|
||||
"DreamState",
|
||||
"DreamTopic",
|
||||
"DreamUnit",
|
||||
"EmbNode",
|
||||
"FileChunk",
|
||||
"FileFrontMatter",
|
||||
"FileLink",
|
||||
"FileNode",
|
||||
"IntegrateOutcome",
|
||||
"ProactiveResult",
|
||||
"Request",
|
||||
"Response",
|
||||
"StreamChunk",
|
||||
"TopicSelectionOutput",
|
||||
]
|
||||
|
|
|
|||
|
|
@ -1,18 +1,17 @@
|
|||
"""Shared auto-dream schemas."""
|
||||
"""Auto-dream schemas."""
|
||||
|
||||
from typing import Literal
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
BUCKETS: tuple[str, ...] = ("procedure", "personal", "wiki")
|
||||
Bucket = Literal["procedure", "personal", "wiki"]
|
||||
from ..enumeration import DreamBucketEnum
|
||||
|
||||
|
||||
class DreamUnit(BaseModel):
|
||||
"""One cross-file memory unit emitted by global extract."""
|
||||
|
||||
name: str = Field(description="Short kebab-case handle for the abstraction.")
|
||||
bucket: str = Field(description="procedure, personal, or wiki; unknown values route to wiki.")
|
||||
bucket: DreamBucketEnum = Field(description="Digest bucket; unknown raw values route to wiki before validation.")
|
||||
summary: str = Field(description="Grounded abstraction summary with evidence pointers.")
|
||||
paths: list[str] = Field(default_factory=list, description="Workspace-relative source paths.")
|
||||
|
||||
|
|
@ -61,7 +60,7 @@ class ProactiveResult(BaseModel):
|
|||
|
||||
|
||||
class DreamState(BaseModel):
|
||||
"""Shared state passed across the four dream steps."""
|
||||
"""Shared state passed across the dream steps."""
|
||||
|
||||
date: str = ""
|
||||
dates: list[str] = Field(default_factory=list)
|
||||
|
|
@ -5,7 +5,8 @@ import json
|
|||
from ...base_step import BaseStep
|
||||
from ...file_io import refresh_day_index
|
||||
from ....components import R
|
||||
from .schema import BUCKETS, DreamState
|
||||
from ....enumeration import DreamBucketEnum
|
||||
from ....schema import DreamState
|
||||
from .utils import (
|
||||
clean_paths,
|
||||
daily_dir,
|
||||
|
|
@ -100,7 +101,7 @@ class DreamExtractStep(BaseStep):
|
|||
system_prompt=self.prompt_format(
|
||||
"extract_system_prompt",
|
||||
workspace_dir=str(workspace),
|
||||
buckets=", ".join(BUCKETS),
|
||||
buckets=", ".join(bucket.value for bucket in DreamBucketEnum),
|
||||
),
|
||||
job_tools=list(_TOOLS),
|
||||
)
|
||||
|
|
@ -127,13 +128,15 @@ class DreamExtractStep(BaseStep):
|
|||
continue
|
||||
name = str(raw.get("name") or "").strip()
|
||||
summary = str(raw.get("summary") or "").strip()
|
||||
bucket = str(raw.get("bucket") 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
|
||||
if bucket not in BUCKETS:
|
||||
self.logger.warning(f"[{self.name}] unit {name!r} emitted bucket {bucket!r}; routing to wiki")
|
||||
bucket = "wiki"
|
||||
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)
|
||||
|
|
|
|||
|
|
@ -4,8 +4,7 @@ from pathlib import Path
|
|||
|
||||
from ...base_step import BaseStep
|
||||
from ....components import R
|
||||
from ....schema import FileNode
|
||||
from .schema import DreamState
|
||||
from ....schema import DreamState, FileNode
|
||||
from .utils import state_from_context, store_state, workspace_dir
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -5,7 +5,8 @@ from pathlib import Path
|
|||
|
||||
from ...base_step import BaseStep
|
||||
from ....components import R
|
||||
from .schema import BUCKETS, IntegrateOutcome
|
||||
from ....enumeration import DreamBucketEnum
|
||||
from ....schema import IntegrateOutcome
|
||||
from .utils import llm_available, pack_paths, parse_structured_reply, state_from_context, store_state, workspace_dir
|
||||
|
||||
_TOOLS = ("node_search", "read", "frontmatter_read", "write", "edit", "frontmatter_update")
|
||||
|
|
@ -29,6 +30,8 @@ class DreamIntegrateStep(BaseStep):
|
|||
|
||||
workspace = Path(state.workspace).resolve() if state.workspace else workspace_dir(self)
|
||||
digest_dir = self.config_value("digest_dir")
|
||||
for bucket in DreamBucketEnum:
|
||||
(workspace / digest_dir / bucket.value).mkdir(parents=True, exist_ok=True)
|
||||
for i, unit in enumerate(state.units, start=1):
|
||||
await self._integrate_one(state, unit, i, workspace, digest_dir)
|
||||
state.failed_paths = sorted(set(state.failed_paths))
|
||||
|
|
@ -36,7 +39,10 @@ class DreamIntegrateStep(BaseStep):
|
|||
return self._finish(state, not state.failed_units, answer)
|
||||
|
||||
async def _integrate_one(self, state, unit: dict, index: int, workspace: Path, digest_dir: str) -> None:
|
||||
bucket = unit.get("bucket") if unit.get("bucket") in BUCKETS else "wiki"
|
||||
try:
|
||||
bucket = DreamBucketEnum(str(unit.get("bucket") or "")).value
|
||||
except ValueError:
|
||||
bucket = DreamBucketEnum.WIKI.value
|
||||
paths = [str(p) for p in unit.get("paths", [])]
|
||||
try:
|
||||
result = await self.agent_wrapper.reply(
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@
|
|||
|
||||
from ...base_step import BaseStep
|
||||
from ....components import R
|
||||
from .schema import ProactiveResult
|
||||
from ....schema import ProactiveResult
|
||||
from .utils import load_yaml_topics, today, workspace_dir
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ import yaml
|
|||
|
||||
from .._evolve import now
|
||||
from ...base_step import BaseStep
|
||||
from .schema import DreamState
|
||||
from ....schema import DreamState
|
||||
|
||||
|
||||
def state_from_context(step: BaseStep) -> DreamState:
|
||||
|
|
|
|||
|
|
@ -9,8 +9,8 @@ import yaml
|
|||
from reme.components.file_catalog import BaseFileCatalog
|
||||
from reme.components.file_store import BaseFileStore
|
||||
from reme.components.runtime_context import RuntimeContext
|
||||
from reme.schema import DreamState
|
||||
from reme.steps.evolve.dream.finish import DreamFinishStep
|
||||
from reme.steps.evolve.dream.schema import DreamState
|
||||
from reme.steps.evolve.dream.topics import DreamTopicsStep
|
||||
from reme.steps.evolve.dream.utils import parse_structured_reply, recent_dates, scan_day_files
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue