mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-09-30 01:52:29 +00:00
up
This commit is contained in:
parent
b3774c67db
commit
9a7877a3a4
13 changed files with 992 additions and 50 deletions
|
|
@ -5,6 +5,7 @@ from . import as_llm_formatter
|
|||
from . import as_token_counter
|
||||
from . import client
|
||||
from . import embedding
|
||||
from . import file_catalog
|
||||
from . import file_graph
|
||||
from . import file_parser
|
||||
from . import file_store
|
||||
|
|
@ -31,6 +32,7 @@ __all__ = [
|
|||
"as_token_counter",
|
||||
"client",
|
||||
"embedding",
|
||||
"file_catalog",
|
||||
"file_graph",
|
||||
"file_parser",
|
||||
"file_store",
|
||||
|
|
|
|||
9
reme4/components/file_catalog/__init__.py
Normal file
9
reme4/components/file_catalog/__init__.py
Normal file
|
|
@ -0,0 +1,9 @@
|
|||
"""File catalog """
|
||||
|
||||
from .base_file_catalog import BaseFileCatalog
|
||||
from .local_file_catalog import LocalFileCatalog
|
||||
|
||||
__all__ = [
|
||||
"BaseFileCatalog",
|
||||
"LocalFileCatalog",
|
||||
]
|
||||
56
reme4/components/file_catalog/base_file_catalog.py
Normal file
56
reme4/components/file_catalog/base_file_catalog.py
Normal file
|
|
@ -0,0 +1,56 @@
|
|||
"""Abstract base for file-catalog backends."""
|
||||
|
||||
from abc import abstractmethod
|
||||
from pathlib import Path
|
||||
|
||||
from ..base_component import BaseComponent
|
||||
from ...enumeration import ComponentEnum
|
||||
from ...schema import FileNode
|
||||
|
||||
|
||||
class BaseFileCatalog(BaseComponent):
|
||||
"""Abstract base for file-catalog backends.
|
||||
|
||||
A catalog records FileNode entries keyed by path — a lightweight
|
||||
counterpart to FileStore that drops chunk/embedding/keyword/link
|
||||
machinery and exposes only node upsert / delete / lookup.
|
||||
"""
|
||||
|
||||
component_type = ComponentEnum.FILE_CATALOG
|
||||
|
||||
def __init__(self, catalog_name: str = "default", catalog_version: str = "v1", **kwargs):
|
||||
super().__init__(**kwargs)
|
||||
self.catalog_name: str = catalog_name or self.name
|
||||
self.catalog_version: str = catalog_version
|
||||
self.catalog_path: Path = self.vault_metadata_path / self.component_type.value / self.catalog_name
|
||||
self.catalog_path.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# -- Lifecycle ---------------------------------------------------------
|
||||
|
||||
async def _start(self) -> None:
|
||||
await super()._start()
|
||||
await self.load()
|
||||
|
||||
async def _close(self) -> None:
|
||||
await self.dump()
|
||||
await super()._close()
|
||||
|
||||
async def load(self) -> None:
|
||||
"""Load persisted state. No-op for backends without local files."""
|
||||
|
||||
async def dump(self) -> None:
|
||||
"""Persist state. No-op for backends without local files."""
|
||||
|
||||
# -- CRUD --------------------------------------------------------------
|
||||
|
||||
@abstractmethod
|
||||
async def upsert(self, nodes: list[FileNode]) -> None:
|
||||
"""Insert or update nodes keyed by path."""
|
||||
|
||||
@abstractmethod
|
||||
async def delete(self, path: str | list[str]) -> None:
|
||||
"""Delete nodes by path; missing paths are skipped."""
|
||||
|
||||
@abstractmethod
|
||||
async def get_nodes(self, paths: list[str] | None = None) -> list[FileNode]:
|
||||
"""Return nodes by paths; None = all nodes; [] = []; missing paths are skipped."""
|
||||
56
reme4/components/file_catalog/local_file_catalog.py
Normal file
56
reme4/components/file_catalog/local_file_catalog.py
Normal file
|
|
@ -0,0 +1,56 @@
|
|||
"""In-memory file catalog with JSONL persistence."""
|
||||
|
||||
import aiofiles
|
||||
|
||||
from .base_file_catalog import BaseFileCatalog
|
||||
from ..component_registry import R
|
||||
from ...schema import FileNode
|
||||
|
||||
|
||||
@R.register("local")
|
||||
class LocalFileCatalog(BaseFileCatalog):
|
||||
"""Dict-backed file catalog persisted as JSONL on close."""
|
||||
|
||||
def __init__(self, encoding: str = "utf-8", **kwargs):
|
||||
super().__init__(**kwargs)
|
||||
self.encoding = encoding
|
||||
self._nodes: dict[str, FileNode] = {}
|
||||
self._catalog_file = self.catalog_path / f"nodes_{self.catalog_version}.jsonl"
|
||||
|
||||
async def load(self) -> None:
|
||||
if not self._catalog_file.exists():
|
||||
return
|
||||
try:
|
||||
async with aiofiles.open(self._catalog_file, encoding=self.encoding) as f:
|
||||
async for line in f:
|
||||
line = line.strip()
|
||||
if line:
|
||||
node = FileNode.model_validate_json(line)
|
||||
self._nodes[node.path] = node
|
||||
self.logger.info(f"Loaded {len(self._nodes)} nodes from {self._catalog_file}")
|
||||
except Exception as e:
|
||||
self.logger.exception(f"Failed to load {self._catalog_file}: {e}")
|
||||
|
||||
async def dump(self) -> None:
|
||||
try:
|
||||
tmp = self._catalog_file.with_suffix(".tmp")
|
||||
async with aiofiles.open(tmp, "w", encoding=self.encoding) as f:
|
||||
await f.write("\n".join(n.model_dump_json() for n in self._nodes.values()))
|
||||
tmp.replace(self._catalog_file)
|
||||
self.logger.info(f"Saved {len(self._nodes)} nodes to {self._catalog_file}")
|
||||
except Exception as e:
|
||||
self.logger.exception(f"Failed to write {self._catalog_file}: {e}")
|
||||
|
||||
async def upsert(self, nodes: list[FileNode]) -> None:
|
||||
for node in nodes:
|
||||
self._nodes[node.path] = node
|
||||
|
||||
async def delete(self, path: str | list[str]) -> None:
|
||||
paths = [path] if isinstance(path, str) else path
|
||||
for p in paths:
|
||||
self._nodes.pop(p, None)
|
||||
|
||||
async def get_nodes(self, paths: list[str] | None = None) -> list[FileNode]:
|
||||
if paths is None:
|
||||
return list(self._nodes.values())
|
||||
return [self._nodes[p] for p in paths if p in self._nodes]
|
||||
68
reme4/config/demo.yaml
Normal file
68
reme4/config/demo.yaml
Normal file
|
|
@ -0,0 +1,68 @@
|
|||
service:
|
||||
backend: http
|
||||
|
||||
jobs:
|
||||
- backend: base
|
||||
name: version
|
||||
description: "return reme4 package version"
|
||||
parameters:
|
||||
type: object
|
||||
properties: { }
|
||||
steps:
|
||||
- backend: version_step
|
||||
|
||||
- backend: base
|
||||
name: help
|
||||
description: "list all registered jobs with their metadata"
|
||||
parameters:
|
||||
type: object
|
||||
properties: { }
|
||||
steps:
|
||||
- backend: help_step
|
||||
|
||||
- backend: base
|
||||
name: demo
|
||||
description: "demo job description"
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
query:
|
||||
type: string
|
||||
description: "query"
|
||||
min_score:
|
||||
type: number
|
||||
description: "min score"
|
||||
default: 0.5
|
||||
required:
|
||||
- query
|
||||
steps:
|
||||
- backend: demo_echo_step1
|
||||
- backend: demo_echo_step2
|
||||
|
||||
- backend: stream
|
||||
name: stream_demo
|
||||
description: "stream demo job: repeat query 10x and stream char-by-char"
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
query:
|
||||
type: string
|
||||
description: "query to echo"
|
||||
repeat:
|
||||
type: integer
|
||||
description: "number of times to repeat the query"
|
||||
default: 10
|
||||
interval:
|
||||
type: number
|
||||
description: "seconds between chunks"
|
||||
default: 0.1
|
||||
required:
|
||||
- query
|
||||
steps:
|
||||
- backend: stream_demo_step1
|
||||
- backend: stream_demo_step2
|
||||
|
||||
components:
|
||||
tokenizer:
|
||||
default:
|
||||
backend: regex
|
||||
570
reme4/config/qwenpaw.yaml
Normal file
570
reme4/config/qwenpaw.yaml
Normal file
|
|
@ -0,0 +1,570 @@
|
|||
service:
|
||||
backend: http
|
||||
# backend: mcp
|
||||
|
||||
# Default dev config points vault_dir at ./.reme so `python -m reme4 start`
|
||||
# can be run from the repo root and exercise the full atomic-tool surface
|
||||
# against the seeded test data. Override the `vault_dir=` CLI arg.
|
||||
vault_dir: .reme
|
||||
daily_dir: daily
|
||||
digest_dir: digest
|
||||
resource_dir: resource
|
||||
|
||||
jobs:
|
||||
# ════════════════════════════════════════════════════════════════════
|
||||
# UTILITY — service introspection
|
||||
# ════════════════════════════════════════════════════════════════════
|
||||
- backend: base
|
||||
name: version
|
||||
description: "return reme4 package version"
|
||||
parameters:
|
||||
type: object
|
||||
properties: {}
|
||||
steps:
|
||||
- backend: version_step
|
||||
|
||||
- backend: base
|
||||
name: health_check
|
||||
description: "return a concise health-check snapshot of reme4 components"
|
||||
parameters:
|
||||
type: object
|
||||
properties: {}
|
||||
steps:
|
||||
- backend: health_check_step
|
||||
|
||||
- backend: base
|
||||
name: help
|
||||
description: "list all registered jobs with their metadata"
|
||||
parameters:
|
||||
type: object
|
||||
properties: {}
|
||||
steps:
|
||||
- backend: help_step
|
||||
|
||||
- backend: base
|
||||
name: reindex
|
||||
description: "wipe the file store and rebuild it from the watcher's tracked files"
|
||||
parameters:
|
||||
type: object
|
||||
properties: {}
|
||||
steps:
|
||||
- backend: reindex_step
|
||||
|
||||
- backend: base
|
||||
name: index_changes
|
||||
description: "apply a batch of file changes (added/modified/deleted) into file_store"
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
changes:
|
||||
type: array
|
||||
description: "list of change items"
|
||||
items:
|
||||
type: object
|
||||
properties:
|
||||
change:
|
||||
type: string
|
||||
enum: [added, modified, deleted]
|
||||
description: "type of file change"
|
||||
path:
|
||||
type: string
|
||||
description: "absolute file path"
|
||||
required:
|
||||
- change
|
||||
- path
|
||||
required:
|
||||
- changes
|
||||
steps:
|
||||
- backend: index_changes_step
|
||||
|
||||
# ════════════════════════════════════════════════════════════════════
|
||||
# ATOMIC TOOLS — same surface plugins/reme-{service,expert} expose
|
||||
# ════════════════════════════════════════════════════════════════════
|
||||
|
||||
# ── Retrieve ───────────────────────────────────────────────────────
|
||||
- backend: base
|
||||
name: search
|
||||
description: "Hybrid vault search (vector + BM25, RRF-fused)."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
query:
|
||||
type: string
|
||||
description: "search query"
|
||||
limit:
|
||||
type: integer
|
||||
description: "max results"
|
||||
default: 5
|
||||
min_score:
|
||||
type: number
|
||||
description: "min fused score"
|
||||
default: 0.0
|
||||
required:
|
||||
- query
|
||||
steps:
|
||||
- backend: search_step
|
||||
vector_weight: 0.7
|
||||
candidate_multiplier: 3.0
|
||||
expand_links: true
|
||||
max_links_per_direction: 10
|
||||
|
||||
- backend: base
|
||||
name: traverse
|
||||
description: "Walk the wikilink graph from a seed path."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "seed path (vault-relative)"
|
||||
depth:
|
||||
type: integer
|
||||
description: "hop limit"
|
||||
default: 1
|
||||
direction:
|
||||
type: string
|
||||
description: "forward / backward / both"
|
||||
default: both
|
||||
required:
|
||||
- path
|
||||
steps:
|
||||
- backend: traverse_step
|
||||
|
||||
# ── Read Operations ───────────────────────────────────────────────────────────
|
||||
- backend: base
|
||||
name: list
|
||||
description: "List files under a vault path."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative dir; empty = root"
|
||||
default: ""
|
||||
recursive:
|
||||
type: boolean
|
||||
description: "recurse"
|
||||
default: false
|
||||
limit:
|
||||
type: integer
|
||||
description: "max results"
|
||||
default: 100
|
||||
steps:
|
||||
- backend: list_step
|
||||
|
||||
- backend: base
|
||||
name: read
|
||||
description: "Read a markdown file under the vault."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path; markdown only"
|
||||
start_line:
|
||||
type: integer
|
||||
description: "first line (1-based, inclusive)"
|
||||
end_line:
|
||||
type: integer
|
||||
description: "last line (1-based, inclusive)"
|
||||
required:
|
||||
- path
|
||||
steps:
|
||||
- backend: read_step
|
||||
|
||||
- backend: base
|
||||
name: stat
|
||||
description: "Stat a vault file (size, mtime, exists, is_dir, is_file)."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
required:
|
||||
- path
|
||||
steps:
|
||||
- backend: stat_step
|
||||
|
||||
- backend: base
|
||||
name: frontmatter:read
|
||||
description: "Read a file's YAML frontmatter as a dict."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
required:
|
||||
- path
|
||||
steps:
|
||||
- backend: frontmatter:read_step
|
||||
|
||||
# ── Write Operations──────────────────────────────────────────────────────────
|
||||
- backend: base
|
||||
name: write
|
||||
description: "Write a markdown file (create or overwrite) with name/description frontmatter."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path; markdown only"
|
||||
name:
|
||||
type: string
|
||||
description: "frontmatter name"
|
||||
description:
|
||||
type: string
|
||||
description: "frontmatter description"
|
||||
content:
|
||||
type: string
|
||||
description: "body"
|
||||
required:
|
||||
- path
|
||||
- name
|
||||
- description
|
||||
- content
|
||||
steps:
|
||||
- backend: write_step
|
||||
|
||||
- backend: base
|
||||
name: edit
|
||||
description: "Find-and-replace in a markdown file (all occurrences)."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
old:
|
||||
type: string
|
||||
description: "text to find"
|
||||
new:
|
||||
type: string
|
||||
description: "replacement"
|
||||
default: ""
|
||||
required:
|
||||
- path
|
||||
- old
|
||||
- new
|
||||
steps:
|
||||
- backend: edit_step
|
||||
|
||||
- backend: base
|
||||
name: append
|
||||
description: "Append content to a markdown file."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
content:
|
||||
type: string
|
||||
description: "content to append"
|
||||
required:
|
||||
- path
|
||||
- content
|
||||
steps:
|
||||
- backend: append_step
|
||||
|
||||
- backend: base
|
||||
name: frontmatter:update
|
||||
description: "Merge keys into a file's YAML frontmatter."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
metadata:
|
||||
type: object
|
||||
description: "keys to merge"
|
||||
additionalProperties: true
|
||||
required:
|
||||
- path
|
||||
- metadata
|
||||
steps:
|
||||
- backend: frontmatter_update_step
|
||||
|
||||
- backend: base
|
||||
name: frontmatter:delete
|
||||
description: "Drop keys from a file's YAML frontmatter."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
keys:
|
||||
type: array
|
||||
description: "keys to remove"
|
||||
items:
|
||||
type: string
|
||||
required:
|
||||
- path
|
||||
- keys
|
||||
steps:
|
||||
- backend: frontmatter_delete_step
|
||||
|
||||
# ── File Operations (relocate / cross-realm) ──────────────────────────────
|
||||
- backend: base
|
||||
name: move
|
||||
description: "Move / rename a vault file; rewrites inbound wikilinks by default."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
src_path:
|
||||
type: string
|
||||
description: "vault-relative source"
|
||||
dst_path:
|
||||
type: string
|
||||
description: "vault-relative destination"
|
||||
overwrite:
|
||||
type: boolean
|
||||
description: "overwrite if dst exists"
|
||||
default: false
|
||||
retarget:
|
||||
type: boolean
|
||||
description: "rewrite [[src]] → [[dst]] across the vault"
|
||||
default: true
|
||||
required:
|
||||
- src_path
|
||||
- dst_path
|
||||
steps:
|
||||
- backend: move_step
|
||||
|
||||
- backend: base
|
||||
name: delete
|
||||
description: "Delete a vault file or folder; returns surviving inbound wikilinks."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "vault-relative path"
|
||||
required:
|
||||
- path
|
||||
steps:
|
||||
- backend: delete_step
|
||||
|
||||
- backend: base
|
||||
name: upload
|
||||
description: "Copy a host file into the vault at an explicit destination."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
src_path:
|
||||
type: string
|
||||
description: "host absolute path"
|
||||
dst_path:
|
||||
type: string
|
||||
description: "vault-relative destination"
|
||||
overwrite:
|
||||
type: boolean
|
||||
description: "overwrite if dst exists"
|
||||
default: false
|
||||
required:
|
||||
- src_path
|
||||
- dst_path
|
||||
steps:
|
||||
- backend: upload_step
|
||||
|
||||
- backend: base
|
||||
name: upload_resource
|
||||
description: "Ingest an external-channel asset into resource/<today>/ with provenance."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
path:
|
||||
type: string
|
||||
description: "host source path"
|
||||
channel:
|
||||
type: string
|
||||
description: "channel id (wechat / email / browser / api / ...)"
|
||||
description:
|
||||
type: string
|
||||
description: "what the asset is and how to interpret it"
|
||||
metadata:
|
||||
type: object
|
||||
description: "extra provenance keys (e.g. source)"
|
||||
default: {}
|
||||
required:
|
||||
- path
|
||||
- channel
|
||||
- description
|
||||
steps:
|
||||
- backend: upload_resource_step
|
||||
|
||||
- backend: base
|
||||
name: download
|
||||
description: "Copy a vault file out to the host filesystem."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
src_path:
|
||||
type: string
|
||||
description: "vault-relative source"
|
||||
dst_path:
|
||||
type: string
|
||||
description: "host absolute dest; empty = temp file"
|
||||
default: ""
|
||||
overwrite:
|
||||
type: boolean
|
||||
description: "overwrite if dst exists"
|
||||
default: false
|
||||
required:
|
||||
- src_path
|
||||
steps:
|
||||
- backend: download_step
|
||||
|
||||
# ── Daily Operations (note CRUD + day-index rollup) ───────────────────
|
||||
- backend: base
|
||||
name: daily:read
|
||||
description: "Read daily/<date>/<slug>.md (body + frontmatter)."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
slug:
|
||||
type: string
|
||||
description: "note slug"
|
||||
date:
|
||||
type: string
|
||||
description: "ISO date; empty = today"
|
||||
default: ""
|
||||
required:
|
||||
- slug
|
||||
steps:
|
||||
- backend: daily_read_step
|
||||
|
||||
- backend: base
|
||||
name: daily:write
|
||||
description: "Write daily/<date>/<slug>.md (body + frontmatter); refreshes the day index."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
slug:
|
||||
type: string
|
||||
description: "note slug"
|
||||
body:
|
||||
type: string
|
||||
description: "note body"
|
||||
default: ""
|
||||
frontmatter:
|
||||
type: object
|
||||
description: "frontmatter dict; defaults to {name: <slug>}"
|
||||
default: {}
|
||||
date:
|
||||
type: string
|
||||
description: "ISO date; empty = today"
|
||||
default: ""
|
||||
overwrite:
|
||||
type: boolean
|
||||
description: "false = skip if exists; true = replace"
|
||||
default: false
|
||||
refresh_index:
|
||||
type: boolean
|
||||
description: "refresh daily/<date>.md after write"
|
||||
default: true
|
||||
required:
|
||||
- slug
|
||||
steps:
|
||||
- backend: daily_write_step
|
||||
|
||||
- backend: base
|
||||
name: daily:list
|
||||
description: "List notes under a single day."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
date:
|
||||
type: string
|
||||
description: "ISO date; empty = today"
|
||||
default: ""
|
||||
steps:
|
||||
- backend: daily_list_step
|
||||
|
||||
- backend: base
|
||||
name: daily:reindex
|
||||
description: "Rebuild the day-index page daily/<date>.md."
|
||||
parameters:
|
||||
type: object
|
||||
properties:
|
||||
date:
|
||||
type: string
|
||||
description: "ISO date; empty = today"
|
||||
default: ""
|
||||
steps:
|
||||
- backend: daily_reindex_step
|
||||
|
||||
- backend: background
|
||||
name: watch_file
|
||||
watch_paths:
|
||||
- MEMORY.md
|
||||
- memory
|
||||
suffix_filters:
|
||||
- md
|
||||
steps:
|
||||
- backend: update_store_step
|
||||
- backend: watch_changes_step
|
||||
|
||||
components:
|
||||
tokenizer:
|
||||
default:
|
||||
backend: regex
|
||||
|
||||
embedding_model:
|
||||
default:
|
||||
backend: ${EMBEDDING_BACKEND:-openai}
|
||||
api_key: ${EMBEDDING_API_KEY:-}
|
||||
base_url: ${EMBEDDING_BASE_URL:-https://api.openai.com/v1}
|
||||
model_name: ${EMBEDDING_MODEL_NAME:-text-embedding-v4}
|
||||
dimensions: 1024
|
||||
|
||||
file_graph:
|
||||
default:
|
||||
backend: local
|
||||
|
||||
file_parser:
|
||||
linked:
|
||||
backend: linked
|
||||
supported_extensions:
|
||||
- md
|
||||
chunked:
|
||||
backend: chunked
|
||||
supported_extensions:
|
||||
- txt
|
||||
- html
|
||||
- json
|
||||
- yaml
|
||||
- py
|
||||
default:
|
||||
backend: default
|
||||
|
||||
keyword_index:
|
||||
default:
|
||||
backend: bm25
|
||||
tokenizer: default
|
||||
|
||||
file_store:
|
||||
default:
|
||||
backend: local
|
||||
store_name: local
|
||||
embedding_model: default
|
||||
keyword_index: default
|
||||
file_graph: default
|
||||
|
||||
# as_llm / formatter aren't required for atomic primitives; configure
|
||||
# only if you'll invoke digester/synchronizer or other LLM-driven
|
||||
# paths from the dev server.
|
||||
# as_llm:
|
||||
# default:
|
||||
# backend: ${LLM_BACKEND:-openai}
|
||||
# api_key: ${LLM_API_KEY:-}
|
||||
# model_name: ${LLM_MODEL_NAME:-gpt-4o-mini}
|
||||
# client_kwargs:
|
||||
# base_url: ${LLM_BASE_URL:-https://api.openai.com/v1}
|
||||
#
|
||||
# as_llm_formatter:
|
||||
# default:
|
||||
# backend: ${LLM_BACKEND:-openai}
|
||||
|
|
@ -22,6 +22,8 @@ class ComponentEnum(str, Enum):
|
|||
|
||||
FILE_GRAPH = "file_graph"
|
||||
|
||||
FILE_CATALOG = "file_catalog"
|
||||
|
||||
KEYWORD_INDEX = "keyword_index"
|
||||
|
||||
SERVICE = "service"
|
||||
|
|
|
|||
|
|
@ -1,11 +1,11 @@
|
|||
"""Background steps."""
|
||||
|
||||
from .index_changes import IndexChangesStep
|
||||
from .update_store import UpdateStoreStep
|
||||
from .scan_changes import ScanChangesStep
|
||||
from .update_store_index import UpdateStoreIndexStep
|
||||
from .watch_changes import WatchChangesStep
|
||||
|
||||
__all__ = [
|
||||
"IndexChangesStep",
|
||||
"UpdateStoreStep",
|
||||
"ScanChangesStep",
|
||||
"UpdateStoreIndexStep",
|
||||
"WatchChangesStep",
|
||||
]
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
"""Initial sync: diff watch_paths vs file_store, then index the diff."""
|
||||
"""One-shot scan: diff watch_paths vs file_store and emit changes for indexing."""
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
|
|
@ -6,14 +6,21 @@ from ..base_step import BaseStep
|
|||
from ...components import R
|
||||
|
||||
|
||||
@R.register("update_store_step")
|
||||
class UpdateStoreStep(BaseStep):
|
||||
"""One-shot sync: compute added/modified/deleted vs file_store and index."""
|
||||
@R.register("scan_changes_step")
|
||||
class ScanChangesStep(BaseStep):
|
||||
"""One-shot scan: compute added/modified/deleted vs file_store and dispatch."""
|
||||
|
||||
def __init__(self, recursive: bool = True, dump: bool = True, **kwargs):
|
||||
def __init__(
|
||||
self,
|
||||
recursive: bool = True,
|
||||
dump_store_index: bool = True,
|
||||
dispatch_job: str = "",
|
||||
**kwargs,
|
||||
):
|
||||
super().__init__(**kwargs)
|
||||
self.recursive: bool = recursive
|
||||
self.dump: bool = dump
|
||||
self.dump_store_index: bool = dump_store_index
|
||||
self.dispatch_job: str = dispatch_job
|
||||
|
||||
async def execute(self):
|
||||
assert self.context is not None
|
||||
|
|
@ -55,10 +62,13 @@ class UpdateStoreStep(BaseStep):
|
|||
counts = {"added": len(to_add), "modified": len(to_modify), "deleted": len(to_delete)}
|
||||
|
||||
if changes:
|
||||
self.logger.info(f"[{self.name}] initial sync: {counts}")
|
||||
await self.run_job("index_changes", changes=changes)
|
||||
if self.dump:
|
||||
await self.file_store.dump()
|
||||
self.logger.info(f"[{self.name}] scan: {counts}")
|
||||
if self.dispatch_job:
|
||||
await self.run_job(
|
||||
self.dispatch_job,
|
||||
changes=changes,
|
||||
dump_store_index=self.dump_store_index,
|
||||
)
|
||||
else:
|
||||
self.logger.info(f"[{self.name}] store is up to date")
|
||||
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
"""Index a batch of file changes into file_store."""
|
||||
"""Update store index with a batch of file changes."""
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
|
|
@ -9,14 +9,15 @@ from ...components import R
|
|||
from ...schema import FileChunk, FileNode
|
||||
|
||||
|
||||
@R.register("index_changes_step")
|
||||
class IndexChangesStep(BaseStep):
|
||||
"""Classify raw watcher changes and index them into file_store."""
|
||||
@R.register("update_store_index_step")
|
||||
class UpdateStoreIndexStep(BaseStep):
|
||||
"""Classify raw watcher changes and update the file_store index."""
|
||||
|
||||
async def execute(self):
|
||||
assert self.context is not None
|
||||
# Each item: {"change": Change | "added"|"modified"|"deleted", "path": absolute path}
|
||||
changes: list[dict] = self.context.get("changes") or []
|
||||
dump_store_index: bool = bool(self.context.get("dump_store_index", False))
|
||||
|
||||
buckets: dict[Change, list[str]] = {Change.added: [], Change.modified: [], Change.deleted: []}
|
||||
for item in changes:
|
||||
|
|
@ -76,6 +77,9 @@ class IndexChangesStep(BaseStep):
|
|||
self.logger.exception(f"Failed to delete {len(deleted)} file(s)")
|
||||
results.extend({"change": "deleted", "path": p, "success": False, "error": str(e)} for p in deleted)
|
||||
|
||||
if dump_store_index and results:
|
||||
await self.file_store.dump()
|
||||
|
||||
self.context.response.answer = results
|
||||
self.context.response.success = all(r["success"] for r in results) if results else True
|
||||
return self.context.response
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
"""Long-running awatch loop: convert raw changes into index_changes calls."""
|
||||
"""Long-running awatch loop: convert raw changes into update_store_index calls."""
|
||||
|
||||
import asyncio
|
||||
|
||||
|
|
@ -10,7 +10,7 @@ from ...components import R
|
|||
|
||||
@R.register("watch_changes_step")
|
||||
class WatchChangesStep(BaseStep):
|
||||
"""Watch files and forward each batch of raw changes to the index_changes job."""
|
||||
"""Watch files and forward each batch of raw changes to the update_store_index job."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
|
|
@ -18,6 +18,7 @@ class WatchChangesStep(BaseStep):
|
|||
force_polling: bool = True,
|
||||
debounce: int = 2000,
|
||||
poll_delay_ms: int = 2000,
|
||||
dispatch_job: str = "",
|
||||
**kwargs,
|
||||
):
|
||||
super().__init__(**kwargs)
|
||||
|
|
@ -25,6 +26,7 @@ class WatchChangesStep(BaseStep):
|
|||
self.force_polling: bool = force_polling
|
||||
self.debounce: int = debounce
|
||||
self.poll_delay_ms: int = poll_delay_ms
|
||||
self.dispatch_job: str = dispatch_job
|
||||
|
||||
def _filter(self, _change: Change, path: str) -> bool:
|
||||
suffixes = (self.context.get("suffix_filters") if self.context else None) or ["md"]
|
||||
|
|
@ -62,6 +64,7 @@ class WatchChangesStep(BaseStep):
|
|||
]
|
||||
if changes:
|
||||
self.logger.info(f"Detected {len(changes)} change(s)")
|
||||
await self.run_job("index_changes", changes=changes)
|
||||
if self.dispatch_job:
|
||||
await self.run_job(self.dispatch_job, changes=changes)
|
||||
|
||||
return self.context.response
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
"""Tests for background steps: UpdateStoreStep + WatchChangesStep.
|
||||
"""Tests for background steps: ScanChangesStep + WatchChangesStep.
|
||||
|
||||
Both steps are subclasses of BaseStep. To exercise them without spinning up the
|
||||
full ApplicationContext / index_changes job, we:
|
||||
full ApplicationContext / update_store_index job, we:
|
||||
* pass real (started) file_store/file_parser via the step's kwargs (so the
|
||||
BaseStep _resolve() machinery returns them);
|
||||
* stub run_job() with a small recorder that captures the changes payload.
|
||||
|
|
@ -22,7 +22,7 @@ from reme4.components.file_parser import ChunkedFileParser
|
|||
from reme4.components.file_store import LocalFileStore
|
||||
from reme4.components.runtime_context import RuntimeContext
|
||||
from reme4.schema import Response
|
||||
from reme4.steps.background import UpdateStoreStep, WatchChangesStep
|
||||
from reme4.steps.background import ScanChangesStep, WatchChangesStep
|
||||
|
||||
warnings.filterwarnings("ignore", category=DeprecationWarning, module="jieba")
|
||||
warnings.filterwarnings("ignore", category=DeprecationWarning, module="pkg_resources")
|
||||
|
|
@ -52,7 +52,7 @@ def write_file(path: Path, content: str = "x") -> Path:
|
|||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# UpdateStoreStep
|
||||
# ScanChangesStep
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
|
|
@ -63,12 +63,12 @@ class _RecorderStep:
|
|||
dispatched: int
|
||||
|
||||
def install_recorder(self):
|
||||
"""Install a fake run_job that records dispatched 'index_changes' payloads."""
|
||||
"""Install a fake run_job that records dispatched 'update_store_index' payloads."""
|
||||
self.recorded = []
|
||||
self.dispatched = 0
|
||||
|
||||
async def fake_run_job(name: str, **kwargs: Any):
|
||||
assert name == "index_changes"
|
||||
assert name == "update_store_index"
|
||||
self.recorded = kwargs.get("changes") or []
|
||||
self.dispatched += 1
|
||||
return Response()
|
||||
|
|
@ -77,23 +77,25 @@ class _RecorderStep:
|
|||
self.run_job = fake_run_job # type: ignore[assignment]
|
||||
|
||||
|
||||
class _RecordingUpdateStoreStep(UpdateStoreStep, _RecorderStep):
|
||||
class _RecordingScanChangesStep(ScanChangesStep, _RecorderStep):
|
||||
pass
|
||||
|
||||
|
||||
async def _make_update_step(
|
||||
async def _make_scan_step(
|
||||
watch_paths: list[str] | str = "vault",
|
||||
suffix_filters: list[str] | None = None,
|
||||
recursive: bool = True,
|
||||
dump: bool = True,
|
||||
) -> tuple[_RecordingUpdateStoreStep, RuntimeContext, LocalFileStore, ChunkedFileParser]:
|
||||
dump_store_index: bool = True,
|
||||
dispatch_job: str = "update_store_index",
|
||||
) -> tuple[_RecordingScanChangesStep, RuntimeContext, LocalFileStore, ChunkedFileParser]:
|
||||
fs = LocalFileStore(store_name="test_store", embedding_model="")
|
||||
parser = ChunkedFileParser()
|
||||
await fs.start()
|
||||
await parser.start()
|
||||
step = _RecordingUpdateStoreStep(
|
||||
step = _RecordingScanChangesStep(
|
||||
recursive=recursive,
|
||||
dump=dump,
|
||||
dump_store_index=dump_store_index,
|
||||
dispatch_job=dispatch_job,
|
||||
file_store=fs,
|
||||
file_parser=parser,
|
||||
)
|
||||
|
|
@ -110,7 +112,7 @@ async def _teardown(fs: LocalFileStore, parser: ChunkedFileParser) -> None:
|
|||
await fs.close()
|
||||
|
||||
|
||||
def test_update_store_initial_all_added():
|
||||
def test_scan_changes_initial_all_added():
|
||||
"""First run on a fresh store emits 'added' for every existing file (abs paths)."""
|
||||
|
||||
async def run():
|
||||
|
|
@ -121,7 +123,7 @@ def test_update_store_initial_all_added():
|
|||
vault = cwd / "vault"
|
||||
write_file(vault / "a.md", "alpha")
|
||||
write_file(vault / "b.md", "beta")
|
||||
step, ctx, fs, parser = await _make_update_step()
|
||||
step, ctx, fs, parser = await _make_scan_step()
|
||||
try:
|
||||
resp = await step(ctx)
|
||||
counts = resp.metadata["counts"]
|
||||
|
|
@ -134,12 +136,12 @@ def test_update_store_initial_all_added():
|
|||
assert paths == expected
|
||||
finally:
|
||||
await _teardown(fs, parser)
|
||||
print("✓ test_update_store_initial_all_added passed")
|
||||
print("✓ test_scan_changes_initial_all_added passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_update_store_no_changes_skips_dispatch():
|
||||
def test_scan_changes_no_changes_skips_dispatch():
|
||||
"""A second run over an unchanged store reports zero counts and does not dispatch."""
|
||||
|
||||
async def run():
|
||||
|
|
@ -147,7 +149,7 @@ def test_update_store_no_changes_skips_dispatch():
|
|||
cwd = Path.cwd()
|
||||
vault = cwd / "vault"
|
||||
a = write_file(vault / "a.md", "alpha")
|
||||
seed_step, ctx, fs, parser = await _make_update_step()
|
||||
seed_step, ctx, fs, parser = await _make_scan_step()
|
||||
try:
|
||||
node, chunks = await parser.parse(a)
|
||||
await fs.upsert([(node, chunks)])
|
||||
|
|
@ -158,12 +160,12 @@ def test_update_store_no_changes_skips_dispatch():
|
|||
assert seed_step.dispatched == 0
|
||||
finally:
|
||||
await _teardown(fs, parser)
|
||||
print("✓ test_update_store_no_changes_skips_dispatch passed")
|
||||
print("✓ test_scan_changes_no_changes_skips_dispatch passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_update_store_detects_modify_and_delete():
|
||||
def test_scan_changes_detects_modify_and_delete():
|
||||
"""Second pass distinguishes added/modified/deleted; paths are absolute."""
|
||||
|
||||
async def run():
|
||||
|
|
@ -172,7 +174,7 @@ def test_update_store_detects_modify_and_delete():
|
|||
vault = cwd / "vault"
|
||||
a = write_file(vault / "a.md", "alpha")
|
||||
b = write_file(vault / "b.md", "beta")
|
||||
step, ctx, fs, parser = await _make_update_step()
|
||||
step, ctx, fs, parser = await _make_scan_step()
|
||||
try:
|
||||
# Seed via direct parse/upsert.
|
||||
for p in (a, b):
|
||||
|
|
@ -194,25 +196,25 @@ def test_update_store_detects_modify_and_delete():
|
|||
assert by_kind["deleted"] == str(b)
|
||||
finally:
|
||||
await _teardown(fs, parser)
|
||||
print("✓ test_update_store_detects_modify_and_delete passed")
|
||||
print("✓ test_scan_changes_detects_modify_and_delete passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
def test_update_store_missing_watch_path_silently_skipped():
|
||||
def test_scan_changes_missing_watch_path_silently_skipped():
|
||||
"""Non-existent watch_paths entries are dropped silently."""
|
||||
|
||||
async def run():
|
||||
with tempfile.TemporaryDirectory() as tmpdir, temp_chdir(tmpdir):
|
||||
(Path(tmpdir) / "vault").mkdir()
|
||||
step, ctx, fs, parser = await _make_update_step(watch_paths=["vault", "ghost"])
|
||||
step, ctx, fs, parser = await _make_scan_step(watch_paths=["vault", "ghost"])
|
||||
try:
|
||||
resp = await step(ctx)
|
||||
assert resp.metadata["counts"] == {"added": 0, "modified": 0, "deleted": 0}
|
||||
assert step.dispatched == 0
|
||||
finally:
|
||||
await _teardown(fs, parser)
|
||||
print("✓ test_update_store_missing_watch_path_silently_skipped passed")
|
||||
print("✓ test_scan_changes_missing_watch_path_silently_skipped passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
|
@ -283,11 +285,11 @@ def test_watch_changes_filter_only_passes_md():
|
|||
|
||||
if __name__ == "__main__":
|
||||
print("\n=== Background Steps Tests ===")
|
||||
# UpdateStoreStep
|
||||
test_update_store_initial_all_added()
|
||||
test_update_store_no_changes_skips_dispatch()
|
||||
test_update_store_detects_modify_and_delete()
|
||||
test_update_store_missing_watch_path_silently_skipped()
|
||||
# ScanChangesStep
|
||||
test_scan_changes_initial_all_added()
|
||||
test_scan_changes_no_changes_skips_dispatch()
|
||||
test_scan_changes_detects_modify_and_delete()
|
||||
test_scan_changes_missing_watch_path_silently_skipped()
|
||||
# WatchChangesStep
|
||||
test_watch_changes_requires_stop_event()
|
||||
test_watch_changes_raises_when_no_valid_paths()
|
||||
|
|
|
|||
160
tests4/unittest/test_file_catalog.py
Normal file
160
tests4/unittest/test_file_catalog.py
Normal file
|
|
@ -0,0 +1,160 @@
|
|||
"""Tests for FileCatalog backends."""
|
||||
|
||||
# pylint: disable=protected-access
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
import tempfile
|
||||
|
||||
import pytest
|
||||
|
||||
from reme4.components.file_catalog import LocalFileCatalog
|
||||
from reme4.schema import FileNode
|
||||
|
||||
|
||||
class temp_chdir:
|
||||
"""Context manager to temporarily chdir into a path and restore on exit."""
|
||||
|
||||
def __init__(self, path):
|
||||
self.path = path
|
||||
self.old = None
|
||||
|
||||
def __enter__(self):
|
||||
self.old = os.getcwd()
|
||||
os.chdir(self.path)
|
||||
return self
|
||||
|
||||
def __exit__(self, *exc):
|
||||
os.chdir(self.old)
|
||||
|
||||
|
||||
def make_node(path: str, mtime: float = 1.0) -> FileNode:
|
||||
return FileNode(path=path, st_mtime=mtime)
|
||||
|
||||
|
||||
# All backends should satisfy the same BaseFileCatalog contract.
|
||||
BACKENDS = [LocalFileCatalog]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("backend_cls", BACKENDS)
|
||||
def test_upsert_and_get_nodes(backend_cls):
|
||||
"""upsert stores nodes; get_nodes returns by path or all."""
|
||||
|
||||
async def run():
|
||||
with tempfile.TemporaryDirectory() as tmpdir, temp_chdir(tmpdir):
|
||||
catalog = backend_cls()
|
||||
await catalog.start()
|
||||
|
||||
await catalog.upsert([make_node("a.md"), make_node("b.md")])
|
||||
|
||||
got_all = await catalog.get_nodes()
|
||||
assert {n.path for n in got_all} == {"a.md", "b.md"}
|
||||
|
||||
got_one = await catalog.get_nodes(["a.md"])
|
||||
assert len(got_one) == 1
|
||||
assert got_one[0].path == "a.md"
|
||||
|
||||
got_missing = await catalog.get_nodes(["nope.md"])
|
||||
assert got_missing == []
|
||||
|
||||
await catalog.close()
|
||||
print(f"✓ test_upsert_and_get_nodes[{backend_cls.__name__}] passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
@pytest.mark.parametrize("backend_cls", BACKENDS)
|
||||
def test_upsert_replaces_existing(backend_cls):
|
||||
"""Re-upserting a node with the same path overwrites the prior entry."""
|
||||
|
||||
async def run():
|
||||
with tempfile.TemporaryDirectory() as tmpdir, temp_chdir(tmpdir):
|
||||
catalog = backend_cls()
|
||||
await catalog.start()
|
||||
|
||||
await catalog.upsert([make_node("a.md", mtime=1.0)])
|
||||
await catalog.upsert([make_node("a.md", mtime=2.0)])
|
||||
|
||||
nodes = await catalog.get_nodes(["a.md"])
|
||||
assert len(nodes) == 1
|
||||
assert nodes[0].st_mtime == 2.0
|
||||
|
||||
await catalog.close()
|
||||
print(f"✓ test_upsert_replaces_existing[{backend_cls.__name__}] passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
@pytest.mark.parametrize("backend_cls", BACKENDS)
|
||||
def test_delete_single_and_list(backend_cls):
|
||||
"""delete accepts both a single path and a list; missing paths are no-ops."""
|
||||
|
||||
async def run():
|
||||
with tempfile.TemporaryDirectory() as tmpdir, temp_chdir(tmpdir):
|
||||
catalog = backend_cls()
|
||||
await catalog.start()
|
||||
|
||||
await catalog.upsert([make_node("a.md"), make_node("b.md"), make_node("c.md")])
|
||||
|
||||
await catalog.delete("a.md")
|
||||
assert {n.path for n in await catalog.get_nodes()} == {"b.md", "c.md"}
|
||||
|
||||
await catalog.delete(["b.md", "ghost.md"])
|
||||
assert {n.path for n in await catalog.get_nodes()} == {"c.md"}
|
||||
|
||||
await catalog.close()
|
||||
print(f"✓ test_delete_single_and_list[{backend_cls.__name__}] passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
@pytest.mark.parametrize("backend_cls", BACKENDS)
|
||||
def test_get_nodes_empty_inputs(backend_cls):
|
||||
"""get_nodes([]) returns []; get_nodes(None) returns all."""
|
||||
|
||||
async def run():
|
||||
with tempfile.TemporaryDirectory() as tmpdir, temp_chdir(tmpdir):
|
||||
catalog = backend_cls()
|
||||
await catalog.start()
|
||||
|
||||
await catalog.upsert([make_node("a.md")])
|
||||
assert await catalog.get_nodes([]) == []
|
||||
assert len(await catalog.get_nodes(None)) == 1
|
||||
|
||||
await catalog.close()
|
||||
print(f"✓ test_get_nodes_empty_inputs[{backend_cls.__name__}] passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
@pytest.mark.parametrize("backend_cls", BACKENDS)
|
||||
def test_persistence_roundtrip(backend_cls):
|
||||
"""close() dumps; a fresh instance loads the same nodes from disk."""
|
||||
|
||||
async def run():
|
||||
with tempfile.TemporaryDirectory() as tmpdir, temp_chdir(tmpdir):
|
||||
c1 = backend_cls()
|
||||
await c1.start()
|
||||
await c1.upsert([make_node("a.md", mtime=10.0), make_node("b.md", mtime=20.0)])
|
||||
await c1.close()
|
||||
|
||||
c2 = backend_cls()
|
||||
await c2.start()
|
||||
nodes = sorted(await c2.get_nodes(), key=lambda n: n.path)
|
||||
assert [n.path for n in nodes] == ["a.md", "b.md"]
|
||||
assert [n.st_mtime for n in nodes] == [10.0, 20.0]
|
||||
await c2.close()
|
||||
print(f"✓ test_persistence_roundtrip[{backend_cls.__name__}] passed")
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
print("\n=== FileCatalog Tests ===")
|
||||
for backend in BACKENDS:
|
||||
test_upsert_and_get_nodes(backend)
|
||||
test_upsert_replaces_existing(backend)
|
||||
test_delete_single_and_list(backend)
|
||||
test_get_nodes_empty_inputs(backend)
|
||||
test_persistence_roundtrip(backend)
|
||||
print("\n所有测试通过!")
|
||||
Loading…
Add table
Reference in a new issue