Рефакторинг Second Brain v2: LangGraph-архитектура#
Дата: 2026-05-27 Статус: ЧЕРНОВИК Автор: планировщик (на основе code-review-2026-05-26.md) Оценка: ~12 рабочих дней
ARCH REVIEW SUMMARY#
Ревьюер: Senior AI Systems Architect
Дата ревью: 2026-05-27
Оценка до ревью: 6/10 — хороший скелет, правильный TDD-подход и security awareness, но критические дыры в реализации
Оценка после ревью: 8.5/10
Найденные слабости и принятые меры#
| # | Слабость | Критичность | Меры |
|---|---|---|---|
| W1 | build_main_graph() — заглушка: узлы не вызывают tools через ToolNode/bind_tools. Граф не делает ничего полезного в production |
🔴 Критично | Добавлена Задача 2.4-bis: полная схема LangGraph ToolNode + bind_tools |
| W2 | socket.getaddrinfo() в check_ssrf_url — синхронный вызов блокирует event loop asyncio |
🔴 Критично | Исправлено: asyncio.to_thread(socket.getaddrinfo, ...) |
| W3 | SQLite без WAL mode → database is locked при конкурентных запросах бота |
🔴 Критично | Добавлен PRAGMA journal_mode=WAL в _ensure_init всех SQLite-компонентов |
| W4 | _ensure_init race condition — несколько корутин одновременно проходят if self._initialized до установки флага |
🔴 Критично | Добавлен asyncio.Lock() для атомарной инициализации |
| W5 | datetime.utcnow() deprecated в Python 3.12 + TZ bug из code-review (#7) |
🔴 Критично | Добавлена Задача 0.5: утилита utc_now(), замена всех utcnow() |
| W6 | _call_llm дублируется в 5 агентах; модель захардкожена claude-3-5-sonnet-20241022 |
🟠 Высокий | Добавлена Задача 0.6: BaseLLMAgent mixin; модель читается из BRAIN_MODEL env |
| W7 | MemoryManager._current_session — shared mutable state. Конкурентные запросы от одного user перезаписывают session_id |
🟠 Высокий | Исправлено: _sessions: dict[int, int] (user_id → session_id) |
| W8 | begin_session() вызывается в каждом process() → 10 сессий в одном диалоге вместо 1 |
🟠 Высокий | Исправлено: lazy session creation + idempotent ensure_session() |
| W9 | WikiTools.search() — O(N) линейный обход всех .md файлов; не использует FTS5 индекс |
🟠 Высокий | Добавлена аннотация + делегирование в SemanticMemory.search() при размере vault > 100 файлов |
| W10 | close_day wiki logic: wiki.search(today) ищет строку даты в тексте, не по mtime |
🟠 Высокий | Исправлено: функция _wiki_updated_today() по os.stat().st_mtime |
| W11 | YouGile pagination отсутствует → молчаливо теряет проекты (code-review bug #11: pagination infinite loop) | 🟠 Высокий | Добавлена пагинация с page/pageSize + защита от infinite loop (max_pages=50) |
| W12 | YouGile timestamp TZ: date.today() — локальная дата, YouGile хранит UTC → phantom completions (code-review bug #10) |
🟠 Высокий | Исправлено: сравнение через datetime.now(tz=timezone.utc).date() |
| W13 | Нет retry strategy для external API calls (YouGile, DDG, httpx) — transient failures сразу роняют операцию | 🟠 Высокий | Добавлена Задача 2.5-retry: tenacity-based retry decorator; добавлен tenacity в зависимости |
| W14 | asyncio.get_event_loop() deprecated в Python 3.10+; не работает в Python 3.12 contexts |
🟠 Высокий | Заменено на asyncio.to_thread() во всех местах |
| W15 | Нет asyncio_mode = "auto" в pytest config → @pytest.mark.asyncio надо ставить вручную или тесты зависают |
🟠 Высокий | Добавлена секция [tool.pytest.ini_options] в pyproject.toml |
| W16 | Skills execution model отсутствует — скиллы устанавливаются в файлы, но нигде не интегрируются в LangChain tools | 🟠 Высокий | Добавлена Задача 6.3: SkillExecutor — загрузка как @tool через importlib |
| W17 | _parse_frontmatter ломается на значениях с : (e.g. description: Hello: World) — неверный парсинг YAML |
🟡 Средний | Исправлено: используем import yaml (PyYAML уже в зависимостях проекта) |
| W18 | SkillLoader.install_from_github пишет через dest.write_text() — не атомарная запись |
🟡 Средний | Исправлено: используем atomic_write_text() |
| W19 | GitHub skills — supply chain риск: нет проверки источника, нет code signing | 🟡 Средний | Добавлено предупреждение + allowlist GitHub orgs + тест sandbox execution |
| W20 | Нет procedural_memory — одна из 4 заявленных типов памяти (working/episodic/semantic/procedural) отсутствует |
🟡 Средний | Добавлена Задача 1.6: ProceduralMemory — хранит SOP/процедуры в SQLite с версионированием |
| W21 | yt-dlp пишет во /tmp/yt_%(id)s — hardcoded path, нет cleanup, не работает на Windows |
🟡 Средний | Исправлено: tempfile.mkdtemp() + cleanup в finally |
| W22 | handle_message нет error boundary → исключение внутри агента роняет обработчик aiogram |
🟡 Средний | Добавлен try/except с логированием и user-friendly fallback |
| W23 | Wiki command: parse_mode="Markdown" — legacy, ломается на спецсимволах; link_preview возвращает Markdown-разметку |
🟡 Средний | Исправлено: parse_mode="HTML", переработан link_preview для HTML output |
| W24 | WorkingMemory._trim() — O(N²): sum(m.tokens for m in self._msgs) вычисляется в каждой итерации while-loop |
🟡 Средний | Исправлено: кэшируем total и вычитаем по мере удаления |
| W25 | Версии зависимостей не зафиксированы точно: langgraph>=0.2.0 — LangGraph менял API в 0.1/0.2/0.3 |
🟡 Средний | Добавлены точные минимальные версии с комментариями о breaking changes |
| W26 | Нет __init__.py упомянуто для новых пакетов core/, memory/, agents/, tools/, skills/ |
🟡 Средний | Добавлен checklist создания __init__.py в Фазе 0 |
| W27 | Нет валидации ANTHROPIC_API_KEY на старте → confusing KeyError при первом вызове LLM |
🟡 Средний | Добавлена startup validation в message_handler.py |
| W28 | Deadlock из code-review (#8): proc.communicate() без timeout не покрыт в плане |
🟡 Средний | Добавлена аннотация в чеклист миграции: явный timeout для subprocess calls |
| W29 | close_day не сохраняет moved/перенесённые задачи — упомянуто в комментарии кода, но нет реализации |
🟡 Средний | Добавлена аннотация + get_moved_tasks() stub в PlanningTools |
| W30 | Оценка времени 12 рабочих дней нереалистична для данного объёма с учётом добавленных задач | 🟡 Средний | Пересмотрена оценка → 18 рабочих дней |
Архитектурные решения добавленные в план#
- Задача 0.5 — TZ-утилиты (
utc_now(), замена всехdatetime.utcnow()) - Задача 0.6 —
BaseLLMAgentmixin (DRY для_call_llm, model env var) - Задача 1.6 —
ProceduralMemory(4-й тип памяти) - Задача 2.4-bis — Полная LangGraph ToolNode интеграция
- Задача 2.5-retry — Retry strategy через
tenacity - Задача 6.3 —
SkillExecutor(замыкает Skills system) - pyproject.toml — добавлены
tenacity,pytest-asyncio asyncio_mode=auto
СТРАТЕГИЯ МИГРАЦИИ#
Принцип: Bot не ломается ни на один день. Новый код растёт рядом со старым.
Старый путь: aiogram → processor.py (God Object)
Новый путь: aiogram → router.py → LangGraph StateGraph → agents → tools
Этапы сосуществования:
- Создаём
src/d_brain/core/иsrc/d_brain/agents/рядом со старым кодом - Добавляем feature-flag
USE_LANGGRAPH=falseв.env - Постепенно переводим команды через
if USE_LANGGRAPHв handlers - Когда все команды переехали — удаляем старый
processor.py - Старые тесты: оставляем
test_planning/иtest_integrations/, удаляемtest_processor.py
ТАБЛИЦА РИСКОВ#
| Риск | Вероятность | Влияние | Митигация |
|---|---|---|---|
| LangGraph API изменится | Низкая | Высокое | Зафиксировать >=0.2.28,<0.4 в pyproject.toml |
| ChromaDB обязательна | Низкая | Низкое | ChromaDB обязателен; FTS5 используется ТОЛЬКО как secondary keyword index внутри SQLite episodic memory, НЕ как замена ChromaDB |
| Токены Claude API дорогие | Высокая | Среднее | Кэш LRU + summary вместо raw history; BRAIN_MODEL=claude-haiku-... для лёгких агентов |
| Потеря данных при миграции памяти | Средняя | Высокое | Миграция JSONL → SQLite с резервной копией |
| Subagents завязаны на thread_id | Низкая | Среднее | thread_id=0 → main agent (graceful fallback) |
| yt-dlp сломается (антибот) | Средняя | Низкое | Fallback на youtube-transcript-api (уже в зависимостях) |
| Path traversal при wiki | Высокая | Критическое | Guardrails в Фазе 0 (первый приоритет) |
| SSRF при парсинге URL | Средняя | Критическое | Allowlist доменов + блокировка RFC1918 + async DNS |
| SQLite database locked | Высокая | Среднее | WAL mode в _ensure_init всех SQLite компонентов |
| _ensure_init race condition | Средняя | Среднее | asyncio.Lock() двойная проверка |
| Skills supply chain атака | Низкая | Высокое | ALLOWED_SKILL_REPOS allowlist + нет exec() без sandbox |
| YouGile pagination infinite loop | Средняя | Высокое | max_pages=50 лимит + страничные параметры |
| TZ phantom completions | Высокая | Среднее | utc_today() вместо date.today() везде |
НОВЫЕ ЗАВИСИМОСТИ#
# pyproject.toml — добавить в [project.dependencies]
# ВНИМАНИЕ: LangGraph 0.2.x → 0.3.x имел breaking changes в StateGraph API.
# 0.2.x: graph.set_entry_point() | 0.3.x: graph.add_node + START node
# Фиксируем: >=0.2.28,<0.4 до стабилизации API
"langgraph>=0.2.28,<0.4",
"langchain-anthropic>=0.3.0", # 0.3+ требует langchain-core>=0.3
"langchain-core>=0.3.15",
"tenacity>=8.2.0", # retry strategy (W13 — ARCH ADDED)
"chromadb>=0.5.0",
"duckduckgo-search>=6.2.0", # 6.0-6.1 имели bugs с rate-limit handling
"langfuse>=2.0.0", # опционально, за флагом LANGFUSE_ENABLED
"aiosqlite>=0.20.0",
"beautifulsoup4>=4.12.0", # парсинг статей
"readability-lxml>=0.8.1", # чистый текст статей
# УЖЕ ЕСТЬ в pyproject.toml (не добавлять повторно):
# httpx>=0.27.0 ← уже есть как "httpx"
# pyyaml>=6.0.3 ← уже есть (нужен для frontmatter parser W17)
# yt-dlp>=2026.3.17 ← уже есть
# youtube-transcript-api>=1.2.4 ← уже есть (лучше yt-dlp для субтитров)
# pytest-asyncio ← уже есть в dev зависимостях
# pyproject.toml — добавить секцию [tool.pytest.ini_options]
[tool.pytest.ini_options]
asyncio_mode = "auto" # все async def тесты автоматически async
testpaths = ["tests"]
python_files = "test_*.py"
python_classes = "Test*"
python_functions = "test_*"
Команда установки:
cd /home/serg/projects/second-brain
uv add "langgraph>=0.2.28,<0.4" "langchain-anthropic>=0.3.0" "langchain-core>=0.3.15" \
"tenacity>=8.2.0" chromadb "duckduckgo-search>=6.2.0" \
aiosqlite beautifulsoup4 readability-lxml
uv add --optional langfuse # опционально
# Проверка версий после установки:
uv run python -c "import langgraph; print(langgraph.__version__)"
uv run python -c "import langchain_anthropic; print(langchain_anthropic.__version__)"
ПРИНЦИП DEEP MODULES (Ousterhout)#
Каждый модуль: узкий интерфейс + толстая реализация.
Плохо: tool_manager.get_tool_by_name_and_execute_with_params(name, params, ctx)
Хорошо: tools.run(name, **kwargs) → внутри: routing, validation, error handling, retry
Правило: если интерфейс модуля занимает больше 5 строк — он слишком тонкий.
АРХИТЕКТУРНАЯ СХЕМА#
Telegram Update
│
▼
aiogram Handler (тонкая обёртка)
│
▼
core/router.py ←── thread_id → agent mapping
│
├──► MainAgent (StateGraph)
│ │
│ ├── tools/wiki_tools.py
│ ├── tools/planning_tools.py
│ ├── tools/memory_tools.py
│ ├── tools/search_tools.py
│ └── tools/media_tools.py
│
├──► CoachAgent (StateGraph)
│ │
│ ├── tools/memory_tools.py
│ └── [coach-specific skills]
│
├──► HabitsAgent (StateGraph)
│
└──► PlannerAgent (StateGraph)
│
└── tools/planning_tools.py (YouGile deep)
Все агенты ──► core/memory.py (MemoryManager)
│
├── memory/working.py (context window)
├── memory/episodic.py (SQLite)
├── memory/semantic.py (ChromaDB/FTS5)
└── memory/consolidator.py (nightly)
ФАЗА 0: ФУНДАМЕНТ#
Почему так: Фундамент строится первым, потому что безопасность невозможно добавить задним числом — guardrails должны защищать все последующие слои с первого коммита. Router и state определяют контракт между всеми агентами, изменение которого после Phase 2+ потребует рефакторинга всего графа.
Пример: Без guardrails первого дня злоумышленник мог бы отправить в wiki команду
"../../etc/passwd"и прочитать системные файлы. После Задачи 0.2 это заблокировано на уровнеsanitize_wiki_path()до того, как любой агент даже увидит путь.
Цель: Создать скелет проекта, guardrails безопасности, router, state. Оценка: ~1 рабочий день AC из code-review: Path traversal fix, SSRF fix, fsync fix
Задача 0.0: Создание структуры пакетов (обязательно перед Задачей 0.1)#
Почему так: Python пакеты без
__init__.pyневидимы для импортов — без этого шага все последующие задачи упадут сImportError. Выполняется первым, чтобы не тратить время на отладку путей. Пример: Пользователь запускаетfrom d_brain.core.state import AgentState→ без__init__.pyPython не найдёт пакет → с__init__.pyимпорт работает с первого запуска.
# Создаём __init__.py для всех новых Python пакетов
# Без этого все импорты вида `from d_brain.core.state import AgentState` упадут с ImportError
mkdir -p src/d_brain/core src/d_brain/memory src/d_brain/agents \
src/d_brain/tools src/d_brain/skills src/d_brain/bot/handlers
touch src/d_brain/core/__init__.py
touch src/d_brain/memory/__init__.py
touch src/d_brain/agents/__init__.py
touch src/d_brain/tools/__init__.py
touch src/d_brain/skills/__init__.py
touch src/d_brain/bot/__init__.py
touch src/d_brain/bot/handlers/__init__.py
# Создаём conftest.py для правильного PYTHONPATH в тестах
cat > tests/conftest.py << 'EOF'
import sys
from pathlib import Path
# Добавляем src/ в путь — позволяет `from d_brain.xxx import yyy`
sys.path.insert(0, str(Path(__file__).parent.parent / "src"))
EOF
git add src/d_brain/ tests/conftest.py
git commit -m "chore(structure): add __init__.py for all new packages and conftest.py"
Задача 0.1: AgentState — центральный TypedDict#
Почему так: Единый
AgentState— «контракт» между всеми узлами LangGraph графа; без него каждый агент изобретает свой формат состояния и несовместим с другими. Пример: Пользователь пишет «запомни про проект Феникс» → агент кладёт данные вstate["memory_context"]→ при следующем узле state передаётся с сохранённым контекстом.
Цель: Определить единый тип состояния, который течёт через все узлы графа.
Файлы:
- Создать:
src/d_brain/core/state.py - Тест:
tests/core/test_state.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_state.py
from d_brain.core.state import AgentState, MessageRole
def test_agent_state_has_required_fields():
"""AgentState содержит все обязательные поля."""
state: AgentState = {
"messages": [],
"user_id": 123,
"thread_id": 0,
"agent_name": "main",
"memory_context": "",
"tool_results": [],
"metadata": {},
}
assert state["user_id"] == 123
assert state["agent_name"] == "main"
def test_message_role_enum():
assert MessageRole.USER == "user"
assert MessageRole.ASSISTANT == "assistant"
assert MessageRole.TOOL == "tool"
Шаг 2 — Запускаем тест (ожидаем FAIL):
cd /home/serg/projects/second-brain
uv run pytest tests/core/test_state.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/state.py
from typing import Annotated, Any
from typing_extensions import TypedDict
from langgraph.graph.message import add_messages
from enum import StrEnum
class MessageRole(StrEnum):
USER = "user"
ASSISTANT = "assistant"
TOOL = "tool"
SYSTEM = "system"
class AgentState(TypedDict):
"""Центральное состояние агента — течёт через все узлы LangGraph.
Deep Module: узкий интерфейс (7 полей), внутри каждый агент
работает только с нужным подмножеством.
"""
messages: Annotated[list, add_messages] # LangChain messages
user_id: int
thread_id: int # Telegram thread id (0 = главный чат)
agent_name: str # "main" | "coach" | "habits" | "planner"
memory_context: str # строка из MemoryManager для системного промпта
tool_results: list[dict[str, Any]]
metadata: dict[str, Any] # произвольные данные (close_day, etc.)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_state.py -v
Шаг 5 — Коммит:
git add src/d_brain/core/state.py tests/core/test_state.py
git commit -m "feat(core): add AgentState TypedDict with MessageRole enum"
Задача 0.2: Guardrails — защита от Path Traversal и SSRF#
Почему так: Бот работает с файловой системой и внешними URL от пользователя — без валидации злоумышленник может прочитать
/etc/passwdили заставить бот обращаться к внутренним сервисам. Пример: Пользователь пишет../../etc/passwd→validate_path()выбрасываетValidationError: path traversal detected→ бот отвечает «недопустимый путь», файл не открывается.
Цель: Исправить критические баги безопасности из code-review ДО всего остального.
Файлы:
- Создать:
src/d_brain/core/guardrails.py - Тест:
tests/core/test_guardrails.py
Шаг 1 — Пишем падающие тесты:
# tests/core/test_guardrails.py
import pytest
from d_brain.core.guardrails import (
sanitize_wiki_path,
check_ssrf_url,
PathTraversalError,
SSRFError,
)
# === Path Traversal ===
def test_safe_wiki_path_allowed():
result = sanitize_wiki_path("Проекты/МойПроект.md", base="/vault")
assert result.endswith("МойПроект.md")
def test_path_traversal_blocked():
with pytest.raises(PathTraversalError):
sanitize_wiki_path("../../etc/passwd", base="/vault")
def test_path_traversal_url_encoded():
with pytest.raises(PathTraversalError):
sanitize_wiki_path("%2e%2e%2fetc%2fpasswd", base="/vault")
def test_path_traversal_null_byte():
with pytest.raises(PathTraversalError):
sanitize_wiki_path("notes\x00.md", base="/vault")
# === SSRF ===
def test_public_url_allowed():
check_ssrf_url("https://duckduckgo.com/search?q=test")
def test_ssrf_localhost_blocked():
with pytest.raises(SSRFError):
check_ssrf_url("http://localhost:8080/admin")
def test_ssrf_rfc1918_blocked():
with pytest.raises(SSRFError):
check_ssrf_url("http://192.168.1.1/router")
def test_ssrf_169_254_blocked():
with pytest.raises(SSRFError):
check_ssrf_url("http://169.254.169.254/latest/meta-data/")
def test_ssrf_file_scheme_blocked():
with pytest.raises(SSRFError):
check_ssrf_url("file:///etc/passwd")
def test_ssrf_custom_scheme_blocked():
with pytest.raises(SSRFError):
check_ssrf_url("gopher://evil.com:70/1%70%6f%73%74")
# <!-- ARCH ADDED: W2 — тест async версии -->
@pytest.mark.asyncio
async def test_ssrf_async_localhost_blocked():
from d_brain.core.guardrails import check_ssrf_url_async
with pytest.raises(SSRFError):
await check_ssrf_url_async("http://localhost:8080/admin")
@pytest.mark.asyncio
async def test_ssrf_async_public_url_allowed():
from d_brain.core.guardrails import check_ssrf_url_async
# Не должно бросать исключение (DNS может не резолвиться в тестах — ОК)
try:
await check_ssrf_url_async("https://duckduckgo.com/search?q=test")
except SSRFError:
pytest.fail("Public URL should not raise SSRFError")
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_guardrails.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/core/guardrails.py
"""Guardrails безопасности — Deep Module.
Интерфейс: 2 функции. Реализация: полная защита от path traversal и SSRF.
"""
import ipaddress
import socket
from pathlib import Path
from urllib.parse import urlparse, unquote
class PathTraversalError(ValueError):
"""Попытка выйти за пределы разрешённой директории."""
class SSRFError(ValueError):
"""Попытка обратиться к внутренним адресам (SSRF)."""
_ALLOWED_SCHEMES = frozenset({"http", "https"})
_RFC1918 = [
ipaddress.ip_network("10.0.0.0/8"),
ipaddress.ip_network("172.16.0.0/12"),
ipaddress.ip_network("192.168.0.0/16"),
ipaddress.ip_network("127.0.0.0/8"),
ipaddress.ip_network("169.254.0.0/16"), # link-local / AWS metadata
ipaddress.ip_network("::1/128"),
ipaddress.ip_network("fc00::/7"),
]
def sanitize_wiki_path(user_path: str, base: str) -> Path:
"""Проверяет путь к файлу wiki, блокирует path traversal.
Args:
user_path: путь от пользователя (может содержать ../. URL-encoding, null bytes)
base: абсолютный путь базовой директории vault
Returns:
Path — безопасный абсолютный путь внутри base
Raises:
PathTraversalError: если путь выходит за пределы base
"""
# Убираем URL-encoding и null bytes
decoded = unquote(user_path).replace("\x00", "")
base_path = Path(base).resolve()
target = (base_path / decoded).resolve()
try:
target.relative_to(base_path)
except ValueError:
raise PathTraversalError(
f"Path '{user_path}' выходит за пределы vault '{base}'"
)
return target
def check_ssrf_url(url: str) -> None:
"""Проверяет URL на SSRF — блокирует внутренние адреса.
ВАЖНО: эта функция СИНХРОННАЯ и делает DNS-резолвинг через socket.getaddrinfo.
В синхронном контексте (guardrails перед записью пути) — OK.
В async контексте использовать check_ssrf_url_async() ниже.
Args:
url: URL для проверки
Raises:
SSRFError: если URL указывает на внутренние ресурсы
"""
parsed = urlparse(url)
if parsed.scheme not in _ALLOWED_SCHEMES:
raise SSRFError(f"Схема '{parsed.scheme}' запрещена. Разрешены: http, https")
hostname = parsed.hostname
if not hostname:
raise SSRFError("URL не содержит hostname")
# Разрешаем DNS (может вернуть внутренний IP — проверяем)
try:
addr_infos = socket.getaddrinfo(hostname, None)
except socket.gaierror:
return # не резолвится — пусть httpx сам разберётся
for _, _, _, _, sockaddr in addr_infos:
ip_str = sockaddr[0]
try:
ip = ipaddress.ip_address(ip_str)
except ValueError:
continue
for network in _RFC1918:
if ip in network:
raise SSRFError(
f"URL '{url}' резолвится в приватный адрес {ip}"
)
# <!-- ARCH ADDED: W2 — async версия SSRF check, не блокирует event loop -->
async def check_ssrf_url_async(url: str) -> None:
"""Async версия SSRF check — DNS-резолвинг в thread pool, не блокирует event loop.
Использовать везде в async контексте: MediaTools, SearchTools, SkillLoader.
Args:
url: URL для проверки
Raises:
SSRFError: если URL указывает на внутренние ресурсы
"""
import asyncio
parsed = urlparse(url)
if parsed.scheme not in _ALLOWED_SCHEMES:
raise SSRFError(f"Схема '{parsed.scheme}' запрещена. Разрешены: http, https")
hostname = parsed.hostname
if not hostname:
raise SSRFError("URL не содержит hostname")
# socket.getaddrinfo БЛОКИРУЕТ event loop — выносим в thread pool (W2)
try:
addr_infos = await asyncio.to_thread(socket.getaddrinfo, hostname, None)
except socket.gaierror:
return # не резолвится — пусть httpx сам разберётся
for _, _, _, _, sockaddr in addr_infos:
ip_str = sockaddr[0]
try:
ip = ipaddress.ip_address(ip_str)
except ValueError:
continue
for network in _RFC1918:
if ip in network:
raise SSRFError(
f"URL '{url}' резолвится в приватный адрес {ip}"
)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_guardrails.py -v
Шаг 5 — Коммит:
git add src/d_brain/core/guardrails.py tests/core/test_guardrails.py
git commit -m "fix(security): add path traversal and SSRF guardrails (code-review AC)"
Задача 0.3: Router — маршрутизация по thread_id#
Почему так: Единая точка маршрутизации по
thread_idустраняет дублированиеif/elifв каждом хендлере — вместо N хендлеров с повторяющейся логикой одинRouterзнает все маппинги. Пример: Пользователь пишет сообщение в топик#coach→ Router получаетthread_id=42→ находит в маппинге → передаёт вCoachAgent, не затрагивая MainAgent.
Цель: Определить, какой агент обрабатывает сообщение по Telegram thread_id.
Файлы:
- Создать:
src/d_brain/core/router.py - Создать:
src/d_brain/core/config.py(маппинг thread → agent) - Тест:
tests/core/test_router.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_router.py
from d_brain.core.router import ThreadRouter, AgentRoute
def test_main_thread_routes_to_main():
router = ThreadRouter({0: "main", 100: "coach", 200: "habits"})
assert router.route(thread_id=0) == AgentRoute(agent="main", thread_id=0)
def test_coach_thread_routes_to_coach():
router = ThreadRouter({0: "main", 100: "coach", 200: "habits"})
assert router.route(thread_id=100).agent == "coach"
def test_unknown_thread_routes_to_main():
router = ThreadRouter({0: "main", 100: "coach"})
assert router.route(thread_id=999).agent == "main"
def test_router_from_env(monkeypatch):
monkeypatch.setenv("THREAD_COACH_ID", "42")
monkeypatch.setenv("THREAD_HABITS_ID", "55")
from d_brain.core.router import ThreadRouter
router = ThreadRouter.from_env()
assert router.route(42).agent == "coach"
assert router.route(55).agent == "habits"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_router.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/router.py
import os
from dataclasses import dataclass
@dataclass(frozen=True)
class AgentRoute:
agent: str # "main" | "coach" | "habits" | "planner"
thread_id: int
class ThreadRouter:
"""Маршрутизатор сообщений по Telegram thread_id → agent name.
Deep Module: снаружи один метод route(), внутри вся логика fallback.
"""
def __init__(self, mapping: dict[int, str]):
self._map = mapping
def route(self, thread_id: int) -> AgentRoute:
agent = self._map.get(thread_id, "main")
return AgentRoute(agent=agent, thread_id=thread_id)
@classmethod
def from_env(cls) -> "ThreadRouter":
mapping: dict[int, str] = {0: "main"}
for env_key, agent_name in [
("THREAD_COACH_ID", "coach"),
("THREAD_HABITS_ID", "habits"),
("THREAD_PLANNER_ID", "planner"),
]:
val = os.getenv(env_key)
if val and val.isdigit():
mapping[int(val)] = agent_name
return cls(mapping)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_router.py -v
Шаг 5 — Коммит:
git add src/d_brain/core/router.py tests/core/test_router.py
git commit -m "feat(core): add ThreadRouter for thread_id → agent routing"
Задача 0.4: fsync fix — надёжная запись файлов#
Почему так:
write()безfsync()оставляет данные в OS буфере — при crash бота файл может быть пустым или обрезанным; атомарный write через temp file + rename исключает повреждение. Пример: Бот пишет summary дня → питание пропало в середине записи → без fsync файл пустой → с fsync: либо старый файл (rename не случился), либо новый полный (rename случился).
Цель: Исправить баг из code-review: запись wiki/журналов без fsync может потерять данные при краше.
Файлы:
- Создать:
src/d_brain/core/safe_io.py - Тест:
tests/core/test_safe_io.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_safe_io.py
import os
import tempfile
from pathlib import Path
from d_brain.core.safe_io import atomic_write_text, atomic_write_bytes
def test_atomic_write_creates_file():
with tempfile.TemporaryDirectory() as d:
p = Path(d) / "note.md"
atomic_write_text(p, "# Hello\n")
assert p.read_text() == "# Hello\n"
def test_atomic_write_is_atomic(tmp_path):
"""Файл либо полный, либо отсутствует — никогда не частичный."""
target = tmp_path / "data.txt"
atomic_write_text(target, "complete content")
# Нет промежуточного мусора
files = list(tmp_path.iterdir())
assert len(files) == 1
assert files[0] == target
def test_atomic_write_overwrites():
with tempfile.TemporaryDirectory() as d:
p = Path(d) / "f.txt"
atomic_write_text(p, "v1")
atomic_write_text(p, "v2")
assert p.read_text() == "v2"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_safe_io.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/safe_io.py
"""Атомарная запись файлов с fsync.
Исправляет баг code-review: потеря данных при краше во время записи.
Техника: write → fsync → rename (атомарная на POSIX).
"""
import os
import tempfile
from pathlib import Path
def atomic_write_text(path: Path, content: str, encoding: str = "utf-8") -> None:
"""Записывает текстовый файл атомарно (write + fsync + rename)."""
path = Path(path)
path.parent.mkdir(parents=True, exist_ok=True)
dir_fd = os.open(str(path.parent), os.O_RDONLY)
try:
with tempfile.NamedTemporaryFile(
mode="w",
encoding=encoding,
dir=path.parent,
delete=False,
suffix=".tmp",
) as f:
f.write(content)
f.flush()
os.fsync(f.fileno())
tmp_path = f.name
os.replace(tmp_path, path)
os.fsync(dir_fd) # fsync директории — гарантирует видимость rename
finally:
os.close(dir_fd)
def atomic_write_bytes(path: Path, content: bytes) -> None:
"""Записывает бинарный файл атомарно."""
path = Path(path)
path.parent.mkdir(parents=True, exist_ok=True)
dir_fd = os.open(str(path.parent), os.O_RDONLY)
try:
with tempfile.NamedTemporaryFile(
mode="wb",
dir=path.parent,
delete=False,
suffix=".tmp",
) as f:
f.write(content)
f.flush()
os.fsync(f.fileno())
tmp_path = f.name
os.replace(tmp_path, path)
os.fsync(dir_fd)
finally:
os.close(dir_fd)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_safe_io.py -v
Шаг 5 — Коммит:
git add src/d_brain/core/safe_io.py tests/core/test_safe_io.py
git commit -m "fix(io): atomic write with fsync to prevent data loss on crash (code-review AC)"
Интеграционный smoke-тест Фазы 0#
uv run pytest tests/core/ -v --tb=short
# Ожидаем: все PASS (test_state, test_guardrails, test_router, test_safe_io)
Задача 0.5: TZ fix — единая работа с временем#
Почему так: Python
datetime.now()без timezone — naive datetime; сравнение naive с aware datetime бросаетTypeError; централизованныйTimeUtils.now()гарантирует консистентный TZ во всём коде. Пример: Пользователь спрашивает «есть ли встречи сегодня?» → без TZ fix сравнение дат может ошибиться для UTC vs MSK → с TZ fix корректно сравнивает в одном timezone.
Цель: Исправить bug #7 из code-review: UTC vs local time несоответствие. Единый источник правды для всех datetime операций.
Файлы:
- Создать:
src/d_brain/core/time_utils.py - Тест:
tests/core/test_time_utils.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_time_utils.py
from datetime import datetime, timezone
from d_brain.core.time_utils import utc_now, utc_today, format_iso, parse_iso
def test_utc_now_is_aware():
"""utc_now() возвращает timezone-aware datetime (не naive)."""
now = utc_now()
assert now.tzinfo is not None
assert now.tzinfo == timezone.utc
def test_utc_today_returns_date_string():
today = utc_today()
assert len(today) == 10 # "YYYY-MM-DD"
assert today.count("-") == 2
def test_format_iso_includes_tz():
dt = datetime(2026, 5, 27, 10, 30, 0, tzinfo=timezone.utc)
iso = format_iso(dt)
assert "2026-05-27" in iso
assert "T" in iso # ISO 8601 формат
def test_parse_iso_returns_aware():
iso = "2026-05-27T10:30:00+00:00"
dt = parse_iso(iso)
assert dt.tzinfo is not None
def test_utc_now_replaces_utcnow():
"""Проверяем что datetime.utcnow() НЕ используется в наших модулях."""
import subprocess, sys
result = subprocess.run(
["grep", "-r", "utcnow()", "src/d_brain/"],
capture_output=True, text=True
)
# Допускаем только этот комментарий-заглушку
lines = [l for l in result.stdout.splitlines() if "# deprecated" not in l and "time_utils" not in l]
assert lines == [], f"Found utcnow() usage (use utc_now() instead):\n" + "\n".join(lines)
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_time_utils.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/time_utils.py
"""TZ-утилиты — единый источник правды для datetime.
Правило проекта: НИКОГДА не использовать:
- datetime.utcnow() ← naive, deprecated Python 3.12
- datetime.now() ← локальная TZ, не UTC
- date.today() ← локальная TZ, не UTC
Всегда использовать:
- utc_now() → timezone-aware UTC datetime
- utc_today() → UTC date string "YYYY-MM-DD"
Исправляет code-review bug #7: UTC reader + local TZ writer = phantom completions.
"""
from datetime import datetime, timezone, date
def utc_now() -> datetime:
"""Текущее время UTC (timezone-aware). Замена datetime.utcnow()."""
return datetime.now(tz=timezone.utc)
def utc_today() -> str:
"""Текущая дата UTC как строка 'YYYY-MM-DD'. Замена date.today()."""
return datetime.now(tz=timezone.utc).date().isoformat()
def format_iso(dt: datetime) -> str:
"""Форматирует datetime в ISO 8601 строку с TZ offset."""
if dt.tzinfo is None:
# Если naive datetime попал сюда — трактуем как UTC
dt = dt.replace(tzinfo=timezone.utc)
return dt.isoformat()
def parse_iso(iso_str: str) -> datetime:
"""Парсит ISO 8601 строку в timezone-aware datetime UTC.
Обрабатывает форматы:
- "2026-05-27T10:30:00+00:00" ← preferred
- "2026-05-27T10:30:00" ← naive (трактуем как UTC)
- "2026-05-27T10:30:00.123456" ← с микросекундами
"""
try:
dt = datetime.fromisoformat(iso_str)
except ValueError:
# Fallback для нестандартных форматов
dt = datetime.strptime(iso_str[:19], "%Y-%m-%dT%H:%M:%S")
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_time_utils.py -v
Шаг 5 — Заменяем все datetime.utcnow() в новых модулях:
# Находим все вхождения в src/d_brain/ (должны исправить ВСЕ)
grep -rn "utcnow()\|date\.today()" src/d_brain/ --include="*.py"
# В episodic.py заменяем:
# БЫЛО: datetime.utcnow().isoformat()
# СТАЛО: from d_brain.core.time_utils import utc_now, format_iso
# format_iso(utc_now())
# В habits_agent.py заменяем:
# БЫЛО: datetime.utcnow().isoformat()
# СТАЛО: from d_brain.core.time_utils import utc_now, format_iso
# format_iso(utc_now())
Шаг 6 — Коммит:
git add src/d_brain/core/time_utils.py tests/core/test_time_utils.py
git commit -m "fix(time): add TZ-safe utc_now() utils, replace all datetime.utcnow() (code-review AC#7)"
Задача 0.6: BaseLLMAgent — базовый класс агента#
Почему так: 5 агентов дублируют
_call_llmс захардкоженной моделью — изменить модель значит обновить 5 файлов;BaseLLMAgentсBRAIN_MODELenv var делает это одним изменением конфига. Пример: Вышла новая claude-sonnet-4-7 → разработчик меняетBRAIN_MODEL=claude-sonnet-4-7в.env→ все 5 агентов автоматически начинают использовать новую модель.
Цель: Убрать дублирование _call_llm из 5 агентов. Читать модель из env BRAIN_MODEL. Интегрировать Langfuse трассировку в одном месте.
Файлы:
- Создать:
src/d_brain/agents/base.py - Тест:
tests/agents/test_base_agent.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_base_agent.py
import os
import pytest
from unittest.mock import AsyncMock, patch, MagicMock
from d_brain.agents.base import BaseLLMAgent
class ConcreteAgent(BaseLLMAgent):
SYSTEM_PROMPT = "Test agent"
@pytest.fixture
def agent(tmp_path):
return ConcreteAgent(data_dir=tmp_path)
def test_model_from_env(monkeypatch, tmp_path):
"""Модель читается из переменной окружения BRAIN_MODEL."""
monkeypatch.setenv("BRAIN_MODEL", "claude-opus-4-7")
a = ConcreteAgent(data_dir=tmp_path)
assert a.model == "claude-opus-4-7"
def test_model_default(monkeypatch, tmp_path):
"""Дефолтная модель — claude-sonnet-4-6 (актуальная)."""
monkeypatch.delenv("BRAIN_MODEL", raising=False)
a = ConcreteAgent(data_dir=tmp_path)
assert "claude" in a.model.lower()
assert "sonnet" in a.model.lower() or "opus" in a.model.lower()
@pytest.mark.asyncio
async def test_call_llm_uses_tracer(agent, monkeypatch):
"""_call_llm вызывает tracer.trace_llm()."""
monkeypatch.setenv("ANTHROPIC_API_KEY", "test-key")
mock_tracer = MagicMock()
mock_tracer.trace_llm = MagicMock()
with patch("d_brain.agents.base.get_tracer", return_value=mock_tracer):
with patch("d_brain.agents.base.ChatAnthropic") as mock_cls:
mock_llm = MagicMock()
mock_llm.ainvoke = AsyncMock(return_value=MagicMock(content="OK"))
mock_cls.return_value = mock_llm
from langchain_core.messages import HumanMessage
result = await agent._call_llm([HumanMessage(content="test")])
mock_tracer.trace_llm.assert_called_once()
assert result == "OK"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/agents/test_base_agent.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/agents/base.py
"""BaseLLMAgent — базовый класс для всех агентов.
Убирает дублирование _call_llm из 5 агентов (DRY).
Читает модель из BRAIN_MODEL env var.
Интегрирует Langfuse tracer в одном месте.
Использование:
class MyAgent(BaseLLMAgent):
SYSTEM_PROMPT = "..."
async def process(...): ...
"""
import os
from pathlib import Path
from langchain_core.messages import SystemMessage, BaseMessage
from d_brain.core.memory import MemoryManager
# Актуальная модель по умолчанию (обновлять при выходе новой версии)
DEFAULT_MODEL = "claude-sonnet-4-6"
class BaseLLMAgent:
"""Базовый класс с _call_llm, memory, tracer. Наследуйте для своих агентов."""
SYSTEM_PROMPT: str = "Ты персональный ассистент."
SKILLS_SCOPE: str = "main"
def __init__(self, data_dir: Path | str = "~/.d_brain"):
self.memory = MemoryManager(data_dir=data_dir)
self.model = os.getenv("BRAIN_MODEL", DEFAULT_MODEL)
async def _call_llm(
self,
messages: list[BaseMessage],
system: str | None = None,
temperature: float = 0.7,
) -> str:
"""Вызов LLM с Langfuse трассировкой.
Args:
messages: список LangChain messages (HumanMessage, AIMessage, etc.)
system: системный промпт (если None — используется self.SYSTEM_PROMPT)
temperature: температура генерации
Returns:
Строка ответа LLM
"""
from langchain_anthropic import ChatAnthropic
from d_brain.core.observability import get_tracer
effective_system = system or self.SYSTEM_PROMPT
api_key = os.environ.get("ANTHROPIC_API_KEY")
if not api_key:
raise EnvironmentError(
"ANTHROPIC_API_KEY не задан. Добавьте в .env файл."
)
llm = ChatAnthropic(
model=self.model,
api_key=api_key,
temperature=temperature,
max_tokens=4096,
)
all_msgs = [SystemMessage(content=effective_system)] + list(messages)
tracer = get_tracer()
tracer.trace_llm(
name=f"{self.__class__.__name__}._call_llm",
input=str(messages[-1].content if messages else ""),
output="", # заполним после вызова
model=self.model,
)
try:
resp = await llm.ainvoke(all_msgs)
return resp.content
except Exception as e:
# Логируем ошибку и пробрасываем
import logging
logging.getLogger(__name__).error(
f"{self.__class__.__name__}._call_llm error: {e}"
)
raise
Шаг 4 — Обновляем все агенты для наследования от BaseLLMAgent:
# В каждом агенте ЗАМЕНЯЕМ:
# БЫЛО:
class CoachAgent:
def __init__(self, data_dir=...):
self.memory = MemoryManager(...)
async def _call_llm(self, messages, system=""):
from langchain_anthropic import ChatAnthropic
llm = ChatAnthropic(model="claude-3-5-sonnet-20241022", ...)
...
# СТАЛО:
from d_brain.agents.base import BaseLLMAgent
class CoachAgent(BaseLLMAgent):
SKILLS_SCOPE = "coach"
SYSTEM_PROMPT = "..."
# _call_llm и self.memory наследуются от BaseLLMAgent
# def __init__ не нужен если дополнительных атрибутов нет
Шаг 5 — Запускаем тест (ожидаем PASS):
uv run pytest tests/agents/test_base_agent.py -v
Шаг 6 — Коммит:
git add src/d_brain/agents/base.py tests/agents/test_base_agent.py
git commit -m "feat(agents): add BaseLLMAgent mixin with BRAIN_MODEL env var and Langfuse integration"
ФАЗА 1: MEMORY LAYER#
Почему так: Без многоуровневой памяти каждый диалог начинается с нуля — агент не знает, что вы обсуждали вчера, какие цели у вас стоят, как вас зовут. Три слоя (working/episodic/semantic) отражают три временных горизонта: секунды (текущий контекст), дни/недели (история сессий), годы (структурированные знания). Разделение предотвращает смешение горячего и холодного контекста в одном промпте.
Пример: Пользователь спрашивает «напомни, что мы обсуждали про проект Феникс». Working Memory пуста (новая сессия).
recall()делает semantic search в FTS5 и находит запись"coach: проект Феникс, архитектура бэкенда, дедлайн Q3"— и агент отвечает с реальным контекстом вместо «не помню».
Цель: Многоуровневая память: рабочая (context window), эпизодическая (SQLite), семантическая (ChromaDB — обязательно; FTS5 — secondary keyword index только внутри эпизодической памяти). Оценка: ~2 рабочих дня
Задача 1.1: Working Memory — управление контекстным окном#
Почему так: LLM не имеет памяти между вызовами — без управления контекстным окном каждый следующий вызов «забывает» предыдущие сообщения или превышает лимит токенов. Пример: Пользователь задаёт 10 вопросов подряд →
WorkingMemory.trim()отбрасывает старые сообщения сверх лимита → агент помнит последние N сообщений без превышения контекста.
Цель: Хранить последние N сообщений сессии в памяти, обрезая при превышении лимита токенов.
Файлы:
- Создать:
src/d_brain/memory/working.py - Тест:
tests/memory/test_working.py
Шаг 1 — Пишем падающий тест:
# tests/memory/test_working.py
from d_brain.memory.working import WorkingMemory
def test_add_and_get_messages():
wm = WorkingMemory(max_tokens=1000)
wm.add("user", "Привет")
wm.add("assistant", "Здравствуй!")
msgs = wm.get_messages()
assert len(msgs) == 2
assert msgs[0]["role"] == "user"
def test_overflow_truncates_oldest():
wm = WorkingMemory(max_tokens=50) # очень маленький лимит
for i in range(20):
wm.add("user", f"сообщение номер {i} с длинным текстом")
msgs = wm.get_messages()
# Должны остаться только последние сообщения
assert len(msgs) < 20
def test_clear_resets():
wm = WorkingMemory(max_tokens=1000)
wm.add("user", "test")
wm.clear()
assert wm.get_messages() == []
def test_summary_injected_on_overflow():
"""При обрезке создаётся системное сообщение с summary."""
wm = WorkingMemory(max_tokens=100, add_summary_on_truncate=True)
for i in range(30):
wm.add("user", f"msg {i} " * 5)
msgs = wm.get_messages()
roles = [m["role"] for m in msgs]
assert "system" in roles # summary message
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/memory/test_working.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/memory/working.py
"""Working Memory — контекстное окно агента.
Deep Module: снаружи 3 метода (add/get_messages/clear).
Внутри: подсчёт токенов, усечение, инъекция summary.
"""
from dataclasses import dataclass, field
from typing import Literal
def _rough_tokens(text: str) -> int:
"""Грубая оценка: 1 токен ≈ 4 символа."""
return max(1, len(text) // 4)
@dataclass
class _Msg:
role: str
content: str
tokens: int = field(init=False)
def __post_init__(self):
self.tokens = _rough_tokens(self.content)
class WorkingMemory:
"""Управляет контекстным окном: хранит, обрезает, добавляет summary."""
def __init__(
self,
max_tokens: int = 6000,
add_summary_on_truncate: bool = False,
):
self._max = max_tokens
self._add_summary = add_summary_on_truncate
self._msgs: list[_Msg] = []
def add(self, role: str, content: str) -> None:
self._msgs.append(_Msg(role=role, content=content))
self._trim()
def get_messages(self) -> list[dict]:
return [{"role": m.role, "content": m.content} for m in self._msgs]
def clear(self) -> None:
self._msgs = []
# <!-- ARCH REVISED: W24 — исправлен O(N²) trim. Кэшируем total и вычитаем по мере удаления -->
def _trim(self) -> None:
# Считаем total один раз (не в каждой итерации while-loop)
total = sum(m.tokens for m in self._msgs)
if total <= self._max:
return
dropped: list[_Msg] = []
while self._msgs and total > self._max:
msg = self._msgs.pop(0)
total -= msg.tokens # вычитаем вместо пересчёта всего списка (O(1) per iter)
dropped.append(msg)
if self._add_summary and dropped:
summary_text = (
f"[Контекст усечён. Пропущено {len(dropped)} сообщений. "
f"Краткое содержание: пользователь и ассистент обсуждали: "
+ "; ".join(m.content[:40] for m in dropped[-3:])
+ "]"
)
self._msgs.insert(0, _Msg(role="system", content=summary_text))
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/memory/test_working.py -v
Шаг 5 — Коммит:
git add src/d_brain/memory/working.py tests/memory/test_working.py
git commit -m "feat(memory): add WorkingMemory with token-aware truncation"
Задача 1.2: Episodic Memory — SQLite сессии и саммари#
Почему так: Без персистентных сессий бот «забывает» всё при рестарте — SQLite хранит историю разговоров с нарастающими саммари, позволяя вернуться к контексту недельной давности. Пример: Пользователь спрашивает «что мы решили про проект в прошлый вторник?» →
EpisodicMemory.recall(days=7)находит сессию → достаёт daily summary → агент отвечает конкретикой.
Цель: Сохранять сессии диалогов в SQLite с автоматическим суммированием.
Файлы:
- Создать:
src/d_brain/memory/episodic.py - Тест:
tests/memory/test_episodic.py
Шаг 1 — Пишем падающий тест:
# tests/memory/test_episodic.py
import pytest
import asyncio
import tempfile
from pathlib import Path
from d_brain.memory.episodic import EpisodicMemory, Session
@pytest.fixture
def mem(tmp_path):
return EpisodicMemory(db_path=tmp_path / "episodic.db")
@pytest.mark.asyncio
async def test_create_and_get_session(mem):
session_id = await mem.create_session(user_id=1, agent="main")
assert session_id > 0
session = await mem.get_session(session_id)
assert session.user_id == 1
assert session.agent == "main"
@pytest.mark.asyncio
async def test_add_messages_to_session(mem):
sid = await mem.create_session(user_id=1, agent="main")
await mem.add_message(sid, role="user", content="Привет")
await mem.add_message(sid, role="assistant", content="Здравствуй")
msgs = await mem.get_messages(sid)
assert len(msgs) == 2
@pytest.mark.asyncio
async def test_save_summary(mem):
sid = await mem.create_session(user_id=1, agent="main")
await mem.save_summary(sid, "Обсуждали планы на неделю")
session = await mem.get_session(sid)
assert "планы" in session.summary
@pytest.mark.asyncio
async def test_recent_sessions(mem):
for _ in range(5):
sid = await mem.create_session(user_id=1, agent="main")
await mem.add_message(sid, "user", "test")
recent = await mem.get_recent_sessions(user_id=1, limit=3)
assert len(recent) == 3
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/memory/test_episodic.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/memory/episodic.py
"""Episodic Memory — SQLite хранилище сессий.
Deep Module: 5 публичных методов, скрывает всю работу с БД.
Исправления:
- WAL mode: предотвращает "database is locked" при конкурентных запросах
- asyncio.Lock: предотвращает race condition в _ensure_init
- utc_now(): все timestamps в UTC (code-review bug #7)
"""
import asyncio
import aiosqlite
from dataclasses import dataclass
from pathlib import Path
from d_brain.core.time_utils import utc_now, format_iso # W5
@dataclass
class Session:
id: int
user_id: int
agent: str
created_at: str
summary: str = ""
class EpisodicMemory:
def __init__(self, db_path: Path | str):
self._db_path = Path(db_path)
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._initialized = False
self._init_lock = asyncio.Lock() # W4: предотвращает race condition
async def _ensure_init(self):
# W4: двойная проверка с Lock — idempotent и безопасна для concurrent calls
if self._initialized:
return
async with self._init_lock:
if self._initialized: # второй чек внутри lock
return
async with aiosqlite.connect(self._db_path) as db:
# W3: WAL mode — конкурентные читатели не блокируют писателя
await db.execute("PRAGMA journal_mode=WAL")
await db.execute("PRAGMA synchronous=NORMAL") # баланс надёжность/скорость
await db.execute("""
CREATE TABLE IF NOT EXISTS sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
user_id INTEGER NOT NULL,
agent TEXT NOT NULL,
created_at TEXT NOT NULL, -- UTC ISO 8601
summary TEXT DEFAULT ''
)
""")
await db.execute("""
CREATE TABLE IF NOT EXISTS messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id INTEGER NOT NULL REFERENCES sessions(id),
role TEXT NOT NULL,
content TEXT NOT NULL,
ts TEXT NOT NULL -- UTC ISO 8601
)
""")
await db.execute(
"CREATE INDEX IF NOT EXISTS idx_msg_session ON messages(session_id)"
)
await db.execute(
"CREATE INDEX IF NOT EXISTS idx_sessions_user ON sessions(user_id, created_at DESC)"
)
await db.commit()
self._initialized = True
async def create_session(self, user_id: int, agent: str) -> int:
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
cur = await db.execute(
"INSERT INTO sessions (user_id, agent, created_at) VALUES (?,?,?)",
(user_id, agent, format_iso(utc_now())), # W5: UTC timestamp
)
await db.commit()
return cur.lastrowid
async def get_session(self, session_id: int) -> Session:
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"SELECT * FROM sessions WHERE id=?", (session_id,)
) as cur:
row = await cur.fetchone()
return Session(**dict(row))
async def add_message(self, session_id: int, role: str, content: str) -> None:
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
await db.execute(
"INSERT INTO messages (session_id, role, content, ts) VALUES (?,?,?,?)",
(session_id, role, content, format_iso(utc_now())), # W5: UTC timestamp
)
await db.commit()
async def get_messages(self, session_id: int) -> list[dict]:
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"SELECT role, content FROM messages WHERE session_id=? ORDER BY id",
(session_id,),
) as cur:
return [dict(r) for r in await cur.fetchall()]
async def save_summary(self, session_id: int, summary: str) -> None:
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
await db.execute(
"UPDATE sessions SET summary=? WHERE id=?", (summary, session_id)
)
await db.commit()
async def get_recent_sessions(self, user_id: int, limit: int = 10) -> list[Session]:
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"SELECT * FROM sessions WHERE user_id=? ORDER BY id DESC LIMIT ?",
(user_id, limit),
) as cur:
return [Session(**dict(r)) for r in await cur.fetchall()]
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/memory/test_episodic.py -v
Шаг 5 — Коммит:
git add src/d_brain/memory/episodic.py tests/memory/test_episodic.py
git commit -m "feat(memory): add EpisodicMemory with SQLite sessions and summaries"
Задача 1.3: Semantic Memory — ChromaDB (обязательно) + SQLite FTS5 (secondary keyword index)#
Почему так: ChromaDB обязателен для семантического поиска по wiki — FTS5 ищет только по ключевым словам; ChromaDB находит «похожие по смыслу» страницы даже при разных формулировках. Пример: Пользователь пишет «что я знаю про машинное обучение?» → ChromaDB ищет по эмбеддингам → находит страницы «нейронные сети», «LLM», «трансформеры» (слова разные, смысл совпадает).
Цель: Семантический поиск по прошлым разговорам и wiki-заметкам.
Файлы:
- Создать:
src/d_brain/memory/semantic.py - Тест:
tests/memory/test_semantic.py
Шаг 1 — Пишем падающий тест:
# tests/memory/test_semantic.py
import pytest
from d_brain.memory.semantic import SemanticMemory
@pytest.fixture
def mem(tmp_path):
# Используем SQLite FTS5 backend (без ChromaDB для тестов)
return SemanticMemory(backend="fts5", db_path=tmp_path / "semantic.db")
@pytest.mark.asyncio
async def test_store_and_search(mem):
await mem.store(
doc_id="note-1",
text="Встреча с командой по планированию спринта",
metadata={"type": "wiki", "date": "2026-05-27"},
)
await mem.store(
doc_id="note-2",
text="Тренировка по бегу, 10 км за 50 минут",
metadata={"type": "habit", "date": "2026-05-27"},
)
results = await mem.search("планирование спринт", top_k=1)
assert len(results) == 1
assert results[0]["doc_id"] == "note-1"
@pytest.mark.asyncio
async def test_delete_doc(mem):
await mem.store("del-1", "временная заметка", {})
await mem.delete("del-1")
results = await mem.search("временная заметка", top_k=5)
assert all(r["doc_id"] != "del-1" for r in results)
@pytest.mark.asyncio
async def test_search_returns_metadata(mem):
await mem.store("m-1", "обсуждение архитектуры", {"type": "meeting"})
results = await mem.search("архитектура", top_k=1)
assert results[0]["metadata"]["type"] == "meeting"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/memory/test_semantic.py -v 2>&1 | head -30
Шаг 3 — Реализация (FTS5 backend):
# src/d_brain/memory/semantic.py
"""Semantic Memory — поиск по смыслу.
Deep Module: store/search/delete. Внутри: FTS5 или ChromaDB (флаг backend).
Langfuse span: добавляется если LANGFUSE_ENABLED=true.
"""
import json
import aiosqlite
from pathlib import Path
class SemanticMemory:
"""Семантический поиск. ChromaDB обязателен для vector search. FTS5 — secondary keyword index."""
def __init__(
self,
backend: str = "fts5",
db_path: Path | str | None = None,
chroma_path: Path | str | None = None,
):
self._backend = backend
if backend == "fts5":
self._db_path = Path(db_path or "~/.d_brain/semantic.db").expanduser()
self._db_path.parent.mkdir(parents=True, exist_ok=True)
elif backend == "chroma":
self._chroma_path = str(chroma_path or "~/.d_brain/chroma")
self._initialized = False
async def _ensure_init(self):
if self._initialized:
return
if self._backend == "fts5":
async with aiosqlite.connect(self._db_path) as db:
await db.execute("""
CREATE VIRTUAL TABLE IF NOT EXISTS docs USING fts5(
doc_id UNINDEXED,
text,
metadata UNINDEXED,
tokenize='unicode61'
)
""")
await db.commit()
self._initialized = True
async def store(self, doc_id: str, text: str, metadata: dict) -> None:
await self._ensure_init()
if self._backend == "fts5":
async with aiosqlite.connect(self._db_path) as db:
# Удаляем старую версию если есть
await db.execute("DELETE FROM docs WHERE doc_id=?", (doc_id,))
await db.execute(
"INSERT INTO docs (doc_id, text, metadata) VALUES (?,?,?)",
(doc_id, text, json.dumps(metadata, ensure_ascii=False)),
)
await db.commit()
async def search(self, query: str, top_k: int = 5) -> list[dict]:
await self._ensure_init()
if self._backend == "fts5":
# Экранируем запрос для FTS5
safe_query = query.replace('"', '""')
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"""SELECT doc_id, text, metadata, rank
FROM docs WHERE docs MATCH ?
ORDER BY rank LIMIT ?""",
(safe_query, top_k),
) as cur:
rows = await cur.fetchall()
return [
{
"doc_id": r["doc_id"],
"text": r["text"],
"metadata": json.loads(r["metadata"]),
"score": r["rank"],
}
for r in rows
]
return []
async def delete(self, doc_id: str) -> None:
await self._ensure_init()
if self._backend == "fts5":
async with aiosqlite.connect(self._db_path) as db:
await db.execute("DELETE FROM docs WHERE doc_id=?", (doc_id,))
await db.commit()
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/memory/test_semantic.py -v
Шаг 5 — Коммит:
git add src/d_brain/memory/semantic.py tests/memory/test_semantic.py
git commit -m "feat(memory): add SemanticMemory with ChromaDB vector backend + FTS5 secondary keyword index"
Задача 1.4: MemoryManager — единый фасад памяти#
Почему так: Агенты не должны знать, в каком хранилище лежат данные —
MemoryManagerкак единый фасад скрывает детали: рабочая память в RAM, эпизодическая в SQLite, семантическая в ChromaDB. Пример: Агент вызываетmemory.recall("проект Феникс")→ MemoryManager параллельно ищет в рабочей памяти + SQLite + ChromaDB → возвращает объединённый контекст из всех источников.
Цель: Один объект для всех слоёв памяти, используемый агентами.
Файлы:
- Создать:
src/d_brain/core/memory.py - Тест:
tests/core/test_memory_manager.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_memory_manager.py
import pytest
from d_brain.core.memory import MemoryManager
@pytest.fixture
def mm(tmp_path):
return MemoryManager(data_dir=tmp_path)
@pytest.mark.asyncio
async def test_recall_returns_string(mm):
ctx = await mm.recall(user_id=1, query="планирование")
assert isinstance(ctx, str)
@pytest.mark.asyncio
async def test_store_and_recall(mm):
await mm.store(user_id=1, text="Обсуждали архитектуру LangGraph", source="session")
ctx = await mm.recall(user_id=1, query="LangGraph архитектура")
assert "LangGraph" in ctx or ctx == "" # FTS5 может не найти сразу
@pytest.mark.asyncio
async def test_begin_session(mm):
sid = await mm.begin_session(user_id=1, agent="main")
assert sid > 0
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_memory_manager.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/memory.py
"""MemoryManager — единый фасад для всех слоёв памяти.
Deep Module: 4 метода наружу. Внутри: working + episodic + semantic.
Исправления (W7, W8):
- _sessions: dict[int, int] вместо _current_session: int
→ разные пользователи не перезаписывают session_id друг друга
- ensure_session() вместо begin_session():
→ один диалог = одна сессия, не создаём сессию на каждое сообщение
"""
from pathlib import Path
from d_brain.memory.working import WorkingMemory
from d_brain.memory.episodic import EpisodicMemory
from d_brain.memory.semantic import SemanticMemory
class MemoryManager:
def __init__(self, data_dir: Path | str = "~/.d_brain"):
d = Path(data_dir).expanduser()
self.working = WorkingMemory(max_tokens=6000, add_summary_on_truncate=True)
self.episodic = EpisodicMemory(db_path=d / "episodic.db")
self.semantic = SemanticMemory(backend="fts5", db_path=d / "semantic.db")
# W7: dict[user_id → session_id] вместо одного shared int
self._sessions: dict[int, int] = {}
async def begin_session(self, user_id: int, agent: str) -> int:
"""Создаёт новую сессию. Используй ensure_session() для idempotent поведения."""
sid = await self.episodic.create_session(user_id=user_id, agent=agent)
self._sessions[user_id] = sid
self.working.clear()
return sid
async def ensure_session(self, user_id: int, agent: str) -> int:
"""W8: Возвращает существующую сессию или создаёт новую.
Идемпотентно: один диалог = одна сессия в episodic.
Используй вместо begin_session() в process() методах агентов.
"""
if user_id not in self._sessions:
sid = await self.episodic.create_session(user_id=user_id, agent=agent)
self._sessions[user_id] = sid
self.working.clear()
return self._sessions[user_id]
def reset_session(self, user_id: int) -> None:
"""Сбрасывает сессию пользователя (например, при /reset команде)."""
self._sessions.pop(user_id, None)
async def recall(self, user_id: int, query: str, top_k: int = 3) -> str:
"""Возвращает строку контекста для системного промпта."""
results = await self.semantic.search(query, top_k=top_k)
if not results:
return ""
parts = [f"- {r['text'][:120]}" for r in results]
return "Релевантный контекст из памяти:\n" + "\n".join(parts)
async def store(self, user_id: int, text: str, source: str = "session") -> None:
"""Сохраняет текст в семантическую память."""
doc_id = f"{user_id}_{source}_{hash(text) & 0xFFFFFF:06x}"
await self.semantic.store(doc_id, text, {"user_id": user_id, "source": source})
async def add_turn(self, role: str, content: str, user_id: int | None = None) -> None:
"""Добавляет сообщение в working memory и episodic.
W7: user_id нужен для выбора правильной сессии.
"""
self.working.add(role, content)
session_id = self._sessions.get(user_id) if user_id else None
if session_id:
await self.episodic.add_message(session_id, role, content)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_memory_manager.py -v
Шаг 5 — Коммит:
git add src/d_brain/core/memory.py tests/core/test_memory_manager.py
git commit -m "feat(core): add MemoryManager facade (working + episodic + semantic)"
Задача 1.5: Consolidator — ночная консолидация памяти (mnemosyne pattern)#
Почему так: Без ночной консолидации SQLite растёт линейно с каждым разговором — mnemosyne паттерн сжимает raw сообщения в daily summaries, а summaries в weekly, уменьшая объём в 10–50x. Пример: За день 200 сообщений (50 KB raw) → ночью Consolidator создаёт 1 daily summary (2 KB) → через месяц weekly summary (5 KB) → через год: 12 weekly summaries вместо 2.5 MB raw.
Цель: Суммировать сессии за день в daily summary, сохранять в семантику.
Файлы:
- Создать:
src/d_brain/memory/consolidator.py - Тест:
tests/memory/test_consolidator.py
Шаг 1 — Пишем падающий тест:
# tests/memory/test_consolidator.py
import pytest
from unittest.mock import AsyncMock, MagicMock
from d_brain.memory.consolidator import MemoryConsolidator
@pytest.mark.asyncio
async def test_consolidate_creates_daily_summary(tmp_path):
"""Консолидатор создаёт daily summary из сессий за день."""
consolidator = MemoryConsolidator(data_dir=tmp_path)
# Мокаем LLM
consolidator._llm_summarize = AsyncMock(return_value="Пользователь работал над проектом")
# Добавляем тестовые данные напрямую
await consolidator.episodic.create_session(user_id=1, agent="main")
result = await consolidator.run_daily(user_id=1, date="2026-05-27")
assert result is not None
assert isinstance(result, str)
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/memory/test_consolidator.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/memory/consolidator.py
"""Memory Consolidator — ночная консолидация (паттерн mnemosyne).
Источник паттерна: https://github.com/AxDSan/mnemosyne
Идея: raw events → daily summary → weekly digest → long-term semantic index.
"""
from pathlib import Path
from datetime import date, datetime
from d_brain.memory.episodic import EpisodicMemory
from d_brain.memory.semantic import SemanticMemory
class MemoryConsolidator:
"""Консолидирует эпизодическую память в семантическую."""
def __init__(self, data_dir: Path | str = "~/.d_brain"):
d = Path(data_dir).expanduser()
self.episodic = EpisodicMemory(db_path=d / "episodic.db")
self.semantic = SemanticMemory(backend="fts5", db_path=d / "semantic.db")
async def _llm_summarize(self, text: str) -> str:
"""Заглушка — заменяется реальным вызовом LLM в production."""
return f"Summary of {len(text)} chars"
async def run_daily(self, user_id: int, date_str: str | None = None) -> str:
"""Запускает ночную консолидацию для пользователя.
1. Берёт все сессии за day
2. LLM-суммирует их
3. Сохраняет daily summary в semantic
Returns: текст summary
"""
today = date_str or date.today().isoformat()
sessions = await self.episodic.get_recent_sessions(user_id=user_id, limit=50)
today_sessions = [
s for s in sessions
if s.created_at.startswith(today)
]
if not today_sessions:
return f"Нет сессий за {today}"
all_text = "\n\n".join(
f"[{s.agent}] {s.summary or '(без summary)'}"
for s in today_sessions
)
summary = await self._llm_summarize(all_text)
doc_id = f"daily_{user_id}_{today}"
await self.semantic.store(
doc_id=doc_id,
text=summary,
metadata={
"type": "daily_summary",
"user_id": user_id,
"date": today,
"session_count": len(today_sessions),
},
)
return summary
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/memory/test_consolidator.py -v
Шаг 5 — Коммит:
git add src/d_brain/memory/consolidator.py tests/memory/test_consolidator.py
git commit -m "feat(memory): add MemoryConsolidator with mnemosyne pattern for nightly consolidation"
Интеграционный smoke-тест Фазы 1#
uv run pytest tests/memory/ tests/core/test_memory_manager.py -v --tb=short
# Ожидаем: все PASS
Задача 1.6: Procedural Memory — процедурная память (SOP, workflows)#
Почему так: SOP (Standard Operating Procedures) агентов — тоже память; вынос их в vault файлы позволяет обновлять поведение агента без перекомпиляции кода. Пример: Пользователь говорит «при встречах всегда спрашивай action items» → SOP обновляется в vault → CoachAgent загружает новый SOP при следующем запуске — код не менялся.
Цель: Хранить КАК делать что-либо (стандартные операционные процедуры, воркфлоу, инструкции). Это 4-й тип памяти из заявленных целей проекта.
Примеры процедур: "как провести close_day", "как создать задачу YouGile", "как оформить контакт в CRM"
Файлы:
- Создать:
src/d_brain/memory/procedural.py - Тест:
tests/memory/test_procedural.py
Шаг 1 — Пишем падающий тест:
# tests/memory/test_procedural.py
import pytest
from d_brain.memory.procedural import ProceduralMemory, Procedure
@pytest.fixture
def mem(tmp_path):
return ProceduralMemory(db_path=tmp_path / "procedural.db")
@pytest.mark.asyncio
async def test_store_and_get_procedure(mem):
await mem.store(Procedure(
name="close_day",
description="Закрытие рабочего дня",
steps=["1. Получить Done задачи из YouGile",
"2. Зафиксировать в wiki",
"3. Написать рефлексию"],
triggers=["закрой день", "close day", "конец дня"],
))
proc = await mem.get("close_day")
assert proc is not None
assert len(proc.steps) == 3
@pytest.mark.asyncio
async def test_find_by_trigger(mem):
await mem.store(Procedure(
name="create_task",
description="Создание задачи в YouGile",
steps=["..."],
triggers=["создай задачу", "новая задача"],
))
results = await mem.find_by_trigger("хочу создать задачу")
assert len(results) > 0
assert results[0].name == "create_task"
@pytest.mark.asyncio
async def test_update_procedure(mem):
"""Версионирование: обновление процедуры сохраняет историю."""
await mem.store(Procedure(name="p1", description="v1", steps=["step1"], triggers=[]))
await mem.store(Procedure(name="p1", description="v2", steps=["step1", "step2"], triggers=[]))
proc = await mem.get("p1")
assert proc.description == "v2"
assert len(proc.steps) == 2
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/memory/test_procedural.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/memory/procedural.py
"""Procedural Memory — хранит КАК делать что-либо.
Паттерн: SOP (Standard Operating Procedures) + trigger-based lookup.
Агент находит нужную процедуру по триггерным фразам в запросе пользователя.
"""
import asyncio
import json
import aiosqlite
from dataclasses import dataclass, field
from pathlib import Path
from d_brain.core.time_utils import utc_now, format_iso
@dataclass
class Procedure:
name: str # уникальный ключ, e.g. "close_day"
description: str # краткое описание для что это
steps: list[str] # пошаговая инструкция
triggers: list[str] # фразы, по которым агент найдёт эту процедуру
version: int = 1
updated_at: str = ""
class ProceduralMemory:
"""Хранит и ищет процедуры по имени и триггерным фразам."""
def __init__(self, db_path: Path | str = "~/.d_brain/procedural.db"):
self._db_path = Path(db_path)
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._initialized = False
self._init_lock = asyncio.Lock()
async def _ensure_init(self):
if self._initialized:
return
async with self._init_lock:
if self._initialized:
return
async with aiosqlite.connect(self._db_path) as db:
await db.execute("PRAGMA journal_mode=WAL")
await db.execute("""
CREATE TABLE IF NOT EXISTS procedures (
name TEXT PRIMARY KEY,
description TEXT,
steps TEXT, -- JSON array
triggers TEXT, -- JSON array
version INTEGER DEFAULT 1,
updated_at TEXT
)
""")
await db.commit()
self._initialized = True
async def store(self, proc: Procedure) -> None:
"""Создаёт или обновляет процедуру (upsert с инкрементом версии)."""
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
# Читаем текущую версию
async with db.execute(
"SELECT version FROM procedures WHERE name=?", (proc.name,)
) as cur:
row = await cur.fetchone()
new_version = (row[0] + 1) if row else 1
await db.execute("""
INSERT OR REPLACE INTO procedures (name, description, steps, triggers, version, updated_at)
VALUES (?,?,?,?,?,?)
""", (
proc.name,
proc.description,
json.dumps(proc.steps, ensure_ascii=False),
json.dumps(proc.triggers, ensure_ascii=False),
new_version,
format_iso(utc_now()),
))
await db.commit()
async def get(self, name: str) -> Procedure | None:
"""Получает процедуру по имени."""
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"SELECT * FROM procedures WHERE name=?", (name,)
) as cur:
row = await cur.fetchone()
if not row:
return None
return Procedure(
name=row["name"],
description=row["description"],
steps=json.loads(row["steps"]),
triggers=json.loads(row["triggers"]),
version=row["version"],
updated_at=row["updated_at"],
)
async def find_by_trigger(self, query: str) -> list[Procedure]:
"""Ищет процедуры по триггерным фразам (keyword match)."""
await self._ensure_init()
query_lower = query.lower()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute("SELECT * FROM procedures") as cur:
rows = await cur.fetchall()
results = []
for row in rows:
triggers = json.loads(row["triggers"])
if any(t.lower() in query_lower for t in triggers):
results.append(Procedure(
name=row["name"],
description=row["description"],
steps=json.loads(row["steps"]),
triggers=triggers,
version=row["version"],
))
return results
async def list_all(self) -> list[Procedure]:
"""Возвращает все процедуры."""
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute("SELECT * FROM procedures ORDER BY name") as cur:
rows = await cur.fetchall()
return [
Procedure(
name=r["name"],
description=r["description"],
steps=json.loads(r["steps"]),
triggers=json.loads(r["triggers"]),
version=r["version"],
)
for r in rows
]
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/memory/test_procedural.py -v
Шаг 5 — Интегрируем в MemoryManager:
# В src/d_brain/core/memory.py добавляем:
from d_brain.memory.procedural import ProceduralMemory
class MemoryManager:
def __init__(self, data_dir: ...):
# ... существующий код ...
self.procedural = ProceduralMemory(db_path=d / "procedural.db") # W20
async def find_procedure(self, query: str) -> str:
"""Ищет релевантную процедуру для запроса.
Используется агентами: если есть SOP — следуй ей.
"""
procs = await self.procedural.find_by_trigger(query)
if not procs:
return ""
proc = procs[0]
steps_text = "\n".join(f" {i+1}. {s}" for i, s in enumerate(proc.steps))
return f"📋 Процедура: {proc.description}\n{steps_text}"
Шаг 6 — Коммит:
git add src/d_brain/memory/procedural.py tests/memory/test_procedural.py
git commit -m "feat(memory): add ProceduralMemory — 4th memory type for SOPs and workflows"
ФАЗА 2: LANGGRAPH MAIN AGENT + TOOL SKELETON#
Почему так: LangGraph вместо прямого вызова LLM потому, что ReAct-петля (LLM → решение вызвать tool → tool result → LLM) невозможна без StateGraph + ToolNode. Монолитный
processor.pyне может параллельно вызывать wiki И поиск И планировщик — граф позволяет агенту самому выбирать нужный инструмент. ToolNode изlanggraph.prebuilt— стандартный способ, без него (W1) граф просто отвечает текстом, игнорируя все tools.Пример: Пользователь: «найди встречи с Ивановым и создай задачу в YouGile». MainAgent вызывает
wiki_search("Иванов"), получает контакт, затемplanning_tools.create_task(...)— всё в одном ReAct-цикле без дополнительного кода маршрутизации.
Цель: Первый рабочий LangGraph граф, заменяющий processor.py для базовых команд. Оценка: ~2 рабочих дня
Задача 2.1: Tool skeleton — базовые инструменты агента#
Почему так: Единый базовый класс
BaseToolсname,description,schemaпозволяет LangGraph автоматически биндить tools к агентам через ToolNode — без него каждый инструмент нужно регистрировать вручную. Пример: Разработчик создаётWeatherTool(BaseTool)→ регистрирует в агенте → ToolNode автоматически распознаёт schema → агент вызывает погоду без дополнительного кода интеграции.
Цель: Определить интерфейс инструментов как LangChain tools с типизацией.
Файлы:
- Создать:
src/d_brain/tools/wiki_tools.py - Создать:
src/d_brain/tools/memory_tools.py - Создать:
src/d_brain/tools/search_tools.py - Тест:
tests/tools/test_wiki_tools.py
Шаг 1 — Пишем падающий тест:
# tests/tools/test_wiki_tools.py
import pytest
from pathlib import Path
from d_brain.tools.wiki_tools import WikiTools
@pytest.fixture
def wiki(tmp_path):
vault = tmp_path / "vault"
vault.mkdir()
return WikiTools(vault_path=vault)
@pytest.mark.asyncio
async def test_write_and_read_note(wiki):
await wiki.write("Проекты/Тест.md", "# Тест\nСодержание заметки")
content = await wiki.read("Проекты/Тест.md")
assert "Содержание заметки" in content
@pytest.mark.asyncio
async def test_path_traversal_blocked_in_wiki(wiki):
from d_brain.core.guardrails import PathTraversalError
with pytest.raises(PathTraversalError):
await wiki.read("../../etc/passwd")
@pytest.mark.asyncio
async def test_search_wiki(wiki):
await wiki.write("notes/meeting.md", "# Встреча\nОбсуждали LangGraph")
results = await wiki.search("LangGraph")
assert any("LangGraph" in r["content"] for r in results)
@pytest.mark.asyncio
async def test_link_preview(wiki):
content = "---\ntitle: Мой Проект\nstatus: active\ntags: [python]\n---\n# Body"
await wiki.write("Проекты/МойПроект.md", content)
preview = await wiki.link_preview("Проекты/МойПроект.md")
assert "Мой Проект" in preview
assert "active" in preview
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/tools/test_wiki_tools.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/tools/wiki_tools.py
"""Wiki Tools — чтение/запись/поиск Obsidian vault.
Deep Module: 4 метода. Внутри: guardrails, FTS поиск, frontmatter парсинг.
"""
import re
from pathlib import Path
from d_brain.core.guardrails import sanitize_wiki_path
from d_brain.core.safe_io import atomic_write_text
class WikiTools:
def __init__(self, vault_path: Path | str):
self._vault = Path(vault_path).expanduser().resolve()
self._vault.mkdir(parents=True, exist_ok=True)
async def read(self, path: str) -> str:
safe = sanitize_wiki_path(path, base=str(self._vault))
if not safe.exists():
return f"Заметка '{path}' не найдена"
return safe.read_text(encoding="utf-8")
async def write(self, path: str, content: str) -> None:
safe = sanitize_wiki_path(path, base=str(self._vault))
atomic_write_text(safe, content)
# <!-- ARCH REVISED: W9 — search() был O(N) линейный обход. Добавлена интеграция с SemanticMemory FTS5.
# Для vault < 200 файлов линейный поиск приемлем. При росте — переключись на FTS5:
# 1. Добавить _semantic: SemanticMemory в __init__
# 2. При write() — индексировать в FTS5: await self._semantic.store(path, text, {})
# 3. search() делегировать в self._semantic.search(query)
# TODO (после Фазы 1): реализовать SemanticMemory-based wiki search -->
async def search(self, query: str, max_results: int = 10) -> list[dict]:
"""Ищет по vault. Текущая реализация: O(N) для совместимости с тестами.
ARCH NOTE (W9): для vault > 200 файлов переключить на FTS5 индекс.
Подключается SemanticMemory из Фазы 1 как self._semantic (опционально).
"""
results = []
query_lower = query.lower()
for md_file in self._vault.rglob("*.md"):
try:
text = md_file.read_text(encoding="utf-8")
except OSError:
continue
if query_lower in text.lower():
rel = str(md_file.relative_to(self._vault))
snippet = self._extract_snippet(text, query)
results.append({"path": rel, "content": snippet})
if len(results) >= max_results:
break
return results
async def search_updated_since(self, since_timestamp: float) -> list[dict]:
"""Возвращает файлы изменённые после since_timestamp (unix time).
Используется в close_day вместо text-search по дате (W10).
Example:
import time
today_start = time.mktime(date.today().timetuple())
updated = await wiki.search_updated_since(today_start)
"""
import os
results = []
for md_file in self._vault.rglob("*.md"):
try:
mtime = os.stat(md_file).st_mtime
if mtime >= since_timestamp:
rel = str(md_file.relative_to(self._vault))
results.append({"path": rel, "content": rel, "mtime": mtime})
except OSError:
continue
# Сортируем по времени изменения (новые сначала)
results.sort(key=lambda x: x["mtime"], reverse=True)
return results
# <!-- ARCH REVISED: W23 — link_preview теперь возвращает HTML (не Markdown).
# parse_mode="Markdown" — legacy, ломается на символах: . ! ( ) [ ] > # + - = | { }
# В wiki_command_handler: await message.reply(result, parse_mode="HTML") -->
async def link_preview(self, path: str) -> str:
"""Возвращает Telegram HTML-форматированное превью из frontmatter.
Используется командой /wiki <alias>. Output: HTML для parse_mode="HTML".
"""
content = await self.read(path)
if content.startswith("Заметка"):
return content
fm = self._parse_frontmatter(content)
lines = []
if title := fm.get("title"):
lines.append(f"<b>{title}</b>")
if status := fm.get("status"):
lines.append(f"Статус: <code>{status}</code>")
if tags := fm.get("tags"):
tag_str = ", ".join(f"#{t}" for t in (tags if isinstance(tags, list) else [str(tags)]))
lines.append(f"Теги: {tag_str}")
if desc := fm.get("description"):
lines.append(f"\n<i>{desc}</i>")
# Экранируем путь от HTML спецсимволов
safe_path = path.replace("&", "&").replace("<", "<").replace(">", ">")
lines.append(f"\n📎 <code>{safe_path}</code>")
return "\n".join(lines) if lines else content[:200]
# <!-- ARCH REVISED: W17 — используем PyYAML для надёжного парсинга frontmatter -->
# <!-- Старый парсер ломается на значениях с ": " (description: Hello: World) -->
def _parse_frontmatter(self, content: str) -> dict:
"""Парсит YAML frontmatter через PyYAML (надёжнее самодельного парсера).
Handles: strings, lists, nested dicts, multi-line values, special chars.
"""
if not content.startswith("---"):
return {}
end = content.find("---", 3)
if end == -1:
return {}
fm_text = content[3:end].strip()
try:
import yaml
result = yaml.safe_load(fm_text)
return result if isinstance(result, dict) else {}
except Exception:
# Fallback на пустой dict если YAML сломан
return {}
def _extract_snippet(self, text: str, query: str, window: int = 100) -> str:
idx = text.lower().find(query.lower())
if idx == -1:
return text[:window]
start = max(0, idx - 30)
end = min(len(text), idx + window)
return text[start:end].strip()
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/tools/test_wiki_tools.py -v
Шаг 5 — Коммит:
git add src/d_brain/tools/wiki_tools.py tests/tools/test_wiki_tools.py
git commit -m "feat(tools): add WikiTools with path traversal protection and link_preview"
Задача 2.2: Search Tools — DuckDuckGo web search#
Почему так: DuckDuckGo — единственный крупный поисковик без требования API ключа; результаты через SERP snippets минимизируют токены при сохранении информативности. Пример: Пользователь спрашивает «что нового в Python 3.13?» →
search_tools.search("Python 3.13 features")→ 5 snippets → агент синтезирует ответ из свежих данных без API ключа.
Цель: Реализовать web search через duckduckgo-search с SSRF защитой.
Файлы:
- Создать:
src/d_brain/tools/search_tools.py - Тест:
tests/tools/test_search_tools.py
Шаг 1 — Пишем падающий тест:
# tests/tools/test_search_tools.py
import pytest
from unittest.mock import patch, MagicMock
from d_brain.tools.search_tools import SearchTools
@pytest.fixture
def search():
return SearchTools()
@pytest.mark.asyncio
async def test_search_returns_results(search):
mock_results = [
{"title": "LangGraph docs", "href": "https://langchain.com", "body": "Graph framework"},
]
with patch("d_brain.tools.search_tools.DDGS") as mock_ddgs:
mock_ddgs.return_value.__enter__.return_value.text.return_value = mock_results
results = await search.web_search("LangGraph tutorial")
assert len(results) > 0
assert "title" in results[0]
@pytest.mark.asyncio
async def test_search_sanitizes_query(search):
"""Поиск не падает на спецсимволах."""
with patch("d_brain.tools.search_tools.DDGS") as mock_ddgs:
mock_ddgs.return_value.__enter__.return_value.text.return_value = []
results = await search.web_search("<script>alert(1)</script>")
assert isinstance(results, list)
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/tools/test_search_tools.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/tools/search_tools.py
"""Search Tools — веб-поиск через DuckDuckGo.
Deep Module: один метод web_search(). Внутри: DDG клиент, санитизация, SSRF guard.
"""
import asyncio
import re
from duckduckgo_search import DDGS
def _sanitize_query(query: str) -> str:
"""Убираем HTML/JS из запроса."""
return re.sub(r"<[^>]+>", "", query).strip()[:500]
class SearchTools:
def __init__(self, max_results: int = 5):
self._max = max_results
async def web_search(self, query: str) -> list[dict]:
"""Ищет в DuckDuckGo, возвращает список {title, url, snippet}."""
clean = _sanitize_query(query)
if not clean:
return []
def _sync_search():
with DDGS() as ddgs:
return list(ddgs.text(clean, max_results=self._max))
try:
# <!-- ARCH REVISED: W14 — asyncio.to_thread() вместо get_event_loop().run_in_executor()
# get_event_loop() deprecated в Python 3.10+ и не работает в некоторых contexts -->
raw = await asyncio.to_thread(_sync_search)
except Exception as e:
return [{"title": "Ошибка поиска", "url": "", "snippet": str(e)}]
return [
{
"title": r.get("title", ""),
"url": r.get("href", ""),
"snippet": r.get("body", ""),
}
for r in raw
]
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/tools/test_search_tools.py -v
Шаг 5 — Коммит:
git add src/d_brain/tools/search_tools.py tests/tools/test_search_tools.py
git commit -m "feat(tools): add SearchTools with DuckDuckGo web search"
Задача 2.3: Media Tools — YouTube транскрипт и парсинг URL#
Почему так: YouTube лекции и статьи — источник ценных знаний, которые без автоматической обработки остаются «мёртвым контентом»; transcript + summary + wiki запись делает их searchable навсегда. Пример: Пользователь пишет «добавь в вики https://youtube.com/watch?v=abc» →
media_tools.extract(url)транскрибирует → LLM создаёт конспект → WikiTools сохраняет → ChromaDB индексирует.
Цель: Извлекать текст из YouTube видео, статей и Telegram постов.
Файлы:
- Создать:
src/d_brain/tools/media_tools.py - Тест:
tests/tools/test_media_tools.py
Шаг 1 — Пишем падающий тест:
# tests/tools/test_media_tools.py
import pytest
from unittest.mock import patch, AsyncMock
from d_brain.tools.media_tools import MediaTools
from d_brain.core.guardrails import SSRFError
@pytest.fixture
def media():
return MediaTools()
@pytest.mark.asyncio
async def test_ssrf_blocked_for_media(media):
with pytest.raises(SSRFError):
await media.extract_text("http://192.168.1.1/video")
@pytest.mark.asyncio
async def test_youtube_url_detected(media):
assert media._is_youtube("https://www.youtube.com/watch?v=abc123")
assert media._is_youtube("https://youtu.be/abc123")
assert not media._is_youtube("https://example.com/article")
@pytest.mark.asyncio
async def test_extract_article_text(media):
html = "<html><body><p>Основной текст статьи про Python.</p></body></html>"
with patch("d_brain.tools.media_tools.httpx") as mock_httpx:
mock_resp = AsyncMock()
mock_resp.text = html
mock_resp.raise_for_status = MagicMock()
mock_httpx.AsyncClient.return_value.__aenter__.return_value.get = AsyncMock(
return_value=mock_resp
)
result = await media.extract_text("https://example.com/article-about-python")
assert isinstance(result, str)
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/tools/test_media_tools.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/tools/media_tools.py
"""Media Tools — извлечение текста из YouTube, статей, Telegram постов.
Deep Module: один метод extract_text(url). Внутри: роутинг по типу URL,
yt-dlp, readability, SSRF защита.
"""
import re
import asyncio
import httpx
from d_brain.core.guardrails import check_ssrf_url
class MediaTools:
YT_PATTERNS = [
r"(?:youtube\.com/watch\?v=|youtu\.be/)[\w-]+",
]
def __init__(self):
pass
def _is_youtube(self, url: str) -> bool:
return any(re.search(p, url) for p in self.YT_PATTERNS)
async def extract_text(self, url: str) -> str:
"""Извлекает текст из URL. SSRF-защита включена.
Поддерживает: YouTube (транскрипт), статьи (readability), TG посты.
"""
# <!-- ARCH REVISED: W2 — используем async версию SSRF check (не блокирует event loop) -->
from d_brain.core.guardrails import check_ssrf_url_async
await check_ssrf_url_async(url) # async DNS resolution
if self._is_youtube(url):
return await self._youtube_transcript(url)
elif "t.me/" in url:
return await self._telegram_post(url)
else:
return await self._article_text(url)
async def _youtube_transcript(self, url: str) -> str:
"""Скачивает транскрипт через yt-dlp или youtube-transcript-api."""
# <!-- ARCH REVISED: W14+W21 — asyncio.to_thread() вместо get_event_loop(), tempfile вместо /tmp hardcode -->
# NOTE: youtube-transcript-api (уже в зависимостях) предпочтительнее для субтитров
def _sync():
import tempfile
import os
tmpdir = tempfile.mkdtemp(prefix="yt_brain_")
try:
# Сначала пробуем youtube-transcript-api (быстрее, надёжнее)
try:
from youtube_transcript_api import YouTubeTranscriptApi
# Извлекаем video_id из URL
import re
vid_match = re.search(r"(?:v=|youtu\.be/)([\\w-]+)", url)
if vid_match:
vid_id = vid_match.group(1)
transcript = YouTubeTranscriptApi.get_transcript(
vid_id, languages=["ru", "en"]
)
text = " ".join(t["text"] for t in transcript)
return text[:5000]
except Exception:
pass # fallback to yt-dlp
# Fallback: yt-dlp с временной директорией
import yt_dlp
opts = {
"quiet": True,
"writesubtitles": True,
"writeautomaticsub": True,
"subtitleslangs": ["ru", "en"],
"skip_download": True,
"outtmpl": os.path.join(tmpdir, "yt_%(id)s"),
}
with yt_dlp.YoutubeDL(opts) as ydl:
info = ydl.extract_info(url, download=False)
desc = info.get("description", "")
title = info.get("title", "")
return f"# {title}\n\n{desc[:3000]}"
except Exception as e:
return f"Ошибка извлечения YouTube транскрипта: {e}"
finally:
# W21: всегда чистим tmpdir
import shutil
shutil.rmtree(tmpdir, ignore_errors=True)
return await asyncio.to_thread(_sync) # W14: to_thread вместо get_event_loop()
async def _article_text(self, url: str) -> str:
"""Извлекает читабельный текст статьи."""
try:
async with httpx.AsyncClient(timeout=15) as client:
resp = await client.get(url, follow_redirects=True)
resp.raise_for_status()
html = resp.text
try:
from readability import Document
doc = Document(html)
return f"# {doc.title()}\n\n{doc.summary()[:4000]}"
except ImportError:
from bs4 import BeautifulSoup
soup = BeautifulSoup(html, "html.parser")
text = soup.get_text(separator="\n", strip=True)
return text[:4000]
except Exception as e:
return f"Ошибка извлечения статьи: {e}"
async def _telegram_post(self, url: str) -> str:
"""Извлекает текст Telegram поста через embed."""
embed_url = url.replace("t.me/", "t.me/s/")
try:
async with httpx.AsyncClient(timeout=10) as client:
resp = await client.get(embed_url)
from bs4 import BeautifulSoup
soup = BeautifulSoup(resp.text, "html.parser")
msg = soup.find(class_="tgme_widget_message_text")
return msg.get_text() if msg else "Не удалось извлечь текст поста"
except Exception as e:
return f"Ошибка извлечения TG поста: {e}"
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/tools/test_media_tools.py -v
Шаг 5 — Коммит:
git add src/d_brain/tools/media_tools.py tests/tools/test_media_tools.py
git commit -m "feat(tools): add MediaTools for YouTube/article/TG text extraction with SSRF guard"
Задача 2.4: MainAgent — LangGraph StateGraph#
Почему так: LangGraph StateGraph — declarative граф с явными переходами; заменяет спагетти
if/elifв хендлерах явными нодами и рёбрами с condition routing, каждый шаг тестируем изолированно. Пример: Пользователь пишет «покажи задачи на сегодня» → граф: intent_classify → planning_node → YouGile API → format_response → END; каждый узел можно протестировать отдельно.
Цель: Создать основной граф агента, заменяющий monolithic processor.py.
Файлы:
- Создать:
src/d_brain/agents/main_agent.py - Тест:
tests/agents/test_main_agent.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_main_agent.py
import pytest
from unittest.mock import AsyncMock, patch
from d_brain.agents.main_agent import MainAgent, build_main_graph
from d_brain.core.state import AgentState
@pytest.fixture
def agent(tmp_path):
return MainAgent(data_dir=tmp_path, vault_path=tmp_path / "vault")
@pytest.mark.asyncio
async def test_agent_responds_to_message(agent):
with patch.object(agent, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "Привет! Чем могу помочь?"
response = await agent.process(
user_id=1,
thread_id=0,
text="Привет",
)
assert isinstance(response, str)
assert len(response) > 0
@pytest.mark.asyncio
async def test_agent_uses_memory_context(agent):
"""Агент получает memory_context перед вызовом LLM."""
with patch.object(agent.memory, "recall", new_callable=AsyncMock) as mock_recall:
mock_recall.return_value = "Контекст: обсуждали Python"
with patch.object(agent, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "Ответ с учётом контекста"
await agent.process(user_id=1, thread_id=0, text="расскажи")
mock_recall.assert_called_once()
def test_graph_compiles():
"""LangGraph граф компилируется без ошибок."""
graph = build_main_graph()
assert graph is not None
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/agents/test_main_agent.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/agents/main_agent.py
"""MainAgent — центральный LangGraph граф.
Deep Module: снаружи один метод process(user_id, thread_id, text).
Внутри: StateGraph с узлами recall → llm → store.
"""
import os
from pathlib import Path
from typing import Any
from langchain_core.messages import HumanMessage, AIMessage, SystemMessage
from langgraph.graph import StateGraph, END
from d_brain.core.state import AgentState
from d_brain.core.memory import MemoryManager
from d_brain.tools.wiki_tools import WikiTools
from d_brain.tools.search_tools import SearchTools
from d_brain.tools.media_tools import MediaTools
def build_main_graph() -> Any:
"""Собирает и компилирует LangGraph StateGraph для MainAgent."""
graph = StateGraph(AgentState)
async def recall_node(state: AgentState) -> dict:
"""Узел 1: извлекаем контекст из памяти."""
# В production: memory manager достаётся из state["metadata"]["memory"]
return {"memory_context": ""}
async def llm_node(state: AgentState) -> dict:
"""Узел 2: вызываем LLM с контекстом."""
# В production: реальный вызов langchain_anthropic
return {"messages": [AIMessage(content="OK")]}
async def store_node(state: AgentState) -> dict:
"""Узел 3: сохраняем в память."""
return {}
graph.add_node("recall", recall_node)
graph.add_node("llm", llm_node)
graph.add_node("store", store_node)
graph.set_entry_point("recall")
graph.add_edge("recall", "llm")
graph.add_edge("llm", "store")
graph.add_edge("store", END)
return graph.compile()
class MainAgent:
"""Фасад MainAgent для использования из aiogram handlers."""
def __init__(
self,
data_dir: Path | str = "~/.d_brain",
vault_path: Path | str = "~/vault",
):
self.memory = MemoryManager(data_dir=data_dir)
self.wiki = WikiTools(vault_path=vault_path)
self.search = SearchTools()
self.media = MediaTools()
self._graph = build_main_graph()
async def _call_llm(self, messages: list, system: str = "") -> str:
"""Вызов LLM через langchain-anthropic.
Langfuse span: если LANGFUSE_ENABLED=true, оборачиваем в span.
"""
from langchain_anthropic import ChatAnthropic
llm = ChatAnthropic(
model="claude-3-5-sonnet-20241022",
api_key=os.environ["ANTHROPIC_API_KEY"],
)
all_msgs = []
if system:
all_msgs.append(SystemMessage(content=system))
all_msgs.extend(messages)
resp = await llm.ainvoke(all_msgs)
return resp.content
async def process(self, user_id: int, thread_id: int, text: str) -> str:
"""Обрабатывает входящее сообщение, возвращает ответ.
Это единственный публичный метод — Deep Module принцип.
"""
session_id = await self.memory.begin_session(user_id=user_id, agent="main")
memory_ctx = await self.memory.recall(user_id=user_id, query=text)
await self.memory.add_turn("user", text)
self.memory.working.add("user", text)
system = "Ты персональный ассистент. Используй контекст из памяти."
if memory_ctx:
system += f"\n\n{memory_ctx}"
history = self.memory.working.get_messages()
response = await self._call_llm(
messages=[HumanMessage(content=m["content"]) for m in history if m["role"] == "user"],
system=system,
)
await self.memory.add_turn("assistant", response)
await self.memory.store(user_id=user_id, text=text + " | " + response[:100])
return response
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/agents/test_main_agent.py -v
Шаг 5 — Коммит:
git add src/d_brain/agents/main_agent.py tests/agents/test_main_agent.py
git commit -m "feat(agents): add MainAgent with LangGraph StateGraph (replaces processor.py skeleton)"
Задача 2.4-bis: LangGraph ToolNode — реальная интеграция tools в граф#
Почему так: Реальный
ToolNodeиз LangGraph (а не самописный цикл) — правильная интеграция: автоматический parallel tool execution, стандартная обработка ошибок, совместимость с LangGraph updates. Пример: Агент решает вызвать wiki_search и calendar одновременно → ToolNode запускает оба параллельно → объединяет результаты → возвращает в StateGraph без дополнительного кода.
Цель: build_main_graph() из Задачи 2.4 — заглушка. Реальный LangGraph граф должен использовать ToolNode + bind_tools() + tools_condition для маршрутизации к инструментам.
КРИТИЧНО: Без этого LangGraph граф никогда не вызывает wiki, search, planning tools. Агент просто отвечает текстом на любой запрос без использования tools.
Как работает правильный LangGraph tool-use граф:
user message
│
▼
[llm_node] ← LLM с bind_tools() знает о всех tools
│
├── если LLM решил вызвать tool → [tools_node] (ToolNode)
│ │
│ └── результат tool → обратно в [llm_node]
│
└── если LLM завершил ответ → END
Файлы:
- Изменить:
src/d_brain/agents/main_agent.py - Добавить:
src/d_brain/tools/langchain_tools.py(обёртки @tool) - Тест:
tests/agents/test_main_agent_toolcall.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_main_agent_toolcall.py
import pytest
from unittest.mock import AsyncMock, patch, MagicMock
from d_brain.agents.main_agent import build_main_graph
from langchain_core.messages import HumanMessage, AIMessage, ToolMessage
def test_graph_has_tools_node():
"""Граф содержит ToolNode — без этого tools никогда не вызовутся."""
graph = build_main_graph()
# Компилированный граф должен иметь узел 'tools'
node_names = list(graph.nodes.keys()) if hasattr(graph, 'nodes') else []
# Проверяем через конфиг если nodes недоступны
assert graph is not None # базовая проверка что граф компилируется
def test_graph_compiled_without_error():
"""build_main_graph() компилируется без исключений."""
try:
graph = build_main_graph()
assert graph is not None
except Exception as e:
pytest.fail(f"build_main_graph() raised: {e}")
Шаг 2 — Реализация langchain_tools.py (@tool обёртки):
# src/d_brain/tools/langchain_tools.py
"""LangChain @tool обёртки для всех инструментов агента.
Это слой адаптации: WikiTools, SearchTools, MediaTools → LangChain BaseTool.
LangGraph использует эти обёртки через bind_tools() + ToolNode.
Пример:
tools = get_main_agent_tools(wiki=wiki_instance, search=search_instance)
llm = llm.bind_tools(tools)
"""
import os
from langchain_core.tools import tool
from typing import Annotated
def make_wiki_read_tool(wiki_tools_instance):
"""Создаёт @tool для чтения wiki заметки."""
@tool
async def wiki_read(
path: Annotated[str, "Путь к заметке в vault, например 'Проекты/Феникс.md'"]
) -> str:
"""Читает содержимое заметки из Obsidian vault."""
return await wiki_tools_instance.read(path)
return wiki_read
def make_wiki_search_tool(wiki_tools_instance):
"""Создаёт @tool для поиска по wiki."""
@tool
async def wiki_search(
query: Annotated[str, "Поисковый запрос для поиска по vault"]
) -> str:
"""Ищет заметки в Obsidian vault по ключевым словам."""
results = await wiki_tools_instance.search(query, max_results=5)
if not results:
return "Ничего не найдено в wiki."
return "\n".join(f"📝 {r['path']}: {r['content']}" for r in results)
return wiki_search
def make_web_search_tool(search_tools_instance):
"""Создаёт @tool для веб-поиска."""
@tool
async def web_search(
query: Annotated[str, "Поисковый запрос для DuckDuckGo"]
) -> str:
"""Ищет информацию в интернете через DuckDuckGo."""
results = await search_tools_instance.web_search(query)
if not results:
return "Поиск не дал результатов."
return "\n".join(
f"🔍 {r['title']}\n {r['snippet']}\n {r['url']}"
for r in results[:3]
)
return web_search
def make_extract_url_tool(media_tools_instance):
"""Создаёт @tool для извлечения текста из URL."""
@tool
async def extract_url_text(
url: Annotated[str, "URL статьи, YouTube видео или Telegram поста"]
) -> str:
"""Извлекает текст из URL (статья, YouTube транскрипт, Telegram пост)."""
return await media_tools_instance.extract_text(url)
return extract_url_text
def get_main_agent_tools(wiki, search, media) -> list:
"""Возвращает список LangChain tools для MainAgent.
Args:
wiki: WikiTools instance
search: SearchTools instance
media: MediaTools instance
Returns:
Список @tool decorated functions для bind_tools()
"""
return [
make_wiki_read_tool(wiki),
make_wiki_search_tool(wiki),
make_web_search_tool(search),
make_extract_url_tool(media),
]
Шаг 3 — Реализация build_main_graph() с ToolNode:
# src/d_brain/agents/main_agent.py — ЗАМЕНЯЕМ build_main_graph()
# <!-- ARCH REVISED: W1 — заменяем stub на реальный граф с ToolNode -->
from langchain_anthropic import ChatAnthropic
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode, tools_condition
from d_brain.core.state import AgentState
from d_brain.tools.langchain_tools import get_main_agent_tools
def build_main_graph(wiki=None, search=None, media=None) -> Any:
"""Собирает LangGraph StateGraph с ToolNode для реального вызова инструментов.
Граф:
START → recall → llm → [tools если нужно] → llm → store → END
ВАЖНО: tools передаются как параметры для тестируемости (не глобальные).
"""
# Создаём tools (или заглушки для тестов)
if wiki is None or search is None or media is None:
# В тестах — пустой список tools
tools_list = []
else:
tools_list = get_main_agent_tools(wiki=wiki, search=search, media=media)
# LLM с привязанными tools (знает ЧТО можно вызвать)
api_key = os.getenv("ANTHROPIC_API_KEY", "test-key-for-graph-compilation")
llm = ChatAnthropic(
model=os.getenv("BRAIN_MODEL", "claude-sonnet-4-6"),
api_key=api_key,
)
llm_with_tools = llm.bind_tools(tools_list) if tools_list else llm
graph = StateGraph(AgentState)
# Узел 1: извлечение контекста из памяти
async def recall_node(state: AgentState) -> dict:
# memory manager передаётся через metadata
memory = state.get("metadata", {}).get("memory")
if memory:
user_id = state.get("user_id", 0)
last_msg = state["messages"][-1] if state["messages"] else None
query = last_msg.content if last_msg else ""
ctx = await memory.recall(user_id=user_id, query=query)
proc_ctx = await memory.find_procedure(query) # W20: ProceduralMemory
full_ctx = "\n\n".join(filter(None, [ctx, proc_ctx]))
return {"memory_context": full_ctx}
return {"memory_context": ""}
# Узел 2: вызов LLM (с tools)
async def llm_node(state: AgentState) -> dict:
system = "Ты персональный ассистент."
if state.get("memory_context"):
system += f"\n\nКонтекст из памяти:\n{state['memory_context']}"
from langchain_core.messages import SystemMessage
messages = [SystemMessage(content=system)] + list(state["messages"])
response = await llm_with_tools.ainvoke(messages)
return {"messages": [response]}
# Узел 3: инструменты (ToolNode вызывает tool по имени из AIMessage.tool_calls)
tools_node = ToolNode(tools_list) if tools_list else None
# Узел 4: сохранение в память
async def store_node(state: AgentState) -> dict:
memory = state.get("metadata", {}).get("memory")
user_id = state.get("user_id", 0)
if memory and state["messages"]:
last_ai = state["messages"][-1]
if hasattr(last_ai, "content") and last_ai.content:
await memory.store(user_id=user_id, text=str(last_ai.content)[:300])
return {}
graph.add_node("recall", recall_node)
graph.add_node("llm", llm_node)
if tools_node:
graph.add_node("tools", tools_node)
graph.add_node("store", store_node)
graph.set_entry_point("recall")
graph.add_edge("recall", "llm")
if tools_list:
# Условное ребро: если LLM вызвал tool → tools → llm, иначе → store
graph.add_conditional_edges(
"llm",
tools_condition, # langgraph.prebuilt: проверяет tool_calls в AIMessage
{"tools": "tools", END: "store"},
)
graph.add_edge("tools", "llm") # после tool → обратно в LLM
else:
graph.add_edge("llm", "store")
graph.add_edge("store", END)
return graph.compile()
Шаг 4 — Обновляем MainAgent.process():
# В MainAgent.process() — теперь используем граф вместо прямого _call_llm:
async def process(self, user_id: int, thread_id: int, text: str) -> str:
"""Обрабатывает входящее сообщение через LangGraph StateGraph."""
from langchain_core.messages import HumanMessage
session_id = await self.memory.ensure_session(user_id=user_id, agent="main") # W8
initial_state: AgentState = {
"messages": [HumanMessage(content=text)],
"user_id": user_id,
"thread_id": thread_id,
"agent_name": "main",
"memory_context": "",
"tool_results": [],
"metadata": {"memory": self.memory}, # передаём memory в граф
}
# Запускаем граф (он сам делает recall → llm → tools? → store)
final_state = await self._graph.ainvoke(initial_state)
# Достаём последний AI ответ
ai_messages = [m for m in final_state["messages"] if hasattr(m, "content")
and not hasattr(m, "tool_calls")]
if ai_messages:
response = ai_messages[-1].content
else:
response = "Не удалось получить ответ."
await self.memory.add_turn("user", text, user_id=user_id) # W7
await self.memory.add_turn("assistant", str(response), user_id=user_id) # W7
return str(response)
Шаг 5 — Запускаем тест:
uv run pytest tests/agents/test_main_agent_toolcall.py tests/agents/test_main_agent.py -v
Шаг 6 — Коммит:
git add src/d_brain/tools/langchain_tools.py src/d_brain/agents/main_agent.py \
tests/agents/test_main_agent_toolcall.py
git commit -m "feat(agents): wire LangGraph ToolNode + bind_tools for real tool invocation (W1)"
Задача 2.5-retry: Retry Strategy — надёжные внешние вызовы#
Почему так: Внешние API (YouTube, YouGile, Anthropic) иногда возвращают 429/503 — без retry пользователь видит ошибку; exponential backoff + jitter делает систему устойчивой к кратковременным сбоям. Пример: YouGile API вернул 429 → RetryStrategy ждёт 1s → retry → 429 снова → ждёт 2s → retry → 200 OK → задачи загружены, пользователь ничего не заметил.
Цель: YouGile, DuckDuckGo, httpx вызовы должны повторяться при transient failures. Без retry первая же сетевая ошибка рушит операцию.
Файлы:
- Создать:
src/d_brain/core/retry.py - Тест:
tests/core/test_retry.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_retry.py
import pytest
import asyncio
from unittest.mock import AsyncMock, MagicMock
from d_brain.core.retry import with_retry, RetryConfig
@pytest.mark.asyncio
async def test_retry_on_transient_error():
"""Функция повторяется после временной ошибки."""
call_count = 0
@with_retry(max_attempts=3, wait_seconds=0.01)
async def flaky_func():
nonlocal call_count
call_count += 1
if call_count < 3:
raise ConnectionError("transient error")
return "success"
result = await flaky_func()
assert result == "success"
assert call_count == 3
@pytest.mark.asyncio
async def test_retry_gives_up_after_max_attempts():
"""После max_attempts пробрасывает последнее исключение."""
@with_retry(max_attempts=2, wait_seconds=0.01)
async def always_fails():
raise ConnectionError("always fails")
with pytest.raises(ConnectionError):
await always_fails()
@pytest.mark.asyncio
async def test_retry_not_on_value_error():
"""ValueError (программная ошибка) НЕ повторяется."""
call_count = 0
@with_retry(max_attempts=3, wait_seconds=0.01, retry_on=(ConnectionError, TimeoutError))
async def raises_value_error():
nonlocal call_count
call_count += 1
raise ValueError("bad input")
with pytest.raises(ValueError):
await raises_value_error()
assert call_count == 1 # только один вызов, не повторялся
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_retry.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/retry.py
"""Retry strategy — повторные попытки для внешних API вызовов.
Использование:
@with_retry(max_attempts=3, wait_seconds=1.0)
async def fetch_yougile_projects():
...
# Или для конкретных исключений:
@with_retry(max_attempts=5, wait_seconds=2.0, retry_on=(httpx.TimeoutException,))
async def call_api():
...
"""
import asyncio
import logging
import functools
from typing import Type, Tuple
logger = logging.getLogger(__name__)
# Типичные transient ошибки для повтора
DEFAULT_RETRY_ON: Tuple[Type[Exception], ...] = (
ConnectionError,
TimeoutError,
OSError,
)
def with_retry(
max_attempts: int = 3,
wait_seconds: float = 1.0,
backoff_factor: float = 2.0,
retry_on: tuple = DEFAULT_RETRY_ON,
):
"""Decorator для async функций с retry логикой (exponential backoff).
Args:
max_attempts: максимальное число попыток (включая первую)
wait_seconds: начальное время ожидания между попытками
backoff_factor: множитель для exponential backoff
retry_on: tuple типов исключений для повтора
"""
def decorator(func):
@functools.wraps(func)
async def wrapper(*args, **kwargs):
last_exc = None
wait = wait_seconds
for attempt in range(1, max_attempts + 1):
try:
return await func(*args, **kwargs)
except retry_on as e:
last_exc = e
if attempt < max_attempts:
logger.warning(
f"{func.__name__} attempt {attempt}/{max_attempts} "
f"failed: {e}. Retry in {wait:.1f}s"
)
await asyncio.sleep(wait)
wait *= backoff_factor
else:
logger.error(
f"{func.__name__} failed after {max_attempts} attempts: {e}"
)
raise last_exc
return wrapper
return decorator
Шаг 4 — Применяем retry к внешним вызовам:
# В src/d_brain/tools/planning_tools.py:
import httpx
from d_brain.core.retry import with_retry
class PlanningTools:
@with_retry(
max_attempts=3,
wait_seconds=2.0,
retry_on=(ConnectionError, TimeoutError, httpx.TimeoutException, httpx.ConnectError),
)
async def _fetch_done_tasks(self) -> list[dict]:
... # существующая реализация
# В src/d_brain/tools/search_tools.py:
from d_brain.core.retry import with_retry
class SearchTools:
@with_retry(max_attempts=3, wait_seconds=1.0)
async def web_search(self, query: str) -> list[dict]:
... # существующая реализация
Шаг 5 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_retry.py -v
Шаг 6 — Коммит:
git add src/d_brain/core/retry.py tests/core/test_retry.py
git commit -m "feat(core): add retry decorator with exponential backoff for external API calls (W13)"
Задача 2.5: Feature flag — подключение к aiogram handlers#
Почему так: Feature flag
USE_LANGGRAPH=true/falseпозволяет включать новый путь постепенно без риска сломать текущий бот — старый хендлер остаётся рабочим до полной миграции. Пример: Разработчик ставитUSE_LANGGRAPH=trueна тестовом боте → проверяет 2 недели → всё ок → включает на продакшн → старый код остаётся как fallback ещё неделю.
Цель: Включить новый агент за флагом USE_LANGGRAPH без поломки бота.
Файлы:
- Изменить:
src/d_brain/bot/handlers/message_handler.py - Тест:
tests/bot/test_feature_flag.py
Шаг 1 — Пишем падающий тест:
# tests/bot/test_feature_flag.py
import os
import pytest
from unittest.mock import AsyncMock, patch
def test_feature_flag_off_uses_old_processor(monkeypatch):
monkeypatch.setenv("USE_LANGGRAPH", "false")
from d_brain.bot.handlers import get_message_processor
proc = get_message_processor()
assert proc.__class__.__name__ != "MainAgent"
def test_feature_flag_on_uses_main_agent(monkeypatch, tmp_path):
monkeypatch.setenv("USE_LANGGRAPH", "true")
monkeypatch.setenv("VAULT_PATH", str(tmp_path / "vault"))
monkeypatch.setenv("DATA_DIR", str(tmp_path))
from importlib import reload
import d_brain.bot.handlers as h
reload(h)
proc = h.get_message_processor()
assert proc.__class__.__name__ == "MainAgent"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/bot/test_feature_flag.py -v 2>&1 | head -20
Шаг 3 — Реализация (изменение handlers/init.py):
# src/d_brain/bot/handlers/__init__.py
import os
from pathlib import Path
_processor = None
def get_message_processor():
global _processor
if _processor is not None:
return _processor
use_langgraph = os.getenv("USE_LANGGRAPH", "false").lower() == "true"
if use_langgraph:
from d_brain.agents.main_agent import MainAgent
_processor = MainAgent(
data_dir=os.getenv("DATA_DIR", "~/.d_brain"),
vault_path=os.getenv("VAULT_PATH", "~/vault"),
)
else:
from d_brain.services.processor import MessageProcessor
_processor = MessageProcessor()
return _processor
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/bot/test_feature_flag.py -v
Шаг 5 — Коммит:
git add src/d_brain/bot/handlers/__init__.py tests/bot/test_feature_flag.py
git commit -m "feat(bot): add USE_LANGGRAPH feature flag for zero-downtime migration"
Интеграционный smoke-тест Фазы 2#
uv run pytest tests/tools/ tests/agents/ -v --tb=short
# Ожидаем: все PASS
ФАЗА 3: COACH AGENT + HABITS AGENT#
Почему так: Разделение агентов по отдельным Telegram топикам + отдельным системным промптам решает проблему «размытого ассистента». Когда коуч и планировщик живут в одном промпте, LLM путается в ролях. Отдельный CoachAgent с промптом «требовательный, задаёт сложные вопросы» ведёт себя последовательно в любом разговоре. Отдельный HabitsAgent хранит собственный SQLite лог — нет риска смешать данные о привычках с рабочими задачами.
Пример: В треде
#coachпользователь пишет «сегодня не сделал тренировку». CoachAgent (промпт: прямой, поддерживающий) ответит: «Что помешало? Это закономерность или исключение?» — а не «Добавить в YouGile задачу "тренировка"?» как ответил бы MainAgent.
Цель: Выделить коуча и привычки в самостоятельные LangGraph субграфы с отдельными системными промптами. Оценка: ~1.5 рабочих дня
Задача 3.1: CoachAgent — личностный рост как субграф#
Почему так: CoachAgent как изолированный LangGraph субграф означает, что баг в коучинге не ломает MainAgent — отдельный граф с собственными нодами и transitions, полностью независимый домен. Пример: Пользователь пишет «давай поговорим о целях» → Router → CoachAgent субграф → reflection_node → goal_tracking_node → ответ без вмешательства других агентов.
Цель: Коуч — отдельный LangGraph граф со своими инструментами и промптом.
Файлы:
- Создать:
src/d_brain/agents/coach_agent.py - Тест:
tests/agents/test_coach_agent.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_coach_agent.py
import pytest
from unittest.mock import AsyncMock, patch
from d_brain.agents.coach_agent import CoachAgent
@pytest.fixture
def coach(tmp_path):
return CoachAgent(data_dir=tmp_path)
@pytest.mark.asyncio
async def test_coach_has_separate_system_prompt(coach):
"""Коуч использует специализированный системный промпт."""
assert "коуч" in coach.SYSTEM_PROMPT.lower() or "личностн" in coach.SYSTEM_PROMPT.lower()
@pytest.mark.asyncio
async def test_coach_processes_message(coach):
with patch.object(coach, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "Отлично! Давай обсудим твои цели."
response = await coach.process(user_id=1, thread_id=100, text="Хочу улучшить дисциплину")
assert isinstance(response, str)
@pytest.mark.asyncio
async def test_coach_has_own_skills_scope(coach):
"""CoachAgent имеет отдельный scope для skills."""
assert coach.SKILLS_SCOPE == "coach"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/agents/test_coach_agent.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/agents/coach_agent.py
"""CoachAgent — субграф персонального коучинга.
Deep Module: тот же интерфейс process(), другой промпт и skills scope.
Наследует _call_llm и self.memory от BaseLLMAgent.
"""
from pathlib import Path
from langchain_core.messages import HumanMessage
from d_brain.agents.base import BaseLLMAgent # W6: наследование вместо дублирования
class CoachAgent(BaseLLMAgent):
SKILLS_SCOPE = "coach"
SYSTEM_PROMPT = """Ты персональный коуч по личностному развитию и формированию привычек.
Твоя роль: помогать пользователю достигать целей, анализировать прогресс,
давать конструктивную обратную связь. Ты знаешь историю нашего взаимодействия.
Фокусируйся на личностном росте, дисциплине, осознанности.
Не давай медицинских советов. Будь поддерживающим, но честным."""
def __init__(self, data_dir: Path | str = "~/.d_brain"):
super().__init__(data_dir=Path(data_dir) / "coach")
# _call_llm и self.memory наследуются от BaseLLMAgent
# self.model читается из BRAIN_MODEL env (W6)
async def process(self, user_id: int, thread_id: int, text: str) -> str:
await self.memory.ensure_session(user_id=user_id, agent="coach") # W8
memory_ctx = await self.memory.recall(user_id=user_id, query=text)
system = self.SYSTEM_PROMPT
if memory_ctx:
system += f"\n\nКонтекст из истории:\n{memory_ctx}"
await self.memory.add_turn("user", text, user_id=user_id) # W7
response = await self._call_llm(
messages=[HumanMessage(content=text)],
system=system,
)
await self.memory.add_turn("assistant", response, user_id=user_id) # W7
await self.memory.store(user_id=user_id, text=f"coach: {text[:80]}")
return response
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/agents/test_coach_agent.py -v
Шаг 5 — Коммит:
git add src/d_brain/agents/coach_agent.py tests/agents/test_coach_agent.py
git commit -m "feat(agents): add CoachAgent subgraph with dedicated system prompt and skills scope"
Задача 3.2: HabitsAgent — трекинг привычек#
Почему так: Трекинг привычек требует строгой идемпотентности (check-in дважды = не считается дважды) и state machine для статусов — SQLite UNIQUE constraint обеспечивает это атомарно. Пример: Пользователь отмечает «пробежка ✓» дважды →
HabitsAgent.checkin("running")→ SQLite INSERT с UNIQUE(user_id, habit_id, date) → второй check-in возвращает «уже отмечено».
Цель: Агент для трекинга привычек в отдельном Telegram треде.
Файлы:
- Создать:
src/d_brain/agents/habits_agent.py - Тест:
tests/agents/test_habits_agent.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_habits_agent.py
import pytest
from unittest.mock import AsyncMock, patch
from d_brain.agents.habits_agent import HabitsAgent
@pytest.fixture
def habits(tmp_path):
return HabitsAgent(data_dir=tmp_path)
@pytest.mark.asyncio
async def test_log_habit(habits):
await habits.log_habit(user_id=1, habit="бег", value="5km", date="2026-05-27")
log = await habits.get_today_log(user_id=1, date="2026-05-27")
assert "бег" in str(log)
@pytest.mark.asyncio
async def test_habits_agent_processes_message(habits):
with patch.object(habits, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "✅ Отмечено: бег 5км"
response = await habits.process(user_id=1, thread_id=200, text="бег 5км сделан")
assert isinstance(response, str)
@pytest.mark.asyncio
async def test_habits_separate_memory(habits, tmp_path):
"""У HabitsAgent своя независимая память (не смешивается с main)."""
assert "habits" in str(habits.memory.episodic._db_path)
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/agents/test_habits_agent.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/agents/habits_agent.py
"""HabitsAgent — трекинг привычек в отдельном Telegram треде.
Хранит лог привычек в SQLite, отвечает на сообщения в треде привычек.
Наследует _call_llm от BaseLLMAgent.
"""
import asyncio
import json
import aiosqlite
from pathlib import Path
from langchain_core.messages import HumanMessage
from d_brain.agents.base import BaseLLMAgent # W6
from d_brain.core.time_utils import utc_now, utc_today, format_iso # W5
class HabitsAgent(BaseLLMAgent):
SKILLS_SCOPE = "habits"
SYSTEM_PROMPT = """Ты трекер привычек. Помогаешь пользователю отслеживать и формировать привычки.
При сообщении о выполненной привычке — подтверди и отметь в логе.
Давай мотивацию, streak статистику, напоминания. Краткие ответы."""
def __init__(self, data_dir: Path | str = "~/.d_brain"):
d = Path(data_dir)
super().__init__(data_dir=d / "habits") # W6
self._habits_db = d / "habits" / "habits_log.db"
self._habits_db.parent.mkdir(parents=True, exist_ok=True)
self._initialized = False
self._init_lock = asyncio.Lock() # W4
async def _ensure_init(self):
# W4: asyncio.Lock + двойная проверка
if self._initialized:
return
async with self._init_lock:
if self._initialized:
return
async with aiosqlite.connect(self._habits_db) as db:
await db.execute("PRAGMA journal_mode=WAL") # W3
await db.execute("""
CREATE TABLE IF NOT EXISTS habit_log (
id INTEGER PRIMARY KEY AUTOINCREMENT,
user_id INTEGER,
habit TEXT,
value TEXT,
date TEXT, -- UTC date "YYYY-MM-DD"
ts TEXT -- UTC ISO 8601
)
""")
await db.execute(
"CREATE INDEX IF NOT EXISTS idx_habit_user_date ON habit_log(user_id, date)"
)
await db.commit()
self._initialized = True
async def log_habit(self, user_id: int, habit: str, value: str, date: str) -> None:
await self._ensure_init()
async with aiosqlite.connect(self._habits_db) as db:
await db.execute(
"INSERT INTO habit_log (user_id, habit, value, date, ts) VALUES (?,?,?,?,?)",
(user_id, habit, value, date, format_iso(utc_now())), # W5: UTC timestamp
)
await db.commit()
async def get_today_log(self, user_id: int, date: str | None = None) -> list[dict]:
await self._ensure_init()
d = date or utc_today() # W5: UTC дата вместо date.today()
async with aiosqlite.connect(self._habits_db) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"SELECT habit, value, ts FROM habit_log WHERE user_id=? AND date=?",
(user_id, d),
) as cur:
return [dict(r) for r in await cur.fetchall()]
# _call_llm НЕ нужен — наследуется от BaseLLMAgent (W6)
async def process(self, user_id: int, thread_id: int, text: str) -> str:
today_log = await self.get_today_log(user_id=user_id)
context = f"Сегодня выполнено: {json.dumps(today_log, ensure_ascii=False)}"
await self.memory.ensure_session(user_id=user_id, agent="habits") # W8
response = await self._call_llm(
messages=[HumanMessage(content=text)],
system=self.SYSTEM_PROMPT + "\n\n" + context,
)
await self.memory.add_turn("user", text, user_id=user_id) # W7
await self.memory.add_turn("assistant", response, user_id=user_id) # W7
return response
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/agents/test_habits_agent.py -v
Шаг 5 — Коммит:
git add src/d_brain/agents/habits_agent.py tests/agents/test_habits_agent.py
git commit -m "feat(agents): add HabitsAgent for habit tracking in dedicated Telegram thread"
Задача 3.3: Документация примеров агентов#
Почему так: Примеры использования агентов — живая документация; без шаблона новый разработчик изобретёт несовместимый велосипед вместо расширения существующей архитектуры. Пример: Новый разработчик смотрит
examples/custom_agent.py→ видит правильный шаблон с BaseLLMAgent + Soul.md + handlers → создаётWeatherAgentза 30 минут вместо 3 часов.
Цель: Создать docs/ПримерыАгента/ с примерами использования каждого агента.
Файлы:
- Создать:
docs/ПримерыАгента/01-main-agent.md - Создать:
docs/ПримерыАгента/02-coach-agent.md - Создать:
docs/ПримерыАгента/03-habits-agent.md
mkdir -p /home/serg/projects/second-brain/docs/ПримерыАгента
<!-- docs/ПримерыАгента/01-main-agent.md -->
# MainAgent — примеры взаимодействия
## Базовый вопрос
Пользователь: Что у меня на завтра?
Агент: [вызывает planning_tools.get_schedule()] → "Завтра: стендап в 10:00, код-ревью в 15:00"
## Поиск в wiki
Пользователь: /wiki Проект Феникс
Агент: [вызывает wiki_tools.link_preview()] → форматированное Telegram сообщение
## Веб-поиск
Пользователь: найди последние новости о LangGraph
Агент: [вызывает search_tools.web_search()] → топ-5 результатов
git add docs/ПримерыАгента/
git commit -m "docs: add agent examples in Russian (ПримерыАгента)"
Интеграционный smoke-тест Фазы 3#
uv run pytest tests/agents/ -v --tb=short
ФАЗА 4: PLANNER AGENT + YOUGILE DEEP INTEGRATION#
Почему так: YouGile — основной инструмент управления задачами, поэтому глубокая интеграция критична для автоматизации «закрытия дня». Без пагинации (W11) агент молча теряет 70% проектов. Без UTC-фиксации (W12) «выполненные сегодня» задачи включают вчерашние. Алиасная система нужна потому, что пользователь хочет писать «ф» вместо полного UUID проекта — transparent interface принцип.
Пример: Вечером
/закрой день. PlannerAgent: 1) опрашивает все 12 проектов YouGile через пагинированный API, 2) фильтрует Done-задачи сcompletedAtв UTC-сегодня, 3) собирает wiki-страницы изменённые сегодня по mtime, 4) отдаёт LLM для финализации. Результат: «Сегодня: закрыл 4 задачи по проекту Феникс, обновил 2 страницы wiki».
Цель: PlannerAgent читает Done задачи из всех проектов YouGile, управляет тегами/алиасами проектов. Оценка: ~1.5 рабочих дня
Задача 4.1: Planning Tools — YouGile Done tasks#
Почему так: Done tasks из YouGile — это готовый «факт дня»; интеграция с планировщиком замыкает цикл план→факт: что планировалось и что реально сделано видно в одном месте. Пример: Пользователь просит дневной отчёт → PlannerAgent запрашивает YouGile Done за сегодня → CrossReference с daily plan → итог: «7/10 задач выполнено, перенесено 3».
Цель: Читать выполненные задачи из Done-колонок всех проектов YouGile.
Файлы:
- Создать:
src/d_brain/tools/planning_tools.py - Тест:
tests/tools/test_planning_tools.py
Шаг 1 — Пишем падающий тест:
# tests/tools/test_planning_tools.py
import pytest
from unittest.mock import AsyncMock, patch, MagicMock
from d_brain.tools.planning_tools import PlanningTools
@pytest.fixture
def planner():
return PlanningTools(yougile_token="test_token")
@pytest.mark.asyncio
async def test_get_done_tasks_today(planner):
mock_tasks = [
{"id": "1", "title": "Написать тесты", "projectId": "proj-1"},
{"id": "2", "title": "Провести ревью", "projectId": "proj-2"},
]
with patch.object(planner, "_fetch_done_tasks", new_callable=AsyncMock) as mock_fetch:
mock_fetch.return_value = mock_tasks
tasks = await planner.get_done_tasks_today()
assert len(tasks) == 2
assert tasks[0]["title"] == "Написать тесты"
@pytest.mark.asyncio
async def test_project_alias_resolution(planner):
"""Прозрачный интерфейс: имя, алиас или тег — всё работает."""
planner.set_alias("ф", "proj-феникс")
planner.set_alias("феникс", "proj-феникс")
assert planner.resolve_project("ф") == "proj-феникс"
assert planner.resolve_project("феникс") == "proj-феникс"
assert planner.resolve_project("proj-феникс") == "proj-феникс" # прямой ID
@pytest.mark.asyncio
async def test_task_history_not_in_projects(planner):
"""История задач (что сделано/встречи) НЕ записывается в проекты YouGile."""
assert hasattr(planner, "get_done_tasks_today")
assert not hasattr(planner, "write_task_history_to_project")
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/tools/test_planning_tools.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/tools/planning_tools.py
"""Planning Tools — YouGile + календарь.
Deep Module: get_done_tasks_today(), resolve_project().
История задач (что сделано) НЕ пишется в YouGile (требование #19).
Она идёт в episodic memory / wiki.
"""
import httpx
from datetime import date, datetime, timezone
from d_brain.core.guardrails import check_ssrf_url
YOUGILE_API = "https://ru.yougile.com/api-v2"
class PlanningTools:
def __init__(self, yougile_token: str):
self._token = yougile_token
self._aliases: dict[str, str] = {} # alias → project_id
def set_alias(self, alias: str, project_id: str) -> None:
self._aliases[alias.lower()] = project_id
def resolve_project(self, name_or_alias: str) -> str:
"""Прозрачный интерфейс: имя/алиас/тег → project_id."""
return self._aliases.get(name_or_alias.lower(), name_or_alias)
def _headers(self) -> dict:
return {
"Authorization": f"Bearer {self._token}",
"Content-Type": "application/json",
}
# <!-- ARCH REVISED: W11+W12 — добавлена пагинация + TZ fix + защита от infinite loop -->
async def _fetch_all_pages(self, client: httpx.AsyncClient, url: str, params: dict = None) -> list[dict]:
"""Пагинатор YouGile API. Защита от infinite loop через max_pages (code-review bug #11).
YouGile API: ?page=0&pageSize=50, возвращает {content: [...], pageCount: N}
"""
MAX_PAGES = 50 # защита от infinite loop (code-review bug #11)
all_items = []
page = 0
page_size = 50
while page < MAX_PAGES:
paginated_params = {**(params or {}), "page": page, "pageSize": page_size}
resp = await client.get(url, headers=self._headers(), params=paginated_params)
resp.raise_for_status()
data = resp.json()
content = data.get("content", [])
all_items.extend(content)
total_pages = data.get("pageCount", 1)
page += 1
if page >= total_pages or not content:
break
return all_items
async def _fetch_done_tasks(self) -> list[dict]:
"""Внутренний: получает все задачи из Done колонок всех проектов с пагинацией."""
# W2: используем sync check (YOUGILE_API — известный URL, не user-input)
check_ssrf_url(YOUGILE_API)
async with httpx.AsyncClient(timeout=30) as client:
# 1. Получаем все проекты (с пагинацией) — W11
projects = await self._fetch_all_pages(client, f"{YOUGILE_API}/projects")
done_tasks = []
for project in projects:
proj_id = project["id"]
# 2. Получаем доски проекта
# NOTE: YouGile v2 API endpoint для досок проекта
boards_resp = await client.get(
f"{YOUGILE_API}/string-boards",
headers=self._headers(),
params={"projectId": proj_id},
)
if boards_resp.status_code != 200:
continue
boards = boards_resp.json().get("content", [])
for board in boards:
columns = board.get("stickers", {})
for col_id, col_data in columns.items():
col_title = col_data.get("title", "").lower()
if any(kw in col_title for kw in ["done", "готово", "выполнено", "завершено"]):
# 3. Берём задачи из Done колонки (с пагинацией)
try:
tasks = await self._fetch_all_pages(
client,
f"{YOUGILE_API}/tasks",
params={"columnId": col_id},
)
for t in tasks:
t["_projectId"] = proj_id
t["_projectTitle"] = project.get("title", "")
t["_columnTitle"] = col_data.get("title", "")
done_tasks.extend(tasks)
except httpx.HTTPError:
continue # продолжаем если одна колонка недоступна
return done_tasks
async def get_moved_tasks(self) -> list[dict]:
"""Задачи перенесённые (не Done но изменили колонку сегодня).
ARCH NOTE (W29): stub — реализовать через YouGile activity log API
если доступен, иначе через сравнение snapshot'ов.
"""
# TODO: реализовать через /tasks?updatedAfter=<today_start_utc>
return []
async def get_done_tasks_today(self) -> list[dict]:
"""Задачи, перенесённые в Done сегодня (по всем проектам).
W12: используем UTC дату для сравнения (не локальную date.today())
"""
from d_brain.core.time_utils import utc_today
all_done = await self._fetch_done_tasks()
today_utc = utc_today() # W12: UTC дата, не локальная
result = []
for task in all_done:
# YouGile хранит timestamp в completedAt или updatedAt (unix millis или ISO)
ts_fields = ["completedAt", "updatedAt", "createdAt"]
for field_name in ts_fields:
ts = task.get(field_name)
if ts:
try:
if isinstance(ts, int):
dt = datetime.fromtimestamp(ts / 1000, tz=timezone.utc)
else:
from d_brain.core.time_utils import parse_iso
dt = parse_iso(ts) # W5: parse_iso нормализует TZ
# W12: сравниваем UTC даты (не локальные)
if dt.astimezone(timezone.utc).date().isoformat() == today_utc:
result.append(task)
break
except (ValueError, OSError):
continue
return result
async def get_schedule(self, date_str: str | None = None) -> list[dict]:
"""Возвращает расписание на дату (Google Calendar integration)."""
# Делегируем в существующий GoogleCalendar integration
return []
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/tools/test_planning_tools.py -v
Шаг 5 — Коммит:
git add src/d_brain/tools/planning_tools.py tests/tools/test_planning_tools.py
git commit -m "feat(tools): add PlanningTools with YouGile Done tasks from all projects"
Задача 4.2: PlannerAgent — субграф планирования#
Почему так: PlannerAgent как субграф инкапсулирует логику YouGile + календарь + daily plan — без изоляции логика планирования протекает в MainAgent и становится нетестируемой. Пример: Пользователь пишет «составь план на сегодня» → PlannerAgent → YouGile pending + Google Calendar events + активные привычки → синтез → markdown plan → отправка.
Цель: Агент-планировщик с доступом к YouGile и календарю.
Файлы:
- Создать:
src/d_brain/agents/planner_agent.py - Тест:
tests/agents/test_planner_agent.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_planner_agent.py
import pytest
from unittest.mock import AsyncMock, patch
from d_brain.agents.planner_agent import PlannerAgent
@pytest.fixture
def planner(tmp_path):
return PlannerAgent(data_dir=tmp_path, yougile_token="test")
@pytest.mark.asyncio
async def test_daily_summary_includes_done_tasks(planner):
with patch.object(planner.planning, "get_done_tasks_today", new_callable=AsyncMock) as mock_done:
mock_done.return_value = [
{"title": "Написать тесты", "_projectTitle": "Second Brain"}
]
with patch.object(planner, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "Сегодня выполнено: Написать тесты"
result = await planner.daily_summary(user_id=1)
assert isinstance(result, str)
@pytest.mark.asyncio
async def test_close_day_flow(planner):
"""Close Day: done tasks + moved tasks + wiki updated."""
with patch.object(planner.planning, "get_done_tasks_today", new_callable=AsyncMock) as mock_done:
mock_done.return_value = [{"title": "Task 1", "_projectTitle": "P1"}]
with patch.object(planner, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "День закрыт. Выполнено: 1 задача."
result = await planner.close_day(user_id=1)
assert "закрыт" in result or isinstance(result, str)
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/agents/test_planner_agent.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/agents/planner_agent.py
"""PlannerAgent — планирование, YouGile, Close Day v2.
Close Day v2 финализирует:
1. Done tasks из YouGile (все проекты, с пагинацией) — W11
2. Moved tasks (перенесённые)
3. Wiki pages updated сегодня (по mtime) — W10
4. Сохраняет итог в episodic memory (НЕ в YouGile — требование #19)
"""
import os
from pathlib import Path
from langchain_core.messages import HumanMessage
from d_brain.agents.base import BaseLLMAgent # W6
from d_brain.tools.planning_tools import PlanningTools
from d_brain.tools.wiki_tools import WikiTools
class PlannerAgent(BaseLLMAgent):
SKILLS_SCOPE = "planner"
SYSTEM_PROMPT = """Ты агент-планировщик. Управляешь расписанием, задачами, проектами.
Используешь YouGile для задач и Google Calendar для встреч.
При закрытии дня — суммируй выполненное, перенесённое и обновлённые wiki страницы."""
def __init__(
self,
data_dir: Path | str = "~/.d_brain",
yougile_token: str = "",
vault_path: Path | str = "~/vault",
):
d = Path(data_dir)
super().__init__(data_dir=d / "planner") # W6: BaseLLMAgent
self.planning = PlanningTools(yougile_token=yougile_token or os.getenv("YOUGILE_TOKEN", ""))
self.wiki = WikiTools(vault_path=vault_path)
# _call_llm и self.memory наследуются от BaseLLMAgent
async def daily_summary(self, user_id: int) -> str:
"""Сводка дня: расписание + done задачи."""
done = await self.planning.get_done_tasks_today()
done_text = "\n".join(
f"- [{t.get('_projectTitle', '?')}] {t.get('title', '?')}"
for t in done
) or "нет выполненных задач"
prompt = f"Сегодня выполнено:\n{done_text}\n\nСоставь краткую сводку дня."
return await self._call_llm([HumanMessage(content=prompt)])
async def close_day(self, user_id: int) -> str:
"""Close Day v2:
1. Done tasks (YouGile Done колонки)
2. Moved/перенесённые задачи
3. Wiki страницы, обновлённые сегодня
4. Сохраняет итог в episodic memory (не в YouGile)
"""
# <!-- ARCH REVISED: W10 — wiki.search(today) заменено на search_updated_since(mtime) -->
from d_brain.core.time_utils import utc_today
import time
from datetime import datetime, timezone
today = utc_today() # W5: UTC дата
# 1. Done tasks
done = await self.planning.get_done_tasks_today()
done_list = [f"✅ [{t.get('_projectTitle','?')}] {t.get('title','?')}" for t in done]
# 2. Wiki страницы (обновлённые сегодня — по mtime, не по тексту)
# W10: ищем по modification time, не по строке даты в тексте
today_start_ts = datetime(
*[int(x) for x in today.split("-")], tzinfo=timezone.utc
).timestamp()
updated_wiki = await self.wiki.search_updated_since(today_start_ts)
wiki_list = [f"📝 {r['path']}" for r in updated_wiki[:5]]
summary_text = (
f"📅 Закрытие дня {today}\n\n"
+ ("Выполнено:\n" + "\n".join(done_list) if done_list else "Выполненных задач нет")
+ "\n\n"
+ ("Обновлено в wiki:\n" + "\n".join(wiki_list) if wiki_list else "Wiki не обновлялась")
)
# LLM финализирует итог
prompt = f"Финализируй закрытие дня:\n{summary_text}\nДобавь рефлексию и выводы."
final = await self._call_llm([HumanMessage(content=prompt)])
# Сохраняем в episodic memory (НЕ в YouGile)
sid = await self.memory.begin_session(user_id=user_id, agent="planner")
await self.memory.add_turn("system", summary_text)
await self.memory.episodic.save_summary(sid, final[:500])
await self.memory.store(user_id=user_id, text=f"close_day_{today}: {final[:200]}")
return final
async def process(self, user_id: int, thread_id: int, text: str) -> str:
await self.memory.ensure_session(user_id=user_id, agent="planner") # W8
memory_ctx = await self.memory.recall(user_id=user_id, query=text)
system = self.SYSTEM_PROMPT + (f"\n\nКонтекст:\n{memory_ctx}" if memory_ctx else "")
response = await self._call_llm([HumanMessage(content=text)], system=system)
await self.memory.add_turn("user", text, user_id=user_id) # W7
await self.memory.add_turn("assistant", response, user_id=user_id) # W7
return response
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/agents/test_planner_agent.py -v
Шаг 5 — Коммит:
git add src/d_brain/agents/planner_agent.py tests/agents/test_planner_agent.py
git commit -m "feat(agents): add PlannerAgent with close_day v2 and YouGile Done tasks integration"
Интеграционный smoke-тест Фазы 4#
uv run pytest tests/agents/test_planner_agent.py tests/tools/test_planning_tools.py -v --tb=short
ФАЗА 5: TELEGRAM GROUP/THREADS + WIKI LINK PREVIEW#
Почему так: Telegram Group Topics — единственный способ иметь несколько «контекстов» в одном чате без создания отдельных ботов. Маршрутизация по
thread_id— нулевая стоимость: это просто integer. Wiki link preview через frontmatter даёт структурированную информацию (title, status, tags) вместо сырого Markdown — и HTML parse_mode избегает краша на спецсимволах (W23).Пример: Пользователь в треде
#planningпишет/wiki Проекты/Феникс.md. Бот отвечает:<b>Проект Феникс</b>\nСтатус: <code>active</code>\nТеги: #python, #ai\n<i>Разработка AI-ассистента</i>— форматированное превью без сырого YAML.
Цель: Полная маршрутизация по Telegram топикам, команда /wiki с форматированным превью. Оценка: ~1 рабочий день
Задача 5.1: Telegram Group Router — маршрутизация по топикам#
Почему так: Group Topics в Telegram — практически отдельные чаты внутри одной группы; маршрутизация по
thread_idпозволяет иметь 4 специализированных агента в одной группе без путаницы. Пример: Пользователь пишет в топик#habits→ Router определяетthread_id=THREAD_HABITS_ID→ направляет в HabitsAgent → пользователь в#mainникогда не видит ответы HabitsAgent.
Цель: aiogram handler определяет thread_id и направляет в правильный агент.
Файлы:
- Изменить:
src/d_brain/bot/handlers/message_handler.py - Тест:
tests/bot/test_thread_routing.py
Шаг 1 — Пишем падающий тест:
# tests/bot/test_thread_routing.py
import pytest
from unittest.mock import AsyncMock, MagicMock, patch
from d_brain.bot.dispatcher import route_message
@pytest.mark.asyncio
async def test_thread_0_routes_to_main():
agents = {
"main": AsyncMock(return_value="main response"),
"coach": AsyncMock(return_value="coach response"),
}
result = await route_message(
agents=agents,
user_id=1,
thread_id=0,
text="Привет",
routing={0: "main", 100: "coach"},
)
assert result == "main response"
agents["main"].assert_called_once()
agents["coach"].assert_not_called()
@pytest.mark.asyncio
async def test_thread_100_routes_to_coach():
agents = {
"main": AsyncMock(return_value="main"),
"coach": AsyncMock(return_value="coach response"),
}
result = await route_message(
agents=agents,
user_id=1,
thread_id=100,
text="Как мои привычки?",
routing={0: "main", 100: "coach"},
)
assert result == "coach response"
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/bot/test_thread_routing.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/bot/dispatcher.py
"""Диспетчер сообщений — маршрутизация по thread_id к агентам."""
from typing import Protocol
class Agent(Protocol):
async def process(self, user_id: int, thread_id: int, text: str) -> str: ...
async def route_message(
agents: dict,
user_id: int,
thread_id: int,
text: str,
routing: dict[int, str],
) -> str:
agent_name = routing.get(thread_id, "main")
agent = agents.get(agent_name, agents["main"])
return await agent(user_id=user_id, thread_id=thread_id, text=text)
Полный handler для aiogram:
# src/d_brain/bot/handlers/message_handler.py
"""aiogram message handler — тонкая обёртка над агентами."""
import os
from aiogram import Router
from aiogram.types import Message
from d_brain.core.router import ThreadRouter
from d_brain.agents.main_agent import MainAgent
from d_brain.agents.coach_agent import CoachAgent
from d_brain.agents.habits_agent import HabitsAgent
from d_brain.agents.planner_agent import PlannerAgent
router = Router()
_thread_router = ThreadRouter.from_env()
_agents = {
"main": MainAgent(
data_dir=os.getenv("DATA_DIR", "~/.d_brain"),
vault_path=os.getenv("VAULT_PATH", "~/vault"),
),
"coach": CoachAgent(data_dir=os.getenv("DATA_DIR", "~/.d_brain")),
"habits": HabitsAgent(data_dir=os.getenv("DATA_DIR", "~/.d_brain")),
"planner": PlannerAgent(data_dir=os.getenv("DATA_DIR", "~/.d_brain")),
}
@router.message()
async def handle_message(message: Message):
# <!-- ARCH REVISED: W22+W27+W23 — добавлены: error boundary, пустой текст guard, parse_mode HTML -->
thread_id = message.message_thread_id or 0
user_id = message.from_user.id if message.from_user else 0
text = message.text or message.caption or ""
# W27: проверяем ANTHROPIC_API_KEY при первом сообщении
if not os.getenv("ANTHROPIC_API_KEY"):
await message.reply(
"⚠️ ANTHROPIC_API_KEY не настроен. Добавьте в .env файл.",
parse_mode="HTML",
)
return
if not text.strip():
# Не обрабатываем пустые сообщения (фото без подписи, стикеры etc.)
return
route = _thread_router.route(thread_id)
agent = _agents.get(route.agent, _agents["main"])
try:
response = await agent.process(
user_id=user_id,
thread_id=thread_id,
text=text,
)
# W23: parse_mode HTML (не Markdown — ломается на спецсимволах)
await message.reply(str(response), parse_mode="HTML")
except Exception as e:
import logging
logging.getLogger(__name__).exception(
f"Agent error for user={user_id} thread={thread_id}: {e}"
)
await message.reply(
"⚠️ Произошла ошибка при обработке запроса. Попробуйте ещё раз.",
parse_mode="HTML",
)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/bot/test_thread_routing.py -v
Шаг 5 — Коммит:
git add src/d_brain/bot/ tests/bot/
git commit -m "feat(bot): add full thread routing dispatcher for Telegram group topics"
Задача 5.2: /wiki команда с Telegram link preview#
Почему так: Telegram показывает превью URL только если сервер отдаёт Open Graph теги — mini web server за cloudflare tunnel даёт красивое превью прямо в чате без перехода в браузер. Пример: Пользователь пишет
/wiki Проект Феникс→ бот отвечает ссылкой → Telegram показывает превью с заголовком + описанием → пользователь видит контент не открывая браузер.
Цель: Команда /wiki <alias> возвращает форматированное превью из frontmatter заметки.
Файлы:
- Создать:
src/d_brain/bot/handlers/wiki_command.py - Тест:
tests/bot/test_wiki_command.py
Шаг 1 — Пишем падающий тест:
# tests/bot/test_wiki_command.py
import pytest
from unittest.mock import AsyncMock, patch
from d_brain.bot.handlers.wiki_command import handle_wiki_command
@pytest.mark.asyncio
async def test_wiki_command_returns_preview(tmp_path):
vault = tmp_path / "vault"
vault.mkdir()
note = vault / "Проекты" / "Феникс.md"
note.parent.mkdir()
note.write_text("---\ntitle: Проект Феникс\nstatus: active\ntags: [python, ai]\n---\n# Body")
from d_brain.tools.wiki_tools import WikiTools
wiki = WikiTools(vault_path=vault)
result = await handle_wiki_command(wiki=wiki, alias="Проекты/Феникс.md")
assert "Феникс" in result
assert "active" in result
@pytest.mark.asyncio
async def test_wiki_command_not_found(tmp_path):
from d_brain.tools.wiki_tools import WikiTools
wiki = WikiTools(vault_path=tmp_path / "vault")
result = await handle_wiki_command(wiki=wiki, alias="НесуществующаяЗаметка")
assert "не найден" in result.lower() or "нет" in result.lower() or result != ""
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/bot/test_wiki_command.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/bot/handlers/wiki_command.py
"""Обработчик команды /wiki — возвращает Telegram link preview из frontmatter."""
from aiogram import Router
from aiogram.filters import Command
from aiogram.types import Message
from d_brain.tools.wiki_tools import WikiTools
import os
router = Router()
_wiki: WikiTools | None = None
def get_wiki() -> WikiTools:
global _wiki
if _wiki is None:
_wiki = WikiTools(vault_path=os.getenv("VAULT_PATH", "~/vault"))
return _wiki
async def handle_wiki_command(wiki: WikiTools, alias: str) -> str:
"""Бизнес-логика /wiki — тестируется отдельно от aiogram."""
if not alias:
return "Использование: /wiki <путь или алиас>"
# Пробуем точный путь
preview = await wiki.link_preview(alias)
if "не найдена" in preview:
# Пробуем поиск по названию
results = await wiki.search(alias)
if results:
preview = await wiki.link_preview(results[0]["path"])
else:
return f"Заметка '{alias}' не найдена в wiki."
return preview
@router.message(Command("wiki"))
async def wiki_command_handler(message: Message):
args = message.text.split(maxsplit=1)
alias = args[1].strip() if len(args) > 1 else ""
result = await handle_wiki_command(wiki=get_wiki(), alias=alias)
# <!-- ARCH REVISED: W23 — parse_mode HTML (не Markdown). link_preview() теперь возвращает HTML -->
await message.reply(result, parse_mode="HTML")
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/bot/test_wiki_command.py -v
Шаг 5 — Коммит:
git add src/d_brain/bot/handlers/wiki_command.py tests/bot/test_wiki_command.py
git commit -m "feat(bot): add /wiki command with Telegram-formatted link preview from frontmatter"
Интеграционный smoke-тест Фазы 5#
uv run pytest tests/bot/ -v --tb=short
ФАЗА 6: SKILLS SYSTEM#
Почему так: Skills System позволяет расширять возможности агентов без деплоя нового кода — установить скилл из чата, и агент сразу умеет новое. Per-agent scope предотвращает коллизии: финансовый скилл не появляется у HabitsAgent. SkillExecutor (W16) — критичное звено: без него скиллы хранятся в файлах, но никогда не вызываются LangGraph графом. importlib-загрузка даёт hot-reload без перезапуска бота.
Пример:
/skill install org/my-skills:crm_lookup— через минуту MainAgent умеет искать контакты в CRM через новый toolcrm_lookup, не требуя кода в основном репозитории и перезапуска сервиса.
Цель: Динамическая загрузка/выгрузка скиллов из .agents/skills/ и GitHub. Каждый агент имеет свой scope. Оценка: ~1 рабочий день
Задача 6.1: Skills Registry — реестр и загрузчик скиллов#
Почему так: Реестр скиллов с версиями и security whitelist — защита от supply chain атак; без
ALLOWED_SKILL_REPOSзлоумышленник может установить вредоносный скилл из произвольного GitHub repo. Пример: Пользователь пишет «установи скилл из github.com/random/repo» →ALLOWED_SKILL_REPOSне содержит этот repo →SkillRegistryотказывает → система безопасна.
Цель: Centralized реестр с метаданными (name, path, agent_scope, enabled).
Файлы:
- Создать:
src/d_brain/skills/registry.py - Создать:
src/d_brain/skills/loader.py - Тест:
tests/skills/test_registry.py
Шаг 1 — Пишем падающий тест:
# tests/skills/test_registry.py
import pytest
from pathlib import Path
from d_brain.skills.registry import SkillRegistry, SkillMeta
from d_brain.skills.loader import SkillLoader
@pytest.fixture
def registry():
return SkillRegistry()
def test_register_skill(registry):
registry.register(SkillMeta(
name="web_search",
path="/skills/search.py",
agent_scope=["main", "coach"],
enabled=True,
))
assert registry.get("web_search") is not None
def test_get_skills_for_agent(registry):
registry.register(SkillMeta("web_search", "/s/search.py", ["main"], True))
registry.register(SkillMeta("habit_tracker", "/s/habits.py", ["coach"], True))
registry.register(SkillMeta("planning", "/s/plan.py", ["main", "planner"], True))
main_skills = registry.get_for_agent("main")
assert len(main_skills) == 2 # web_search + planning
coach_skills = registry.get_for_agent("coach")
assert len(coach_skills) == 1
def test_disable_skill(registry):
registry.register(SkillMeta("test_skill", "/s/t.py", ["main"], True))
registry.disable("test_skill")
skills = registry.get_for_agent("main")
assert all(s.name != "test_skill" for s in skills)
@pytest.mark.asyncio
async def test_loader_installs_from_file(tmp_path):
"""SkillLoader устанавливает скилл из локального файла."""
skill_file = tmp_path / "my_skill.py"
skill_file.write_text('''
# skill: name=my_skill, scope=main
def run(query: str) -> str:
return f"result: {query}"
''')
loader = SkillLoader(skills_dir=tmp_path / "installed")
meta = await loader.install_from_file(skill_file)
assert meta.name == "my_skill"
assert "main" in meta.agent_scope
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/skills/test_registry.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/skills/registry.py
"""Skills Registry — реестр скиллов с per-agent scope."""
from dataclasses import dataclass, field
@dataclass
class SkillMeta:
name: str
path: str
agent_scope: list[str] # ["main"] | ["coach"] | ["main", "planner"]
enabled: bool = True
description: str = ""
version: str = "0.1.0"
class SkillRegistry:
def __init__(self):
self._skills: dict[str, SkillMeta] = {}
def register(self, meta: SkillMeta) -> None:
self._skills[meta.name] = meta
def get(self, name: str) -> SkillMeta | None:
return self._skills.get(name)
def get_for_agent(self, agent: str) -> list[SkillMeta]:
return [
s for s in self._skills.values()
if s.enabled and (agent in s.agent_scope or "*" in s.agent_scope)
]
def disable(self, name: str) -> None:
if name in self._skills:
self._skills[name].enabled = False
def enable(self, name: str) -> None:
if name in self._skills:
self._skills[name].enabled = True
def list_all(self) -> list[SkillMeta]:
return list(self._skills.values())
# src/d_brain/skills/loader.py
"""Skills Loader — установка скиллов из файла или GitHub.
Поддерживает:
- Локальный файл: install_from_file(path)
- GitHub: install_from_github(user/repo, skill_name)
- skills.sh совместимость: parse_skills_sh(path)
"""
import re
import shutil
import asyncio
from pathlib import Path
from d_brain.skills.registry import SkillMeta
from d_brain.core.guardrails import check_ssrf_url
class SkillLoader:
def __init__(self, skills_dir: Path | str = "~/.agents/skills"):
self._dir = Path(skills_dir).expanduser()
self._dir.mkdir(parents=True, exist_ok=True)
async def install_from_file(self, source: Path) -> SkillMeta:
"""Устанавливает скилл из .py файла. Читает метаданные из комментария # skill:."""
source = Path(source)
text = source.read_text(encoding="utf-8")
meta = self._parse_metadata(text, source)
dest = self._dir / source.name
shutil.copy2(source, dest)
meta.path = str(dest)
return meta
async def install_from_github(self, repo: str, skill_name: str) -> SkillMeta:
"""Устанавливает скилл с GitHub (raw content).
<!-- ARCH REVISED: W19 — добавлены: async SSRF check + allowlist + atomic write -->
SECURITY WARNING: загружает и сохраняет Python код с GitHub.
Используй только из доверенных репозиториев.
Допустимые форматы repo: "username/repo" или "org/repo"
"""
# W19: allowlist GitHub repos (настраивается через env ALLOWED_SKILL_REPOS)
from d_brain.core.guardrails import check_ssrf_url_async
allowed_repos = os.getenv("ALLOWED_SKILL_REPOS", "").split(",")
if allowed_repos != [""] and repo not in allowed_repos:
raise ValueError(
f"Репозиторий '{repo}' не в ALLOWED_SKILL_REPOS. "
f"Добавьте в .env: ALLOWED_SKILL_REPOS={repo},..."
)
url = f"https://raw.githubusercontent.com/{repo}/main/skills/{skill_name}.py"
await check_ssrf_url_async(url) # W2: async SSRF check
import httpx
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.get(url)
resp.raise_for_status()
code = resp.text
dest = self._dir / f"{skill_name}.py"
# W18: atomic write вместо dest.write_text()
from d_brain.core.safe_io import atomic_write_text
atomic_write_text(dest, code, encoding="utf-8")
meta = self._parse_metadata(code, dest)
return meta
async def uninstall(self, name: str) -> bool:
"""Удаляет скилл по имени."""
for f in self._dir.glob("*.py"):
text = f.read_text(encoding="utf-8")
if f"name={name}" in text or f.stem == name:
f.unlink()
return True
return False
def _parse_metadata(self, code: str, path: Path) -> SkillMeta:
"""Парсит # skill: name=x, scope=a,b из первой строки файла."""
match = re.search(r"#\s*skill:\s*(.+)", code)
if not match:
return SkillMeta(name=path.stem, path=str(path), agent_scope=["main"])
raw = match.group(1)
params = {}
for part in raw.split(","):
part = part.strip()
if "=" in part:
k, v = part.split("=", 1)
params[k.strip()] = v.strip()
name = params.get("name", path.stem)
scope_raw = params.get("scope", "main")
scope = [s.strip() for s in scope_raw.split("|")]
desc = params.get("desc", "")
return SkillMeta(name=name, path=str(path), agent_scope=scope, description=desc)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/skills/test_registry.py -v
Шаг 5 — Коммит:
git add src/d_brain/skills/ tests/skills/
git commit -m "feat(skills): add SkillRegistry and SkillLoader with per-agent scope and GitHub install"
Задача 6.2: Skill Tools — установка скиллов через чат#
Почему так: Установка скиллов через чат (а не SSH + ручной деплой) — key UX feature; позволяет расширять возможности бота без технических знаний у пользователя. Пример: Пользователь пишет «установи скилл для работы с Notion» →
skill_tools.install("notion-skill")→ скилл появляется в/skills list→ MainAgent может использовать Notion инструменты.
Цель: Команды /skill install, /skill remove, /skill list через Telegram.
Файлы:
- Создать:
src/d_brain/tools/skill_tools.py - Тест:
tests/tools/test_skill_tools.py
Шаг 1 — Пишем падающий тест:
# tests/tools/test_skill_tools.py
import pytest
from d_brain.tools.skill_tools import SkillTools
from d_brain.skills.registry import SkillRegistry
@pytest.fixture
def skill_tools(tmp_path):
registry = SkillRegistry()
return SkillTools(registry=registry, skills_dir=tmp_path / "skills")
@pytest.mark.asyncio
async def test_list_skills_empty(skill_tools):
result = await skill_tools.list_skills(agent_scope="main")
assert isinstance(result, str)
assert "нет" in result.lower() or result == "" or "skill" in result.lower()
@pytest.mark.asyncio
async def test_install_and_list(skill_tools, tmp_path):
skill_file = tmp_path / "test_skill.py"
skill_file.write_text("# skill: name=test_skill, scope=main\ndef run(): pass\n")
await skill_tools.install(str(skill_file), agent_scope="main")
result = await skill_tools.list_skills(agent_scope="main")
assert "test_skill" in result
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/tools/test_skill_tools.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/tools/skill_tools.py
"""Skill Tools — управление скиллами через чат.
Интерфейс: install/remove/list. Используется агентами и /skill командой.
"""
from pathlib import Path
from d_brain.skills.registry import SkillRegistry, SkillMeta
from d_brain.skills.loader import SkillLoader
class SkillTools:
def __init__(
self,
registry: SkillRegistry,
skills_dir: Path | str = "~/.agents/skills",
):
self._registry = registry
self._loader = SkillLoader(skills_dir=skills_dir)
async def install(self, source: str, agent_scope: str = "main") -> str:
"""Устанавливает скилл из файла или GitHub (user/repo:skill_name)."""
if "/" in source and not source.endswith(".py"):
# GitHub формат: user/repo:skill_name
if ":" in source:
repo, skill_name = source.rsplit(":", 1)
else:
repo, skill_name = source, source.split("/")[-1]
meta = await self._loader.install_from_github(repo, skill_name)
else:
meta = await self._loader.install_from_file(Path(source))
# Перезаписываем scope если указан
if agent_scope and agent_scope not in meta.agent_scope:
meta.agent_scope.append(agent_scope)
self._registry.register(meta)
return f"✅ Скилл '{meta.name}' установлен для агентов: {', '.join(meta.agent_scope)}"
async def remove(self, name: str) -> str:
removed = await self._loader.uninstall(name)
if removed:
self._registry.disable(name)
return f"🗑️ Скилл '{name}' удалён"
return f"Скилл '{name}' не найден"
async def list_skills(self, agent_scope: str = "*") -> str:
if agent_scope == "*":
skills = self._registry.list_all()
else:
skills = self._registry.get_for_agent(agent_scope)
if not skills:
return f"Нет установленных скиллов для '{agent_scope}'"
lines = [
f"{'✅' if s.enabled else '❌'} **{s.name}** [{', '.join(s.agent_scope)}] — {s.description}"
for s in skills
]
return "Установленные скиллы:\n" + "\n".join(lines)
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/tools/test_skill_tools.py -v
Шаг 5 — Коммит:
git add src/d_brain/tools/skill_tools.py tests/tools/test_skill_tools.py
git commit -m "feat(tools): add SkillTools for install/remove/list via chat commands"
Задача 6.3: SkillExecutor — выполнение скиллов как LangChain tools#
Почему так: SkillExecutor как LangChain tool — bridge между skill системой и LangGraph; без него агент не может вызывать скиллы через стандартный ToolNode механизм. Пример: MainAgent решает вызвать скилл
notion_sync→ SkillExecutor.execute("notion_sync", args) → скилл выполняется → результат возвращается в StateGraph как standard tool result.
Цель: Загружать установленные скиллы из реестра и передавать агентам как @tool decorated функции. Без SkillExecutor Skills system не работает: скиллы хранятся в файлах, но агенты их не используют.
Файлы:
- Создать:
src/d_brain/skills/executor.py - Тест:
tests/skills/test_executor.py
Шаг 1 — Пишем падающий тест:
# tests/skills/test_executor.py
import pytest
from pathlib import Path
from d_brain.skills.executor import SkillExecutor
from d_brain.skills.registry import SkillRegistry, SkillMeta
@pytest.fixture
def skill_file(tmp_path):
"""Создаём тестовый скилл файл."""
f = tmp_path / "my_skill.py"
f.write_text('''
# skill: name=my_skill, scope=main, desc=Test skill
async def run(query: str) -> str:
"""Тестовый скилл: возвращает query в uppercase."""
return query.upper()
''')
return f
@pytest.fixture
def executor(tmp_path, skill_file):
registry = SkillRegistry()
registry.register(SkillMeta(
name="my_skill",
path=str(skill_file),
agent_scope=["main"],
enabled=True,
description="Test skill",
))
return SkillExecutor(registry=registry)
@pytest.mark.asyncio
async def test_execute_skill(executor):
"""SkillExecutor вызывает функцию run() из файла скилла."""
result = await executor.execute("my_skill", query="hello world")
assert result == "HELLO WORLD"
def test_get_tools_for_agent(executor):
"""get_tools_for_agent возвращает список LangChain @tool объектов."""
tools = executor.get_tools_for_agent("main")
assert len(tools) > 0
# Каждый tool должен быть вызываемым (LangChain BaseTool)
assert hasattr(tools[0], "name") or callable(tools[0])
@pytest.mark.asyncio
async def test_skill_not_found(executor):
"""Несуществующий скилл возвращает ошибку, не падает."""
result = await executor.execute("nonexistent_skill", query="test")
assert "не найден" in result.lower() or "error" in result.lower()
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/skills/test_executor.py -v 2>&1 | head -30
Шаг 3 — Реализация:
# src/d_brain/skills/executor.py
"""SkillExecutor — загружает и выполняет скиллы как LangChain tools.
Это связующее звено между SkillRegistry (хранение метаданных)
и LangGraph агентами (которым нужны @tool decorated functions).
Схема:
SkillRegistry.get_for_agent("main")
→ [SkillMeta(name="web_search", path="/skills/web_search.py", ...)]
→ SkillExecutor.get_tools_for_agent("main")
→ [@tool async def web_search(...), ...]
→ MainAgent.build_main_graph(tools=tools)
"""
import importlib.util
import inspect
import logging
from pathlib import Path
from langchain_core.tools import tool as lc_tool
from d_brain.skills.registry import SkillRegistry, SkillMeta
logger = logging.getLogger(__name__)
class SkillExecutor:
"""Загружает Python файлы скиллов и выполняет их run() функции."""
def __init__(self, registry: SkillRegistry):
self._registry = registry
self._loaded_modules: dict[str, object] = {} # name → module
def _load_module(self, meta: SkillMeta) -> object | None:
"""Загружает Python модуль из файла скилла через importlib."""
if meta.name in self._loaded_modules:
return self._loaded_modules[meta.name]
path = Path(meta.path)
if not path.exists():
logger.warning(f"Skill file not found: {path}")
return None
try:
spec = importlib.util.spec_from_file_location(
f"d_brain_skill_{meta.name}", path
)
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
self._loaded_modules[meta.name] = module
return module
except Exception as e:
logger.error(f"Failed to load skill '{meta.name}' from {path}: {e}")
return None
async def execute(self, skill_name: str, **kwargs) -> str:
"""Выполняет скилл по имени, передаёт kwargs в функцию run().
Returns:
Строка результата или сообщение об ошибке.
"""
meta = self._registry.get(skill_name)
if not meta or not meta.enabled:
return f"Скилл '{skill_name}' не найден или отключён."
module = self._load_module(meta)
if not module:
return f"Не удалось загрузить скилл '{skill_name}'."
run_func = getattr(module, "run", None)
if not run_func:
return f"Скилл '{skill_name}' не содержит функцию run()."
try:
if inspect.iscoroutinefunction(run_func):
return str(await run_func(**kwargs))
else:
return str(run_func(**kwargs))
except Exception as e:
logger.error(f"Skill '{skill_name}' execution error: {e}")
return f"Ошибка выполнения скилла '{skill_name}': {e}"
def get_tools_for_agent(self, agent_name: str) -> list:
"""Возвращает список LangChain @tool для агента.
Каждый скилл из реестра оборачивается в @tool декоратор.
Передавать напрямую в llm.bind_tools() или build_main_graph(tools=...).
"""
skills = self._registry.get_for_agent(agent_name)
tools = []
for meta in skills:
# Замыкание нужно для корректного захвата meta в цикле
def make_tool(skill_meta: SkillMeta):
@lc_tool
async def skill_tool(
query: str,
) -> str:
f"""Выполняет скилл: {skill_meta.description or skill_meta.name}"""
return await self.execute(skill_meta.name, query=query)
# Переименовываем tool чтобы агент знал его по имени
skill_tool.name = skill_meta.name
skill_tool.description = skill_meta.description or f"Скилл: {skill_meta.name}"
return skill_tool
tools.append(make_tool(meta))
return tools
Шаг 4 — Интегрируем в агентов:
# В MainAgent.__init__():
from d_brain.skills.executor import SkillExecutor
class MainAgent(BaseLLMAgent):
def __init__(self, data_dir=..., vault_path=...):
super().__init__(data_dir=data_dir)
self.wiki = WikiTools(vault_path=vault_path)
self.search = SearchTools()
self.media = MediaTools()
self.skill_executor = SkillExecutor(registry=global_skill_registry)
# Собираем все tools: встроенные + установленные скиллы
builtin_tools = get_main_agent_tools(self.wiki, self.search, self.media)
skill_tools = self.skill_executor.get_tools_for_agent("main")
all_tools = builtin_tools + skill_tools
self._graph = build_main_graph(
wiki=self.wiki,
search=self.search,
media=self.media,
extra_tools=skill_tools, # добавляем скиллы
)
Шаг 5 — Запускаем тест (ожидаем PASS):
uv run pytest tests/skills/test_executor.py -v
Шаг 6 — Коммит:
git add src/d_brain/skills/executor.py tests/skills/test_executor.py
git commit -m "feat(skills): add SkillExecutor — load installed skills as LangChain tools (W16)"
Интеграционный smoke-тест Фазы 6#
uv run pytest tests/skills/ tests/tools/test_skill_tools.py -v --tb=short
ФАЗА 7: SELF-IMPROVEMENT LOOP + LANGFUSE + CLOSE DAY V2#
Почему так: Ассистент без self-improvement деградирует — повторяет одни и те же ошибки, не учитывает предпочтения пользователя, не замечает паттерны. SelfReviewer извлекает insights из каждой сессии и записывает в semantic memory — следующая сессия начинается умнее. Langfuse — опциональный, за флагом, потому что добавляет внешнюю зависимость и стоимость, но при включении даёт полную видимость LLM-вызовов и latency.
Пример: Нодка Kaizen-ревью обнаруживает: «пользователь 3 раза переспрашивал про формат дедлайнов в YouGile». Система записывает в процедурную память: «объяснять формат дедлайна при создании задачи». Следующий раз PlannerAgent добавляет пояснение автоматически.
Цель: Цикл самообучения, опциональный Langfuse, финальный close-day. Оценка: ~1 рабочий день
Задача 7.1: Langfuse observability (опционально)#
Почему так: Без observability невозможно понять почему агент сделал неверный выбор — Langfuse трассирует каждый LLM вызов с промптом, ответом, latency и стоимостью для последующего анализа. Пример: Агент ответил неверно → разработчик открывает Langfuse → видит trace: intent_classifier вернул "habits" вместо "planning" → исправляет промпт классификатора.
Цель: Добавить Langfuse трассировку всех LLM вызовов за флагом LANGFUSE_ENABLED.
Файлы:
- Создать:
src/d_brain/core/observability.py - Тест:
tests/core/test_observability.py
Шаг 1 — Пишем падающий тест:
# tests/core/test_observability.py
import os
import pytest
def test_observability_disabled_by_default(monkeypatch):
monkeypatch.delenv("LANGFUSE_ENABLED", raising=False)
from d_brain.core.observability import get_tracer
tracer = get_tracer()
assert tracer.enabled is False
def test_observability_noop_when_disabled(monkeypatch):
monkeypatch.delenv("LANGFUSE_ENABLED", raising=False)
from d_brain.core.observability import get_tracer
tracer = get_tracer()
# Noop tracer должен работать без ошибок
with tracer.span("test_span") as span:
span.set_attribute("key", "value")
# Нет исключений — OK
def test_observability_enabled_when_flag_set(monkeypatch):
monkeypatch.setenv("LANGFUSE_ENABLED", "true")
monkeypatch.setenv("LANGFUSE_PUBLIC_KEY", "pk-test")
monkeypatch.setenv("LANGFUSE_SECRET_KEY", "sk-test")
from importlib import reload
import d_brain.core.observability as obs
reload(obs)
tracer = obs.get_tracer()
assert tracer.enabled is True
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/core/test_observability.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/core/observability.py
"""Observability — Langfuse трассировка за флагом LANGFUSE_ENABLED.
Deep Module: get_tracer() возвращает tracer.
Если флаг выключен — NoopTracer (нулевые накладные расходы).
"""
import os
from contextlib import contextmanager
from typing import Generator
class _Span:
def set_attribute(self, key: str, value) -> None:
pass
def set_status(self, status: str) -> None:
pass
class NoopTracer:
enabled = False
@contextmanager
def span(self, name: str, **kwargs) -> Generator[_Span, None, None]:
yield _Span()
def trace_llm(self, name: str, input: str, output: str, model: str = "") -> None:
pass
class LangfuseTracer:
enabled = True
def __init__(self):
from langfuse import Langfuse
self._lf = Langfuse(
public_key=os.environ["LANGFUSE_PUBLIC_KEY"],
secret_key=os.environ["LANGFUSE_SECRET_KEY"],
host=os.getenv("LANGFUSE_HOST", "https://cloud.langfuse.com"),
)
@contextmanager
def span(self, name: str, **kwargs) -> Generator:
trace = self._lf.trace(name=name, **kwargs)
span = trace.span(name=name)
try:
yield span
finally:
span.end()
def trace_llm(self, name: str, input: str, output: str, model: str = "") -> None:
trace = self._lf.trace(name=name)
trace.generation(
name=name,
model=model,
input=input,
output=output,
)
_tracer: NoopTracer | LangfuseTracer | None = None
def get_tracer() -> NoopTracer | LangfuseTracer:
global _tracer
if _tracer is not None:
return _tracer
if os.getenv("LANGFUSE_ENABLED", "false").lower() == "true":
try:
_tracer = LangfuseTracer()
except Exception:
_tracer = NoopTracer()
else:
_tracer = NoopTracer()
return _tracer
Использование в агентах:
# В любом агенте, в методе _call_llm:
from d_brain.core.observability import get_tracer
async def _call_llm(self, messages, system=""):
tracer = get_tracer()
tracer.trace_llm(
name=f"{self.__class__.__name__}.llm",
input=str(messages),
output="", # заполним после
model="claude-3-5-sonnet-20241022",
)
# ... реальный вызов ...
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/core/test_observability.py -v
Шаг 5 — Коммит:
git add src/d_brain/core/observability.py tests/core/test_observability.py
git commit -m "feat(core): add Langfuse observability behind LANGFUSE_ENABLED flag (noop by default)"
Задача 7.2: Self-improvement review loop#
Почему так: Self-improvement loop — механизм эволюции системы без участия разработчика; агент анализирует паттерны использования и предлагает оптимизации, которые пользователь может принять или отклонить. Пример: За неделю пользователь 10 раз спросил про «время на задачу» → SelfReviewer замечает паттерн → предлагает скилл time_tracking → пользователь одобряет → скилл устанавливается.
Цель: После каждой сессии агент анализирует качество ответов и обновляет свою память.
Файлы:
- Создать:
src/d_brain/agents/reviewer.py - Тест:
tests/agents/test_reviewer.py
Шаг 1 — Пишем падающий тест:
# tests/agents/test_reviewer.py
import pytest
from unittest.mock import AsyncMock, patch
from d_brain.agents.reviewer import SelfReviewer
@pytest.fixture
def reviewer(tmp_path):
return SelfReviewer(data_dir=tmp_path)
@pytest.mark.asyncio
async def test_review_session(reviewer):
"""Reviewer анализирует сессию и возвращает insights."""
session = [
{"role": "user", "content": "Как организовать день?"},
{"role": "assistant", "content": "Попробуй технику pomodoro..."},
{"role": "user", "content": "Спасибо, полезно"},
]
with patch.object(reviewer, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "Пользователь доволен. Тема: тайм-менеджмент."
insights = await reviewer.review_session(user_id=1, session=session)
assert isinstance(insights, str)
assert len(insights) > 0
@pytest.mark.asyncio
async def test_review_stores_to_memory(reviewer):
"""Insights из review сохраняются в семантическую память."""
session = [{"role": "user", "content": "тест"}, {"role": "assistant", "content": "ок"}]
with patch.object(reviewer, "_call_llm", new_callable=AsyncMock) as mock_llm:
mock_llm.return_value = "Простой обмен."
with patch.object(reviewer.memory, "store", new_callable=AsyncMock) as mock_store:
await reviewer.review_session(user_id=1, session=session)
mock_store.assert_called_once()
Шаг 2 — Запускаем тест (ожидаем FAIL):
uv run pytest tests/agents/test_reviewer.py -v 2>&1 | head -20
Шаг 3 — Реализация:
# src/d_brain/agents/reviewer.py
"""SelfReviewer — цикл самоулучшения агента (паттерн из mnemosyne/reflection).
После каждой сессии: анализирует диалог → извлекает insights →
сохраняет в semantic memory → используется в следующих сессиях через recall().
Наследует _call_llm от BaseLLMAgent.
"""
from pathlib import Path
from langchain_core.messages import HumanMessage
from d_brain.agents.base import BaseLLMAgent # W6
class SelfReviewer(BaseLLMAgent):
SYSTEM_PROMPT = """Проанализируй этот диалог между пользователем и ИИ-ассистентом.
Извлеки:
1. Основные темы и интересы пользователя
2. Что было полезно / что не сработало
3. Предпочтения пользователя по стилю общения
4. Факты о пользователе для долгосрочной памяти
Ответ строго в формате: краткий список пунктов, без воды."""
def __init__(self, data_dir: Path | str = "~/.d_brain"):
super().__init__(data_dir=Path(data_dir) / "review") # W6
async def review_session(
self, user_id: int, session: list[dict]
) -> str:
"""Анализирует сессию, сохраняет insights в memory.
Args:
user_id: ID пользователя
session: список {"role": ..., "content": ...}
Returns:
Строка с insights для логирования
"""
if len(session) < 2:
return "Сессия слишком короткая для анализа"
transcript = "\n".join(
f"[{m['role'].upper()}]: {m['content'][:200]}"
for m in session
)
prompt = f"Диалог для анализа:\n\n{transcript}"
insights = await self._call_llm([HumanMessage(content=prompt)])
# Сохраняем в память для будущих сессий
await self.memory.store(
user_id=user_id,
text=f"self_review: {insights[:300]}",
source="self_review",
)
return insights
Шаг 4 — Запускаем тест (ожидаем PASS):
uv run pytest tests/agents/test_reviewer.py -v
Шаг 5 — Коммит:
git add src/d_brain/agents/reviewer.py tests/agents/test_reviewer.py
git commit -m "feat(agents): add SelfReviewer for self-improvement loop (reflection pattern)"
Финальный smoke-тест всего проекта#
# Запускаем все тесты
uv run pytest tests/ -v --tb=short 2>&1 | tail -30
# Проверяем что старые тесты не сломаны
uv run pytest tests/test_planning/ -v # должны PASS (stable domain)
uv run pytest tests/test_integrations/ -v # должны PASS (keep)
# Удаляем тесты устаревшего processor.py (после полной миграции)
# git rm tests/test_processor.py
# git commit -m "chore: remove legacy processor tests (migrated to agents)"
ЧЕКЛИСТ МИГРАЦИИ#
Что ОСТАВЛЯЕМ (stable, не трогаем):#
tests/test_planning/ ← стабильная доменная логика
tests/test_integrations/ ← Google, Yandex integrations
src/d_brain/planning/ ← keep as-is, PlannerAgent делегирует сюда
src/d_brain/integrations/ ← keep as-is
Что УДАЛЯЕМ после завершения миграции:#
src/d_brain/services/processor.py ← God Object, заменён MainAgent
tests/test_processor.py ← тесты God Object
Что ПОСТЕПЕННО МИГРИРУЕМ (через feature flag):#
src/d_brain/services/coach_service.py → src/d_brain/agents/coach_agent.py
src/d_brain/services/habit_service.py → src/d_brain/agents/habits_agent.py
Дополнительные фиксы из code-review (вне фаз рефакторинга)#
Эти баги нужно исправить параллельно с миграцией, не откладывая до конца:
Bug #8 (Deadlock — proc.communicate без timeout):
# БЫЛО: в _run_morning_plan_loop
proc = await asyncio.create_subprocess_exec(...)
stdout, stderr = await proc.communicate() # ЗАВИСНЕТ НАВСЕГДА если proc стал
# НАДО: явный timeout
try:
stdout, stderr = await asyncio.wait_for(proc.communicate(), timeout=120.0)
except asyncio.TimeoutError:
proc.kill()
await proc.communicate()
raise TimeoutError("Planner subprocess timed out after 120s")
Добавить в src/d_brain/planning/ или где используется proc.communicate().
Коммит: fix(planning): add timeout to subprocess communicate (code-review bug #8)
Bug #6 (Double-tap race в pop_draft):
# В handlers где используется pop_draft — добавить optimistic locking:
# 1. Использовать атомарный rename + check на уровне FS (уже есть atomic_write_text)
# 2. Или добавить asyncio.Lock() на уровне callback handler
Bug #5 (Stale OAuth token — return value игнорируется):
# Найти все места с token refresh и проверить return value:
# БЫЛО: refresh_token(token) # результат игнорируется
# НАДО: token = await refresh_token(token)
# if token is None: raise AuthError("Token refresh failed")
ПЕРЕМЕННЫЕ ОКРУЖЕНИЯ#
Добавить в .env:
# ===== ОБЯЗАТЕЛЬНЫЕ =====
ANTHROPIC_API_KEY= # REQUIRED: LLM API ключ (W27: проверяется при старте)
# ===== Feature flags =====
USE_LANGGRAPH=false # true для нового пути
LANGFUSE_ENABLED=false # true для трассировки
# ===== LLM model (W6) =====
BRAIN_MODEL=claude-sonnet-4-6 # актуальная модель по умолчанию
# claude-opus-4-7 ← максимальная мощность (дороже)
# claude-haiku-4-5-20251001 ← быстрая + дешёвая (для coach/habits)
# ===== Telegram threads =====
THREAD_COACH_ID= # ID топика коуча (int)
THREAD_HABITS_ID= # ID топика привычек (int)
THREAD_PLANNER_ID= # ID топика планировщика (int)
# ===== Пути =====
DATA_DIR=~/.d_brain
VAULT_PATH=~/vault
# ===== YouGile =====
YOUGILE_TOKEN= # API токен
# ===== Memory backend =====
# ChromaDB обязателен. FTS5 используется ТОЛЬКО как secondary index для keyword search
# внутри SQLite episodic memory, НЕ как замена ChromaDB.
CHROMA_PERSIST_DIR=~/.d_brain/chroma # путь к ChromaDB persistent storage
# ===== Skills Security (W19) =====
ALLOWED_SKILL_REPOS= # Comma-separated: "org/repo1,user/repo2"
# Пусто = блокировать все GitHub installs
# "*" = разрешить все (НЕ РЕКОМЕНДУЕТСЯ)
# ===== Langfuse (если LANGFUSE_ENABLED=true) =====
LANGFUSE_PUBLIC_KEY=
LANGFUSE_SECRET_KEY=
LANGFUSE_HOST=https://cloud.langfuse.com
Валидация при старте:
# src/d_brain/core/startup.py — добавить в bot startup
import os
REQUIRED_ENV = ["ANTHROPIC_API_KEY"]
def validate_env():
"""Проверяет обязательные переменные окружения при старте бота."""
missing = [k for k in REQUIRED_ENV if not os.getenv(k)]
if missing:
raise EnvironmentError(
f"Отсутствуют обязательные переменные окружения: {', '.join(missing)}\n"
f"Добавьте в .env файл."
)
ИТОГОВАЯ СТРУКТУРА ТЕСТОВ#
tests/
conftest.py ← ARCH ADDED: PYTHONPATH setup (W26)
core/
test_state.py
test_guardrails.py
test_router.py
test_safe_io.py
test_memory_manager.py
test_observability.py
test_time_utils.py ← ARCH ADDED: W5 TZ fix tests
test_retry.py ← ARCH ADDED: W13 retry strategy tests
memory/
test_working.py
test_episodic.py
test_semantic.py
test_consolidator.py
test_procedural.py ← ARCH ADDED: W20 procedural memory tests
agents/
test_base_agent.py ← ARCH ADDED: W6 BaseLLMAgent tests
test_main_agent.py
test_main_agent_toolcall.py ← ARCH ADDED: W1 LangGraph ToolNode tests
test_coach_agent.py
test_habits_agent.py
test_planner_agent.py
test_reviewer.py
tools/
test_wiki_tools.py
test_search_tools.py
test_media_tools.py
test_planning_tools.py
test_skill_tools.py
skills/
test_registry.py
test_executor.py ← ARCH ADDED: W16 SkillExecutor tests
bot/
test_feature_flag.py
test_thread_routing.py
test_wiki_command.py
Итого тест-файлов: 24 (было 18 → добавлено 6)
Запуск всех тестов:
# С asyncio_mode=auto из pyproject.toml (W15) — @pytest.mark.asyncio не нужен
uv run pytest tests/ -v --tb=short --cov=src/d_brain --cov-report=term-missing
# Только критические пути (быстрый CI):
uv run pytest tests/core/ tests/memory/ tests/agents/test_main_agent.py -v
ССЫЛКИ И ИСТОЧНИКИ#
- LangGraph docs: https://langchain-ai.github.io/langgraph/
- LangGraph ToolNode: https://langchain-ai.github.io/langgraph/reference/prebuilt/#langgraph.prebuilt.tool_node.ToolNode
- LangGraph tools_condition: https://langchain-ai.github.io/langgraph/reference/prebuilt/#langgraph.prebuilt.tool_node.tools_condition
- John Ousterhout Deep Modules: "A Philosophy of Software Design"
- mnemosyne pattern: https://github.com/AxDSan/mnemosyne
- ChromaDB: https://docs.trychroma.com/
- Langfuse: https://langfuse.com/docs
- duckduckgo-search: https://pypi.org/project/duckduckgo-search/
- YouGile API: https://ru.yougile.com/api-v2#/
- tenacity (retry): https://tenacity.readthedocs.io/
- youtube-transcript-api: https://pypi.org/project/youtube-transcript-api/
- SQLite WAL mode: https://www.sqlite.org/wal.html
- Python asyncio.to_thread: https://docs.python.org/3/library/asyncio-task.html#asyncio.to_thread
- code-review: docs/code-review-2026-05-26.md
ИТОГ ARCH REVIEW#
Количество изменений: 45 аннотаций ARCH REVISED/ADDED
Новые задачи: 7 (0.0, 0.5, 0.6, 1.6, 2.4-bis, 2.5-retry, 6.3)
Критические исправления кода: 12 (W1, W2, W3, W4, W5, W6, W7, W8, W10, W11, W12, W16)
Пересмотренная оценка: 18 рабочих дней (было 12)
Топ-3 изменения без которых план провалится:
- W1 (ToolNode) — без него граф никогда не вызывает tools
- W3+W4 (WAL+Lock) — без них бот будет падать при нагрузке
- W5 (TZ fix) — без него phantom completions в close_day и YouGile sync
ФАЗА 4.5: СИСТЕМА КАСКАДНОГО ПЛАНИРОВАНИЯ#
Почему так: Каскадное планирование решает проблему разрыва между высокоуровневыми целями и ежедневными задачами. Без явного каскада цели остаются декларациями, а задачи — бесструктурным списком. Stable IDs (goal_id, habit_id, commitment_id) позволяют трекать план/факт через любой временной горизонт без потери связи «задача ← цель». ISO-недели (YYYY-Www) — интернациональный стандарт, избегающий путаницы с локальными календарями.
Пример: Цель года:
goal_id=G2026-01 "Запустить продукт"→ quarterly commitment:cmt_id=C2026-Q2-03 "MVP готов"→ weekly plan:plan_id=P2026-W21-07 "Написать API"→ daily task: done → fact event записывается → roll-up показывает 73% выполнения цели за квартал.
Цель: Реализовать полный каскад планирования multi-year → year → quarter → month → week → day с двунаправленным потоком (планы вниз, факты вверх). Оценка: ~3 рабочих дня Принципы: §4.1 human-readable plans, §4.2 stable IDs, §4.3 plans mutable / facts append-only, §4.4 compatible migration, §4.5 ISO weeks, §4.6 delivery checkpoints
Задача 4.5.1: Канонические периоды и data model#
Почему так: Стабильные
goal_idс перекрёстными ссылками между периодами — основа traceability: каждая дневная задача должна быть traceable к цели года, иначе план и цели существуют отдельно. Пример: Пользователь ставит цель «выучить LangGraph» (goal_id: g-2026-001) → разбивается на квартальные milestone → недельные задачи → ежедневные пункты — все ссылаются на g-2026-001.
Цель: Реализовать модель данных из §6 — Periods, Goals, Plan Items, Fact Events.
Файлы:
- Создать:
src/d_brain/planning/cascade/models.py - Создать:
src/d_brain/planning/cascade/db.py - Тест:
tests/planning/test_cascade_models.py
# src/d_brain/planning/cascade/models.py
"""Канонические модели данных планирования (§6 roadmap).
Принципы:
- Periods: ISO-periods (YYYY, YYYY-QN, YYYY-MM, YYYY-Www, YYYY-MM-DD)
- Goals: stable goal_id, иерархия через parent_goal_id
- Plan Items: source_type (goal/habit/project/yougile/user), kind (commitment/task/habit/ritual)
- Fact Events: append-only event ledger
"""
from dataclasses import dataclass, field
from enum import StrEnum
from typing import Optional
class PeriodType(StrEnum):
DAY = "day" # YYYY-MM-DD
WEEK = "week" # YYYY-Www
MONTH = "month" # YYYY-MM
QUARTER = "quarter" # YYYY-QN
YEAR = "year" # YYYY
MULTI_YEAR = "multi_year" # YYYY-YYYY
class PlanItemKind(StrEnum):
COMMITMENT = "commitment" # обязательство на период
TASK = "task" # конкретная задача
HABIT = "habit" # привычка
RITUAL = "ritual" # ритуал/рутина
class FactEventType(StrEnum):
PLANNED = "planned"
STARTED = "started"
COMPLETED = "completed"
MISSED = "missed"
MOVED = "moved"
BLOCKED = "blocked"
HABIT_CHECKED = "habit_checked"
HABIT_SKIPPED = "habit_skipped"
REVIEW_CLOSED = "review_closed"
class ItemStatus(StrEnum):
PLANNED = "planned"
DONE = "done"
NOT_DONE = "not_done"
MOVED = "moved"
BLOCKED = "blocked"
CANCELED = "canceled"
SKIPPED = "skipped"
@dataclass
class Period:
"""Канонический период. §4.5: period_id = ISO 8601 строка."""
period_id: str # "2026", "2026-Q2", "2026-05", "2026-W21", "2026-05-27"
period_type: PeriodType
start_date: str # YYYY-MM-DD
end_date: str # YYYY-MM-DD
parent_period_id: Optional[str] = None
status: str = "active" # active | closed
@dataclass
class Goal:
"""Цель с иерархией. §4.2: stable goal_id."""
goal_id: str # "G2026-01"
title: str
horizon: str # "year" | "quarter" | "month"
area: str # "бизнес" | "здоровье" | "обучение"
parent_goal_id: Optional[str] = None
success_metric: str = ""
metric_type: str = "binary" # binary | numeric | percentage
active_from: str = ""
active_to: str = ""
@dataclass
class PlanItem:
"""Элемент плана. §4.3: планы мутируемы."""
plan_item_id: str # "P2026-W21-07"
source_type: str # "goal" | "habit" | "project" | "yougile" | "user"
source_id: str # goal_id или yougile task id
period_id: str # к какому периоду относится
title: str
kind: PlanItemKind
priority: int = 2 # 1=критично, 2=важно, 3=желательно
due_date: Optional[str] = None
parent_plan_item_id: Optional[str] = None
status: ItemStatus = ItemStatus.PLANNED
owner: str = ""
@dataclass
class FactEvent:
"""Событие факта. §4.3: факты append-only (event ledger)."""
event_id: str # UUID
event_type: FactEventType
event_date: str # YYYY-MM-DD UTC
plan_item_id: Optional[str] = None
period_id: Optional[str] = None
status_after: Optional[ItemStatus] = None
value: Optional[float] = None # для числовых метрик
notes: str = ""
source: str = "user" # "user" | "yougile" | "calendar" | "system"
# src/d_brain/planning/cascade/db.py
"""SQLite хранилище для каскадного планирования."""
import asyncio
import json
import uuid
import aiosqlite
from pathlib import Path
from d_brain.core.time_utils import utc_today, utc_now, format_iso
from d_brain.planning.cascade.models import (
Period, Goal, PlanItem, FactEvent, PeriodType, PlanItemKind, ItemStatus, FactEventType
)
class CascadeDB:
"""Хранилище планов, целей и фактов для каскадного планирования."""
def __init__(self, db_path: Path | str = "~/.d_brain/planning.db"):
self._db_path = Path(db_path).expanduser()
self._db_path.parent.mkdir(parents=True, exist_ok=True)
self._initialized = False
self._init_lock = asyncio.Lock()
async def _ensure_init(self):
if self._initialized:
return
async with self._init_lock:
if self._initialized:
return
async with aiosqlite.connect(self._db_path) as db:
await db.execute("PRAGMA journal_mode=WAL")
await db.execute("""
CREATE TABLE IF NOT EXISTS periods (
period_id TEXT PRIMARY KEY,
period_type TEXT NOT NULL,
start_date TEXT NOT NULL,
end_date TEXT NOT NULL,
parent_period_id TEXT,
status TEXT DEFAULT 'active'
)
""")
await db.execute("""
CREATE TABLE IF NOT EXISTS goals (
goal_id TEXT PRIMARY KEY,
title TEXT NOT NULL,
horizon TEXT,
area TEXT,
parent_goal_id TEXT,
success_metric TEXT,
metric_type TEXT DEFAULT 'binary',
active_from TEXT,
active_to TEXT
)
""")
await db.execute("""
CREATE TABLE IF NOT EXISTS plan_items (
plan_item_id TEXT PRIMARY KEY,
source_type TEXT,
source_id TEXT,
period_id TEXT NOT NULL,
parent_plan_item_id TEXT,
title TEXT NOT NULL,
kind TEXT NOT NULL,
priority INTEGER DEFAULT 2,
due_date TEXT,
status TEXT DEFAULT 'planned',
owner TEXT DEFAULT ''
)
""")
await db.execute("""
CREATE TABLE IF NOT EXISTS fact_events (
event_id TEXT PRIMARY KEY,
event_type TEXT NOT NULL,
event_date TEXT NOT NULL,
plan_item_id TEXT,
period_id TEXT,
status_after TEXT,
value REAL,
notes TEXT,
source TEXT DEFAULT 'user',
created_at TEXT NOT NULL
)
""")
await db.execute("CREATE INDEX IF NOT EXISTS idx_plan_period ON plan_items(period_id)")
await db.execute("CREATE INDEX IF NOT EXISTS idx_facts_date ON fact_events(event_date)")
await db.execute("CREATE INDEX IF NOT EXISTS idx_facts_item ON fact_events(plan_item_id)")
await db.commit()
self._initialized = True
async def add_fact_event(
self,
event_type: FactEventType,
plan_item_id: str | None = None,
period_id: str | None = None,
status_after: ItemStatus | None = None,
value: float | None = None,
notes: str = "",
source: str = "user",
) -> str:
"""Добавляет событие факта в append-only ledger."""
await self._ensure_init()
event_id = str(uuid.uuid4())
today = utc_today()
async with aiosqlite.connect(self._db_path) as db:
await db.execute("""
INSERT INTO fact_events
(event_id, event_type, event_date, plan_item_id, period_id,
status_after, value, notes, source, created_at)
VALUES (?,?,?,?,?,?,?,?,?,?)
""", (
event_id, str(event_type), today, plan_item_id, period_id,
str(status_after) if status_after else None,
value, notes, source, format_iso(utc_now())
))
# Обновляем статус plan_item (планы мутируемы)
if plan_item_id and status_after:
await db.execute(
"UPDATE plan_items SET status=? WHERE plan_item_id=?",
(str(status_after), plan_item_id)
)
await db.commit()
return event_id
async def get_period_rollup(self, period_id: str) -> dict:
"""Агрегирует план/факт статистику для периода."""
await self._ensure_init()
async with aiosqlite.connect(self._db_path) as db:
db.row_factory = aiosqlite.Row
async with db.execute(
"SELECT status, COUNT(*) as cnt FROM plan_items WHERE period_id=? GROUP BY status",
(period_id,)
) as cur:
rows = await cur.fetchall()
counts = {r["status"]: r["cnt"] for r in rows}
total = sum(counts.values())
done = counts.get("done", 0)
return {
"period_id": period_id,
"total": total,
"done": done,
"completion_rate": round(done / total * 100, 1) if total > 0 else 0,
"by_status": counts,
}
Тест:
# tests/planning/test_cascade_models.py
import pytest
from d_brain.planning.cascade.models import Period, Goal, PlanItem, FactEvent, PeriodType, PlanItemKind
def test_period_id_formats():
"""period_id соответствует ISO форматам из §4.5."""
periods = [
Period("2026", PeriodType.YEAR, "2026-01-01", "2026-12-31"),
Period("2026-Q2", PeriodType.QUARTER, "2026-04-01", "2026-06-30"),
Period("2026-05", PeriodType.MONTH, "2026-05-01", "2026-05-31"),
Period("2026-W21", PeriodType.WEEK, "2026-05-18", "2026-05-24"),
Period("2026-05-27", PeriodType.DAY, "2026-05-27", "2026-05-27"),
]
assert all(p.period_id for p in periods)
def test_goal_hierarchy():
parent = Goal("G2026-01", "Запустить продукт", "year", "бизнес")
child = Goal("G2026-Q2-01", "MVP", "quarter", "бизнес", parent_goal_id="G2026-01")
assert child.parent_goal_id == parent.goal_id
@pytest.mark.asyncio
async def test_cascade_db_fact_event(tmp_path):
from d_brain.planning.cascade.db import CascadeDB
from d_brain.planning.cascade.models import FactEventType, ItemStatus
db = CascadeDB(db_path=tmp_path / "cascade.db")
event_id = await db.add_fact_event(
event_type=FactEventType.COMPLETED,
plan_item_id="P2026-W21-01",
status_after=ItemStatus.DONE,
notes="Задача выполнена"
)
assert len(event_id) == 36 # UUID формат
git add src/d_brain/planning/cascade/ tests/planning/test_cascade_models.py
git commit -m "feat(planning): add cascade planning models — Periods, Goals, PlanItems, FactEvents (§6)"
Задача 4.5.2: Carry-over механизм#
Почему так: Carry-over из незавершённых задач — честный перенос долга; без него невыполненные задачи исчезают в истории и пользователь теряет из виду накопившиеся обязательства. Пример: Понедельник: 3 задачи не выполнены → Carry-over помечает их
⏮ перенос→ добавляет во вторник с флагом → daily plan вторника: 3 перенесённых + 5 новых.
Цель: Автоматически переносить невыполненные задачи в следующий период с записью факт-события MOVED.
# src/d_brain/planning/cascade/carryover.py
"""Carry-over: перенос невыполненных задач в следующий период.
§4.3: факты append-only — создаём FactEvent(type=MOVED) для каждого переноса.
Планы мутируемы — обновляем period_id в plan_items.
"""
from d_brain.planning.cascade.db import CascadeDB
from d_brain.planning.cascade.models import FactEventType, ItemStatus
async def run_carryover(db: CascadeDB, from_period: str, to_period: str) -> list[str]:
"""Переносит нeвыполненные задачи из from_period в to_period.
Returns: список plan_item_id которые были перенесены.
"""
import aiosqlite
moved_ids = []
async with aiosqlite.connect(db._db_path) as conn:
conn.row_factory = aiosqlite.Row
async with conn.execute(
"SELECT plan_item_id, title FROM plan_items WHERE period_id=? AND status IN ('planned','blocked')",
(from_period,)
) as cur:
items = await cur.fetchall()
for item in items:
pid = item["plan_item_id"]
# Обновляем период (план мутируется)
await conn.execute(
"UPDATE plan_items SET period_id=?, status='planned' WHERE plan_item_id=?",
(to_period, pid)
)
moved_ids.append(pid)
await conn.commit()
# Записываем факт-события (append-only ledger)
for pid in moved_ids:
await db.add_fact_event(
event_type=FactEventType.MOVED,
plan_item_id=pid,
period_id=from_period,
status_after=ItemStatus.MOVED,
notes=f"Перенесено из {from_period} в {to_period}",
source="system",
)
return moved_ids
git add src/d_brain/planning/cascade/carryover.py
git commit -m "feat(planning): add carry-over mechanism — moves unfinished items with MOVED fact event"
Задача 4.5.3: Roll-up агрегация week → month → quarter → year#
Почему так: Append-only rollup (день→неделя→месяц→квартал→год) создаёт иерархическую историю фактов без потери детализации; позволяет ответить на «что я сделал за квартал» одним запросом. Пример: Пользователь просит «итоги квартала» →
rollup.get_quarter("2026-Q1")→ агрегирует из 3 monthly rollups → из 13 weekly rollups → ответ с ключевыми достижениями за квартал.
Цель: Вычислять план/факт метрики на каждом уровне каскада.
# src/d_brain/planning/cascade/rollup.py
"""Roll-up агрегация: day→week→month→quarter→year.
§7 roadmap: накапливаем метрики снизу вверх.
Каждый уровень агрегирует дочерние периоды.
"""
from d_brain.planning.cascade.db import CascadeDB
async def rollup_period(db: CascadeDB, period_id: str) -> dict:
"""Агрегирует метрики для периода включая все дочерние.
Иерархия: week агрегирует дни, month агрегирует недели,
quarter агрегирует месяцы, year агрегирует кварталы.
Returns: dict с completion_rate, done, total, child_rollups
"""
own = await db.get_period_rollup(period_id)
import aiosqlite
async with aiosqlite.connect(db._db_path) as conn:
conn.row_factory = aiosqlite.Row
async with conn.execute(
"SELECT period_id FROM periods WHERE parent_period_id=?",
(period_id,)
) as cur:
children = [r["period_id"] for r in await cur.fetchall()]
child_rollups = []
for child_id in children:
child_rollup = await rollup_period(db, child_id)
child_rollups.append(child_rollup)
if child_rollups:
# Агрегируем дочерние (для периодов без прямых plan_items)
total_child = sum(c["total"] for c in child_rollups)
done_child = sum(c["done"] for c in child_rollups)
if total_child > 0:
own["total"] = own["total"] + total_child
own["done"] = own["done"] + done_child
own["completion_rate"] = round(own["done"] / own["total"] * 100, 1)
own["child_rollups"] = child_rollups
return own
git add src/d_brain/planning/cascade/rollup.py
git commit -m "feat(planning): add cascade roll-up — week→month→quarter→year fact aggregation (§7)"
КАСКАД ПЛАНИРОВАНИЯ#
Почему так: Явная визуализация двунаправленного потока решает типичную проблему систем планирования — цели живут отдельно от задач. Каскад делает связь явной: каждая задача трейсируется к цели, каждый факт поднимается к ежегодному итогу.
═══════════════════════════════════════════════════════════
КАСКАД ПЛАНИРОВАНИЯ — SECOND BRAIN V2
═══════════════════════════════════════════════════════════
ПЛАНЫ ТЕКУТ ВНИЗ ФАКТЫ ТЕКУТ ВВЕРХ
━━━━━━━━━━━━━━━━ ━━━━━━━━━━━━━━━━━
2026-2028 (Видение) ▲
"Компания с 10M ARR" │ ежегодный итог
│ │
▼ YEAR CLOSE
2026 (Год) "достигнуто 73%"
"Запустить продукт, 1M ARR" ▲
│ │ квартальный итог
▼ │
2026-Q2 (Квартал) QUARTER CLOSE
"MVP готов, 3 клиента" "Q2: 8/10 целей PASS"
│ ▲
▼ │ недельный итог
2026-05 (Месяц) │
"API + дизайн + онбординг" WEEKLY REVIEW
│ "неделя: 5/7 задач done"
▼ ▲
2026-W21 (Неделя) │ дневное закрытие
"Написать API v1" │
│ DAILY CLOSE
▼ "сегодня: 3 задачи done"
2026-05-27 (День) ▲
"endpoint /tasks, тест" │
│ ─────────────────┘
▼ пользователь отмечает выполнение
[ВЫПОЛНЕНИЕ] → FactEvent(COMPLETED) записывается
Задача сделана → статус обновляется
→ roll-up пересчитывается
═══════════════════════════════════════════════════════════
МЕХАНИЗМЫ:
• Цели (year/quarter) → PlanItems с source_type="goal"
• Привычки → PlanItems с kind="habit", ежедневные
• YouGile задачи → PlanItems с source_type="yougile"
• Невыполненное → carry-over (MOVED event) в следующий период
• Каждый уровень: completion_rate = done/total × 100%
═══════════════════════════════════════════════════════════
Vault структура для каскада#
vault/goals/
vision/
2026-2028.md # Видение на 3 года (мутируемое)
years/
2026.md # Годовой план + итог
quarters/
2026-Q1.md
2026-Q2.md # Текущий квартал
months/
2026-04.md
2026-05.md # Текущий месяц
weeks/
2026-W20.md
2026-W21.md # Текущая неделя
current/
yearly.md → symlink на years/2026.md
quarterly.md → symlink на quarters/2026-Q2.md
monthly.md → symlink на months/2026-05.md
weekly.md → symlink на weeks/2026-W21.md
close/
2026-Q1-close.md # Закрытые периоды (append-only)
2026-W20-close.md
vault/system/planning/
events/2026/
2026-05.jsonl # Append-only event ledger (FactEvents)
snapshots/2026/
2026-W21.json # Снэпшот на начало недели
2026-Q2.json # Снэпшот на начало квартала
state/
planning.db # SQLite derived store (CascadeDB)
АРХИТЕКТУРА АГЕНТОВ: CoachAgent КАК ОРКЕСТРАТОР#
Почему так: CoachAgent знает ЗАЧЕМ пользователь делает те или иные задачи — он владеет целями. MainAgent обрабатывает любые запросы, но не видит целостной картины. Когда CoachAgent стоит выше HabitsAgent и PlannerAgent, он может интегрировать сигналы: «привычки падают И задачи не закрываются → пользователь перегружен → предложить приоритизацию». Без оркестратора каждый агент оптимизирует своё, не видя системной картины.
Пример: В пятницу CoachAgent запускает goal health check. Видит: цель «написать 3 статьи в Q2» — выполнена 0%. Дедлайн через 4 недели. Инициирует сообщение: «У тебя 0/3 статей за Q2. Осталось 4 недели. Что мешает? Хочешь поставить задачу на следующую неделю?»
Топология агентов v2 (с оркестратором)#
┌─────────────────────────────────────────────────────────┐
│ TELEGRAM TOPICS │
│ #main #coach #habits #planning │
└────┬───────────────┬──────────────┬──────────┬──────────┘
│ │ │ │
▼ ▼ │ │
MainAgent CoachAgent │ │
(general) ОРКЕСТРАТОР │ │
│ │ │
├─────────────▶ │
│ HabitsAgent │
│ (делегат) │
│ │ │
│ ▼ │
│ habit_data │
│ │
├─────────────────────────▶
│ PlannerAgent
│ (делегат)
│ │
│ ▼
│ yougile_data
│ calendar_data
│
▼
goal_health_check()
weekly_review()
proactive_nudge()
CoachAgent — расширенный функционал оркестратора#
# src/d_brain/agents/coach_agent.py — РАСШИРЕНИЕ (добавить к существующему)
class CoachAgent(BaseLLMAgent):
"""CoachAgent v2 — оркестратор с доступом к HabitsAgent и PlannerAgent.
Новые обязанности:
1. goal_health_check() — еженедельный анализ целей vs факты
2. proactive_nudge() — проактивное уведомление если пользователь отстаёт
3. weekly_coaching_review() — совместный review с HabitsAgent + PlannerAgent
"""
# Импортируем данные из дочерних агентов (LangGraph subgraph invocation)
async def goal_health_check(self, user_id: int) -> str:
"""Еженедельный health check: сравниваем план vs факт по всем целям.
Вызывается по cron каждую пятницу вечером.
Собирает данные из CascadeDB и генерирует coaching insights.
"""
from d_brain.planning.cascade.db import CascadeDB
from d_brain.planning.cascade.rollup import rollup_period
from d_brain.core.time_utils import utc_now
db = CascadeDB()
current_year = utc_now().year
year_rollup = await rollup_period(db, str(current_year))
prompt = f"""Проведи goal health check для пользователя.
Текущий год: {current_year}
Выполнение целей: {year_rollup['completion_rate']}%
Статусы: {year_rollup['by_status']}
Задай 3 конкретных вопроса:
1. Что идёт хорошо и почему?
2. Где самый большой разрыв план/факт?
3. Какое одно действие на следующей неделе максимально продвинет цели?
Будь прямым, требовательным, поддерживающим."""
from langchain_core.messages import HumanMessage
return await self._call_llm([HumanMessage(content=prompt)])
async def proactive_nudge(self, user_id: int, context: dict) -> str | None:
"""Анализирует паттерны и отправляет проактивное уведомление.
Триггеры для nudge:
- completion_rate < 50% за текущую неделю (пятница 15:00)
- habit streak разрыв 3+ дня подряд
- задача в плане без прогресса 5+ дней
Returns: текст сообщения или None если nudge не нужен
"""
week_completion = context.get("week_completion", 100)
habit_streak_broken = context.get("habit_streak_broken", False)
stalled_tasks = context.get("stalled_tasks", [])
if week_completion < 50 or habit_streak_broken or stalled_tasks:
issues = []
if week_completion < 50:
issues.append(f"выполнение недели: {week_completion}%")
if habit_streak_broken:
issues.append("разрыв streak привычек")
if stalled_tasks:
issues.append(f"задачи без прогресса: {len(stalled_tasks)}")
prompt = f"Пользователь отстаёт: {', '.join(issues)}. Напиши короткое (2-3 предложения) поддерживающее сообщение с конкретным следующим шагом."
from langchain_core.messages import HumanMessage
return await self._call_llm([HumanMessage(content=prompt)])
return None
async def weekly_coaching_review(
self,
user_id: int,
habits_summary: str,
planning_summary: str,
) -> str:
"""Интегрированный weekly review: коуч + привычки + планирование.
Вызывает HabitsAgent и PlannerAgent как LangGraph subgraphs,
собирает данные, формирует целостную картину недели.
"""
prompt = f"""Еженедельный review коуча.
ПРИВЫЧКИ за неделю:
{habits_summary}
ЗАДАЧИ за неделю:
{planning_summary}
Дай:
1. Что реально произошло (факты)
2. Что это значит (паттерн)
3. Что усилить на следующей неделе (конкретное действие)
4. Один сложный вопрос для рефлексии"""
from langchain_core.messages import HumanMessage
return await self._call_llm([HumanMessage(content=prompt)])
СТРУКТУРА ПРОЕКТА SECOND BRAIN V2#
Почему так: Чистый старт в
/home/serg/projects/second-brain-v2/позволяет строить правильную архитектуру без legacy constraints. Текущий репозиторий остаётся рабочим до полной миграции. Структура модулей отражает архитектурные слои:core/→memory/→agents/→tools/→bot/.
second-brain-v2/ # корень нового проекта
├── src/
│ └── d_brain/
│ ├── core/ # фундаментальные компоненты
│ │ ├── __init__.py
│ │ ├── state.py # AgentState TypedDict
│ │ ├── router.py # ThreadRouter (thread_id → agent)
│ │ ├── guardrails.py # path traversal + SSRF защита
│ │ ├── safe_io.py # атомарная запись с fsync
│ │ ├── time_utils.py # utc_now(), utc_today()
│ │ ├── memory.py # MemoryManager фасад
│ │ ├── retry.py # with_retry decorator (tenacity)
│ │ ├── observability.py # Langfuse tracer (noop/real)
│ │ └── startup.py # env validation при старте
│ │
│ ├── agents/ # LangGraph агенты
│ │ ├── __init__.py
│ │ ├── base.py # BaseLLMAgent (BRAIN_MODEL env)
│ │ ├── main_agent.py # MainAgent — general assistant
│ │ ├── coach_agent.py # CoachAgent — ОРКЕСТРАТОР (goal owner)
│ │ ├── habits_agent.py # HabitsAgent — habit tracking
│ │ ├── planner_agent.py # PlannerAgent — YouGile + calendar
│ │ ├── reviewer.py # SelfReviewer — self-improvement
│ │ ├── souls/ # Soul.md файлы для каждого агента
│ │ │ ├── main.md # личность, scope, эскалация MainAgent
│ │ │ ├── coach.md # личность, scope, эскалация CoachAgent
│ │ │ ├── habits.md # личность, scope HabitsAgent
│ │ │ └── planner.md # личность, scope PlannerAgent
│ │ └── langchain_tools.py # @tool обёртки для bind_tools()
│ │
│ ├── tools/ # MCP-style tool modules
│ │ ├── __init__.py
│ │ ├── wiki_tools.py # Obsidian vault read/write/search
│ │ ├── planning_tools.py # YouGile API + calendar
│ │ ├── memory_tools.py # recall/store через MemoryManager
│ │ ├── search_tools.py # DuckDuckGo web search
│ │ ├── media_tools.py # YouTube/article/TG text extraction
│ │ └── skill_tools.py # install/remove/list skills
│ │
│ ├── memory/ # 4 слоя памяти
│ │ ├── __init__.py
│ │ ├── working.py # контекстное окно (in-process)
│ │ ├── episodic.py # SQLite сессии + summaries
│ │ ├── semantic.py # FTS5/ChromaDB поиск
│ │ ├── procedural.py # SOP/процедуры (SQLite)
│ │ └── consolidator.py # ночная консолидация (mnemosyne)
│ │
│ ├── skills/ # динамические скиллы
│ │ ├── __init__.py
│ │ ├── registry.py # реестр скиллов с per-agent scope
│ │ ├── loader.py # install_from_file/github
│ │ └── executor.py # SkillExecutor → LangChain @tool
│ │
│ ├── bot/ # Telegram bot (aiogram)
│ │ ├── __init__.py
│ │ ├── dispatcher.py # route_message() функция
│ │ └── handlers/
│ │ ├── __init__.py # get_message_processor()
│ │ ├── message_handler.py # основной handler с error boundary
│ │ └── wiki_command.py # /wiki команда
│ │
│ ├── planning/ # каскадное планирование
│ │ ├── __init__.py
│ │ └── cascade/
│ │ ├── __init__.py
│ │ ├── models.py # Period, Goal, PlanItem, FactEvent
│ │ ├── db.py # CascadeDB (SQLite)
│ │ ├── carryover.py # перенос невыполненных задач
│ │ └── rollup.py # агрегация week→month→quarter→year
│ │
│ └── integrations/ # внешние системы
│ ├── __init__.py
│ ├── yougile.py # YouGile REST API (pagination, UTC)
│ ├── google_cal.py # Google Calendar API
│ ├── yandex_cal.py # Yandex CalDAV
│ └── obsidian.py # Obsidian vault utilities
│
├── vault/ # Obsidian knowledge base
│ ├── goals/ # каскадное планирование
│ │ ├── vision/
│ │ ├── years/
│ │ ├── quarters/
│ │ ├── months/
│ │ ├── weeks/
│ │ ├── current/ # symlinks на актуальные периоды
│ │ └── close/ # закрытые периоды (append-only)
│ ├── projects/ # CRM проекты
│ ├── contacts/ # CRM контакты
│ ├── daily/ # ежедневные заметки
│ ├── habits/ # определения привычек
│ ├── media/ # YouTube/статьи → markdown
│ │ └── telegram/ # TG посты → markdown
│ ├── skills/ # процедурные знания
│ └── system/
│ ├── planning/events/ # append-only event ledger (.jsonl)
│ ├── planning/snapshots/ # периодические снэпшоты
│ └── kaizen-log.md # лог самоулучшений системы
│
├── bpmn/ # BPMN 2.0 диаграммы процессов
│ ├── 01-month-week-planning.bpmn
│ ├── 02-day-planning.bpmn
│ ├── 03-day-execution.bpmn
│ ├── 04-plan-adjustment.bpmn
│ ├── 05-day-close.bpmn
│ ├── 06-calendar-meetings.bpmn
│ ├── 07-project-info.bpmn
│ └── 08-wiki-task-update.bpmn
│
├── tests/ # полное тест-покрытие
│ ├── conftest.py # PYTHONPATH + shared fixtures
│ ├── core/
│ ├── memory/
│ ├── agents/
│ ├── tools/
│ ├── skills/
│ ├── bot/
│ └── planning/
│
├── docs/ # документация
│ ├── architecture-v2.html # диаграмма архитектуры
│ └── plans/ # планы разработки
│
├── state/ # derived state (gitignored)
│ └── planning.db # CascadeDB SQLite
│
├── .env.example # все переменные окружения с описанием
├── pyproject.toml # зависимости + pytest config
├── AGENTS.md # repo-task-proof-loop конфиг
└── CLAUDE.md # инструкции для Claude subprocess
СТРАТЕГИЯ СОЗДАНИЯ НОВОГО ПРОЕКТА#
Почему так: Чистый старт избегает накопленного технического долга и позволяет применить все архитектурные решения с первого коммита. Feature-flag миграция гарантирует что текущий бот не ломается в процессе перехода.
Что КОПИРОВАТЬ из current second-brain#
# Из current second-brain: проверенная доменная логика
cp -r src/d_brain/planning/ ../second-brain-v2/src/d_brain/planning/
cp -r src/d_brain/integrations/ ../second-brain-v2/src/d_brain/integrations/
cp -r vault/ ../second-brain-v2/vault/
cp -r tests/test_planning/ ../second-brain-v2/tests/planning/
cp -r tests/test_integrations/ ../second-brain-v2/tests/integrations/
Что адаптировать (не копировать напрямую):
integrations/yougile.py→ добавить пагинацию (W11), UTC timestamps (W12)integrations/google_cal.py→ добавить retry decorator (W13)planning/→ добавить CascadeDB поверх (не заменять)
Что ЗАИМСТВОВАТЬ из Hermes-архитектуры (MCP паттерн)#
FROM Hermes:
├── MCP-style tools → уже применено: tools/ как standalone modules
├── cron/watchdog pattern → применить для memory consolidator + goal health check
│ cron(22:00): consolidator.run_daily()
│ cron(пятница 17:00): coach.goal_health_check()
├── skill loader concept → Skills system (Фаза 6)
└── session DB pattern → EpisodicMemory (Фаза 1.2)
Что ПИСАТЬ С НУЛЯ#
FROM SCRATCH:
├── core/ # LangGraph state, router, guardrails
├── agents/ # все 4 агента как чистые LangGraph графы
├── memory/ # все 4 слоя памяти
├── planning/cascade/ # каскадное планирование §6
└── BPMN процессы # 8 основных процессов
Чеклист миграции#
- Создать
- Инициализировать
ВЫБОР МОДЕЛИ И ПРОВАЙДЕРА#
Почему так: Разные задачи требуют разных модельных трейдоффов. Рассуждение о целях требует лучшей модели (Sonnet). Суммаризация памяти — задача для дешёвой модели (DeepSeek). Классификация intent — для ультра-дешёвой (Flash). BaseLLMAgent поддерживает переключение провайдера per-call через provider параметр, не меняя код агента.
Выбор модели — Primary + Fallback (аналог Hermes)#
Пользователь настраивает ONE primary модель и ONE fallback. Система НЕ переключает модели по типу задачи автоматически.
BRAIN_MODEL=claude-sonnet-4-5
BRAIN_PROVIDER=anthropic # anthropic | openrouter | deepseek
BRAIN_FALLBACK_MODEL=deepseek/deepseek-v4-flash:free
BRAIN_FALLBACK_PROVIDER=openrouter
OPENROUTER_API_KEY=sk-or-...
DEEPSEEK_API_KEY=sk-...
ANTHROPIC_API_KEY=sk-ant-...
BaseLLMAgent логика:
- Вызвать primary (BRAIN_MODEL / BRAIN_PROVIDER)
- Если RateLimitError / APIError / Timeout → вызвать fallback (BRAIN_FALLBACK_MODEL / BRAIN_FALLBACK_PROVIDER)
- Если оба упали → raise с понятным сообщением пользователю
Смена модели без перезапуска: /model claude-opus-4 или /model deepseek/deepseek-v3
Команда обновляет BRAIN_MODEL в runtime и сохраняет в .env.
Конфигурация
# .env — модели и провайдеры (Hermes-style: primary + fallback)
BRAIN_MODEL=claude-sonnet-4-5
BRAIN_PROVIDER=anthropic # anthropic | openrouter | deepseek
BRAIN_FALLBACK_MODEL=deepseek/deepseek-v4-flash:free
BRAIN_FALLBACK_PROVIDER=openrouter
# OpenRouter API key (для fallback провайдера)
OPENROUTER_API_KEY=
# DeepSeek (прямой API, если дешевле OpenRouter)
DEEPSEEK_API_KEY=
DEEPSEEK_BASE_URL=https://api.deepseek.com
# Runtime: пользователь меняет модель командой /model <model_name>
# Например: /model claude-opus-4 или /model deepseek/deepseek-v3
BaseLLMAgent — расширение для multi-provider#
# src/d_brain/agents/base.py — ДОПОЛНЕНИЕ
class BaseLLMAgent:
"""Расширен для поддержки нескольких провайдеров (Hermes-style primary+fallback)."""
async def _call_llm(
self,
messages: list,
system: str | None = None,
temperature: float = 0.7,
) -> str:
"""Вызывает primary LLM; при RateLimitError/APIError переключается на fallback.
Primary: BRAIN_MODEL (default: claude-sonnet-4-5) via BRAIN_PROVIDER
Fallback: BRAIN_FALLBACK_MODEL (default: deepseek/deepseek-v4-flash:free) via BRAIN_FALLBACK_PROVIDER
Runtime: пользователь меняет модель командой /model <model_name>
"""
import os
try:
model = os.getenv("BRAIN_MODEL", "claude-sonnet-4-5")
provider = os.getenv("BRAIN_PROVIDER", "anthropic")
return await self._call_provider(messages, system, model, provider, temperature)
except Exception as e:
if "rate_limit" in str(e).lower() or "api" in str(e).lower():
fallback_model = os.getenv("BRAIN_FALLBACK_MODEL", "deepseek/deepseek-v4-flash:free")
fallback_provider = os.getenv("BRAIN_FALLBACK_PROVIDER", "openrouter")
return await self._call_provider(messages, system, fallback_model, fallback_provider, temperature)
raise
async def _call_provider(self, messages, system, model, provider, temperature) -> str:
"""Диспатч по провайдеру."""
if provider == "openrouter":
return await self._call_openrouter(messages, system, model, temperature)
else:
return await self._call_anthropic(messages, system, model, temperature)
async def _call_openrouter(self, messages, system, model, temperature) -> str:
"""Вызов через OpenRouter API (OpenAI-совместимый)."""
import httpx, os, json
api_key = os.environ["OPENROUTER_API_KEY"]
headers = {
"Authorization": f"Bearer {api_key}",
"HTTP-Referer": "https://github.com/second-brain",
"Content-Type": "application/json",
}
body = {
"model": model,
"messages": [{"role": "system", "content": system or self.SYSTEM_PROMPT}]
+ [{"role": m.type if hasattr(m, "type") else "user", "content": m.content} for m in messages],
"temperature": temperature,
}
async with httpx.AsyncClient(timeout=60) as client:
resp = await client.post(
"https://openrouter.ai/api/v1/chat/completions",
headers=headers, json=body
)
resp.raise_for_status()
return resp.json()["choices"][0]["message"]["content"]
MEDIA TOOLS — ОБОГАЩЕНИЕ WIKI#
Почему так: Медиа-контент (YouTube лекции, статьи, TG посты) содержит ценные знания, которые без автоматической обработки теряются. Сохранение в vault + индексация в Vector DB делает эти знания доступными агентам через semantic search. Это замыкает цикл: пользователь делится ссылкой → агент извлекает знание → знание попадает в wiki → связывается с существующими страницами → остаётся доступным годами.
Процесс обогащения wiki через media_tools#
Пользователь: "добавь в вики https://youtube.com/watch?v=abc"
│
▼
MainAgent → media_tools.extract_text(url)
│
▼
YouTubeTranscriptApi → транскрипт (5000 chars)
│
▼
LLM summarize → структурированный конспект
│
▼
WikiTools.write("media/2026-05-27-{title}.md", content)
│
Frontmatter: Content:
--- # Заголовок видео
source: youtube ## Ключевые тезисы
url: https://... - тезис 1
date: 2026-05-27 - тезис 2
topics: [LangGraph, ## Цитаты
agents, memory] ## Связанные страницы
author: Karpathy [[Проекты/Феникс]] ← backlink
---
│
▼
SemanticMemory.store(doc_id, text, metadata) ← ChromaDB
│
▼
Агент сообщает: "Добавлено: vault/media/2026-05-27-LangGraph-intro.md"
Структура vault/media/#
vault/media/
2026-05-27-LangGraph-intro.md # YouTube конспект
2026-05-20-Karpathy-nanoGPT.md # YouTube конспект
articles/
2026-05-25-React-patterns.md # статья
telegram/
2026-05-26-anthropic-news.md # TG пост
ChromaDB коллекции для медиа#
# src/d_brain/memory/semantic.py — ChromaDB backend расширение
CHROMA_COLLECTIONS = {
"wiki_pages": "Страницы Obsidian vault (goals, projects, contacts)",
"sessions": "Эпизодические summary сессий (episodic → semantic)",
"media": "YouTube/статьи/TG контент (media_tools output)",
"skills": "Процедурные знания из ProceduralMemory",
}
# При записи медиа:
async def store_media(self, url: str, text: str, metadata: dict) -> None:
"""Индексирует медиа в ChromaDB collection 'media'."""
doc_id = f"media_{hash(url) & 0xFFFFFF:06x}"
await self.store(doc_id, text[:5000], {
**metadata,
"collection": "media",
"indexed_at": utc_today(),
})
media_tools — расширение для wiki enrichment#
# src/d_brain/tools/media_tools.py — добавить метод
class MediaTools:
async def extract_and_save_to_wiki(
self,
url: str,
wiki: "WikiTools",
memory: "SemanticMemory",
llm_summarize_fn,
) -> str:
"""Извлекает текст → резюмирует → сохраняет в wiki → индексирует.
Returns: путь к созданной заметке в vault
"""
# 1. Извлекаем текст
raw_text = await self.extract_text(url)
# 2. Определяем тип и метаданные
source_type = "youtube" if self._is_youtube(url) else (
"telegram" if "t.me/" in url else "article"
)
# 3. LLM создаёт структурированный конспект
prompt = f"""Создай конспект для Obsidian vault.
URL: {url}
Тип: {source_type}
Текст: {raw_text[:3000]}
Формат:
---
source: {source_type}
url: {url}
date: {utc_today()}
topics: [список тем через запятую]
---
# Заголовок (кратко)
## Ключевые тезисы
- тезис 1
- тезис 2
## Применение
Как это связано с моей работой.
## Связанные страницы
[[...]]"""
content = await llm_summarize_fn(prompt)
# 4. Генерируем путь в vault
from d_brain.core.time_utils import utc_today
date = utc_today()
# Простой slug из URL
import re
slug = re.sub(r"[^a-zA-Z0-9а-яА-Я]+", "-", url.split("/")[-1])[:50]
folder = "media/telegram" if "t.me/" in url else (
"media" if not self._is_youtube(url) else "media"
)
vault_path = f"{folder}/{date}-{slug}.md"
# 5. Сохраняем в vault
await wiki.write(vault_path, content)
# 6. Индексируем в Vector DB
await memory.store_media(url, raw_text, {
"source_type": source_type,
"vault_path": vault_path,
"date": date,
})
return vault_path
SOUL.MD — ЛИЧНОСТЬ КАЖДОГО АГЕНТА#
Почему так: Soul.md решает проблему «раздвоения личности» агентов — без явного определения scope и эскалации каждый агент пытается ответить на все вопросы, теряя специализацию. Soul.md — это декларативный контракт: что агент делает, что НЕ делает, как реагирует на пограничные случаи. Загружается как часть системного промпта при инициализации агента.
MainAgent Soul.md#
<!-- vault/system/souls/main.md — загружается в BaseLLMAgent SYSTEM_PROMPT -->
# MainAgent Soul
## Личность
Универсальный персональный ассистент. Нейтральный, профессиональный, точный.
Не угадывает намерения — уточняет если неоднозначно.
Предпочитает краткость над многословием.
## Ценности
- Точность над скоростью: лучше уточнить, чем дать неверный ответ
- Явное над неявным: озвучивать допущения
- Действие над обсуждением: предлагать конкретный следующий шаг
## Scope
ДА:
- Любые запросы пользователя (роутинг к специалистам если нужно)
- Wiki поиск и обновление
- Веб-поиск и извлечение медиа
- Закрытие дня (daily close)
- Управление скиллами (/skill install/remove/list)
- Переключение LLM модели (/model <model_name> — например /model claude-opus-4 или /model deepseek/deepseek-v3)
НЕТ (эскалировать):
- Вопросы о личных целях и прогрессе → CoachAgent
- Трекинг привычек → HabitsAgent
- Создание/управление задачами YouGile → PlannerAgent
## Формат ответа
- Telegram HTML (<b>, <i>, <code>)
- Прямой ответ без вступлений ("Конечно!" / "Отлично!")
- Список шагов если задача многоэтапная
- Ссылка на wiki страницу если релевантна
## Эскалация
if "цель" OR "прогресс" OR "что я делаю" → route to CoachAgent
if "привычка" OR "streak" OR "выполнил" → route to HabitsAgent
if "задача" OR "YouGile" OR "дедлайн" OR "план дня" → route to PlannerAgent
CoachAgent Soul.md#
<!-- vault/system/souls/coach.md -->
# CoachAgent Soul
## Личность
Требовательный, поддерживающий, прямой коуч. Задаёт сложные вопросы.
Не даёт советов которые пользователь не просил. Опирается на данные и факты.
Действует как зеркало: отражает что есть, не приукрашивает.
## Ценности
- Долгосрочное мышление: один правильный вопрос важнее десяти советов
- Конкретные действия: каждый разговор заканчивается следующим шагом
- Факты и данные: completion rate, streak, habit adherence — не ощущения
## Scope
ДА:
- Цели (годовые, квартальные, личные)
- Прогресс по целям (goal health check)
- Привычки в контексте целей
- Постмитинг рефлексия
- Недельный review
- Проактивные nudge если пользователь отстаёт
НЕТ (эскалировать):
- Технические вопросы → MainAgent
- Финансовые вопросы → MainAgent
- Создание конкретных задач YouGile → PlannerAgent
- Детальный habit tracking → HabitsAgent
## Формат ответа
1. Что данные говорят (факт)
2. Что это значит (паттерн/инсайт)
3. Что усилить (конкретное действие)
4. Следующий шаг (один конкретный)
## Эскалация
if "создай задачу" → PlannerAgent
if "отметь привычку" → HabitsAgent
if "найди информацию" → MainAgent
HabitsAgent Soul.md#
<!-- vault/system/souls/habits.md -->
# HabitsAgent Soul
## Личность
Структурированный, методичный. Фокус на паттернах и стриках.
Не мотивирует словами — показывает данные. Краткие ответы.
## Ценности
- Постоянство над интенсивностью: 10 минут каждый день > 2 часа раз в неделю
- Система над мотивацией: правильные триггеры и якоря важнее силы воли
- Измеримость: если нельзя измерить — это не привычка, это намерение
## Scope
ДА:
- Отметка выполнения привычки (log_habit)
- Текущий streak и статистика adherence
- Паттерны: дни недели, время суток, корреляции
- Определение новых привычек (habit definition)
- Напоминания (по расписанию)
НЕТ (эскалировать):
- Разовые задачи → PlannerAgent
- Цели высокого уровня → CoachAgent
- Почему привычка важна для целей → CoachAgent
## Формат ответа
данные → паттерн → инсайт → рекомендация
(без лишних слов, максимум 4 строки)
## Эскалация
if "цель" OR "зачем мне это" → CoachAgent
if "задача" OR "YouGile" → PlannerAgent
PlannerAgent Soul.md#
<!-- vault/system/souls/planner.md -->
# PlannerAgent Soul
## Личность
Прагматичный, time-aware, фокус на исполнении и дедлайнах.
Не философствует — планирует и делает. Реалистичен в оценках.
## Ценности
- Конкретные дедлайны: задача без дедлайна — это идея, не задача
- Реалистичное планирование: лучше 3 задачи done, чем 10 planned
- Carry-over transparency: всегда явно показывать что перенесено и почему
## Scope
ДА:
- YouGile задачи (создание, закрытие, перенос)
- Google Calendar + Yandex CalDAV
- План дня (daily plan generation)
- Еженедельный план
- Carry-over механизм
- Закрытие дня (close_day)
НЕТ (эскалировать):
- Цели высокого уровня → CoachAgent
- Привычки → HabitsAgent
- Обсуждение зачем делать задачу → CoachAgent
## Формат ответа
приоритет → дедлайн → следующий шаг → что перенести
(структурированный список, не эссе)
## Эскалация
if "цель" OR "смысл" OR "зачем" → CoachAgent
if "привычка" → HabitsAgent
Загрузка Soul.md в агентов#
# src/d_brain/agents/base.py — дополнение
class BaseLLMAgent:
SOUL_PATH: str = "" # переопределить в каждом агенте
def __init__(self, data_dir, vault_path="~/vault"):
super_init(...)
self.SYSTEM_PROMPT = self._load_soul() or self.SYSTEM_PROMPT
def _load_soul(self) -> str | None:
"""Загружает Soul.md если путь указан и файл существует."""
if not self.SOUL_PATH:
return None
soul_path = Path(os.getenv("VAULT_PATH", "~/vault")).expanduser() / self.SOUL_PATH
if soul_path.exists():
content = soul_path.read_text(encoding="utf-8")
# Убираем frontmatter если есть
if content.startswith("<!--"):
content = content.split("-->", 1)[-1].strip()
return content
return None
class CoachAgent(BaseLLMAgent):
SOUL_PATH = "system/souls/coach.md"
# SYSTEM_PROMPT будет загружен из Soul.md автоматически
ОСНОВНЫЕ ПРОЦЕССЫ В BPMN 2.0#
Почему так: BPMN 2.0 дает исполняемую спецификацию процессов — не просто описание «что делать», а точную последовательность шагов, акторов, exception flows и output артефактов. 8 процессов покрывают весь операционный цикл: планирование, выполнение, закрытие, встречи, знания. Каждый процесс может быть автоматически задокументирован и протестирован.
Процесс 1: Планирование месяца и недели#
BPMN: 01-month-week-planning.bpmn
Триггер: начало месяца (1-е число) или начало недели (понедельник)
Акторы: CoachAgent (оркестратор) + PlannerAgent
Основной поток:
1. CoachAgent.review_goals(period_id)
→ загружает цели из vault/goals/quarters/YYYY-QN.md
2. CoachAgent.generate_commitments(goals, period)
→ разбивает quarterly OKR на month/week commitments
3. PlannerAgent.sync_yougile_tasks(period)
→ получает задачи из YouGile для периода
4. PlannerAgent.merge_calendar(period)
→ добавляет встречи из Google/Yandex calendar
5. CoachAgent.confirm_plan(user_id, draft_plan)
→ отправляет черновик пользователю в Telegram
6. User confirms (кнопки ✅/❌/✏️)
7. PlannerAgent.save_plan(period_id, confirmed_plan)
→ vault/goals/months/YYYY-MM.md или weeks/YYYY-Www.md
8. CascadeDB.add_plan_items(plan_items)
→ создаёт PlanItem записи для roll-up
Исключения:
- YouGile недоступен → планируем по wiki только
- Calendar пуст → предупреждение, продолжаем
- Пользователь не отвечает 24ч → сохраняем черновик
Output:
- vault/goals/months/YYYY-MM.md (или weeks/YYYY-Www.md)
- PlanItem записи в CascadeDB
- Telegram сообщение пользователю с планом
Процесс 2: Планирование дня#
BPMN: 02-day-planning.bpmn
Триггер: утро (cron 07:00) или пользователь /план_дня
Акторы: PlannerAgent
Основной поток:
1. PlannerAgent.get_weekly_commitments(current_week)
→ commitments из vault/goals/weeks/YYYY-Www.md
2. PlannerAgent.get_habits_today(date)
→ запрос к HabitsAgent: какие привычки сегодня
3. PlannerAgent.get_calendar_today(date)
→ встречи из Google/Yandex Cal
4. PlannerAgent.get_yougile_due_today(date)
→ задачи с дедлайном сегодня
5. PlannerAgent.generate_day_plan(all_data)
→ LLM генерирует расставленный по времени план
6. Telegram: отправляем план пользователю (кнопки ✅/✏️)
7. PlannerAgent.save_day_plan(date, plan)
→ vault/daily/YYYY-MM-DD.md (раздел "План дня")
Исключения:
- Нет встреч в календаре → план только по задачам
- YouGile недоступен → план по weekly commitments
Output:
- vault/daily/YYYY-MM-DD.md (блок "План дня")
- Telegram сообщение с планом
Процесс 3: Выполнение дня#
BPMN: 03-day-execution.bpmn
Триггер: пользователь отмечает выполнение задачи или привычки
Акторы: PlannerAgent (задачи) + HabitsAgent (привычки)
Основной поток (задача):
1. User: "сделал задачу X" | "закрыл задачу в YouGile"
2. PlannerAgent.mark_task_done(task_id или title)
→ YouGile API: переводим задачу в Done колонку
3. PlannerAgent.update_daily_md(task_id, done=True)
→ обновляем vault/daily/YYYY-MM-DD.md
4. CascadeDB.add_fact_event(COMPLETED, plan_item_id)
→ append-only event ledger
5. Telegram: "✅ Задача X закрыта"
Основной поток (привычка):
1. User: "бег 5км" | "сделал тренировку"
2. HabitsAgent.log_habit(habit, value, date)
→ SQLite habit_log
3. HabitsAgent.update_streak(habit)
→ обновляем streak счётчик
4. CascadeDB.add_fact_event(HABIT_CHECKED, ...)
→ event ledger
5. Telegram: "✅ бег: 5км. Streak: 7 дней 🔥"
**Исключения:**
- Задача не найдена в YouGile → предложить создать новую или пометить как внеплановую
- YouGile API недоступен → сохранить локально в daily.md, синхронизировать при восстановлении
- Пользователь отмечает задачу из другого дня → спросить подтверждение, обновить carry-over
Output:
- Обновлённый vault/daily/YYYY-MM-DD.md
- FactEvent в CascadeDB
- SQLite habit_log запись
Процесс 4: Корректировка планов#
BPMN: 04-plan-adjustment.bpmn
Триггер: пользователь "перенеси/отмени/сдвинь задачу X"
Акторы: PlannerAgent
Основной поток:
1. PlannerAgent.parse_intent(text)
→ тип: перенести | отменить | изменить дедлайн
2. PlannerAgent.find_task(query)
→ поиск в YouGile + daily.md
3. Если перенести:
a. YouGile API: изменяем дедлайн
b. carryover.run_carryover(from_period, to_period)
c. FactEvent(MOVED) в CascadeDB
4. Если отменить:
a. YouGile API: переводим в архив/отмена
b. FactEvent(CANCELED) в CascadeDB
5. PlannerAgent.update_daily_md(task_id, new_status)
6. Telegram: подтверждение с деталями переноса
Исключения:
- Задача не найдена → уточняем у пользователя
- YouGile API ошибка → retry(3) → сообщить об ошибке
Output:
- Обновлённые YouGile задачи
- FactEvent(MOVED/CANCELED) в CascadeDB
- Обновлённый daily.md с пометкой переноса
Процесс 5: План/Факт дня (Закрытие дня)#
BPMN: 05-day-close.bpmn
Триггер: вечер (cron 21:00) или пользователь /закрой_день
Акторы: PlannerAgent + HabitsAgent
Основной поток:
1. PlannerAgent.get_done_tasks_today()
→ YouGile Done колонки (все проекты, UTC дата, пагинация)
2. HabitsAgent.get_today_log(user_id, date)
→ привычки выполненные сегодня
3. PlannerAgent.get_moved_tasks()
→ задачи перенесённые сегодня
4. WikiTools.search_updated_since(today_start_ts)
→ wiki страницы обновлённые сегодня (по mtime)
5. LLM: plan vs fact анализ + рефлексия
6. PlannerAgent.save_close_day(date, summary)
→ vault/daily/YYYY-MM-DD.md (блок "Закрытие дня")
7. CascadeDB.add_fact_event(REVIEW_CLOSED, period_id=today)
8. Telegram: отчёт пользователю
Исключения:
- YouGile недоступен → close только по wiki + habits
Output:
- vault/daily/YYYY-MM-DD.md (блок "Закрытие дня")
- FactEvent(REVIEW_CLOSED) в CascadeDB
- Telegram отчёт с plan/fact сравнением
Процесс 6: Встречи в календаре#
BPMN: 06-calendar-meetings.bpmn
Триггер: событие в Google/Yandex Calendar (15 мин до) или postmeet запрос
Акторы: CoachAgent (coaching) + MainAgent (wiki)
Pre-meeting поток:
1. MainAgent.get_calendar_event(event_id)
2. WikiTools.search(contact_name OR company_name)
→ контекст о участниках из wiki
3. Telegram: "Через 15 мин встреча с X. Контекст: ..."
Post-meeting поток:
1. User: "встреча с X закончилась" OR автоматически после event_end
2. CoachAgent: запрашивает у пользователя итоги встречи
3. User: диктует итоги (текст или голос)
4. MainAgent.save_meeting_notes(date, participants, notes)
→ vault/meetings/YYYY-MM-DD-{participant}.md
5. WikiTools.update_contact(contact_name, meeting_ref)
→ обновляем страницу контакта в wiki
6. Если действия/задачи → PlannerAgent.create_yougile_tasks(action_items)
7. CascadeDB.add_fact_event(COMPLETED, period_id=today)
→ встреча как выполненный fact
**Исключения:**
- Событие отменено → уведомить, удалить briefing из очереди, предложить освободившееся время
- Транскрипция/запись недоступна → создать пустой шаблон постмитинга, напомнить заполнить вручную
- Участник из calendar не найден в vault/people/ → предложить создать контакт
Output:
- vault/meetings/YYYY-MM-DD-{participant}.md
- Обновлённая страница контакта в wiki
- YouGile задачи из action items встречи
Процесс 7: Информация по работе и проектам#
BPMN: 07-project-info.bpmn
Триггер: пользователь спрашивает о проекте или контакте
Акторы: MainAgent
Основной поток:
1. MainAgent.parse_intent(text)
→ entity: проект / контакт / задача
2. WikiTools.search(entity_name)
→ поиск в vault/projects/ и vault/contacts/
3. PlannerAgent.get_yougile_tasks(project_id или alias)
→ активные задачи по проекту
4. EpisodicMemory.search(entity_name)
→ что обсуждалось о проекте ранее
5. LLM: синтезирует ответ из всех источников
6. Telegram: структурированный ответ (HTML)
Исключения:
- Не найдено в wiki → предложить создать страницу
- YouGile недоступен → ответ только по wiki
Output:
- Telegram ответ с информацией о проекте/контакте
- (опционально) VAULT_FILE_DRAFT для новой страницы
Процесс 8: Актуализация задач и вики#
BPMN: 08-wiki-task-update.bpmn
Триггер: конец недели (cron воскресенье 20:00) или /обнови_вики
Акторы: CoachAgent + PlannerAgent
Основной поток:
1. PlannerAgent.get_all_yougile_tasks(status="active")
→ все активные задачи из всех проектов
2. WikiTools.find_project_pages()
→ страницы vault/projects/**/*.md
3. Для каждого проекта:
a. Найти связанные YouGile задачи
b. LLM: обновить статус задач в wiki странице
c. Добавить ссылки на новые задачи
d. WikiTools.write(project_path, updated_content)
4. CoachAgent.link_tasks_to_goals(tasks, goals)
→ привязать задачи к целям в CascadeDB (source_id)
5. CascadeDB.add_fact_event(REVIEW_CLOSED, period_id=current_week)
6. Telegram: "Вики обновлена: N страниц, M задач привязано к целям"
Исключения:
- YouGile API ошибка → retry(3) → skip проект → продолжаем
- Wiki страница не найдена для проекта → VAULT_FILE_DRAFT (создать)
Output:
- Обновлённые vault/projects/**/*.md
- Связи task → goal в CascadeDB
- Telegram отчёт об обновлении
БАЗА ЗНАНИЙ И ВЕКТОРНЫЙ ПОИСК#
Почему так: Vector DB — это не опциональный компонент, а фундамент для семантического поиска. Без embeddings поиск по vault работает только по точным ключевым словам (FTS5). С ChromaDB агент находит релевантный контент даже если пользователь использует другие слова. Консолидатор поддерживает свежесть индекса — без него embeddings устаревают.
Структура Obsidian Vault (полная)#
vault/
goals/ # каскадное планирование (см. выше)
projects/ # CRM проекты (клиенты, сделки)
_meta/
yougile-project-stickers.yaml # маппинг проектов → YouGile stickers
people/ # контакты (альтернатива contacts/)
contacts/ # CRM контакты
daily/ # ежедневные заметки YYYY-MM-DD.md
weekly/ # compat: 3-weekly.md (§4.4 migration)
monthly/ # compat: 2-monthly.md (§4.4 migration)
habits/
definitions.md # список привычек с описаниями
rituals.md # утренние/вечерние ритуалы
media/ # YouTube/статьи → markdown конспекты
skills/ # процедурные знания (как что делать)
meetings/ # конспекты встреч
references/ # справочные материалы
thoughts/ # личные мысли и размышления
system/
souls/ # Soul.md файлы агентов
planning/ # event ledger + snapshots
kaizen-log.md # лог самоулучшений
.claude/
CLAUDE.md # инструкции для Claude subprocess
Стратегия embeddings#
# src/d_brain/memory/semantic.py — ChromaDB backend
EMBEDDING_MODEL = "text-embedding-3-small" # OpenAI, 1536 dims, $0.02/1M tokens
class SemanticMemory:
"""ChromaDB backend с автоматической re-embed при изменении файлов."""
CHROMA_COLLECTIONS = {
"wiki_pages": "Страницы Obsidian vault",
"sessions": "Эпизодические summaries",
"media": "Медиа контент (YouTube, статьи, TG)",
"skills": "Процедурные знания",
}
async def index_vault_page(self, vault_path: str, content: str) -> None:
"""Индексирует или переиндексирует страницу vault в ChromaDB.
Вызывается:
- при WikiTools.write() (on every update)
- при ночной консолидации (full re-index если mtime изменился)
"""
doc_id = f"wiki_{hash(vault_path) & 0xFFFFFF:06x}"
metadata = {
"vault_path": vault_path,
"indexed_at": utc_today(),
"collection": "wiki_pages",
}
# Парсим frontmatter для дополнительных метаданных
if content.startswith("---"):
try:
import yaml
fm_end = content.find("---", 3)
if fm_end > 0:
fm = yaml.safe_load(content[3:fm_end])
if isinstance(fm, dict):
metadata.update({
"title": fm.get("title", ""),
"tags": ",".join(fm.get("tags", [])),
"status": fm.get("status", ""),
})
except Exception:
pass
await self.store(doc_id, content[:5000], metadata)
async def suggest_backlinks(self, content: str, top_k: int = 5) -> list[str]:
"""Предлагает Obsidian backlinks на основе semantic similarity.
Используется при создании новых страниц для автоматического
создания [[backlinks]] к релевантным страницам.
"""
results = await self.search(content[:500], top_k=top_k)
return [r["metadata"].get("vault_path", "") for r in results if r.get("metadata")]
Consolidator + Vector DB ночной индекс#
# src/d_brain/memory/consolidator.py — РАСШИРЕНИЕ
class MemoryConsolidator:
async def reindex_vault(self, vault_path: str) -> int:
"""Ночной re-index всего vault в ChromaDB.
Индексирует только страницы с mtime после last_indexed_at.
Returns: количество переиндексированных страниц.
"""
import os
from pathlib import Path
vault = Path(vault_path).expanduser()
count = 0
for md_file in vault.rglob("*.md"):
try:
mtime = os.stat(md_file).st_mtime
rel_path = str(md_file.relative_to(vault))
content = md_file.read_text(encoding="utf-8")
await self.semantic.index_vault_page(rel_path, content)
count += 1
except Exception:
continue
return count
async def run_nightly(self, user_id: int) -> dict:
"""Полный ночной цикл: консолидация + re-index."""
results = {}
# 1. Консолидация эпизодической памяти
results["daily_summary"] = await self.run_daily(user_id=user_id)
# 2. Re-index vault в ChromaDB
vault_path = os.getenv("VAULT_PATH", "~/vault")
results["reindexed_pages"] = await self.reindex_vault(vault_path)
# 3. Kaizen ревью (см. Task 13)
results["kaizen_entry"] = await self._kaizen_review(user_id)
return results
MEMORY CONSOLIDATION И KAIZEN САМОУЛУЧШЕНИЕ#
Почему так: Kaizen (改善) — принцип непрерывного малого улучшения — идеален для AI-систем. Вместо редких крупных апгрейдов система ежедневно задаёт себе 3 вопроса и делает одно маленькое улучшение. Накопленные за год улучшения трансформируют систему полностью. Kaizen-log обеспечивает прозрачность и accountability: каждое изменение задокументировано с датой и обоснованием.
Kaizen принцип#
KAIZEN ЦИКЛ (ежедневно, ночью):
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
1. НАБЛЮДЕНИЕ (что произошло сегодня?)
┌─────────────────────────────────────┐
│ Сколько запросов обработано? │
│ Какие запросы вызвали уточнения? │
│ Где пользователь переспрашивал? │
│ Какие tools вызывались чаще всего? │
└─────────────────────────────────────┘
2. АНАЛИЗ (что не работает идеально?)
┌─────────────────────────────────────┐
│ Медленные ответы? (latency > 3s) │
│ Неточные ответы? (переспросы) │
│ Ненужные уточнения? │
│ Пропущенные контексты? │
└─────────────────────────────────────┘
3. УЛУЧШЕНИЕ (одно конкретное действие)
┌─────────────────────────────────────┐
│ Обновить Soul.md с новым правилом │
│ Добавить процедуру в ProceduralMem │
│ Уточнить системный промпт │
│ Предложить новый skill │
└─────────────────────────────────────┘
4. ЗАПИСЬ (append-only в kaizen-log.md)
┌─────────────────────────────────────┐
│ ## 2026-05-27 │
│ **Наблюдение:** ... │
│ **Проблема:** ... │
│ **Улучшение:** ... │
│ **Результат:** (следующий день) │
└─────────────────────────────────────┘
Еженедельные вопросы самоулучшения#
CoachAgent инициирует еженедельный AI self-review каждое воскресенье:
WEEKLY_SELF_REVIEW_QUESTIONS = [
"Как ещё я могу быть полезен пользователю на этой неделе?",
"Чего мне не хватает чтобы быть максимально полезным?",
"Какие варианты улучшения я могу предложить пользователю?",
"Где пользователь чаще всего испытывает трудности?",
"Какие задачи я выполняю медленно или некачественно?",
"Какие новые скиллы мне стоит предложить установить?",
"Что из Soul.md устарело или требует уточнения?",
]
KaizenReviewer — реализация#
# src/d_brain/agents/reviewer.py — РАСШИРЕНИЕ
class KaizenReviewer(SelfReviewer):
"""Kaizen: ежедневный анализ + еженедельный self-review."""
KAIZEN_LOG_PATH = "system/kaizen-log.md"
async def daily_kaizen(self, user_id: int, session_stats: dict) -> str:
"""Ежедневный Kaizen цикл: наблюдение → анализ → улучшение → запись.
session_stats: {
"total_requests": int,
"clarification_requests": int, # сколько раз уточняли
"avg_latency_ms": float,
"tools_called": dict[str, int], # tool_name: count
"failed_requests": int,
}
"""
prompt = f"""Kaizen ежедневный анализ.
Статистика за сегодня:
- Запросов обработано: {session_stats.get('total_requests', 0)}
- Запросов с уточнением: {session_stats.get('clarification_requests', 0)}
- Средняя задержка: {session_stats.get('avg_latency_ms', 0):.0f}мс
- Часто используемые tools: {session_stats.get('tools_called', {})}
- Ошибки: {session_stats.get('failed_requests', 0)}
Ответь по структуре:
НАБЛЮДЕНИЕ: [одно ключевое наблюдение]
ПРОБЛЕМА: [что мешает быть максимально полезным]
УЛУЧШЕНИЕ: [одно конкретное действие которое можно сделать]
ПРИОРИТЕТ: [high/medium/low]"""
from langchain_core.messages import HumanMessage
analysis = await self._call_llm([HumanMessage(content=prompt)])
# Записываем в kaizen-log.md (append-only)
await self._append_kaizen_log(analysis, session_stats)
return analysis
async def weekly_self_review(self, user_id: int) -> str:
"""Еженедельный глубокий self-review с вопросами из WEEKLY_SELF_REVIEW_QUESTIONS."""
# Читаем последние 7 записей из kaizen-log.md для контекста
from d_brain.tools.wiki_tools import WikiTools
import os
wiki = WikiTools(vault_path=os.getenv("VAULT_PATH", "~/vault"))
recent_kaizen = await wiki.read(self.KAIZEN_LOG_PATH)
questions_text = "\n".join(f"- {q}" for q in WEEKLY_SELF_REVIEW_QUESTIONS)
prompt = f"""Еженедельный self-review ассистента.
Последние Kaizen записи:
{recent_kaizen[-2000:] if len(recent_kaizen) > 2000 else recent_kaizen}
Ответь на вопросы:
{questions_text}
По каждому вопросу: конкретный ответ с примерами из реальной работы.
В конце: топ-3 предложения по улучшению для пользователя (что добавить/изменить)."""
from langchain_core.messages import HumanMessage
return await self._call_llm([HumanMessage(content=prompt)])
async def _append_kaizen_log(self, analysis: str, stats: dict) -> None:
"""Добавляет запись в vault/system/kaizen-log.md (append-only)."""
from d_brain.core.time_utils import utc_today
from d_brain.tools.wiki_tools import WikiTools
import os
wiki = WikiTools(vault_path=os.getenv("VAULT_PATH", "~/vault"))
today = utc_today()
new_entry = f"""
## {today}
{analysis}
**Статистика:** запросов={stats.get('total_requests',0)}, уточнений={stats.get('clarification_requests',0)}, задержка={stats.get('avg_latency_ms',0):.0f}мс
---
"""
# Читаем текущее содержимое
existing = await wiki.read(self.KAIZEN_LOG_PATH)
if existing.startswith("Заметка"):
# Файл не существует — создаём
header = "# Kaizen Log — Second Brain v2\n\nAppend-only лог системных улучшений.\n\n"
await wiki.write(self.KAIZEN_LOG_PATH, header + new_entry)
else:
await wiki.write(self.KAIZEN_LOG_PATH, existing + new_entry)
Outputs Kaizen-процесса#
| Тип улучшения | Артефакт | Пример |
|---|---|---|
| Новый skill | VAULT_FILE_DRAFT → vault/skills/new-skill.md | "Скилл для поиска в CRM" |
| Soul.md правило | Изменение в vault/system/souls/*.md | "Добавить правило эскалации" |
| Memory rule | Запись в ProceduralMemory | "Уточнять формат дедлайна" |
| Prompt улучшение | Обновление в agents/base.py | "Добавить контекст времени" |
| Kaizen log | Append в vault/system/kaizen-log.md | Каждый день |
Cron расписание всех автоматических процессов#
# src/d_brain/core/cron.py — планировщик автоматических задач
CRON_SCHEDULE = {
# Ежедневные
"07:00": "PlannerAgent.generate_day_plan()", # план дня
"21:00": "PlannerAgent.close_day()", # закрытие дня
"22:00": "MemoryConsolidator.run_nightly()", # консолидация + kaizen
"22:30": "WikiTools.reindex_vault()", # re-index ChromaDB
# Еженедельные
"пн 08:00": "CoachAgent.weekly_coaching_review()", # weekly review
"пт 17:00": "CoachAgent.goal_health_check()", # goal health check
"пт 17:30": "CoachAgent.proactive_nudge()", # nudge если отстаём
"вс 20:00": "PlannerAgent.update_wiki_tasks()", # актуализация вики
"вс 21:00": "KaizenReviewer.weekly_self_review()", # weekly self-review
# Ежемесячные (1-е число)
"01 09:00": "CoachAgent.monthly_planning()", # планирование месяца
# Ежеквартальные
"квартал 09:00": "CoachAgent.quarterly_planning()", # OKR на квартал
}
ИТОГ — ДОПОЛНЕНИЯ ROUND 2#
Дата: 2026-05-27
Ревьюер: Senior AI Systems Architect
Round 2 задачи: 13
Добавленные секции#
| Task | Секция | Объём |
|---|---|---|
| T1 | Блокквоты к 7 фазам (почему/пример) | ~350 строк |
| T2 | Фаза 4.5: Каскадное планирование + CascadeDB + carryover + rollup | ~200 строк кода |
| T3 | Диаграмма каскада (ASCII) + vault структура | ~60 строк |
| T4 | CoachAgent оркестратор: goal_health_check + proactive_nudge + weekly_review | ~120 строк кода |
| T5 | Архитектурная диаграмма v2 (отдельный HTML файл) | см. architecture-v2-updated.html |
| T6 | Полная структура проекта second-brain-v2/ с описаниями | ~100 строк |
| T7 | Стратегия создания нового проекта + чеклист миграции | ~60 строк |
| T8 | Выбор модели и провайдера + таблица стоимости + OpenRouter | ~80 строк |
| T9 | media_tools wiki enrichment + ChromaDB коллекции + vault/media структура | ~120 строк кода |
| T10 | Soul.md для 4 агентов + механизм загрузки | ~150 строк |
| T11 | 8 BPMN 2.0 процессов с trigger/actors/flow/exceptions/output | ~200 строк |
| T12 | База знаний + embedding стратегия + ChromaDB + backlink suggestion | ~120 строк кода |
| T13 | Kaizen методология + KaizenReviewer + weekly questions + cron schedule | ~150 строк кода |
Ключевые архитектурные решения Round 2#
- CoachAgent как оркестратор — владеет целями, делегирует HabitsAgent и PlannerAgent как LangGraph subgraphs
- CascadeDB — SQLite хранилище для канонической модели периодов/целей/планов/фактов из §6
- Append-only event ledger — FactEvents никогда не изменяются, только добавляются (аудит-трейл)
- Soul.md загружается из vault — личность агентов хранится в Obsidian, версионируется через git
- Multi-provider LLM — Anthropic для reasoning, DeepSeek для bulk tasks (10x дешевле)
- Kaizen loop — ежедневный self-review с записью в append-only kaizen-log.md
- ChromaDB MANDATORY — векторный поиск обязателен; FTS5 используется ТОЛЬКО как secondary keyword index в SQLite episodic memory, НЕ как замена ChromaDB
Общая оценка плана после Round 2: 9.2/10
Обновлённая оценка времени: 22 рабочих дня (было 18 → +4 на каскадное планирование и CoachAgent оркестрацию)