refactor: split config.py into the beaver_agent package (vault, prompts, skills, policy, hands, agents, frontends, texts, memory, jobs)

This commit is contained in:
hh
2026-09-02 00:24:31 +02:00
parent f161d6a46d
commit 3b2f5057da
38 changed files with 1302 additions and 935 deletions
+142
View File
@@ -0,0 +1,142 @@
"""Руки: какие MCP есть у сетапа и кому из агентов какая рука выдана.
Секреты MCP - только через env подпроцесса, никогда argv: процесс модели
видит `ps` всего контейнера. Рука, чей env не задан, просто не появляется.
"""
from __future__ import annotations
import os
from pathlib import Path
from beaver_gateway.agents.base import ExposedMcp
from beaver_gateway.mcp.types import McpServer, McpServerT
from beaver_agent.hands.calendars import calendars
from beaver_agent.hands.homeassistant import HomeAssistant
from beaver_agent.hands.komodo import Komodo
from beaver_agent.hands.vibegram import Vibegram
from beaver_agent.vault import TZ
env = os.environ.get
obsidian = McpServer.stdio(
name="obsidian-fs",
command=["bunx", "-y", "@modelcontextprotocol/server-filesystem", "/vault"],
lenient=True,
)
"""Файловые тулзы для тех, у кого нет своих (Raycast); у Claude vault смонтирован."""
firefly = (
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,
)
if env("FIREFLY_PAT")
else None
)
firefly_exposed = (ExposedMcp(name="firefly", deny=("delete_*",)),) if firefly else ()
telegram = (
McpServer.http(name="telegram", url=os.environ["BEAVERGRAM_MCP"])
if env("BEAVERGRAM_MCP")
else None
)
"""beavergram: история чатов Бобра в Telegram."""
calendar_servers = calendars(env("CALENDAR_MCPS", ""))
calendar_exposed = tuple(ExposedMcp(name=m.name) for m in calendar_servers)
KOMODO = (
Komodo(
url=os.environ["KOMODO_URL"],
key=os.environ["KOMODO_KEY"],
secret=os.environ["KOMODO_SECRET"],
)
if env("KOMODO_URL")
else None
)
"""Флот: перечислимые действия, ключ остаётся в gateway."""
komodo = McpServer.python_tool(name="komodo", tools=[KOMODO.komodo]) if KOMODO else None
komodo_exposed = (ExposedMcp(name="komodo"),) if KOMODO else ()
HA = (
HomeAssistant(url=os.environ["HA_URL"], token=os.environ["HA_TOKEN"])
if env("HA_URL") and env("HA_TOKEN")
else None
)
"""Дом: HA живёт на хосте, из контейнера - host.docker.internal."""
ha = McpServer.python_tool(name="ha", tools=[HA.ha]) if HA else None
ha_exposed = (ExposedMcp(name="ha"),) if HA else ()
VIBEGRAM_REPO = Path("/vibegram") / env("VIBEGRAM_REPO", "x/room").split("/")[-1]
VIBEGRAM = (
Vibegram(
hub=env("VIBEGRAM_HUB", "https://vibegram.example.com"),
token=os.environ["VIBEGRAM_TOKEN"],
nick=env("VIBEGRAM_NICK", "beaver-test"),
repo=VIBEGRAM_REPO if VIBEGRAM_REPO.exists() else None,
tz=TZ,
)
if env("VIBEGRAM_TOKEN")
else None
)
"""Комната агентов: диспетчер сидит в ней сам; клон репы комнаты - том `vibegram`."""
vibegram = (
McpServer.python_tool(name="vibegram", tools=[VIBEGRAM.vibegram])
if VIBEGRAM
else None
)
vibegram_exposed = (ExposedMcp(name="vibegram"),) if VIBEGRAM else ()
t3code = (
McpServer.http(name="t3code", url=env("T3CODE_MCP", "http://t3code-mcp:8000/mcp"))
if "t3" in env("COMPOSE_PROFILES", "").split(",")
else None
)
"""t3code-mcp - сосед по compose (профиль t3); машины и allowlist в его t3code.toml."""
t3code_exposed = (ExposedMcp(name="t3code"),) if t3code else ()
msos = (
McpServer.http(name="msos", url=os.environ["MSOS_MCP"]) if env("MSOS_MCP") else None
)
"""Таскер рабочих проектов (скилл `задачи`)."""
msos_exposed = (ExposedMcp(name="msos"),) if msos else ()
mcps: list[McpServerT] = [
m
for m in (
obsidian,
firefly,
telegram,
*calendar_servers,
komodo,
ha,
vibegram,
t3code,
msos,
)
if m is not None
]
CLAUDE = (
*firefly_exposed,
ExposedMcp(name="telegram"),
*calendar_exposed,
*komodo_exposed,
*ha_exposed,
)
"""Руки любого Claude-агента: firefly без delete, телега, календари, флот, дом."""
DISPATCHER = (*CLAUDE, *t3code_exposed, *vibegram_exposed, *msos_exposed)
"""Диспетчер ещё ходит в код, в комнату агентов и в таскер; глубоким это не даётся."""
RAYCAST = (
ExposedMcp(name="obsidian-fs"),
*firefly_exposed,
ExposedMcp(name="telegram"),
*calendar_exposed,
)
+19
View File
@@ -0,0 +1,19 @@
"""Календари из `CALENDAR_MCPS=name=url,name=url`; в url - приватный фид с токеном."""
from __future__ import annotations
from beaver_gateway.mcp.types import HttpMcp, McpServer
def calendars(raw: str) -> list[HttpMcp]:
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():
msg = f"CALENDAR_MCPS: запись #{position} не вида `name=url`"
raise ValueError(msg)
servers.append(McpServer.http(name=f"calendar-{name.strip()}", url=url.strip()))
return servers
+169
View File
@@ -0,0 +1,169 @@
"""Home Assistant как python_tool с перечислимыми действиями.
Токен HA живёт в gateway и в контекст модели не попадает. Действия:
найти сущности, посмотреть состояние, вызвать сервис, перечислить сервисы
домена. Ничего про конфигурацию, пользователей и интеграции - только
REST `/api/states` и `/api/services`.
"""
from __future__ import annotations
import json
from dataclasses import dataclass
from typing import Any, Literal
import aiohttp
__all__ = ["Action", "HomeAssistant"]
Action = Literal["search", "state", "call", "services"]
SEARCH_LIMIT = 60
ATTR_LIMIT = 25
# Не влезает в ответ и не нужно модели: гигантские списки и бинарники.
NOISY_ATTRS = frozenset({"entity_picture", "icon", "supported_color_modes", "options"})
@dataclass(frozen=True, slots=True)
class HomeAssistant:
url: str
token: str
timeout: float = 20.0
max_chars: int = 8_000
async def ha(
self,
action: Action,
*,
query: str | None = None,
entity: str | None = None,
domain: str | None = None,
service: str | None = None,
data: dict[str, Any] | None = None,
) -> str:
"""Home Assistant - умный дом: свет, кнопки ИК-пульта, сцены, скрипты, датчики.
action:
- search: сущности по подстроке в id или имени (query; пусто - все),
сгруппированы по домену. Начни отсюда: точные entity_id нужны везде.
- state: состояние и атрибуты одной сущности (entity).
- call: вызвать сервис (domain, service, entity опционально, data -
доп. поля: brightness, color_temp, …). Примеры: light/turn_on,
light/turn_off, script/turn_on, scene/turn_on, button/press.
- services: какие сервисы есть у домена (domain) и их поля.
Делай то, о чём попросили: «выключи свет» - это одна сущность или
группа, не весь дом.
"""
match action:
case "search":
return await self._search(query or "")
case "state":
return await self._state(_need(entity, "entity"))
case "call":
return await self._call(
_need(domain, "domain"), _need(service, "service"), entity, data
)
case "services":
return await self._services(_need(domain, "domain"))
async def _search(self, query: str) -> str:
states = await self._get("/api/states")
q = query.lower()
hits = [
s
for s in states
if q in s["entity_id"].lower()
or q in str(s.get("attributes", {}).get("friendly_name", "")).lower()
]
if not hits:
return "ничего не нашлось"
by_domain: dict[str, list[dict[str, Any]]] = {}
for s in sorted(hits, key=lambda s: s["entity_id"]):
by_domain.setdefault(s["entity_id"].split(".", 1)[0], []).append(s)
lines = []
shown = 0
for dom, items in by_domain.items():
lines.append(f"{dom}:")
for s in items:
if shown >= SEARCH_LIMIT:
lines.append(f" … ещё {len(hits) - shown}, уточни query")
return "\n".join(lines)
name = s.get("attributes", {}).get("friendly_name") or ""
lines.append(f" {s['entity_id']:40} {s['state']:12} {name}")
shown += 1
return "\n".join(lines)
async def _state(self, entity: str) -> str:
s = await self._get(f"/api/states/{entity}")
attrs = {
k: v for k, v in (s.get("attributes") or {}).items() if k not in NOISY_ATTRS
}
lines = [f"{s['entity_id']}: {s['state']} (с {s.get('last_changed', '?')})"]
for k, v in list(attrs.items())[:ATTR_LIMIT]:
lines.append(f" {k}: {json.dumps(v, ensure_ascii=False)[:200]}")
return self._clip("\n".join(lines))
async def _call(
self, domain: str, service: str, entity: str | None, data: dict[str, Any] | None
) -> str:
body: dict[str, Any] = dict(data or {})
if entity:
body["entity_id"] = entity
changed = await self._post(f"/api/services/{domain}/{service}", body)
if not changed:
return f"{domain}.{service}: выполнено, состояния не менялись"
lines = [f"{domain}.{service}: выполнено, изменилось:"]
lines.extend(f" {s['entity_id']}: {s['state']}" for s in changed[:20])
return "\n".join(lines)
async def _services(self, domain: str) -> str:
for block in await self._get("/api/services"):
if block.get("domain") == domain:
lines = [f"{domain}:"]
for name, spec in sorted((block.get("services") or {}).items()):
fields = ", ".join((spec.get("fields") or {}).keys())
lines.append(
f" {name}: {spec.get('description') or ''}".rstrip()
+ (f" [{fields}]" if fields else "")
)
return self._clip("\n".join(lines))
return f"домена {domain} нет"
async def _get(self, path: str) -> Any:
return await self._request("GET", path, None)
async def _post(self, path: str, body: dict[str, Any]) -> Any:
return await self._request("POST", path, body)
async def _request(
self, method: str, path: str, body: dict[str, Any] | None
) -> Any:
headers = {"Authorization": f"Bearer {self.token}"}
async with (
aiohttp.ClientSession() as http,
http.request(
method,
f"{self.url.rstrip('/')}{path}",
json=body,
headers=headers,
timeout=aiohttp.ClientTimeout(total=self.timeout),
) as response,
):
text = await response.text()
if response.status >= 400:
msg = f"Home Assistant {method} {path}: {response.status} {text[:300]}"
raise RuntimeError(msg)
return json.loads(text) if text else None
def _clip(self, text: str) -> str:
if len(text) <= self.max_chars:
return text
return text[: self.max_chars] + "\n…[обрезано]"
def _need(value: str | None, name: str) -> str:
if not value:
msg = f"для этого действия нужен параметр {name}"
raise ValueError(msg)
return value
+292
View File
@@ -0,0 +1,292 @@
"""Komodo как python_tool с перечислимыми действиями.
Ключ живёт в gateway и в контекст модели не попадает. В перечислении нет
`Prune*`, `Destroy*`, `DeleteServer` и терминала хоста; `exec` - это
`docker exec` в контейнер, не шелл на сервере, и контейнеры из
``exec_deny`` (periphery с docker.sock - это и есть хост) ему не даются.
"""
from __future__ import annotations
import asyncio
import fnmatch
import json
import re
import time
from dataclasses import dataclass
from datetime import UTC, datetime
from typing import Any, Literal
import aiohttp
__all__ = ["Action", "Komodo"]
Action = Literal[
"status",
"stacks",
"containers",
"logs",
"search_logs",
"updates",
"update",
"deploy",
"restart",
"exec",
]
ANSI = re.compile(r"\x1b\[[0-9;?]*[ -/]*[@-~]")
EXIT_MARK = "__KOMODO_EXIT_CODE:"
NOT_RUNNING_LIMIT = 15
@dataclass(frozen=True, slots=True)
class Komodo:
url: str
key: str
secret: str
exec_deny: tuple[str, ...] = ("*periphery*", "*komodo*")
terminal: str = "beaver"
shell: str = "env PAGER=cat GIT_PAGER=cat TERM=dumb sh"
"""Terminal is a tty: without ``PAGER=cat`` psql/git open a pager and
wait for a key forever, wedging the terminal."""
timeout: float = 120.0
wait: float = 180.0
max_chars: int = 12_000
async def komodo(
self,
action: Action,
*,
stack: str | None = None,
service: str | None = None,
container: str | None = None,
command: str | None = None,
terms: str | None = None,
tail: int = 100,
update: str | None = None,
server: str | None = None,
) -> str:
"""Прод через Komodo: посмотреть, а по просьбе Бобра - подействовать.
action:
- status: серверы и стеки не в running.
- stacks: все стеки - сервер, состояние, есть ли незадеплоенные коммиты.
- containers: контейнеры сервера (server; без него - все серверы).
- logs: логи стека (stack, service опционально, tail строк).
- search_logs: поиск по логам стека (stack, terms - слова через пробел).
- updates: последние операции - кто, что, статус.
- update: подробности одной операции по id (update).
- deploy / restart: стек целиком или один service; ждёт исхода до 3 мин.
- exec: команда внутри контейнера (container, command; server, если имя
контейнера есть на нескольких серверах). Это `docker exec` от
пользователя контейнера, свежий шелл на каждый вызов, без tty-пейджера.
Лимит ~2 мин: долгое запускай в фон с выводом в файл и читай файл.
Деплой и рестарт - только когда об этом попросили, не «заодно».
"""
match action:
case "status":
return await self._status()
case "stacks":
return await self._stacks()
case "containers":
return await self._containers(server)
case "logs":
return await self._logs(_need(stack, "stack"), service, tail)
case "search_logs":
return await self._search(
_need(stack, "stack"), _need(terms, "terms"), service
)
case "updates":
return await self._updates()
case "update":
return await self._update(_need(update, "update"))
case "deploy":
return await self._run("DeployStack", _need(stack, "stack"), service)
case "restart":
return await self._run("RestartStack", _need(stack, "stack"), service)
case "exec":
return await self._exec(
_need(container, "container"), _need(command, "command"), server
)
async def execute(self, operation: str, params: dict[str, Any]) -> dict[str, Any]:
return await self._json("execute", {"type": operation, "params": params})
async def read(self, request: str, params: dict[str, Any]) -> Any:
return await self._json("read", {"type": request, "params": params})
async def _status(self) -> str:
servers, stacks = await asyncio.gather(
self.read("ListServers", {}), self.read("ListStacks", {})
)
lines = [f"{s['name']:24} {s['info']['state']}" for s in servers]
bad = [s for s in stacks if s["info"]["state"] != "running"]
lines.append(f"стеков: {len(stacks)}, не running: {len(bad)}")
lines.extend(
f" {s['info']['state']:12} {s['name']} @{s['info'].get('server_name')}"
for s in bad[:NOT_RUNNING_LIMIT]
)
return "\n".join(lines)
async def _stacks(self) -> str:
stacks = await self.read("ListStacks", {})
lines = []
for s in sorted(stacks, key=lambda s: s["name"]):
info = s["info"]
behind = info.get("latest_hash") and info.get("deployed_hash") != info.get(
"latest_hash"
)
lines.append(
f"{info['state']:12} {s['name']:28} @{info.get('server_name')}"
+ (" (есть незадеплоенные коммиты)" if behind else "")
)
return "\n".join(lines)
async def _containers(self, server: str | None) -> str:
names = (
[server]
if server
else [s["name"] for s in await self.read("ListServers", {})]
)
lines = []
for name in names:
containers = await self.read("ListDockerContainers", {"server": name})
lines.append(f"{name}:")
lines.extend(
f" {c['state']:10} {c['name']:36} {c.get('status', '')}"
for c in sorted(containers, key=lambda c: c["name"])
)
return "\n".join(lines)
async def _logs(self, stack: str, service: str | None, tail: int) -> str:
log = await self.read(
"GetStackLog",
{"stack": stack, "services": [service] if service else [], "tail": tail},
)
return self._clip(_log_text(log))
async def _search(self, stack: str, terms: str, service: str | None) -> str:
log = await self.read(
"SearchStackLog",
{
"stack": stack,
"services": [service] if service else [],
"terms": terms.split(),
},
)
return self._clip(_log_text(log))
async def _updates(self) -> str:
updates = (await self.read("ListUpdates", {}))["updates"][:15]
return "\n".join(_update_line(u) for u in updates) or "операций нет"
async def _update(self, update_id: str) -> str:
return self._clip(
_update_detail(await self.read("GetUpdate", {"id": update_id}))
)
async def _run(self, operation: str, stack: str, service: str | None) -> str:
receipt = await self.execute(
operation, {"stack": stack, "services": [service] if service else []}
)
update_id = receipt.get("id") or receipt.get("_id", {}).get("$oid")
if not update_id:
return f"{operation} {stack}: Komodo вернул {receipt!r}"
deadline = time.monotonic() + self.wait
while time.monotonic() < deadline:
await asyncio.sleep(5)
current = await self.read("GetUpdate", {"id": update_id})
if current.get("status") == "Complete":
return self._clip(_update_detail(current))
return f"{operation} {stack}: ещё идёт, update={update_id} - проверь позже"
async def _exec(self, container: str, command: str, server: str | None) -> str:
if any(fnmatch.fnmatchcase(container, p) for p in self.exec_deny):
return f"exec в {container} запрещён: это контейнер инфраструктуры"
server = server or await self._server_of(container)
if server is None:
return f"контейнер {container} не найден ни на одном сервере"
body = {
"target": {
"type": "Container",
"params": {"server": server, "container": container},
},
"terminal": self.terminal,
"command": command,
"init": {"command": self.shell, "recreate": "Always"},
}
status, text = await self._post("terminal/execute", body)
if status != 200:
return f"exec {container}: Komodo {status}: {text[:500]}"
output, _, code = ANSI.sub("", text).rpartition(EXIT_MARK)
header = f"[{container}@{server}] exit {code.strip() or '?'}"
return f"{header}\n{self._clip(output.strip())}"
async def _server_of(self, container: str) -> str | None:
for s in await self.read("ListServers", {}):
containers = await self.read("ListDockerContainers", {"server": s["name"]})
if any(c["name"] == container for c in containers):
return s["name"]
return None
async def _json(self, route: str, body: dict[str, Any]) -> Any:
status, text = await self._post(route, body)
if status != 200:
msg = f"Komodo {route} {body.get('type')}: {status} {text[:500]}"
raise RuntimeError(msg)
return json.loads(text)
async def _post(self, route: str, body: dict[str, Any]) -> tuple[int, str]:
headers = {"X-Api-Key": self.key, "X-Api-Secret": self.secret}
async with (
aiohttp.ClientSession() as http,
http.post(
f"{self.url.rstrip('/')}/{route}",
json=body,
headers=headers,
timeout=aiohttp.ClientTimeout(total=self.timeout),
) as response,
):
return response.status, await response.text()
def _clip(self, text: str) -> str:
if len(text) <= self.max_chars:
return text
head, tail = self.max_chars * 2 // 3, self.max_chars // 3
cut = len(text) - head - tail
return f"{text[:head]}\n…[обрезано {cut} символов]…\n{text[-tail:]}"
def _need(value: str | None, name: str) -> str:
if not value:
msg = f"для этого действия нужен параметр {name}"
raise ValueError(msg)
return value
def _log_text(log: dict[str, Any]) -> str:
if "error" in log:
return f"Komodo: {log['error']}"
return ANSI.sub("", (log.get("stdout") or "") + (log.get("stderr") or "")).strip()
def _ts(ms: int) -> str:
return datetime.fromtimestamp(ms / 1000, tz=UTC).strftime("%m-%d %H:%M")
def _update_line(u: dict[str, Any]) -> str:
ok = "ok" if u.get("success") else "FAIL"
head = f"{_ts(u['start_ts'])} {u['operation']:20} {u['status']:10} {ok}"
return f"{head} user={u.get('username')} id={u.get('id')}"
def _update_detail(u: dict[str, Any]) -> str:
lines = [f"{u['operation']} {u['status']} {'ok' if u.get('success') else 'FAIL'}"]
for stage in u.get("logs", []):
out = ANSI.sub(
"", (stage.get("stdout") or "") + " " + (stage.get("stderr") or "")
)
lines.append(f"-- {stage['stage']} {'ok' if stage['success'] else 'FAIL'}")
lines.append(re.sub("<[^>]+>", "", out).strip()[:1500])
return "\n".join(lines)
+229
View File
@@ -0,0 +1,229 @@
"""алерты Komodo для диспетчера - что доехать до мастера, когда и как.
Вебхук `/hooks/komodo` получает всё, что шлёт Komodo (типы алертов у
алертера не перечислены намеренно). Сюда мастеру нужна малая часть, и не
сразу: `ServerUnreachable` и смены состояния стеков - флап, который через
минуту сам отбивается («отбой» от Komodo приходит следом), а срочный
инжект прерывает тёрн и стоит токенов. Поэтому такие алерты выдерживаются
(`HOLD`), и по истечении выдержки состояние перепроверяется по API: если
всё уже running - алерт молча выбрасывается; если нет - мастер получает
срочный инжект. Информационные типы (обновления образов, расписания,
дрейф синка) не доезжают вообще - их читает Бобёр в healthbot.
Komodo создаёт `StackStateChange`/`ContainerStateChange` сразу с
`resolved=true` (`resolved_ts == ts`): это факт «состояние сменилось», а
не тревога, которую снимают. У них смотрим на `to`, не на `resolved`.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import UTC, datetime, timedelta
from typing import Any, Literal
__all__ = ["HOLD", "Alert", "AlertMemory", "Verdict", "describe", "parse", "still_bad"]
HOLD = timedelta(minutes=5)
# Не для диспетчера: Бобёр видит их в healthbot, действия по ним - его.
INFORMATIONAL = frozenset(
{
"StackImageUpdateAvailable",
"DeploymentImageUpdateAvailable",
"StackAutoUpdated",
"DeploymentAutoUpdated",
"ResourceSyncPendingUpdates",
"ScheduleRun",
"Test",
"Custom",
}
)
# Флапает: ждать HOLD и перепроверить, прежде чем будить мастера.
FLAPPY = frozenset(
{"ServerUnreachable", "StackStateChange", "ContainerStateChange", "ServerCpu"}
)
STATE_CHANGE = frozenset({"StackStateChange", "ContainerStateChange"})
HEALTHY_STATES = frozenset({"running", "healthy"})
FAILURE = frozenset(
{"ProcedureFailed", "ActionFailed", "BuildFailed", "RepoBuildFailed"}
)
Action = Literal["drop", "hold", "inject", "resolved"]
@dataclass(frozen=True, slots=True)
class Alert:
kind: str
level: str
resolved: bool
name: str
server: str | None
data: dict[str, Any]
@property
def key(self) -> str:
return f"{self.kind}:{self.server or '-'}:{self.name}"
@property
def bad(self) -> bool:
"""Тревога, а не снятие.
У смен состояния - по `to`, у прочих - по `resolved`.
"""
if self.kind in STATE_CHANGE:
return str(self.data.get("to", "")).lower() not in HEALTHY_STATES
return not self.resolved
@dataclass(frozen=True, slots=True)
class Verdict:
action: Action
urgency: Literal["urgent", "normal"] = "normal"
@dataclass(slots=True)
class AlertMemory:
"""Что выдерживается и о чём мастеру уже сказали.
Повторы во время выдержки не дублируются, а «отбой» доезжает только
после тревоги, которую мастер видел.
"""
holding: set[str] = field(default_factory=set)
told: dict[str, datetime] = field(default_factory=dict)
def verdict(self, alert: Alert) -> Verdict:
if alert.kind in INFORMATIONAL:
return Verdict("drop")
if not alert.bad:
if alert.key in self.holding:
# Снялось, пока ждали: ни тревоги, ни отбоя.
self.holding.discard(alert.key)
return Verdict("drop")
if self.told.pop(alert.key, None) is not None:
return Verdict("resolved")
return Verdict("drop")
if alert.kind in FLAPPY:
if alert.key in self.holding:
return Verdict("drop")
self.holding.add(alert.key)
return Verdict("hold")
return Verdict("inject", "normal")
def after_hold(self, alert: Alert, *, still_bad: bool) -> Verdict:
self.holding.discard(alert.key)
if not still_bad:
return Verdict("drop")
return Verdict("inject", "urgent")
def mark_told(self, alert: Alert) -> None:
self.told[alert.key] = datetime.now(UTC)
def parse(payload: dict[str, Any]) -> Alert:
data = payload.get("data") or {}
inner = data.get("data") if isinstance(data, dict) else None
inner = inner if isinstance(inner, dict) else {}
return Alert(
kind=str(data.get("type") or "Unknown")
if isinstance(data, dict)
else "Unknown",
level=str(payload.get("level") or "").upper(),
resolved=bool(payload.get("resolved")),
name=str(inner.get("name") or inner.get("id") or "?"),
server=inner.get("server_name"),
data=inner,
)
def _gb(value: Any) -> str:
try:
return f"{float(value):.1f} ГБ"
except (TypeError, ValueError):
return str(value)
def _pct(value: Any) -> str:
try:
return f"{float(value):.0f}%"
except (TypeError, ValueError):
return str(value)
def _err(value: Any) -> str:
if isinstance(value, dict):
return str(value.get("error") or value)
return str(value or "без деталей")
def describe(alert: Alert, *, held: timedelta | None = None) -> str:
"""Одна-две строки для инжекта: что, где, с каких пор."""
d = alert.data
name, server = alert.name, alert.server
where = f" на {server}" if server else ""
since = f" уже {int(held.total_seconds() // 60)} мин" if held else ""
match alert.kind:
case "ServerUnreachable":
return f"сервер {name} не отвечает{since}: {_err(d.get('err'))}"
case "StackStateChange":
return f"стек {name}{where}: {d.get('from')}{d.get('to')}{since}"
case "ContainerStateChange":
return f"контейнер {name}{where}: {d.get('from')}{d.get('to')}{since}"
case "ServerCpu":
return f"CPU на {name}: {_pct(d.get('percentage'))}{since}"
case "ServerMem":
return (
f"память на {name}: {_gb(d.get('used_gb'))} из {_gb(d.get('total_gb'))}"
)
case "ServerDisk":
return (
f"диск на {name} ({d.get('path', '/')}): "
f"{_gb(d.get('used_gb'))} из {_gb(d.get('total_gb'))}"
)
case "ServerVersionMismatch":
return (
f"версии разъехались на {name}: periphery {d.get('version')}, "
f"core {d.get('core_version')} - чинится только по ssh, это к Бобру"
)
case k if k in FAILURE:
what = {"ProcedureFailed": "процедура", "ActionFailed": "действие"}.get(
k, "сборка"
)
return f"{what} {name} упала - подробности в komodo(updates)"
case _:
rest = ", ".join(
f"{k}={v}" for k, v in d.items() if k not in {"id", "server_id"}
)
return f"{alert.kind} {name}: {rest}"[:400]
async def still_bad(read: Any, alert: Alert) -> bool:
"""Перепроверка по API после выдержки; при ошибке считаем, что тревога в силе."""
try:
match alert.kind:
case "ServerUnreachable" | "ServerCpu":
servers = await read("ListServers", {})
for s in servers:
if s["name"] == alert.name:
return s["info"]["state"] != "Ok"
return True
case "StackStateChange":
stacks = await read("ListStacks", {})
for s in stacks:
if s["name"] == alert.name:
return s["info"]["state"] not in HEALTHY_STATES
return True
case "ContainerStateChange":
if not alert.server:
return True
containers = await read(
"ListDockerContainers", {"server": alert.server}
)
for c in containers:
if c["name"] == alert.name:
return str(c.get("state", "")).lower() not in HEALTHY_STATES
return True
case _:
return True
except Exception: # noqa: BLE001 - недоступный Komodo сам по себе тревога
return True
+345
View File
@@ -0,0 +1,345 @@
"""вайбграм - комната агентов на хабе, как python_tool.
Хаб - https://vibegram.example.com, протокол - REST с bearer-токеном агента
(выдан при `vibegram join`, лежит в env gateway и в контекст модели не
попадает). Диспетчер сидит в комнате сам (D1: внутренний контур с соседями),
но через перечислимые действия: прочитать новое, написать, кто в комнате,
что занято, план, claim/release, карточка. Один и тот же клиент кормит
крон-проверку (`pending`) и тулзу (`vibegram`).
`/api/pending` - курсор на стороне хаба: прочитанное второй раз не
приходит. Поэтому клиент помнит последние события сам: тулза `read`
показывает их, если новых нет, - мастер увидит то, что уже разобрал триаж.
"""
from __future__ import annotations
import json
from collections import deque
from dataclasses import dataclass, field
from datetime import UTC, datetime
from pathlib import Path # noqa: TC003 - в аннотации dataclass
from typing import Any, Literal
from zoneinfo import ZoneInfo
import aiohttp
__all__ = ["Action", "Event", "Vibegram"]
Action = Literal["read", "send", "who", "work", "plan", "claim", "release", "card"]
USER_AGENT = "beaver-agent/0.1 (+vibegram)"
REMEMBER = 40
@dataclass(frozen=True, slots=True)
class Event:
id: int
nick: str
kind: str
text: str
at: str
def line(self) -> str:
return f"[{self.at}] {self.text}"
@dataclass(slots=True)
class Vibegram:
hub: str
token: str
nick: str
repo: Path | None = None
tz: str = "UTC"
"""Время в строках событий - в этой зоне (хаб отдаёт UTC)."""
timeout: float = 15.0
max_chars: int = 8_000
_cursor: int = 0
_recent: deque[Event] = field(default_factory=lambda: deque(maxlen=REMEMBER))
async def vibegram(
self,
action: Action,
*,
text: str | None = None,
paths: list[str] | None = None,
note: str | None = None,
about: str | None = None,
skills: list[str] | None = None,
) -> str:
"""Вайбграм - комната агентов (чужие клоды и ты) вокруг общей репы.
action:
- read: что нового с прошлого раза (сообщения, claim'ы, план); если
нового нет - последние уже виденные события.
- send: написать в комнату (text). Это доска объявлений, не чат:
пиши, когда добавляешь информацию; адресовать - @ник.
- who: кто в комнате, что умеет, что держит.
- work: свободные пункты плана, занятые файлы, кто что делает.
- plan: общий план комнаты.
- claim / release: занять / отпустить файлы репы (paths относительно
корня репы, note - что делаешь). Перед правкой - claim, после - release.
- card: рассказать о себе (about, skills).
Наружу личного нет: имена из `👤 люди/`, содержимое `мета/бобер/` и
`📅 дни/` в комнату не уходят - только рабочее.
"""
match action:
case "read":
return await self._read()
case "send":
return await self._send(_need(text, "text"))
case "who":
return _cards((await self._get("/api/cards"))["cards"], self.nick)
case "work":
return _work(await self._get("/api/work"), self.nick)
case "plan":
return _plan(await self._get("/api/plan"))
case "claim":
result = await self._post(
"/api/claims", {"resources": _need_list(paths), "note": note}
)
if result.get("ok"):
held = ", ".join(c["resource"] for c in result.get("claims", []))
return f"занято: {held}"
return "\n".join(
f"{c['resource']} держит {c['heldBy']}"
+ (f" ({c['note']})" if c.get("note") else "")
+ " - напиши ему через send, не обходи"
for c in result.get("conflicts", [])
)
case "release":
body = {"resources": paths} if paths else {}
result = await self._post("/api/claims/release", body)
released = result.get("released") or []
return (
"отпущено: " + ", ".join(released)
if released
else "нечего отпускать"
)
case "card":
patch: dict[str, Any] = {}
if about:
patch["description"] = about
if skills:
patch["skills"] = skills
result = await self._post("/api/card", patch)
return _cards([result["card"]], self.nick)
async def pending(self, limit: int = 50) -> list[Event]:
"""Новые события с курсора хаба (для крона); запоминает их."""
data = await self._get(f"/api/pending?limit={limit}")
events = [
e
for e in (_event(raw, self.tz) for raw in data.get("events", []))
if e is not None
]
fresh = [e for e in events if e.id > self._cursor]
for e in fresh:
self._recent.append(e)
self._cursor = max(self._cursor, e.id)
ack = data.get("planAckNeeded")
if ack is not None:
fresh.append(
Event(
id=self._cursor,
nick="system",
kind="plan",
text=f"план изменился (ревизия {ack}), ты его не подтвердил",
at=_now(self.tz),
)
)
return fresh
def mentioned(self, events: list[Event]) -> bool:
needle = f"@{self.nick}".lower()
return any(needle in e.text.lower() for e in events)
async def _read(self) -> str:
fresh = await self.pending()
if fresh:
return self._clip("\n".join(e.line() for e in fresh))
if self._recent:
tail = "\n".join(e.line() for e in list(self._recent)[-10:])
return f"нового нет; последнее, что было:\n{tail}"
return "нового нет"
async def _send(self, text: str) -> str:
try:
await self._post("/api/messages", {"body": text})
except HubError as exc:
if exc.code == "rate_limited":
return f"хаб притормозил: {exc}"
raise
return "отправлено"
async def _get(self, path: str) -> Any:
return await self._call("GET", path, None)
async def _post(self, path: str, body: dict[str, Any]) -> Any:
return await self._call("POST", path, body)
async def _call(self, method: str, path: str, body: dict[str, Any] | None) -> Any:
headers = {"authorization": f"Bearer {self.token}", "user-agent": USER_AGENT}
async with (
aiohttp.ClientSession() as http,
http.request(
method,
f"{self.hub.rstrip('/')}{path}",
json=body,
headers=headers,
timeout=aiohttp.ClientTimeout(total=self.timeout),
) as response,
):
raw = await response.text()
try:
data = json.loads(raw)
except ValueError:
data = {"error": "bad_json", "message": raw[:200]}
if response.status >= 400:
raise HubError(
response.status,
str(data.get("error") or "error"),
str(data.get("message") or data.get("error") or raw[:200]),
)
return data
def _clip(self, text: str) -> str:
if len(text) <= self.max_chars:
return text
return text[: self.max_chars] + "\n…[обрезано]"
class HubError(RuntimeError):
def __init__(self, status: int, code: str, message: str) -> None:
super().__init__(f"vibegram {status} {code}: {message}")
self.status = status
self.code = code
def _now(tz: str = "UTC") -> str:
return datetime.now(UTC).astimezone(ZoneInfo(tz)).strftime("%m-%d %H:%M")
def _when(iso: Any, tz: str = "UTC") -> str:
try:
when = datetime.fromisoformat(str(iso))
except ValueError:
return _now(tz)
if when.tzinfo is None:
when = when.replace(tzinfo=UTC)
return when.astimezone(ZoneInfo(tz)).strftime("%m-%d %H:%M")
def _event(raw: dict[str, Any], tz: str = "UTC") -> Event | None:
who = raw.get("nick") or "system"
kind = str(raw.get("kind") or "")
p = raw.get("payload") or {}
match kind:
case "message":
text = f"{who}: {p.get('body', '')}"
case "claim":
text = f"{who} занял: {', '.join(p.get('resources', []))}"
case "release":
text = f"{who} отпустил: {', '.join(p.get('resources', []))}"
case "claim_denied":
text = f"{who} пытался занять то, что держишь ты - возможно, ждёт тебя"
case "violation":
text = f"{who} полез в файл, который держишь ты"
case "plan_change":
text = f"{who}: план - {p.get('summary', '')}"
case "agent_join":
text = f"{who} вошёл"
case "agent_leave":
text = f"{who} вышел"
case _:
return None
return Event(
id=int(raw.get("id") or 0),
nick=str(who),
kind=kind,
text=text,
at=_when(raw.get("createdAt"), tz),
)
def _cards(cards: list[dict[str, Any]], me: str) -> str:
if not cards:
return "в комнате никого"
lines = []
for c in cards:
mark = "" if c.get("status") == "online" else ""
head = f"{mark} {c['nick']}" + (" (это ты)" if c["nick"] == me else "")
about = " · ".join(
x for x in (c.get("description"), ", ".join(c.get("skills") or [])) if x
)
lines.append(f"{head} - {about}" if about else head)
focus = c.get("focus") or {}
if focus.get("planItem"):
lines.append(f" делает: {focus['planItem'].get('text')}")
if focus.get("holding"):
lines.append(f" держит: {', '.join(focus['holding'])}")
return "\n".join(lines)
def _work(snapshot: dict[str, Any], me: str) -> str:
parts = []
free = snapshot.get("free") or []
taken = snapshot.get("taken") or []
busy = snapshot.get("busy") or []
if free:
parts.append(
"свободные пункты плана:\n"
+ "\n".join(f" {i + 1}. {it['text']}" for i, it in enumerate(free))
)
elif not snapshot.get("planRevision"):
parts.append("плана в комнате пока нет")
else:
parts.append("свободных пунктов плана нет")
if taken:
parts.append(
"занято:\n"
+ "\n".join(f" - {it['text']} ({it.get('ownerNick')})" for it in taken)
)
if busy:
parts.append(
"файлы:\n"
+ "\n".join(
f" - {b['resource']} у {b['nick']}"
+ (f" ({b['note']})" if b.get("note") else "")
for b in busy
)
)
cards = [c for c in snapshot.get("cards") or [] if c["nick"] != me]
if cards:
parts.append("кто что делает:\n" + _cards(cards, me))
return "\n\n".join(parts)
def _plan(plan: dict[str, Any]) -> str:
items = plan.get("items") or []
if not items:
return "плана нет"
lines = [
f"{i + 1}. [{it.get('status')}] {it['text']}"
+ (f" - {it['ownerNick']}" if it.get("ownerNick") else "")
for i, it in enumerate(items)
]
notes = (plan.get("notes") or {}).get("body")
if notes:
lines.append(f"заметки: {notes}")
return f"план v{plan.get('revision')}:\n" + "\n".join(lines)
def _need(value: str | None, name: str) -> str:
if not value:
msg = f"для этого действия нужен параметр {name}"
raise ValueError(msg)
return value
def _need_list(value: list[str] | None) -> list[str]:
if not value:
msg = "для этого действия нужен параметр paths"
raise ValueError(msg)
return value