Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@

## 4. Wiki Restructuring

- [ ] **Unify wiki system** — merge `file_wiki.py`, `user_wiki.py`, `agent_wiki.py` into one `WikiManager` with layer-based separation. Current 3 separate classes duplicate logic. Goal: single FTS5 index, shared sync logic, configurable per-layer behavior.
- [x] **Unify wiki system** — merge `file_wiki.py`, `user_wiki.py`, `agent_wiki.py` into one `WikiManager` with layer-based separation. Current 3 separate classes duplicate logic. Goal: single FTS5 index, shared sync logic, configurable per-layer behavior.

## 4. Bug Fixes & Cleanup

Expand Down Expand Up @@ -161,5 +161,5 @@

---

**Completed:** 35/65 items
**Last updated:** 2026-07-03
**Completed:** 36/65 items
**Last updated:** 2026-07-04
2 changes: 2 additions & 0 deletions features/backup.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
from typing import Any, Optional

from config import config
from shared.path_safety import safe_resolve


class BackupManager:
Expand Down Expand Up @@ -50,6 +51,7 @@ async def restore(self, backup_name: str) -> dict[str, Any]:

restored = []
for db_file in manifest.get("files", []):
safe_resolve(self.base_dir, db_file) # raises ValueError if traversal
backup_file = src / db_file
if backup_file.exists():
shutil.copy2(backup_file, self.base_dir / db_file)
Expand Down
10 changes: 7 additions & 3 deletions features/backup_cron.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
from typing import Any, Optional

from config import config
from shared.path_safety import safe_resolve

logger = logging.getLogger(__name__)

Expand All @@ -36,7 +37,9 @@ def __init__(self, base_dir: Optional[str] = None):
def _load_state(self):
if self._state_file.exists():
try:
state = json.loads(self._state_file.read_text(encoding="utf-8"))
from shared.saga_crypto import read_state_legacy_or_encrypted

state = read_state_legacy_or_encrypted(self._state_file)
self._last_backup = state.get("last_backup", 0.0)
self._last_wiki_sync = state.get("last_wiki_sync", 0.0)
except Exception:
Expand Down Expand Up @@ -131,10 +134,10 @@ def _cleanup_old(self):
def _sync_wiki(self):
"""Synchronize wiki files with disk."""
try:
from wiki.file_wiki import FileWiki
from wiki.manager import WikiManager

for layer in ["user", "agent"]:
fw = FileWiki(layer=layer)
fw = WikiManager(layer=layer)
raw = fw.reindex_all()
result: dict[str, Any] = asyncio.run(raw) if asyncio.iscoroutine(raw) else raw
if isinstance(result, dict) and result.get("indexed", 0) > 0:
Expand Down Expand Up @@ -162,6 +165,7 @@ def restore(self, backup_name: str) -> dict[str, Any]:

restored = []
for db_file in manifest.get("files", []):
safe_resolve(self.base_dir, db_file) # raises ValueError if traversal
if db_file.endswith("/"):
# Restore wiki directory
src_wiki = src / db_file
Expand Down
6 changes: 3 additions & 3 deletions features/dashboard.py
Original file line number Diff line number Diff line change
Expand Up @@ -171,10 +171,10 @@ def __init__(self, mm=None, data_dir: Optional[str] = None):

async def get_stats(self, user_id: str = "default") -> dict[str, Any]:
from graph.epistemic import EpistemicGraph
from wiki.file_wiki import FileWiki
from wiki.manager import WikiManager

uw = FileWiki(layer="user")
aw = FileWiki(layer="agent")
uw = WikiManager(layer="user")
aw = WikiManager(layer="agent")
ug = EpistemicGraph(layer="user")

um = self.mm.user_memory(user_id)
Expand Down
2 changes: 2 additions & 0 deletions features/import_export.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from typing import Any, Optional

from shared.connection import AsyncConnectionManager, connection_manager
from shared.path_safety import safe_resolve


class ImportExport:
Expand Down Expand Up @@ -58,6 +59,7 @@ async def export_user(self, user_id: str) -> str:
return str(filepath)

async def import_user(self, filepath: str, target_user_id: Optional[str] = None) -> dict[str, int]:
safe_resolve(self.export_dir, filepath) # raises ValueError if traversal
data = json.loads(Path(filepath).read_text(encoding="utf-8"))
user_id = target_user_id or data.get("user_id", "default")
imported = {"core_memory": 0, "episodes": 0}
Expand Down
69 changes: 46 additions & 23 deletions lifecycle/emotion_trigger.py
Original file line number Diff line number Diff line change
Expand Up @@ -158,47 +158,70 @@ class EmotionTrigger:
def should_save(self, message: str, emotional_state: Optional[dict] = None, state_delta: Optional[dict] = None) -> tuple[bool, str, float]:
msg_lower = message.lower()

# 1. Phrase patterns (high priority)
all_patterns = PHRASE_PATTERNS + PHRASE_PATTERNS_EN
for pattern, emotion, weight in all_patterns:
result = self._check_phrase_patterns(msg_lower)
if result:
return result

result = self._check_emotion_markers(msg_lower)
if result:
return result

result = self._check_emoji(message)
if result:
return result

result = self._check_emotional_state(emotional_state)
if result:
return result

result = self._check_state_shift(state_delta)
if result:
return result

if len(message) > 300:
return True, "long_message", 0.3

if message.count("?") >= 3:
return True, "complex_question", 0.4

if message.count("!") >= 2:
return True, "exclamation", 0.3

return False, "", 0.0

def _check_phrase_patterns(self, msg_lower: str) -> tuple[bool, str, float] | None:
for pattern, emotion, weight in PHRASE_PATTERNS + PHRASE_PATTERNS_EN:
if re.search(pattern, msg_lower):
return True, f"emotion_{emotion}", weight
return None

# 2. Emotion markers
def _check_emotion_markers(self, msg_lower: str) -> tuple[bool, str, float] | None:
high_weight = ("love", "fear", "anger")
for emotion, markers in EMOTION_MARKERS.items():
for marker in markers:
if marker in msg_lower:
weight = 0.7 if emotion in ("love", "fear", "anger") else 0.5
weight = 0.7 if emotion in high_weight else 0.5
return True, f"emotion_{emotion}", weight
return None

# 3. Emoji
def _check_emoji(self, message: str) -> tuple[bool, str, float] | None:
high_weight = ("love", "fear", "anger")
for emotion, emojis in EMOJI_MARKERS.items():
for emoji in emojis:
if emoji in message:
weight = 0.7 if emotion in ("love", "fear", "anger") else 0.5
weight = 0.7 if emotion in high_weight else 0.5
return True, f"emotion_{emotion}", weight
return None

# 4. Emotional state from context
def _check_emotional_state(self, emotional_state: Optional[dict]) -> tuple[bool, str, float] | None:
if emotional_state:
if emotional_state.get("joy", 0) > 0.8 or emotional_state.get("interest", 0) > 0.8:
return True, "high_emotion", 0.6
return None

# 5. State shift
def _check_state_shift(self, state_delta: Optional[dict]) -> tuple[bool, str, float] | None:
if state_delta:
for key, delta in state_delta.items():
if abs(delta) > STATE_SHIFT_THRESHOLD:
return True, f"state_shift_{key}", 0.4

# 6. Long message
if len(message) > 300:
return True, "long_message", 0.3

# 7. Multiple questions
if message.count("?") >= 3:
return True, "complex_question", 0.4

# 8. Exclamation marks
if message.count("!") >= 2:
return True, "exclamation", 0.3

return False, "", 0.0
return None
6 changes: 3 additions & 3 deletions mcp_server/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,15 +35,15 @@
from rag.multi_source import MultiSourceRAG
from shared.cache import MemoryCache
from shared.read_only import read_only_replica
from wiki.file_wiki import FileWiki
from wiki.manager import WikiManager


class AppContext:
def __init__(self):
self.cache = MemoryCache()
self.mm = MemoryManager(cache=self.cache)
self.user_wiki = FileWiki(layer="user")
self.agent_wiki = FileWiki(layer="agent")
self.user_wiki = WikiManager(layer="user")
self.agent_wiki = WikiManager(layer="agent")
self.user_rag = RAGEngine(layer="user")
self.agent_rag = RAGEngine(layer="agent")
self.user_multi = MultiSourceRAG(self.user_rag, self.user_wiki)
Expand Down
18 changes: 18 additions & 0 deletions shared/path_safety.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
"""Path traversal prevention — shared guard for all file-accepting functions."""

import os
from pathlib import Path


def safe_resolve(base: Path, user_input: str) -> Path:
"""Resolve user_input relative to base, raising ValueError if it escapes.

Checks both the base-relative resolution and the real path (follows symlinks).
"""
base_resolved = base.resolve()
target = (base / user_input).resolve()

if not str(target).startswith(str(base_resolved) + os.sep) and target != base_resolved:
raise ValueError(f"Path escapes base directory: {user_input!r}")

return target
52 changes: 30 additions & 22 deletions shared/saga.py
Original file line number Diff line number Diff line change
Expand Up @@ -130,16 +130,15 @@ def _save_state(self):

def _load_state(self, saga_id: str) -> dict | None:
"""Load state from disk (supports encrypted and legacy plain JSON)."""
from shared.saga_crypto import read_state_legacy_or_encrypted

state_file = SAGA_DIR / (saga_id + ".json")
if state_file.exists():
try:
blob = state_file.read_bytes()
if _HAS_ENCRYPTION and is_encrypted_blob(state_file):
return decrypt_json(blob)
return json.loads(blob.decode("utf-8"))
except Exception:
pass
return None
if not state_file.exists():
return None
try:
return read_state_legacy_or_encrypted(state_file)
except Exception:
return None

def _cleanup_state(self):
"""Delete state file after completion."""
Expand Down Expand Up @@ -339,23 +338,32 @@ async def _compensate(self, failed_step: int) -> None:
continue

if isinstance(step.action, Saga):
inner = step.action
for j in range(len(inner._steps) - 1, -1, -1):
inner_step = inner._steps[j]
if inner_step.status == SagaStatus.COMPLETED and inner_step.compensation:
try:
await inner_step.compensation(inner_step.data)
logger.info("Saga '%s' compensated inner step '%s'" % (self.name, inner_step.name))
except Exception as e:
logger.error("Saga '%s' inner compensation failed for '%s': %s" % (self.name, inner_step.name, e))
await self._compensate_inner_saga(step.action)
elif step.compensation:
await self._compensate_step(step)

self._status = SagaStatus.COMPENSATED

async def _compensate_inner_saga(self, inner: "Saga") -> None:
"""Compensate all completed steps of a nested saga in reverse order."""
for j in range(len(inner._steps) - 1, -1, -1):
inner_step = inner._steps[j]
if inner_step.status == SagaStatus.COMPLETED and inner_step.compensation:
try:
await step.compensation(step.data)
logger.info("Saga '%s' compensated step '%s'" % (self.name, step.name))
await inner_step.compensation(inner_step.data)
logger.info("Saga '%s' compensated inner step '%s'" % (self.name, inner_step.name))
except Exception as e:
logger.error("Saga '%s' compensation failed for '%s': %s" % (self.name, step.name, e))
logger.error("Saga '%s' inner compensation failed for '%s': %s" % (self.name, inner_step.name, e))

self._status = SagaStatus.COMPENSATED
async def _compensate_step(self, step: SagaStep) -> None:
"""Run compensation for a single step, logging success or failure."""
if not step.compensation:
return
try:
await step.compensation(step.data)
logger.info("Saga '%s' compensated step '%s'" % (self.name, step.name))
except Exception as e:
logger.error("Saga '%s' compensation failed for '%s': %s" % (self.name, step.name, e))

def get_state(self) -> dict:
return {
Expand Down
4 changes: 2 additions & 2 deletions tests/test_all.py
Original file line number Diff line number Diff line change
Expand Up @@ -91,10 +91,10 @@ async def t():


def test_user_wiki():
from wiki.file_wiki import FileWiki
from wiki.manager import WikiManager

async def t():
w = FileWiki(layer="user")
w = WikiManager(layer="user")
path = await w.add("work_notes", "Day 1", "Started project")
assert path is not None
results = await w.search("project")
Expand Down
63 changes: 63 additions & 0 deletions tests/test_features/test_backup_path_safety.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
"""Tests for backup path traversal prevention."""

import json

import pytest

from features.backup import BackupManager


@pytest.fixture
def bm(tmp_path):
data_dir = tmp_path / "data"
data_dir.mkdir()
return BackupManager(base_dir=str(data_dir))


def test_restore_rejects_traversal_in_manifest(bm):
"""Crafted manifest with ../../ in filenames should be rejected."""
# Create a malicious backup directory with crafted manifest
backup_dir = bm.backup_dir / "malicious"
backup_dir.mkdir()
manifest = {
"files": ["../../etc/crontab", "memory.db"],
"created_at": "2026-01-01T00:00:00",
}
(backup_dir / "manifest.json").write_text(json.dumps(manifest), encoding="utf-8")

import asyncio

with pytest.raises(ValueError, match="escapes base directory"):
asyncio.run(bm.restore("malicious"))


def test_restore_rejects_absolute_path(bm):
"""Manifest with absolute path should be rejected."""
backup_dir = bm.backup_dir / "absolute"
backup_dir.mkdir()
manifest = {
"files": ["/etc/passwd"],
"created_at": "2026-01-01T00:00:00",
}
(backup_dir / "manifest.json").write_text(json.dumps(manifest), encoding="utf-8")

import asyncio

with pytest.raises(ValueError, match="escapes base directory"):
asyncio.run(bm.restore("absolute"))


def test_restore_accepts_valid_files(bm):
"""Valid manifest with normal files should work."""
# Create a real db file
db_file = bm.base_dir / "memory.db"
db_file.write_bytes(b"fake db")

# Create a valid backup
import asyncio

asyncio.run(bm.backup("test_backup"))

# Restore should work
result = asyncio.run(bm.restore("test_backup"))
assert "restored" in result
Loading
Loading