Files
beaver-agent/config.py
T

881 lines
36 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import logging
import os
from datetime import UTC, date, datetime, timedelta
from pathlib import Path
from typing import Any
from zoneinfo import ZoneInfo
from beaver_gateway.agents.base import BaseAgent, ExposedMcp
from beaver_gateway.agents.claude import ClaudeAgent, ClaudeOptions, Prompts
from beaver_gateway.agents.raycast import RaycastAgent, RemoteTool, UserPreferences
from beaver_gateway.core.conversations import (
ConversationTexts,
NewDayContext,
SeedContext,
)
from beaver_gateway.core.distill import DistillContext, Distiller, LineCap
from beaver_gateway.core.injects import INTERRUPTED_TURN, InjectContext
from beaver_gateway.core.prompt import assemble
from beaver_gateway.core.registry import Gateway
from beaver_gateway.core.rotation import HandoutContext, RotationPolicy
from beaver_gateway.core.scheduler import Budget, Job, JobRun
from beaver_gateway.core.turn_record import slugify
from beaver_gateway.core.watch import VaultWatch, WatchRules
from beaver_gateway.frontends.admin import AdminFrontend
from beaver_gateway.frontends.anthropic import AnthropicMessagesFrontend
from beaver_gateway.frontends.api import ApiFrontend
from beaver_gateway.frontends.base import Frontend
from beaver_gateway.frontends.markdown import MarkdownFrontend
from beaver_gateway.frontends.mcp_server import McpServerFrontend
from beaver_gateway.frontends.telegram import TelegramFrontend
from beaver_gateway.mcp.types import HttpMcp, McpServer, McpServerT
from curator import Watched, briefing
from mcps.homeassistant import HomeAssistant
from mcps.komodo import Komodo
from mcps.komodo_alerts import HOLD, AlertMemory, describe, parse, still_bad
from mcps.vibegram import Vibegram
from policy import (
DEEP_DISALLOWED,
DISPATCHER_DISALLOWED,
DISTILLER_DISALLOWED,
TRIAGE_DISALLOWED,
Zones,
bash_zones,
requires_skill,
skill_tracker,
vault_zones,
)
from recall import Recall, ReplyLog, notes_head
TZ = "Europe/Warsaw"
VAULT = Path("/vault")
CHATS_DIR = VAULT / "💬 чаты"
DIARY = VAULT / "📅 дни"
BEAVER = VAULT / "мета" / "бобер"
PROMPTS = BEAVER / "промпты"
GRANULES = PROMPTS / "гранулы"
SKILLS = BEAVER / "скиллы"
DAYS = BEAVER / "дни"
DIGESTS = BEAVER / "выжимки"
INDEX = BEAVER / "индекс.md"
NOTES = BEAVER / "наблюдения"
REPLIES = BEAVER / "реплики"
CURATOR_JOURNAL = BEAVER / "куратор.md"
PORTRAIT = NOTES / "Бобёр - наблюдения.md"
PEOPLE = VAULT / "👤 люди"
BOARDS = VAULT / "📆 доски"
STATE = BEAVER / "состояние.md"
STATE_MAX_LINES = 60
# §4.5: ночное закрытие трогает только чаты с активностью после запуска.
LAUNCH = datetime(2026, 8, 28, tzinfo=ZoneInfo(TZ))
_log = logging.getLogger("beaver_agent.config")
# Гранулы - чистый markdown; xml-теги, в которые они заворачиваются,
# живут здесь. Голос - папка, файл на тег, в этом порядке.
VOICE = tuple(
(tag, GRANULES / "голос" / f"{tag}.md")
for tag in (
"role",
"philosophy",
"user_profile",
"operating_modes",
"how_you_operate",
"interaction_guidelines",
)
)
PROFILE = ("profile", PROMPTS / "профиль.md")
CORRECTIONS = ("corrections", PROMPTS / "поправки.md")
VAULT_MAP = ("vault_map", GRANULES / "карта-vault.md")
def environment(kind: str) -> tuple[str, Path]:
return ("environment", GRANULES / "окружения" / f"{kind}.md")
def overlay(tag: str, name: str) -> tuple[str, Path]:
return (tag, GRANULES / f"{name}.md")
# §4.2: сборки промптов. Окружение kind'а - константная гранула, экземплярное
# (файл, топик) едет первым сообщением. Какие kind агент обслуживает, следует
# из того, какие промпты у него заданы (§3.12).
_BASE = (*VOICE, PROFILE, CORRECTIONS, VAULT_MAP)
_DISTILL = (*_BASE, environment("джоб"), overlay("distiller", "дистиллятор"))
DEEP_PROMPTS = Prompts(
deep=(*_BASE, environment("глубокий"), overlay("deep", "глубокий"))
)
# Диспетчер: тот же голос, окружение по kind, оверлей по роли.
# Форк ветки (слив) думает как дистиллятор в окружении джоба.
DISPATCHER_PROMPTS = Prompts(
master=(*_BASE, environment("мастер"), overlay("dispatcher", "диспетчер")),
branch=(*_BASE, environment("ветка"), overlay("branch", "ветка")),
fork=_DISTILL,
job=_DISTILL,
)
# §4.1: дистиллятор без голоса - иначе выжимки с матом: профиль, окружение
# джоба, оверлей. Форк закрываемого чата и недельная чистка состояния.
_DISTILLER = (PROFILE, environment("джоб"), overlay("distiller", "дистиллятор"))
DISTILLER_PROMPTS = Prompts(fork=_DISTILLER, job=_DISTILLER)
# §4.1: триаж без голоса и без vault - профиль, окружение джоба, оверлей.
TRIAGE_PROMPTS = Prompts(job=(PROFILE, environment("джоб"), overlay("triage", "триаж")))
# Куратор памяти - тот же менеджер бобрения (голос и всё), в окружении джоба.
CURATOR_PROMPTS = Prompts(
job=(*_BASE, environment("джоб"), overlay("curator", "куратор"))
)
QUICK_PROMPT = assemble((*VOICE, PROFILE, overlay("quick", "быстрый")))
# §4.3: наборы скиллов = папки = плагины. Глубокие грузят общие + vault.
DEEP_SKILLS = (SKILLS / "общие", SKILLS / "vault")
DISPATCHER_SKILLS = (SKILLS / "общие", SKILLS / "диспетчер", SKILLS / "vault")
# §3.7 (решение h, 2026-08-30): vault - дом диспетчера. Создавать, править,
# переносить можно везде, кроме `.obsidian` и `мета` вне `бобер`; удалять -
# только своё (`мета/бобер`) и `💬 чаты`, существующие чаты не трогать.
# Дистиллятор остаётся в строгих зонах. firefly пишет только после скилла.
ZONES = Zones(
vault=VAULT,
write=(BEAVER,),
create=(CHATS_DIR,),
edit=(VAULT,),
protected=(VAULT / ".obsidian", VAULT / "мета"),
)
STRICT_ZONES = Zones(vault=VAULT, write=(BEAVER,), create=(CHATS_DIR,))
FIREFLY_WRITES = ("mcp__firefly__store_*", "mcp__firefly__update_*")
VAULT_POLICY = (
skill_tracker(),
vault_zones(ZONES),
bash_zones(ZONES),
requires_skill("firefly", FIREFLY_WRITES),
)
def chat_path(title: str, agent: str, vault: Path) -> Path: # noqa: ARG001
"""Новый файл чата в `💬 чаты/`: помесячная папка, дата и тема."""
today = date.today()
return (
vault / f"{today:%Y-%m}" / f"{today:%Y-%m-%d} - {slugify(title, maxlen=60)}.md"
)
def _calendar_mcps() -> list[HttpMcp]:
raw = os.environ.get("CALENDAR_MCPS", "").strip()
servers: list[HttpMcp] = []
for position, raw_entry in enumerate(raw.split(","), start=1):
entry = raw_entry.strip()
if not entry:
continue
name, sep, url = entry.partition("=")
if not sep or not url.strip():
# Только номер: в url лежит приватный ical-фид с токеном в пути,
# а traceback уезжает в docker logs целиком.
msg = f"CALENDAR_MCPS: запись #{position} не вида `name=url`"
raise ValueError(msg)
servers.append(McpServer.http(name=f"calendar-{name.strip()}", url=url.strip()))
return servers
calendar_mcps = _calendar_mcps()
calendar_exposed = tuple(ExposedMcp(name=m.name) for m in calendar_mcps)
# §4.4: komodo - python_tool с перечислимыми действиями, ключ остаётся в gateway.
KOMODO = (
Komodo(
url=os.environ["KOMODO_URL"],
key=os.environ["KOMODO_KEY"],
secret=os.environ["KOMODO_SECRET"],
)
if os.environ.get("KOMODO_URL")
else None
)
komodo_mcps = (
[McpServer.python_tool(name="komodo", tools=[KOMODO.komodo])] if KOMODO else []
)
komodo_exposed = (ExposedMcp(name="komodo"),) if KOMODO else ()
# §4.4: Home Assistant - python_tool, токен в gateway. HA живёт на хосте
# (network_mode: host), из контейнера - host.docker.internal.
HA = (
HomeAssistant(url=os.environ["HA_URL"], token=os.environ["HA_TOKEN"])
if os.environ.get("HA_URL") and os.environ.get("HA_TOKEN")
else None
)
ha_mcps = [McpServer.python_tool(name="ha", tools=[HA.ha])] if HA else []
ha_exposed = (ExposedMcp(name="ha"),) if HA else ()
# §8.5, §12.D1: вайбграм - комната агентов; диспетчер сидит в ней сам, но
# через перечислимые действия, токен агента в gateway. Клон репы комнаты -
# `vibegram` volume, /vibegram/<repo> ro, тянет сервис vibegram-repo.
VIBEGRAM_REPO = (
Path("/vibegram") / os.environ.get("VIBEGRAM_REPO", "x/room").split("/")[-1]
)
VIBEGRAM = (
Vibegram(
hub=os.environ.get("VIBEGRAM_HUB", "https://vibegram.example.com"),
token=os.environ["VIBEGRAM_TOKEN"],
nick=os.environ.get("VIBEGRAM_NICK", "beaver-test"),
repo=VIBEGRAM_REPO if VIBEGRAM_REPO.exists() else None,
tz=TZ,
)
if os.environ.get("VIBEGRAM_TOKEN")
else None
)
vibegram_mcps = (
[McpServer.python_tool(name="vibegram", tools=[VIBEGRAM.vibegram])]
if VIBEGRAM
else []
)
vibegram_exposed = (ExposedMcp(name="vibegram"),) if VIBEGRAM else ()
# §5: t3code-mcp - сосед по compose (профиль t3), машины и allowlist проектов
# в его t3code.toml; токены t3 живут в его env, gateway их не видит.
T3CODE_MCP = os.environ.get("T3CODE_MCP", "http://t3code-mcp:8000/mcp")
t3code_mcps = (
[McpServer.http(name="t3code", url=T3CODE_MCP)]
if "t3" in os.environ.get("COMPOSE_PROFILES", "").split(",")
else []
)
t3code_exposed = (ExposedMcp(name="t3code"),) if t3code_mcps else ()
# Секреты MCP - только через env подпроцесса (mcp stdio даёт ему белый список
# + это), никогда argv: процесс модели видит `ps` всего контейнера.
mcps: list[McpServerT] = [
McpServer.stdio(
name="obsidian-fs",
command=["bunx", "-y", "@modelcontextprotocol/server-filesystem", "/vault"],
lenient=True,
),
McpServer.stdio(
name="firefly",
command=["bunx", "-y", "@firefly-iii-mcp/local", "--preset", "default"],
env={
"FIREFLY_III_PAT": os.environ["FIREFLY_PAT"],
"FIREFLY_III_BASE_URL": os.environ["FIREFLY_BASE_URL"],
},
lenient=True,
),
McpServer.http(name="telegram", url=os.environ["BEAVERGRAM_MCP"]),
*calendar_mcps,
*komodo_mcps,
*ha_mcps,
*vibegram_mcps,
*t3code_mcps,
]
# Руки глубокого (§4.4): firefly без delete_*, obsidian-fs не даётся - свои
# файловые тулзы есть, vault смонтирован в /vault.
CLAUDE_MCPS = (
ExposedMcp(name="firefly", deny=("delete_*",)),
ExposedMcp(name="telegram"),
*calendar_exposed,
*komodo_exposed,
*ha_exposed,
)
UserPrefsRu = lambda: UserPreferences( # noqa: E731
locale="ru-RU", timezone="Europe/Warsaw", current_date=date.today().isoformat()
)
def dispatcher(name: str, model: str, effort: str | None = None) -> ClaudeAgent:
return ClaudeAgent(
name=name,
model=model,
cwd=VAULT,
prompts=DISPATCHER_PROMPTS,
skill_sets=DISPATCHER_SKILLS,
gateway_tools=("read_conversation", "spawn", "say", "schedule"),
options=ClaudeOptions(effort=effort, disallowed_tools=DISPATCHER_DISALLOWED),
policy=VAULT_POLICY,
# §4.4: t3code и вайбграм - руки диспетчера (мастер, ветки), глубоким
# не даются: наружу и в код ходит только он.
expose_mcps=(*CLAUDE_MCPS, *t3code_exposed, *vibegram_exposed),
)
def deep(name: str, model: str, effort: str | None = None) -> ClaudeAgent:
return ClaudeAgent(
name=name,
model=model,
cwd=VAULT,
prompts=DEEP_PROMPTS,
skill_sets=DEEP_SKILLS,
# §8.4: «ок, обсудили» - тулза, закрытие после ответа.
gateway_tools=("close_chat",),
options=ClaudeOptions(effort=effort, disallowed_tools=DEEP_DISALLOWED),
policy=VAULT_POLICY,
expose_mcps=CLAUDE_MCPS,
)
def distiller(name: str, model: str, effort: str | None = None) -> ClaudeAgent:
"""§4.1: форк закрытого чата - opus medium, без MCP, Read/Write в мета/бобер."""
return ClaudeAgent(
name=name,
model=model,
cwd=VAULT,
prompts=DISTILLER_PROMPTS,
options=ClaudeOptions(
effort=effort,
tools=("Read", "Write"),
disallowed_tools=DISTILLER_DISALLOWED,
),
policy=(vault_zones(STRICT_ZONES),),
)
def curator(name: str, model: str, effort: str | None = None) -> ClaudeAgent:
"""Куратор памяти: opus high, файлы в мета/бобер, без MCP, `inject` в мастер."""
return ClaudeAgent(
name=name,
model=model,
cwd=VAULT,
prompts=CURATOR_PROMPTS,
gateway_tools=("inject",),
options=ClaudeOptions(
effort=effort,
tools=("Read", "Write", "Edit", "Grep", "Glob"),
disallowed_tools=DISTILLER_DISALLOWED,
),
policy=(vault_zones(STRICT_ZONES),),
)
def triage(name: str, model: str, effort: str | None = None) -> ClaudeAgent:
"""§4.1: дешёвый крон-тёрн - sonnet low, без vault, только say/inject."""
return ClaudeAgent(
name=name,
model=model,
cwd=Path("/tmp"), # noqa: S108 - не vault: у триажа нет файлов вообще
prompts=TRIAGE_PROMPTS,
gateway_tools=("say", "inject"),
options=ClaudeOptions(
effort=effort, tools=(), disallowed_tools=TRIAGE_DISALLOWED
),
)
def raycast(name: str, model: str, reasoning_effort: str | None = None) -> RaycastAgent:
return RaycastAgent(
name=name,
model=model,
system_prompt=QUICK_PROMPT,
reasoning_effort=reasoning_effort,
available_native_tools=(RemoteTool.WEB_SEARCH, RemoteTool.READ_PAGE),
user_preferences=UserPrefsRu,
expose_mcps=(
ExposedMcp(name="obsidian-fs"),
ExposedMcp(name="firefly", deny=("delete_*",)),
ExposedMcp(name="telegram"),
*calendar_exposed,
),
)
agents: list[BaseAgent] = [
# §12.D4 пересмотрен 2026-09-01: диспетчер на high - бенч показал, что h
# голосует за находки и суждение, а не за тон; medium держал ~$10/день
# по API-прайсу (cost_usd в usage - кумулятив сессии, суммировать нельзя).
dispatcher("beaver-dispatcher", "claude-opus-5", effort="high"),
deep("beaver-opus-high", "claude-opus-5", effort="high"),
raycast("beaver-gemini-pro-high", "google-gemini-3.1-pro", reasoning_effort="high"),
deep("beaver-fable-high", "claude-fable-5", effort="high"),
deep("beaver-opus-medium", "claude-opus-5", effort="medium"),
deep("beaver-fable-medium", "claude-fable-5", effort="medium"),
deep("beaver-opus-xhigh", "claude-opus-5", effort="xhigh"),
distiller("beaver-distiller", "claude-opus-5", effort="medium"),
curator("beaver-curator", "claude-opus-5", effort="high"),
triage("beaver-triage", "claude-sonnet-5", effort="low"),
raycast("beaver-gemini-pro-low", "google-gemini-3.1-pro", reasoning_effort="low"),
raycast(
"beaver-gemini-flash-high", "google-gemini-3.5-flash", reasoning_effort="high"
),
raycast(
"beaver-gemini-flash-low", "google-gemini-3.5-flash", reasoning_effort="low"
),
]
PUBLIC_BASE_URL = os.environ.get("PUBLIC_BASE_URL", "").rstrip("/")
# §3.8: личка с ботом - домашний фронтенд мастера и веток (General и топики).
# Токен и id - секреты gateway, в env процесса модели они не попадают.
TELEGRAM = (
TelegramFrontend(
token=os.environ["TELEGRAM_BOT_TOKEN"],
user_id=int(os.environ["TELEGRAM_USER_ID"]),
master_agent="beaver-dispatcher",
branch_agent="beaver-dispatcher",
)
if os.environ.get("TELEGRAM_BOT_TOKEN")
else None
)
# Один порт: каждый HTTP-фронтенд живёт под своим путём (/anthropic, /mcp,
# /admin, /api, /md), `/` ведёт в админку. Caddy ничего не срезает.
frontends: list[Frontend] = [
*([TELEGRAM] if TELEGRAM is not None else []),
AnthropicMessagesFrontend(),
McpServerFrontend(),
AdminFrontend(),
# §3.9: /api показывает всё; дефолт master/branch здесь - запасной на случай
# gateway без телеграма, deep - у markdown.
ApiFrontend(
master_agent="beaver-dispatcher",
branch_agent="beaver-dispatcher",
memory_root=BEAVER,
),
MarkdownFrontend(
vault_path=CHATS_DIR,
default_agent="beaver-opus-high",
log_all_chats=True,
chat_path=chat_path,
),
]
def _read(path: Path) -> str | None:
return path.read_text(encoding="utf-8").strip() if path.exists() else None
def _today() -> date:
return datetime.now(UTC).astimezone(ZoneInfo(TZ)).date()
def latest_handout() -> Path | None:
"""§4.5: сид morning берёт последний хендаут, не «вчерашний по календарю»."""
files = sorted(DAYS.glob("????-??-??.md")) if DAYS.exists() else []
return files[-1] if files else None
def seed_body(ctx: SeedContext) -> str | None:
"""§8.2: утренний сид - последний хендаут из `дни/`; остальные сиды - gateway."""
if ctx.seed != "morning":
return None
handout = latest_handout()
if handout is None:
return "Хендаут не приехал."
state = _read(BEAVER / "состояние.md") or ""
signal = "\n".join(state.splitlines()[:5])
parts = [
f"Хендаут ({handout.stem}):\n\n{handout.read_text(encoding='utf-8').strip()}"
]
if signal:
parts.append(
f"Состояние (первые строки, полностью в мета/бобер/состояние.md):\n{signal}"
)
portrait = notes_head(PORTRAIT, ("сейчас", "паттерны"), cap=40)
if portrait:
parts.append(
"Портрет Бобра (шапка "
f"`{PORTRAIT.relative_to(VAULT)}`, лента - там же):\n{portrait}"
)
return "\n\n".join(parts)
DEFAULT_HANDOUT = (
"Этот мастер закрывается ({reason}). Напиши хендаут за {day} в "
"`мета/бобер/дни/{day}.md`: справку на утро, не задание."
)
def handout_prompt(ctx: HandoutContext) -> str:
"""§6.3, §8.3: последний тёрн закрываемого мастера - гранула + факт про дневник."""
day = ctx.day.isoformat()
body = (
(_read(GRANULES / "хендаут.md") or DEFAULT_HANDOUT)
.replace("{day}", day)
.replace("{reason}", ctx.reason)
)
diary = DIARY / f"{day}.md"
if diary.exists():
return f"{body}\n\nДневник за {day} приехал: `📅 дни/{day}.md`, возьми из него."
return f"{body}\n\nДневник за {day} не приехал - так и напиши в разделе «дневник»."
NEW_DAY = (
"Новый день ({day}). Ты - новый мастер того же треда: хендаут и сигнал состояния "
"в сиде выше. Скажи Бобру что-то через say, только если есть что сказать."
)
SAME_DAY = (
"Мастер пересоздан ({why}), день тот же - {day}. Хендаут в сиде выше - справка от "
"предыдущего мастера за сегодня, не утро: продолжай с того места, ничего не "
"приветствуй. Скажи Бобру что-то через say, только если есть что сказать."
)
ROTATION_WHY = {
"возраст": "прошло больше 36 часов",
"транскрипт": "контекст стал слишком большим",
}
def new_day(ctx: NewDayContext) -> str:
"""§4.5: «новый день» - только ночная ротация; остальные - тот же день."""
day = ctx.day.isoformat()
if ctx.reason == "ночь":
return NEW_DAY.format(day=day)
return SAME_DAY.format(day=day, why=ROTATION_WHY.get(ctx.reason, ctx.reason))
DISTILL = (
"Глубокий чат [[{chat}]] закрыт ({reason}), сегодня {day}. Выжимка - `Write` в "
'`{path}`, фронтматтер `source: "[[{chat}]]"`, `date: {day}`; слив - текст '
"ответа, до 5 строк, третье лицо."
)
DISTILL_NO_MEMORY = (
"Глубокий чат [[{chat}]] закрыт ({reason}), сегодня {day}. Память для него "
"выключена: файл не пиши, только слив текстом ответа - до 5 строк, третье лицо."
)
def distill_prompt(ctx: DistillContext) -> str:
"""§8.4: экземплярное для дистиллятора - какой чат, куда файл."""
day = ctx.day.isoformat()
if not ctx.memory:
return DISTILL_NO_MEMORY.format(chat=ctx.chat_name, reason=ctx.reason, day=day)
path = DIGESTS / f"{day} - {slugify(ctx.chat_name, maxlen=60)}.md"
return DISTILL.format(chat=ctx.chat_name, reason=ctx.reason, day=day, path=path)
# Панель и API - это сам Бобёр (single-user); остальные origin - крон, watch,
# komodo, ротация - не он.
BEAVER_ORIGINS = frozenset({"panel", "api"})
def inject_header(ctx: InjectContext) -> str:
if ctx.origin in BEAVER_ORIGINS:
head = (
"[инжект из панели: это Бобёр, ответ он видит в панели; "
"в телегу - только через say]"
)
else:
head = (
f"[инжект: {ctx.origin} - это не Бобёр, отвечать не нужно, "
"голос не обязателен]"
)
return f"{head}\n{INTERRUPTED_TURN}" if ctx.interrupted_turn else head
texts = ConversationTexts(
inject_header=inject_header,
merge_prompt=_read(GRANULES / "слив.md") or ConversationTexts().merge_prompt,
seed=seed_body,
handout=handout_prompt,
new_day=new_day,
distill=distill_prompt,
)
# §4.6: конверт - полный дифф только у дневника за сегодня, остальное именами;
# самый конкретный паттерн побеждает (мета/бобер видно, остальная мета - нет).
watch = VaultWatch(
VAULT,
WatchRules(
full=("📅 дни/{today}.md",),
names=(
"📅 дни/**",
"👤 люди/**",
"📆 планирование/**",
"📆 доски/**",
"💻 проекты/**",
"🧠 мысли/**",
"📶 ресерчи/**",
"мета/бобер/**",
),
ignore=(
"💬 чаты/**",
"мета/**",
"мета/бобер/реплики/**", # пишет сам gateway на каждое сообщение
"**/attachments/**",
".obsidian/**",
),
),
tz=TZ,
)
# §3.3: справка под сообщением - карточка и записки агента по упомянутым людям,
# дни, где они встречались, раз в день сроки из досок; лог реплик Бобра -
# единственное grep-абельное место для сказанного в телеге (§6.0 «спящее»).
recall = Recall(
vault=VAULT,
people_dir=PEOPLE,
notes_dir=NOTES,
diary_dir=DIARY,
boards_dir=BOARDS,
tz=TZ,
)
replies = ReplyLog(REPLIES, tz=TZ)
# §4.5 джобы. Хендлер только ставит работу в очередь и выходит.
async def rotate(run: JobRun) -> None:
for master in await run.rotate():
_log.info("rotated master -> %s", master.external_id)
_vibegram_last_wake: datetime | None = None
_vibegram_backlog: list[str] = []
VIBEGRAM_BRIEF = (
"Новое в вайбграме (комната {room}, ты там {nick}; остальные - чужие агенты):\n"
"{items}\n"
'Будить мастера - inject(conversation="master", urgency="wake", '
"text=резюме до 3 строк: кто, что, чего ждёт). Мастер прочитает подробности "
"через vibegram(read) и ответит через vibegram(send), если решит. "
"Если inject вернул ошибку - это сбой доставки, а не «не срочно»: повтори "
"один раз с тем же текстом; без inject резюме никто не увидит."
)
async def vibegram(run: JobRun) -> None:
"""§8.5: проверка → triage → в мастер ≤ 1/ч, если не ждут ответа."""
global _vibegram_last_wake # noqa: PLW0603 - счётчик «≤ 1/ч», живёт до рестарта
if VIBEGRAM is None:
return
events = await VIBEGRAM.pending()
_vibegram_backlog.extend(e.line() for e in events)
if not _vibegram_backlog:
return
now = datetime.now(UTC)
addressed = VIBEGRAM.mentioned(events)
if (
not addressed
and _vibegram_last_wake is not None
and now - _vibegram_last_wake < timedelta(hours=1)
):
return
_vibegram_last_wake = now
items, _vibegram_backlog[:] = list(_vibegram_backlog), []
await run.spawn_job(
agent="beaver-triage",
title="вайбграм",
text=VIBEGRAM_BRIEF.format(
room=VIBEGRAM_REPO.name,
nick=VIBEGRAM.nick,
items="\n".join(f"- {item}" for item in items[-40:]),
),
)
KOMODO_STACK = os.environ.get("KOMODO_STACK", "beaver-agent")
async def komodo_deploy() -> None:
if KOMODO is None:
_log.warning("deploy hook: KOMODO_URL/KEY/SECRET are not set, nothing deployed")
return
receipt = await KOMODO.execute(
"DeployStack", {"stack": KOMODO_STACK, "services": []}
)
_log.info("deploy hook: komodo accepted %s", receipt.get("id"))
async def deploy(run: JobRun) -> None:
"""§8.6: деплой ждёт conversation.idle(master) - иначе рестарт посреди тёрна."""
master = await run.master()
if master is not None and await run.conversations.busy(master):
await run.retry_in(timedelta(seconds=60))
return
run.background(komodo_deploy())
async def close_idle(run: JobRun) -> None:
"""§4.5, §8.4: ночью, после ротации - глубокие без активности 2 дня, ≤ 3."""
for result in await run.close_idle(kind="deep", days=2, limit=3, since=LAUNCH):
_log.info(
"closed deep chat %s: digest=%s error=%s",
result.conversation.external_id,
result.digest.path if result.digest else None,
result.error,
)
MEMORY = (
"Воскресенье, {day}. Перепиши `{state}` из последних хендаутов: {handouts}. "
"`## сейчас` - только невыводимое из vault, `## мои хвосты` - capability "
"requests в работе и обещания длиннее дня. Меньше {max_lines} строк."
)
async def memory(run: JobRun) -> None:
"""§4.5, §6.2: вс 04:30 - состояние.md из 7 хендаутов, потолок держит gateway."""
handouts = sorted(DAYS.glob("????-??-??.md"))[-7:] if DAYS.exists() else []
if not handouts:
_log.info("memory: no handouts yet, nothing to rewrite")
return
await run.spawn_job(
agent="beaver-distiller",
title="память",
text=MEMORY.format(
day=_today().isoformat(),
state=STATE.relative_to(VAULT),
handouts=", ".join(f"`{h.relative_to(VAULT)}`" for h in handouts),
max_lines=STATE_MAX_LINES,
),
line_cap=LineCap(STATE, max_lines=STATE_MAX_LINES),
)
CURATOR_WATCH = [
Watched("хендауты", DAYS, "????-??-??.md"),
Watched("выжимки", DIGESTS, "*.md"),
Watched("наблюдения", NOTES, "*.md"),
Watched("состояние", BEAVER, "состояние.md"),
Watched("дневник", DIARY, "????-??-??.md"),
Watched("карточки", PEOPLE, "**/*.md"),
]
async def curate(run: JobRun) -> None:
"""§6.0: куратор памяти раз в несколько часов.
Крон собирает «что нового», модель раскладывает по шапкам и лентам;
нечего раскладывать - джоб не спавнится.
"""
text = briefing(
vault=VAULT,
journal=CURATOR_JOURNAL,
replies_dir=REPLIES,
watched=CURATOR_WATCH,
tz=TZ,
)
if text is None:
_log.info("curator: nothing new since the last run")
return
await run.spawn_job(agent="beaver-curator", title="куратор", text=text)
T3CODE_QUESTION_HINT = (
"Если ответ следует из задания - ответь сам: t3_answer(thread_id, "
"{id вопроса: label}). Не знаешь - спроси Бобра через say своими словами, "
"потом t3_answer. Тред стоит, пока не ответят."
)
def t3code_text(p: dict[str, Any]) -> str:
head = (
f"t3code: тред «{p.get('title')}» ({p.get('machine')}, {p.get('project')}, "
f"thread_id {p.get('thread_id')})"
)
event = p.get("event")
if event in ("question", "approval"):
lines = [f"{head} - спрашивает."]
for q in (p.get("question") or {}).get("questions", []):
multi = " (можно несколько)" if q.get("multi_select") else ""
lines.append(
f"- [{q['id']}] {q.get('header') or ''}: {q['question']}{multi}"
)
lines.extend(
f" · {o['label']}"
+ (f" - {o['description']}" if o.get("description") else "")
for o in q.get("options", [])
)
lines.append(T3CODE_QUESTION_HINT)
return "\n".join(lines)
if event == "completed":
lines = [f"{head} - завершён."]
elif event == "interrupted":
lines = [f"{head} - прерван."]
else:
lines = [f"{head} - упал: {p.get('error') or 'без текста ошибки'}."]
if p.get("text"):
lines.append(f"Ответ треда:\n{p['text']}")
if p.get("files"):
lines.append("Файлы: " + ", ".join(f["path"] for f in p["files"]))
return "\n".join(lines)
async def t3code_event(run: JobRun) -> None:
"""§5, §8.6: t3code-mcp следит за тредами по WS и присылает сюда события."""
await run.inject_master(t3code_text(run.payload), urgency="urgent", origin="t3code")
ALERTS = AlertMemory()
KOMODO_HINT = (
"Это флот Бобра; ему такие алерты приходят и так (healthbot), тебя будят, "
"когда через 5 минут само не встало. Посмотри komodo(status/logs), почини, "
"если очевидно (restart стека), иначе - коротко скажи Бобру через say, что "
"видишь. Не деплой и не чини «заодно»."
)
async def komodo_alert(run: JobRun) -> None:
"""§4.5 `/hooks/komodo`: флапы выдерживаются HOLD и перепроверяются."""
payload = run.payload
if payload.get("held"):
alert = parse(payload["alert"])
bad = await still_bad(KOMODO.read, alert) if KOMODO else True
verdict = ALERTS.after_hold(alert, still_bad=bad)
if verdict.action == "drop":
_log.info("komodo alert %s recovered during hold", alert.key)
return
ALERTS.mark_told(alert)
await run.inject_master(
f"Komodo: {describe(alert, held=HOLD)}\n{KOMODO_HINT}",
urgency=verdict.urgency,
origin="komodo",
)
return
alert = parse(payload)
verdict = ALERTS.verdict(alert)
match verdict.action:
case "drop":
_log.info("komodo alert %s dropped (%s)", alert.key, alert.kind)
case "hold":
await run.scheduler.trigger(
run.job,
{
"held": True,
"alert": {k: v for k, v in payload.items() if k != "trigger"},
},
delay=HOLD,
trigger="hold",
)
case "resolved":
await run.inject_master(
f"Komodo: отбой - {describe(alert)}", urgency="normal", origin="komodo"
)
case "inject":
ALERTS.mark_told(alert)
await run.inject_master(
f"Komodo: {describe(alert)}\n{KOMODO_HINT}",
urgency=verdict.urgency,
origin="komodo",
)
jobs = [
Job("ротация", rotate, cron="0 * * * *"),
Job("вайбграм", vibegram, cron="*/10 * * * *", critical=False),
Job("закрытие", close_idle, cron="20 4 * * *", critical=False),
Job("память", memory, cron="30 4 * * 0", critical=False),
# крон в UTC контейнера: 02/10/14/18/22 по Варшаве летом;
# руками - POST /api/jobs/куратор/run
Job("куратор", curate, cron="0 0,8,12,16,20 * * *", critical=False),
Job("deploy", deploy, webhook=True),
# алерты идут пачками (тревога и отбой подряд) - схлопывать нельзя
Job("komodo", komodo_alert, webhook=True, dedupe=False),
# события идут подряд (вопрос и завершение) - схлопывать нельзя
Job("t3code", t3code_event, webhook=True, dedupe=False),
]
gateway = Gateway(
agents=agents,
mcps=mcps,
frontends=frontends,
texts=texts,
jobs=jobs,
# Окно модели - 1M; context rot начинается около половины (эмпирика h по
# счётчику Claude Code), считаем как он - вход последнего API-вызова.
rotation=RotationPolicy(tz=TZ, max_context_tokens=500_000),
watch=watch,
recall=recall.block,
user_sink=replies.write,
budget=Budget(threshold=0.7),
distiller=Distiller(agent="beaver-distiller", dir=DIGESTS, index=INDEX),
tz=TZ,
port=62990,
public_url=PUBLIC_BASE_URL or None,
)