feat(policy,komodo,config): vault zone and firefly skill rules, komodo python_tool with exec, tests
This commit is contained in:
+292
@@ -0,0 +1,292 @@
|
||||
"""§3.7, §4.4: 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 = "sh"
|
||||
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` от
|
||||
пользователя контейнера. Долгое - запускай в фон с выводом в файл.
|
||||
|
||||
Деплой и рестарт - только когда об этом попросили, не «заодно».
|
||||
"""
|
||||
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": "Never"},
|
||||
}
|
||||
status, text = await self._post("terminal/execute", body)
|
||||
if status != 200:
|
||||
body["init"]["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)
|
||||
Reference in New Issue
Block a user