diff --git a/ROADMAP.md b/ROADMAP.md index dde36eac..7c6209ee 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -166,3 +166,44 @@ **Completed:** 38/65 items **Last updated:** 2026-07-05 + +## 17. Migration from agentmemory + +Features to port from agentmemory (rohitg00/agentmemory) for full replacement. + +| Пункт | Фича | Описание | Приоритет | +|-------|------|----------|-----------| +| R1 | Obsidian export | Экспорт памяти в Obsidian markdown с wikilinks | nice-to-have | +| R2 | Mesh sync | P2P синхронизация между инстансами | nice-to-have | +| R3 | Cross-agent sync | Память доступна из Claude Code, Cursor, OpenCode | nice-to-have | +| R4 | Git snapshots | Версионирование памяти через git commit/diff/rollback | nice-to-have | + +### Migration phases + +**Phase 1 (simple, 3-4 days):** +- SHA-256 dedup (5min TTL) +- Circuit breaker (3 errors → open → 30s) +- Token budget (2000 tokens on context_inject) +- Privacy filter (strip secrets before save) + +**Phase 2 (medium, 5-7 days):** +- Slot system (8 pinned memory units) +- Lesson confidence (strengthen/decay) +- Faceted tagging (dimension:value with AND/OR) + +**Phase 3 (complex, 10-14 days):** +- Hybrid search RRF with graph traversal +- Knowledge graph extraction from sessions +- Action graphs with dependencies +- Lease system for multi-agent + +**Phase 4 (hooks, 5-7 days):** +- sync_turn hook (background capture) +- on_memory_write hook (mirror MEMORY.md) +- Diagnostics tool +- Export tool + +--- + +**Completed:** 38/65 items (+4 roadmap) +**Last updated:** 2026-07-07 diff --git a/mcp_server/tools_layer.py b/mcp_server/tools_layer.py index f8bb57fa..8355031d 100644 --- a/mcp_server/tools_layer.py +++ b/mcp_server/tools_layer.py @@ -8,6 +8,7 @@ from shared.constants import DB_NAME import hashlib import logging +import re import time from typing import Any, Optional @@ -24,11 +25,80 @@ StatsResult, ) from mcp_server.registry import _get_ctx, register_tool +from mcp_server.utils.privacy import strip_secrets from shared.metrics import metrics logger = logging.getLogger(__name__) +class _DedupCache: + _doc_ = "SHA-256 dedup with TTL and periodic cleanup." + + def __init__(self, ttl=300, max_size=10000): + self._cache = {} + self._ttl = ttl + self._max_size = max_size + self._last_cleanup = time.time() + + def _cleanup(self): + now = time.time() + if now - self._last_cleanup < 60: + return + self._last_cleanup = now + expired = [k for k, v in self._cache.items() if now - v > self._ttl] + for k in expired: + del self._cache[k] + if len(self._cache) > self._max_size: + oldest = sorted(self._cache.keys(), key=lambda k: self._cache[k])[: len(self._cache) // 4] + for k in oldest: + del self._cache[k] + + def is_duplicate(self, session_id, tool, input_text): + self._cleanup() + key = hashlib.sha256(f"{session_id}:{tool}:{input_text[:500]}".encode()).hexdigest() + now = time.time() + if key in self._cache and now - self._cache[key] < self._ttl: + return True + self._cache[key] = now + return False + + +_dedup_cache = _DedupCache(ttl=300, max_size=10000) + + +# Token budget configuration +DEFAULT_TOKEN_BUDGET = 2000 +CHARS_PER_TOKEN = 4 + + +def _estimate_tokens(text: str) -> int: + if not text: + return 0 + cjk_count = len(re.findall(r"[\u4e00-\u9fff\u3040-\u309f\u30a0-\u30ff]", text)) + remaining_chars = len(text) - cjk_count + non_cjk_tokens = remaining_chars // CHARS_PER_TOKEN + return cjk_count + non_cjk_tokens + + +def _truncate_to_budget(text: str, max_tokens: int) -> tuple: + estimated = _estimate_tokens(text) + if estimated <= max_tokens: + return text, False + char_limit = max_tokens * CHARS_PER_TOKEN + lines = text.split("\\n") + result_lines = [] + current_len = 0 + for line in lines: + line_len = len(line) + 1 + if current_len + line_len > char_limit: + break + result_lines.append(line) + current_len += line_len + truncated = "\\n".join(result_lines) + truncated += "\\n[...truncated to token budget]" + return truncated, True + + def _get_memory(app, layer: str, user_id: str): if layer == "agent": return app.mm.agent_memory(user_id) @@ -135,6 +205,7 @@ async def memory_remember( key: str = "", value: str = "", importance: float = 0.5, + session_id: str = "", ctx: Optional[Context] = None, ) -> dict: """Save a fact to long-term memory (L4 CoreMemory). @@ -145,7 +216,14 @@ async def memory_remember( key: Fact key (e.g. "name", "language", "principle") value: Fact value importance: Importance score 0.0-1.0 (default 0.5) + session_id: Session ID for dedup (optional) """ + # SHA-256 dedup: skip identical calls within 5min window + value = strip_secrets(value) + if session_id and _dedup_cache.is_duplicate(session_id, key, value): + logger.info("Dedup: skipping identical remember key=%s user=%s", key, user_id) + return RememberResult(status="skipped", reason="duplicate_within_ttl").dict() + app = _get_ctx(ctx) layer = _validate_layer(layer) metrics.inc("tool_calls") @@ -802,17 +880,24 @@ async def memory_context_inject( if facts_text: context_parts.append("REMEMBER: " + facts_text) + # Apply token budget + context_text = "\n".join(context_parts) + context_text, was_truncated = _truncate_to_budget(context_text, DEFAULT_TOKEN_BUDGET) + result = { - "context": "\n".join(context_parts), + "context": context_text, "l4_facts_count": len(l4_facts), "l3_episodes_count": len(l3_episodes), "l1_recent_count": len(l1_recent), "wiki_count": len(wiki_entries), + "estimated_tokens": _estimate_tokens(context_text), + "was_truncated": was_truncated, + "token_budget": DEFAULT_TOKEN_BUDGET, } _set_cached(cache_key, result) # Trigger dream_buffer hook for context staging - await _fire_hook("dream_buffer", layer, {"text": "\n".join(context_parts), "user_id": user_id}) + await _fire_hook("dream_buffer", layer, {"text": context_text, "user_id": user_id}) return result diff --git a/mcp_server/utils/__init__.py b/mcp_server/utils/__init__.py new file mode 100644 index 00000000..620868cc --- /dev/null +++ b/mcp_server/utils/__init__.py @@ -0,0 +1,3 @@ +from mcp_server.utils.circuit_breaker import CircuitBreaker + +__all__ = ["CircuitBreaker"] diff --git a/mcp_server/utils/circuit_breaker.py b/mcp_server/utils/circuit_breaker.py new file mode 100644 index 00000000..09127862 --- /dev/null +++ b/mcp_server/utils/circuit_breaker.py @@ -0,0 +1,170 @@ +""" +Circuit Breaker pattern for LLM/embedding calls. + +States: + - closed: normal operation, requests pass through + - open: failures exceeded threshold, requests blocked + - half-open: recovery probe, one request allowed through + +Usage: + breaker = CircuitBreaker(threshold=3, recovery_timeout=30) + + if not breaker.allow_request(): + return cached_result or fallback + + try: + result = await llm_call() + breaker.record_success() + return result + except Exception as e: + breaker.record_failure() + raise +""" + +import logging +import time +from enum import Enum +from typing import Callable, Optional + +logger = logging.getLogger(__name__) + + +class CircuitState(Enum): + CLOSED = "closed" + OPEN = "open" + HALF_OPEN = "half_open" + + +class CircuitBreaker: + """Circuit breaker with configurable threshold and recovery timeout.""" + + def __init__( + self, + threshold: int = 3, + recovery_timeout: float = 30.0, + name: str = "default", + on_state_change: Optional[Callable] = None, + ): + self.threshold = threshold + self.recovery_timeout = recovery_timeout + self.name = name + + self._failures = 0 + self._state = CircuitState.CLOSED + self._opened_at = 0.0 + self._last_failure_at = 0.0 + + self._on_state_change = on_state_change + + # Metrics + self._total_requests = 0 + self._total_failures = 0 + self._total_rejections = 0 + self._state_changes = 0 + + @property + def state(self) -> CircuitState: + if self._state == CircuitState.OPEN: + if time.time() - self._opened_at > self.recovery_timeout: + self._transition_to(CircuitState.HALF_OPEN) + return self._state + + @property + def failures(self) -> int: + return self._failures + + def _transition_to(self, new_state: CircuitState): + old_state = self._state + self._state = new_state + self._state_changes += 1 + logger.info("CircuitBreaker[%s]: %s -> %s", self.name, old_state.value, new_state.value) + if self._on_state_change: + self._on_state_change(self.name, old_state, new_state) + + def record_success(self): + self._total_requests += 1 + self._failures = 0 + if self._state == CircuitState.HALF_OPEN: + self._transition_to(CircuitState.CLOSED) + + def record_failure(self): + self._total_requests += 1 + self._total_failures += 1 + self._failures += 1 + self._last_failure_at = time.time() + + if self._state == CircuitState.HALF_OPEN: + self._transition_to(CircuitState.OPEN) + self._opened_at = time.time() + elif self._failures >= self.threshold: + self._transition_to(CircuitState.OPEN) + self._opened_at = time.time() + + def allow_request(self) -> bool: + self._total_requests += 1 + current_state = self.state + + if current_state == CircuitState.CLOSED: + return True + if current_state == CircuitState.HALF_OPEN: + return True + self._total_rejections += 1 + return False + + def reset(self): + self._failures = 0 + self._state = CircuitState.CLOSED + self._opened_at = 0.0 + + def get_metrics(self) -> dict: + return { + "name": self.name, + "state": self.state.value, + "failures": self._failures, + "threshold": self.threshold, + "recovery_timeout": self.recovery_timeout, + "total_requests": self._total_requests, + "total_failures": self._total_failures, + "total_rejections": self._total_rejections, + "state_changes": self._state_changes, + "last_failure_at": self._last_failure_at, + } + + def __enter__(self): + self._context_allowed = self.allow_request() + return self._context_allowed + + def __exit__(self, exc_type, exc_val, exc_tb): + if exc_type is None: + self.record_success() + else: + self.record_failure() + return False + + +class CircuitBreakerRegistry: + def __init__(self): + self._breakers: dict[str, CircuitBreaker] = {} + + def get(self, name: str, threshold: int = 3, recovery_timeout: float = 30.0, on_state_change: Optional[Callable] = None) -> CircuitBreaker: + if name not in self._breakers: + self._breakers[name] = CircuitBreaker( + threshold=threshold, + recovery_timeout=recovery_timeout, + name=name, + on_state_change=on_state_change, + ) + return self._breakers[name] + + def get_all(self) -> dict[str, CircuitBreaker]: + return dict(self._breakers) + + def get_all_metrics(self) -> dict: + return {name: breaker.get_metrics() for name, breaker in self._breakers.items()} + + def reset_all(self): + for breaker in self._breakers.values(): + breaker.reset() + + +breaker_registry = CircuitBreakerRegistry() diff --git a/mcp_server/utils/privacy.py b/mcp_server/utils/privacy.py new file mode 100644 index 00000000..e0a40e9f --- /dev/null +++ b/mcp_server/utils/privacy.py @@ -0,0 +1,48 @@ +import re +from typing import List, Pattern + + +_CREDENTIAL_PATTERNS: List[Pattern] = [ + re.compile(r"\b(sk-[A-Za-z0-9]{20,})\b"), + re.compile(r"\b(sk-ant-[A-Za-z0-9-]{20,})\b"), + re.compile(r"\b(ghp_[A-Za-z0-9]{36})\b"), + re.compile(r"\b(gho_[A-Za-z0-9]{36})\b"), + re.compile(r"\b(ghs_[A-Za-z0-9]{36})\b"), + re.compile(r"\b(ghr_[A-Za-z0-9]{36})\b"), + re.compile(r"\b(xox[baprs]-[A-Za-z0-9-]{20,})\b"), + re.compile(r"\b(AKIA[0-9A-Z]{16})\b"), + re.compile(r"\b(AIza[0-9A-Za-z_-]{35})\b"), + re.compile(r"\b(sk_live_[0-9a-zA-Z]{24,})\b"), + re.compile(r"\b(pk_live_[0-9a-zA-Z]{24,})\b"), + re.compile(r"\b(sk_test_[0-9a-zA-Z]{24,})\b"), + re.compile(r"\b([0-9]{10}:[A-Za-z0-9_-]{35})\b"), + re.compile(r"\b(Bearer\s+[A-Za-z0-9_\-\.]{20,})\b", re.IGNORECASE), + re.compile(r".*?", re.DOTALL), + re.compile(r".*?", re.DOTALL), + re.compile(r".*?", re.DOTALL), +] + + +def strip_secrets(text: str, replacement: str = "[REDACTED]") -> str: + if not text: + return text + result = text + for pattern in _CREDENTIAL_PATTERNS: + result = pattern.sub(replacement, result) + return result + + +def has_secrets(text: str) -> bool: + if not text: + return False + for pattern in _CREDENTIAL_PATTERNS: + if pattern.search(text): + return True + return False + + +def get_redacted_preview(text: str, max_length: int = 100) -> str: + redacted = strip_secrets(text) + if len(redacted) > max_length: + return redacted[:max_length] + "..." + return redacted