Files
temp/src/cleveragents/actor/loader.py
aditya e476d2de0e feat(actor): extend hierarchical actor YAML schema and loader
Extend actor YAML schema to support hierarchical graphs with explicit
node types (agent, tool, conditional, subgraph), per-node LSP bindings
(lsp_binding with server, languages, auto, capabilities), and tool-source
references (skills, mcp_servers, agent_skills).

Add schema validation for namespaced actor references, duplicate node IDs,
edge target existence, and graph reachability — all nodes must be reachable
from entry_node via explicit edges or conditional node routing targets.

Update loader to report YAML parse errors with precise line/column positions
and schema validation errors with dotted field paths and remediation hints
pointing to docs/reference/actor_config.md.

Add docs/reference/actor_config.md as the practical configuration reference
covering hierarchical graph examples, node type table, topology rules, and
common error cases with fix guidance.

Refresh examples/actors/graph_workflow.yaml to replace deprecated actor_path
with actor_ref. Add benchmarks/actor_yaml_bench.py for schema load overhead.

Tests: 95 Behave scenarios, 10 Robot smoke tests (including hierarchical
loader smoke test), security scan clean, coverage 99% (threshold 97%).

ISSUES CLOSED: #157
2026-02-25 10:30:42 +00:00

303 lines
11 KiB
Python

"""Actor loader for discovering and caching actor YAML configurations.
Discovers actor YAML files from configurable search roots, parses them
against the v3 ``ActorConfigSchema``, applies namespace defaults, and
caches results using content hashing to avoid redundant reloads.
Operations:
| Method | Description |
|--------------------------|------------------------------------------------|
| ``discover()`` | Scan search roots and load/reload actors |
| ``get(name)`` | Retrieve a loaded actor by namespaced name |
| ``list_actors(...)`` | List actors with optional namespace filter |
| ``clear()`` | Drop all cached actors and hashes |
All mutating operations are protected by ``threading.RLock``.
"""
from __future__ import annotations
import hashlib
import logging
import threading
from pathlib import Path
from typing import Any
import yaml
from cleveragents.actor.schema import ActorConfigSchema
from cleveragents.core.exceptions import ValidationError
from cleveragents.tool.registry import ToolRegistry
logger = logging.getLogger(__name__)
_DEFAULT_NAMESPACE = "local"
_YAML_SUFFIXES = frozenset({".yaml", ".yml"})
class _CacheEntry:
"""Internal cache entry storing an actor config alongside its content hash."""
__slots__ = ("config", "content_hash", "load_count", "source_path")
def __init__(
self,
config: ActorConfigSchema,
content_hash: str,
source_path: Path,
) -> None:
self.config = config
self.content_hash = content_hash
self.source_path = source_path
self.load_count = 1
def _compute_hash(content: bytes) -> str:
return hashlib.sha256(content).hexdigest()
def _normalize_name(raw_name: str) -> str:
"""Apply ``local/`` default when the name has no namespace slash."""
if "/" in raw_name:
return raw_name
return f"{_DEFAULT_NAMESPACE}/{raw_name}"
class ActorLoader:
"""Discover, validate, and cache actor YAML configurations.
Parameters
----------
search_roots:
Directories to scan for ``*.yaml`` / ``*.yml`` actor files.
tool_registry:
Optional tool registry used to verify tool references at load time.
"""
def __init__(
self,
search_roots: list[Path],
tool_registry: ToolRegistry | None = None,
) -> None:
self._search_roots = [Path(r) for r in search_roots]
self._tool_registry = tool_registry
self._actors: dict[str, _CacheEntry] = {}
self._path_to_name: dict[Path, str] = {}
self._lock = threading.RLock()
self._warnings: list[str] = []
@property
def warnings(self) -> list[str]:
"""Warnings emitted during the last discovery run."""
with self._lock:
return list(self._warnings)
def discover(self) -> list[ActorConfigSchema]:
"""Scan search roots and load or reload actor configs.
Returns the list of all currently loaded actor configs.
Raises
------
ValidationError
When duplicate actor names are found across files or when
a YAML file fails schema validation.
"""
with self._lock:
self._warnings = []
found_files = self._collect_yaml_files()
self._prune_deleted(found_files)
pending: dict[str, list[tuple[Path, ActorConfigSchema]]] = {}
errors: list[str] = []
for path in found_files:
resolved = path.resolve()
content = resolved.read_bytes()
content_hash = _compute_hash(content)
existing_name = self._path_to_name.get(resolved)
if existing_name and existing_name in self._actors:
entry = self._actors[existing_name]
if entry.content_hash == content_hash:
name = existing_name
config = entry.config
pending.setdefault(name, []).append((resolved, config))
continue
try:
raw = yaml.safe_load(content)
except yaml.YAMLError as exc:
mark = getattr(exc, "problem_mark", None)
location = (
f" at line {mark.line + 1}, column {mark.column + 1}"
if mark is not None
else ""
)
problem = getattr(exc, "problem", str(exc))
errors.append(f"Invalid YAML in {path.name}{location}: {problem}")
continue
if not isinstance(raw, dict):
errors.append(
f"Expected mapping in {path.name}, got {type(raw).__name__}"
)
continue
raw = self._apply_namespace_default(raw)
try:
config = ActorConfigSchema.model_validate(raw)
except Exception as exc:
from pydantic import ValidationError as PydanticValidationError
if isinstance(exc, PydanticValidationError):
field_errors = []
for err in exc.errors():
field_path = ".".join(str(loc) for loc in err["loc"])
field_errors.append(f" {field_path}: {err['msg']}")
detail = "\n".join(field_errors)
hint = (
" Hint: see docs/reference/actor_config.md "
"for the correct schema format."
)
errors.append(
f"Schema validation failed for {path.name}:"
f"\n{detail}\n{hint}"
)
else:
errors.append(
f"Schema validation failed for {path.name}: {exc}"
)
continue
name = config.name
pending.setdefault(name, []).append((resolved, config))
duplicate_msgs: list[str] = []
for name, entries in pending.items():
if len(entries) > 1:
paths_str = ", ".join(str(p) for p, _ in entries)
duplicate_msgs.append(
f"Duplicate actor '{name}' found in: {paths_str}"
)
if duplicate_msgs:
all_errors = errors + duplicate_msgs
raise ValidationError(
"Actor discovery failed:\n" + "\n".join(all_errors),
details={"duplicates": duplicate_msgs, "errors": errors},
)
if errors:
raise ValidationError(
"Actor discovery failed:\n" + "\n".join(errors),
details={"errors": errors},
)
new_actors: dict[str, _CacheEntry] = {}
new_path_map: dict[Path, str] = {}
for name, entries in pending.items():
resolved_path, config = entries[0]
content_hash = _compute_hash(resolved_path.read_bytes())
old_entry = self._actors.get(name)
if old_entry and old_entry.content_hash == content_hash:
new_actors[name] = old_entry
else:
entry = _CacheEntry(
config=config,
content_hash=content_hash,
source_path=resolved_path,
)
if old_entry:
entry.load_count = old_entry.load_count + 1
new_actors[name] = entry
new_path_map[resolved_path] = name
self._resolve_tools(config)
self._actors = new_actors
self._path_to_name = new_path_map
return [e.config for e in self._actors.values()]
def get(self, name: str) -> ActorConfigSchema | None:
"""Retrieve a loaded actor by its namespaced name."""
with self._lock:
entry = self._actors.get(name)
return entry.config if entry else None
def list_actors(
self,
namespace: str | None = None,
) -> list[ActorConfigSchema]:
"""List loaded actors with optional namespace filter.
Parameters
----------
namespace:
If provided, only return actors whose name starts with
``namespace/``.
"""
with self._lock:
configs = [e.config for e in self._actors.values()]
if namespace is not None:
prefix = f"{namespace}/"
configs = [c for c in configs if c.name.startswith(prefix)]
return configs
def clear(self) -> None:
"""Drop all cached actors and content hashes."""
with self._lock:
self._actors.clear()
self._path_to_name.clear()
self._warnings.clear()
def get_load_count(self, name: str) -> int:
"""Return how many times an actor was loaded from disk (for testing)."""
with self._lock:
entry = self._actors.get(name)
return entry.load_count if entry else 0
def _collect_yaml_files(self) -> list[Path]:
files: list[Path] = []
for root in self._search_roots:
if not root.is_dir():
logger.warning("Search root is not a directory: %s", root)
continue
for path in root.rglob("*"):
if path.is_file() and path.suffix in _YAML_SUFFIXES:
files.append(path)
return sorted(files)
def _prune_deleted(self, current_files: list[Path]) -> None:
resolved_set = {p.resolve() for p in current_files}
stale_paths = [p for p in self._path_to_name if p not in resolved_set]
for stale in stale_paths:
name = self._path_to_name.pop(stale, None)
if name:
self._actors.pop(name, None)
@staticmethod
def _apply_namespace_default(raw: dict[str, Any]) -> dict[str, Any]:
name = raw.get("name")
if isinstance(name, str) and "/" not in name:
raw["name"] = _normalize_name(name)
return raw
def _resolve_tools(self, config: ActorConfigSchema) -> None:
if self._tool_registry is None:
return
for tool_ref in config.tools:
if isinstance(tool_ref, str):
spec = self._tool_registry.get(tool_ref)
if spec is None:
msg = (
f"Unresolved tool reference '{tool_ref}' "
f"in actor '{config.name}'"
)
logger.warning(msg)
self._warnings.append(msg)