diff --git a/.env.example b/.env.example index a74feac..5fc08e2 100644 --- a/.env.example +++ b/.env.example @@ -6,8 +6,8 @@ ### Compose ################################################################### # Какие сервисы поднимать: db (postgres), gateway, obsidian (obsidian-headless), -# t3 (t3code-mcp - руки диспетчера на маке и dell) -COMPOSE_PROFILES=db,gateway,obsidian,t3 +# t3 (t3code-mcp - руки диспетчера на маке и dell), vibegram (клон репы комнаты) +COMPOSE_PROFILES=db,gateway,obsidian,t3,vibegram # Ветка/тег beaver-gateway, из которого собирается образ шлюза GATEWAY_REF=main @@ -80,6 +80,29 @@ T3CODE_MCP_REF=main KOMODO_URL= KOMODO_KEY= KOMODO_SECRET= +# Алерты Komodo → /hooks/komodo?token=<токен gateway со scope api> (алертер +# в infra/komodo, Custom endpoint; Komodo не умеет заголовки). Токен минтится +# в admin → Tokens, здесь его хранить не нужно. + +### Home Assistant (python_tool, с S12) ######################################## + +# HA на том же хосте с network_mode: host → host.docker.internal (extra_hosts) +HA_URL=http://host.docker.internal:8123 +# Профиль → Безопасность → долгоживущие токены +HA_TOKEN= + +### Вайбграм (комната агентов, с S12) ######################################### + +# Хаб и токен агента: один раз `vibegram join --nick <ник> --no-hooks` +# из клона репы комнаты на любой машине, токен - в ~/.vibegram/config.json +VIBEGRAM_HUB=https://vibegram.example.com +VIBEGRAM_TOKEN= +VIBEGRAM_NICK=beaver-test +# Репа комнаты (GitHub org/repo) - сервис vibegram-repo держит клон в томе +# `vibegram`, gateway видит его в /vibegram/ только на чтение +VIBEGRAM_REPO=org/room +# Токен GitHub с доступом к репе (приватная); пусто - публичная +VIBEGRAM_REPO_TOKEN= ### Порты хоста и внешний адрес ############################################### diff --git a/README.md b/README.md index 9f26575..f394272 100644 --- a/README.md +++ b/README.md @@ -63,8 +63,9 @@ Claude agents run on the Claude Agent SDK. Prompts are granules under and the order live in `config.py`; skills are the folders under `мета/бобер/скиллы/` (each becomes a plugin, keep `name:` in SKILL.md latin). The vault must be synced before the gateway can start. The model process runs as `beaver-runner` -with a whitelisted environment; the vault is read-only for it except -`мета/бобер` and `💬 чаты`. +with a whitelisted environment; the vault is mounted read-write (it is the +dispatcher's home), and `policy.py` keeps it out of `.obsidian`, `мета` +outside `мета/бобер`, and from deleting anything but its own files. ### 4. Mint a token diff --git a/config.py b/config.py index 0439258..afbd8be 100644 --- a/config.py +++ b/config.py @@ -30,7 +30,10 @@ 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 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, @@ -117,9 +120,18 @@ DEEP_SKILLS = (SKILLS / "общие", SKILLS / "vault") DISPATCHER_SKILLS = (SKILLS / "общие", SKILLS / "диспетчер", SKILLS / "vault") -# §3.7: зоны vault - запись только в мета/бобер, в 💬 чаты только новые файлы; -# firefly пишет только после открытого скилла. Маунт ro - первая линия. -ZONES = Zones(vault=VAULT, write=(BEAVER,), create=(CHATS_DIR,)) +# §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(), @@ -170,6 +182,39 @@ komodo_mcps = ( ) 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/ 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, + ) + 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") @@ -200,6 +245,8 @@ mcps: list[McpServerT] = [ McpServer.http(name="telegram", url=os.environ["BEAVERGRAM_MCP"]), *calendar_mcps, *komodo_mcps, + *ha_mcps, + *vibegram_mcps, *t3code_mcps, ] @@ -210,6 +257,7 @@ CLAUDE_MCPS = ( ExposedMcp(name="telegram"), *calendar_exposed, *komodo_exposed, + *ha_exposed, ) @@ -228,8 +276,9 @@ def dispatcher(name: str, model: str, effort: str | None = None) -> ClaudeAgent: 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), + # §4.4: t3code и вайбграм - руки диспетчера (мастер, ветки), глубоким + # не даются: наружу и в код ходит только он. + expose_mcps=(*CLAUDE_MCPS, *t3code_exposed, *vibegram_exposed), ) @@ -260,7 +309,7 @@ def distiller(name: str, model: str, effort: str | None = None) -> ClaudeAgent: tools=("Read", "Write"), disallowed_tools=DISTILLER_DISALLOWED, ), - policy=(vault_zones(ZONES),), + policy=(vault_zones(STRICT_ZONES),), ) @@ -506,28 +555,44 @@ async def rotate(run: JobRun) -> None: _vibegram_last_wake: datetime | None = None +_vibegram_backlog: list[str] = [] - -def vibegram_check() -> list[str]: - """Заглушка до M7: детерминированная проверка вайбграма, пока ничего нового.""" - return [] +VIBEGRAM_BRIEF = ( + "Новое в вайбграме (комната {room}, ты там {nick}; остальные - чужие агенты):\n" + "{items}\n" + "Будить мастера - inject(master, резюме до 3 строк: кто, что, чего ждёт). " + "Мастер прочитает подробности через vibegram(read) и ответит через " + "vibegram(send), если решит." +) async def vibegram(run: JobRun) -> None: + """§8.5: проверка → triage → в мастер ≤ 1/ч, если не ждут ответа.""" global _vibegram_last_wake # noqa: PLW0603 - счётчик «≤ 1/ч», живёт до рестарта - new = vibegram_check() - if not new: + 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) - if _vibegram_last_wake is not None and now - _vibegram_last_wake < timedelta( - hours=1 + 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="Новое в вайбграме:\n" + "\n".join(f"- {item}" for item in new), + text=VIBEGRAM_BRIEF.format( + room=VIBEGRAM_REPO.name, + nick=VIBEGRAM.nick, + items="\n".join(f"- {item}" for item in items[-40:]), + ), ) @@ -635,10 +700,58 @@ async def t3code_event(run: JobRun) -> None: 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 - summary = payload.get("raw") or {k: v for k, v in payload.items() if k != "trigger"} - await run.inject_master(f"Komodo: {summary}", urgency="urgent", origin="komodo") + 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 = [ @@ -647,7 +760,8 @@ jobs = [ Job("закрытие", close_idle, cron="20 4 * * *", critical=False), Job("память", memory, cron="30 4 * * 0", critical=False), Job("deploy", deploy, webhook=True), - Job("komodo", komodo_alert, webhook=True), + # алерты идут пачками (тревога и отбой подряд) - схлопывать нельзя + Job("komodo", komodo_alert, webhook=True, dedupe=False), # события идут подряд (вопрос и завершение) - схлопывать нельзя Job("t3code", t3code_event, webhook=True, dedupe=False), ] diff --git a/docker-compose.yml b/docker-compose.yml index 889cb0f..56f539a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -55,19 +55,15 @@ services: - ./policy.py:/config/policy.py:ro - ./mcps:/config/mcps:ro - ./config.json:/config/config.json:ro - # §3.7: vault только на чтение, rw - подмонтирования зон агента. - - vault:/vault:ro - - type: volume - source: vault - target: /vault/мета/бобер - volume: - subpath: мета/бобер - - type: volume - source: vault - target: /vault/💬 чаты - volume: - subpath: 💬 чаты + # §3.7: vault - дом диспетчера, rw целиком; границы (`.obsidian`, + # `мета` вне `бобер`, удаление чужого) держит policy.py (2026-08-30). + - vault:/vault + # клон репы комнаты вайбграма (сервис vibegram-repo), только читать + - vibegram:/vibegram:ro - claude-home:/home/beaver-runner/.claude + extra_hosts: + # Home Assistant и прочее с network_mode: host + - "host.docker.internal:host-gateway" healthcheck: test: ["CMD", "python", "-c", "import urllib.request; urllib.request.urlopen('http://127.0.0.1:62990/healthz', timeout=3)"] interval: 30s @@ -83,11 +79,10 @@ services: mkdir -p /home/$$runner/.claude [ -e /home/$$runner/.claude/.claude.json ] || echo '{}' > /home/$$runner/.claude/.claude.json chown -R $$runner /home/$$runner - for d in "/vault/мета/бобер" "/vault/💬 чаты"; do - setfacl -R -m "u:$$runner:rwX" -m "d:u:$$runner:rwX" "$$d" \ - || chown -R $$runner "$$d" \ - || echo "warning: cannot grant $$runner write access to $$d" - done + # ACL, не chown: Obsidian Sync переписывает файлы под root, а default-ACL + # на папках даёт runner'у права и на то, что появится позже. + setfacl -R -m "u:$$runner:rwX" -m "d:u:$$runner:rwX" /vault \ + || echo "warning: cannot grant $$runner write access to /vault" exec python -m beaver_gateway t3code-mcp: @@ -105,9 +100,41 @@ services: - ./t3code.toml:/config/t3code.toml:ro - t3code-state:/data + vibegram-repo: + container_name: beaver-vibegram-repo + image: alpine/git:latest + profiles: [vibegram] + <<: *restart + env_file: .env + volumes: + - vibegram:/vibegram + # Клон репы комнаты и pull раз в 10 минут. Токен - только заголовком + # через GIT_CONFIG_*, в .git/config он не попадает: том читает модель. + entrypoint: + - /bin/sh + - -c + - | + repo="$${VIBEGRAM_REPO:?VIBEGRAM_REPO=org/repo}" + dir="/vibegram/$${repo##*/}" + if [ -n "$$VIBEGRAM_REPO_TOKEN" ]; then + export GIT_CONFIG_COUNT=1 + export GIT_CONFIG_KEY_0="http.https://github.com/.extraheader" + export GIT_CONFIG_VALUE_0="AUTHORIZATION: basic $$(printf 'x-access-token:%s' "$$VIBEGRAM_REPO_TOKEN" | base64 | tr -d '\n')" + fi + while true; do + if [ -d "$$dir/.git" ]; then + git -C "$$dir" pull -q --ff-only || echo "vibegram-repo: pull failed" + else + git clone -q "https://github.com/$$repo.git" "$$dir" || echo "vibegram-repo: clone failed" + fi + chmod -R a+rX "$$dir" 2>/dev/null + sleep "$${VIBEGRAM_PULL_SECONDS:-600}" + done + volumes: t3code-state: postgres-data: vault: + vibegram: obsidian-config: claude-home: diff --git a/mcps/homeassistant.py b/mcps/homeassistant.py new file mode 100644 index 0000000..dbc9c6b --- /dev/null +++ b/mcps/homeassistant.py @@ -0,0 +1,169 @@ +"""§4.4: 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 diff --git a/mcps/komodo_alerts.py b/mcps/komodo_alerts.py new file mode 100644 index 0000000..74e2804 --- /dev/null +++ b/mcps/komodo_alerts.py @@ -0,0 +1,229 @@ +"""§4.5: алерты 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 diff --git a/mcps/vibegram.py b/mcps/vibegram.py new file mode 100644 index 0000000..03cbd1f --- /dev/null +++ b/mcps/vibegram.py @@ -0,0 +1,335 @@ +"""§4.4, §8.5, §12.D1: вайбграм - комната агентов на хабе, как 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 + +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 + 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 map(_event, 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(), + ) + ) + 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() -> str: + return datetime.now(UTC).strftime("%m-%d %H:%M") + + +def _when(iso: Any) -> str: + try: + return datetime.fromisoformat(str(iso)).strftime("%m-%d %H:%M") + except ValueError: + return _now() + + +def _event(raw: dict[str, Any]) -> 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")), + ) + + +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 diff --git a/policy.py b/policy.py index 948b0cd..398bf4d 100644 --- a/policy.py +++ b/policy.py @@ -67,6 +67,7 @@ MUTATING = frozenset( } ) COPYING = frozenset({"cp", "install", "rsync"}) +DELETING = frozenset({"rm", "rmdir", "unlink", "shred", "dd", "truncate"}) INPLACE = frozenset({"sed", "perl"}) REDIRECTS = frozenset({">", ">>", ">|", "&>", "&>>", ">&"}) SEPARATORS = frozenset({";", "&&", "||", "|", "&", "(", ")"}) @@ -74,18 +75,36 @@ SEPARATORS = frozenset({";", "&&", "||", "|", "&", "(", ")"}) @dataclass(frozen=True, slots=True) class Zones: - """Зоны vault: ``write`` - полная запись, ``create`` - только новые файлы.""" + """Зоны vault. + + ``write`` - всё, включая удаление; ``create`` - только новые файлы; + ``edit`` - создавать, править, переносить можно, удалять - нет; + ``protected`` - не трогать. Приоритет: write > protected > create > edit; + путь в vault вне всех зон - отказ. Vault целиком в ``edit`` - это + «дом» диспетчера (решение h, 2026-08-30): заметки, задачи, люди - его; + удаление чужого и `.obsidian`/`мета` остаются за Бобром. + """ vault: Path write: tuple[Path, ...] - create: tuple[Path, ...] + create: tuple[Path, ...] = () + edit: tuple[Path, ...] = () + protected: tuple[Path, ...] = () - def verdict(self, path: Path, *, creating: bool) -> Deny | None: + def verdict( + self, path: Path, *, creating: bool, deleting: bool = False + ) -> Deny | None: if not _under(path, self.vault): return None if any(_under(path, zone) for zone in self.write): return None rel = path.relative_to(self.vault) + for zone in self.protected: + if _under(path, zone): + return Deny( + reason=f"{rel}: «{zone.relative_to(self.vault)}» не трогаем - " + "это зона Бобра" + ) for zone in self.create: if _under(path, zone): if creating and not path.exists(): @@ -94,6 +113,17 @@ class Zones: reason=f"{rel}: в «{zone.relative_to(self.vault)}» можно только " "создавать новые файлы, существующие не трогаем" ) + for zone in self.edit: + if _under(path, zone): + if deleting: + zones = ", ".join( + f"«{z.relative_to(self.vault)}»" for z in self.write + ) + return Deny( + reason=f"{rel}: удалять можно только в {zones}; чужие " + "заметки не удаляем - скажи Бобру, он сам" + ) + return None zones = ", ".join(f"«{z.relative_to(self.vault)}»" for z in self.write) return Deny(reason=f"{rel}: vault только на чтение; писать можно в {zones}") @@ -125,8 +155,10 @@ def bash_zones(zones: Zones) -> PolicyRule: command = call.input.get("command") if not isinstance(command, str): return None - for target, creating in _bash_targets(command): - deny = zones.verdict(call.resolve(target), creating=creating) + for target, creating, deleting in _bash_targets(command): + deny = zones.verdict( + call.resolve(target), creating=creating, deleting=deleting + ) if deny is not None: return Deny(reason=f"Bash: {deny.reason}") return None @@ -134,7 +166,8 @@ def bash_zones(zones: Zones) -> PolicyRule: return rule -def _bash_targets(command: str) -> Iterator[tuple[str, bool]]: +def _bash_targets(command: str) -> Iterator[tuple[str, bool, bool]]: + """``(путь, создаёт, удаляет)`` для каждого пути, который команда трогает.""" for segment in _segments(command): plain: list[str] = [] i = 0 @@ -142,7 +175,7 @@ def _bash_targets(command: str) -> Iterator[tuple[str, bool]]: word = segment[i] if word in REDIRECTS: if i + 1 < len(segment) and _pathlike(segment[i + 1]): - yield segment[i + 1], True + yield segment[i + 1], True, False i += 2 continue plain.append(word) @@ -156,7 +189,7 @@ def _bash_targets(command: str) -> Iterator[tuple[str, bool]]: if word in COPYING: args = args[-1:] for arg in args: - yield arg, False + yield arg, False, word in DELETING break diff --git a/tests/test_komodo_alerts.py b/tests/test_komodo_alerts.py new file mode 100644 index 0000000..f442e66 --- /dev/null +++ b/tests/test_komodo_alerts.py @@ -0,0 +1,80 @@ +from datetime import timedelta + +from mcps.komodo_alerts import HOLD, AlertMemory, describe, parse, still_bad + + +def alert(kind, *, resolved=False, level="CRITICAL", **data): + data.setdefault("name", "ms-agents-proxy-aeza") + return {"level": level, "resolved": resolved, "data": {"type": kind, "data": data}} + + +def test_state_change_is_one_shot_and_read_by_target_state(): + up = parse( + alert( + "StackStateChange", resolved=True, **{"from": "running", "to": "unhealthy"} + ) + ) + down = parse( + alert( + "StackStateChange", resolved=True, **{"from": "unhealthy", "to": "running"} + ) + ) + assert up.bad and not down.bad + assert up.key == down.key + + +def test_flappy_alert_is_held_then_dropped_or_injected(): + mem = AlertMemory() + first = parse( + alert("ServerUnreachable", err={"error": "Timed out waiting for Ping"}) + ) + assert mem.verdict(first).action == "hold" + # повтор той же тревоги, пока держим - тишина + assert mem.verdict(first).action == "drop" + # отбой во время выдержки - тоже тишина, и память чистая + assert ( + mem.verdict(parse(alert("ServerUnreachable", resolved=True))).action == "drop" + ) + assert not mem.holding and not mem.told + # выдержка истекла, всё ещё плохо - срочно, и потом отбой доезжает + mem.verdict(first) + late = mem.after_hold(first, still_bad=True) + assert (late.action, late.urgency) == ("inject", "urgent") + mem.mark_told(first) + assert ( + mem.verdict(parse(alert("ServerUnreachable", resolved=True))).action + == "resolved" + ) + assert not mem.told + + +def test_informational_never_reaches_master_and_failures_go_normal(): + mem = AlertMemory() + assert mem.verdict(parse(alert("StackImageUpdateAvailable"))).action == "drop" + assert mem.verdict(parse(alert("ScheduleRun"))).action == "drop" + failed = mem.verdict(parse(alert("ProcedureFailed", name="Fleet Backup"))) + assert (failed.action, failed.urgency) == ("inject", "normal") + + +def test_describe_mentions_hold(): + a = parse(alert("ServerUnreachable", err={"error": "Timed out waiting for Ping"})) + text = describe(a, held=HOLD) + assert "ms-agents-proxy-aeza" in text and "5 мин" in text and "Timed out" in text + assert describe( + parse(alert("ServerDisk", name="dell", path="/", used_gb=90, total_gb=100)) + ) + + +async def test_still_bad_reads_current_state(): + async def read(kind, params): + assert kind == "ListServers" + return [{"name": "ms-agents-proxy-aeza", "info": {"state": "Ok"}}] + + a = parse(alert("ServerUnreachable")) + assert await still_bad(read, a) is False + + async def boom(kind, params): + raise RuntimeError("komodo down") + + assert await still_bad(boom, a) is True + assert HOLD == timedelta(minutes=5) diff --git a/tests/test_policy.py b/tests/test_policy.py index a5ceb2a..9a1cf7c 100644 --- a/tests/test_policy.py +++ b/tests/test_policy.py @@ -115,3 +115,65 @@ def test_firefly_requires_open_skill(vault): assert gate(call(vault, "mcp__firefly__list_account", state=state)) is None tracker(call(vault, "Skill", state=state, skill="vault:firefly")) assert gate(store) is None + + +@pytest.fixture +def home(vault: Path) -> Zones: + """Зоны диспетчера с 2026-08-30: vault - дом, удаление - только своё.""" + (vault / ".obsidian").mkdir() + (vault / "мета/шаблоны").mkdir() + return Zones( + vault=vault, + write=(vault / "мета/бобер",), + create=(vault / "💬 чаты",), + edit=(vault,), + protected=(vault / ".obsidian", vault / "мета"), + ) + + +def test_home_zones_allow_notes_but_protect_obsidian_and_meta(vault, home): + rule = vault_zones(home) + assert rule(call(vault, "Write", file_path=str(vault / "👤 люди/новый.md"))) is None + assert ( + rule(call(vault, "Edit", file_path=str(vault / "📅 дни/2026-08-29.md"))) is None + ) + assert rule(call(vault, "Write", file_path=str(vault / "мета/бобер/x.md"))) is None + assert ( + rule(call(vault, "Edit", file_path=str(vault / "💬 чаты/старый.md"))) + is not None + ) + deny = rule(call(vault, "Write", file_path=str(vault / ".obsidian/app.json"))) + assert isinstance(deny, Deny) and "Бобра" in deny.reason + assert ( + rule(call(vault, "Edit", file_path=str(vault / "мета/шаблоны/x.md"))) + is not None + ) + + +@pytest.mark.parametrize( + "command", + [ + "mv '👤 люди/x.md' '👤 люди/архив/x.md'", + "echo hi >> '📅 дни/2026-08-29.md'", + "sed -i 's/a/b/' '👤 люди/x.md'", + "rm мета/бобер/дни/old.md", + "mkdir '💻 проекты/новый'", + ], +) +def test_home_bash_edits_allowed(vault, home, command): + assert bash_zones(home)(call(vault, "Bash", command=command)) is None, command + + +@pytest.mark.parametrize( + "command", + [ + "rm '📅 дни/2026-08-29.md'", + "rm -rf 👤\\ люди", + "truncate -s 0 '👤 люди/x.md'", + "echo x > .obsidian/app.json", + "mv '👤 люди/x.md' мета/шаблоны/x.md", + ], +) +def test_home_bash_deletes_and_protected_denied(vault, home, command): + deny = bash_zones(home)(call(vault, "Bash", command=command)) + assert isinstance(deny, Deny), command diff --git a/tests/test_vibegram.py b/tests/test_vibegram.py new file mode 100644 index 0000000..7f2012c --- /dev/null +++ b/tests/test_vibegram.py @@ -0,0 +1,66 @@ +import pytest +from beaver_gateway.mcp.types import McpServer +from beaver_gateway.mcp.wrap import build_python_tool_server +from fastmcp import Client +from fastmcp.exceptions import ToolError + +from mcps.vibegram import Vibegram, _cards, _event, _work + + +def test_events_render_and_unknown_kinds_are_skipped(): + msg = _event( + { + "id": 61, + "nick": "claude-neighbour", + "kind": "message", + "payload": {"body": "@beaver-test привет"}, + "createdAt": "2026-08-30T10:00:00.000Z", + } + ) + assert ( + msg is not None + and msg.line() == "[08-30 10:00] claude-neighbour: @beaver-test привет" + ) + assert _event({"id": 1, "kind": "tree", "payload": {}}) is None + + +def test_mentioned_and_cards(): + v = Vibegram(hub="http://127.0.0.1:9", token="t", nick="beaver-test") + ev = _event( + { + "id": 2, + "nick": "x", + "kind": "message", + "payload": {"body": "эй @Beaver-Test"}, + } + ) + assert ev is not None and v.mentioned([ev]) + text = _cards( + [ + { + "nick": "beaver-test", + "status": "online", + "skills": [], + "focus": {"holding": ["a.md"]}, + } + ], + "beaver-test", + ) + assert "(это ты)" in text and "держит: a.md" in text + assert "плана в комнате пока нет" in _work( + {"free": [], "taken": [], "busy": [], "cards": [], "planRevision": 0}, "me" + ) + + +async def test_tool_schema_lists_actions_and_needs_text_for_send(): + tool = Vibegram(hub="http://127.0.0.1:9", token="t", nick="beaver-test") + server = build_python_tool_server( + McpServer.python_tool(name="vibegram", tools=[tool.vibegram]) + ) + async with Client(server) as client: + (spec,) = await client.list_tools() + assert {"read", "send", "who", "claim", "release"} <= set( + spec.inputSchema["properties"]["action"]["enum"] + ) + with pytest.raises(ToolError, match="text"): + await client.call_tool("vibegram", {"action": "send"})