refactor: split config.py into the beaver_agent package (vault, prompts, skills, policy, hands, agents, frontends, texts, memory, jobs)
This commit is contained in:
@@ -0,0 +1,25 @@
|
||||
"""Джобы: кроны считаются в локальном времени, вебхуки приходят на /hooks/<имя>."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from beaver_gateway.jobs.scheduler import Job
|
||||
|
||||
from beaver_agent.jobs.closing import close_idle
|
||||
from beaver_agent.jobs.deploy import deploy
|
||||
from beaver_agent.jobs.komodo import komodo_alert
|
||||
from beaver_agent.jobs.memory import curate, memory
|
||||
from beaver_agent.jobs.rotation import rotate
|
||||
from beaver_agent.jobs.t3code import t3code_event
|
||||
from beaver_agent.jobs.vibegram import vibegram
|
||||
|
||||
jobs = [
|
||||
Job("ротация", rotate, cron="0 * * * *"),
|
||||
Job("вайбграм", vibegram, cron="*/10 * * * *", critical=False),
|
||||
Job("закрытие", close_idle, cron="20 4 * * *", critical=False),
|
||||
Job("память", memory, cron="30 4 * * 0", critical=False),
|
||||
Job("куратор", curate, cron="0 2,8,12,16,20 * * *", critical=False),
|
||||
Job("deploy", deploy, webhook=True),
|
||||
Job("komodo", komodo_alert, webhook=True, dedupe=False),
|
||||
Job("t3code", t3code_event, webhook=True, dedupe=False),
|
||||
]
|
||||
"""Алерты и события t3code идут пачками (тревога и отбой подряд) - их не схлопывают."""
|
||||
@@ -0,0 +1,23 @@
|
||||
"""Ночью после ротации: глубокие чаты, тихие два дня, закрываются, не больше трёх."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from beaver_agent.vault import LAUNCH
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
_log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def close_idle(run: JobRun) -> None:
|
||||
for result in await run.close_idle(kind="deep", days=2, limit=3, since=LAUNCH):
|
||||
_log.info(
|
||||
"closed deep chat %s: digest=%s error=%s",
|
||||
result.conversation.external_id,
|
||||
result.digest.path if result.digest else None,
|
||||
result.error,
|
||||
)
|
||||
@@ -0,0 +1,33 @@
|
||||
"""`POST /hooks/deploy`: дождаться тишины в мастере и попросить Komodo передеплоить."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
from datetime import timedelta
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from beaver_agent.hands import KOMODO
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
_log = logging.getLogger(__name__)
|
||||
|
||||
STACK = os.environ.get("KOMODO_STACK", "beaver-agent")
|
||||
|
||||
|
||||
async def deploy(run: JobRun) -> None:
|
||||
master = await run.master()
|
||||
if master is not None and await run.conversations.busy(master):
|
||||
await run.retry_in(timedelta(seconds=60))
|
||||
return
|
||||
run.background(_deploy())
|
||||
|
||||
|
||||
async def _deploy() -> None:
|
||||
if KOMODO is None:
|
||||
_log.warning("deploy hook: KOMODO_URL/KEY/SECRET are not set, nothing deployed")
|
||||
return
|
||||
receipt = await KOMODO.execute("DeployStack", {"stack": STACK, "services": []})
|
||||
_log.info("deploy hook: komodo accepted %s", receipt.get("id"))
|
||||
@@ -0,0 +1,72 @@
|
||||
"""`POST /hooks/komodo`: алерты флота - флапы выдерживаются, остальное едет в мастер."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from beaver_agent.hands import KOMODO
|
||||
from beaver_agent.hands.komodo_alerts import (
|
||||
HOLD,
|
||||
AlertMemory,
|
||||
describe,
|
||||
parse,
|
||||
still_bad,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
_log = logging.getLogger(__name__)
|
||||
|
||||
ALERTS = AlertMemory()
|
||||
HINT = (
|
||||
"Это флот Бобра; ему такие алерты приходят и так (healthbot), тебя будят, "
|
||||
"когда через 5 минут само не встало. Посмотри komodo(status/logs), почини, "
|
||||
"если очевидно (restart стека), иначе - коротко скажи Бобру через say, что "
|
||||
"видишь. Не деплой и не чини «заодно»."
|
||||
)
|
||||
|
||||
|
||||
async def komodo_alert(run: JobRun) -> None:
|
||||
payload = run.payload
|
||||
if payload.get("held"):
|
||||
alert = parse(payload["alert"])
|
||||
bad = await still_bad(KOMODO.read, alert) if KOMODO else True
|
||||
verdict = ALERTS.after_hold(alert, still_bad=bad)
|
||||
if verdict.action == "drop":
|
||||
_log.info("komodo alert %s recovered during hold", alert.key)
|
||||
return
|
||||
ALERTS.mark_told(alert)
|
||||
await run.inject_master(
|
||||
f"Komodo: {describe(alert, held=HOLD)}\n{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{HINT}",
|
||||
urgency=verdict.urgency,
|
||||
origin="komodo",
|
||||
)
|
||||
@@ -0,0 +1,66 @@
|
||||
"""Память: воскресная перезапись состояния из хендаутов и куратор по расписанию."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from beaver_gateway.conversations.distill import LineCap
|
||||
|
||||
from beaver_agent.memory.curator import briefing
|
||||
from beaver_agent.memory.watch import CURATOR_WATCH
|
||||
from beaver_agent.vault import (
|
||||
CURATOR_JOURNAL,
|
||||
DAYS,
|
||||
REPLIES,
|
||||
STATE,
|
||||
STATE_MAX_LINES,
|
||||
TZ,
|
||||
VAULT,
|
||||
today,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
_log = logging.getLogger(__name__)
|
||||
|
||||
REWRITE_STATE = (
|
||||
"Воскресенье, {day}. Перепиши `{state}` из последних хендаутов: {handouts}. "
|
||||
"`## сейчас` - только невыводимое из vault, `## мои хвосты` - capability "
|
||||
"requests в работе и обещания длиннее дня. Меньше {max_lines} строк."
|
||||
)
|
||||
|
||||
|
||||
async def memory(run: JobRun) -> None:
|
||||
"""Потолок строк держит gateway: слишком длинный файл откатывается."""
|
||||
handouts = sorted(DAYS.glob("????-??-??.md"))[-7:] if DAYS.exists() else []
|
||||
if not handouts:
|
||||
_log.info("memory: no handouts yet, nothing to rewrite")
|
||||
return
|
||||
await run.spawn_job(
|
||||
agent="beaver-distiller",
|
||||
title="память",
|
||||
text=REWRITE_STATE.format(
|
||||
day=today(),
|
||||
state=STATE.relative_to(VAULT),
|
||||
handouts=", ".join(f"`{h.relative_to(VAULT)}`" for h in handouts),
|
||||
max_lines=STATE_MAX_LINES,
|
||||
),
|
||||
line_cap=LineCap(STATE, max_lines=STATE_MAX_LINES),
|
||||
)
|
||||
|
||||
|
||||
async def curate(run: JobRun) -> None:
|
||||
"""Крон собирает «что нового», модель раскладывает; пусто - джоб не спавнится."""
|
||||
text = briefing(
|
||||
vault=VAULT,
|
||||
journal=CURATOR_JOURNAL,
|
||||
replies_dir=REPLIES,
|
||||
watched=CURATOR_WATCH,
|
||||
tz=TZ,
|
||||
)
|
||||
if text is None:
|
||||
_log.info("curator: nothing new since the last run")
|
||||
return
|
||||
await run.spawn_job(agent="beaver-curator", title="куратор", text=text)
|
||||
@@ -0,0 +1,16 @@
|
||||
"""Раз в час: пора ли сменить мастера (ночь, возраст, размер контекста)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
_log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def rotate(run: JobRun) -> None:
|
||||
for master in await run.rotate():
|
||||
_log.info("rotated master -> %s", master.external_id)
|
||||
@@ -0,0 +1,51 @@
|
||||
"""`POST /hooks/t3code`: тред завершился, упал или спрашивает - срочно в мастер."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
QUESTION_HINT = (
|
||||
"Если ответ следует из задания - ответь сам: t3_answer(thread_id, "
|
||||
"{id вопроса: label}). Не знаешь - спроси Бобра через say своими словами, "
|
||||
"потом t3_answer. Тред стоит, пока не ответят."
|
||||
)
|
||||
|
||||
|
||||
async def t3code_event(run: JobRun) -> None:
|
||||
await run.inject_master(describe(run.payload), urgency="urgent", origin="t3code")
|
||||
|
||||
|
||||
def describe(p: dict[str, Any]) -> str:
|
||||
head = (
|
||||
f"t3code: тред «{p.get('title')}» ({p.get('machine')}, {p.get('project')}, "
|
||||
f"thread_id {p.get('thread_id')})"
|
||||
)
|
||||
event = p.get("event")
|
||||
if event in ("question", "approval"):
|
||||
lines = [f"{head} - спрашивает."]
|
||||
for q in (p.get("question") or {}).get("questions", []):
|
||||
multi = " (можно несколько)" if q.get("multi_select") else ""
|
||||
lines.append(
|
||||
f"- [{q['id']}] {q.get('header') or ''}: {q['question']}{multi}"
|
||||
)
|
||||
lines.extend(
|
||||
f" · {o['label']}"
|
||||
+ (f" - {o['description']}" if o.get("description") else "")
|
||||
for o in q.get("options", [])
|
||||
)
|
||||
lines.append(QUESTION_HINT)
|
||||
return "\n".join(lines)
|
||||
if event == "completed":
|
||||
lines = [f"{head} - завершён."]
|
||||
elif event == "interrupted":
|
||||
lines = [f"{head} - прерван."]
|
||||
else:
|
||||
lines = [f"{head} - упал: {p.get('error') or 'без текста ошибки'}."]
|
||||
if p.get("text"):
|
||||
lines.append(f"Ответ треда:\n{p['text']}")
|
||||
if p.get("files"):
|
||||
lines.append("Файлы: " + ", ".join(f["path"] for f in p["files"]))
|
||||
return "\n".join(lines)
|
||||
@@ -0,0 +1,53 @@
|
||||
"""Раз в 10 минут: новое в комнате → триаж → мастер, не чаще раза в час."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from beaver_agent.hands import VIBEGRAM, VIBEGRAM_REPO
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from beaver_gateway.jobs.scheduler import JobRun
|
||||
|
||||
BRIEF = (
|
||||
"Новое в вайбграме (комната {room}, ты там {nick}; остальные - чужие агенты):\n"
|
||||
"{items}\n"
|
||||
'Будить мастера - inject(conversation="master", urgency="wake", '
|
||||
"text=резюме до 3 строк: кто, что, чего ждёт). Мастер прочитает подробности "
|
||||
"через vibegram(read) и ответит через vibegram(send), если решит. "
|
||||
"Если inject вернул ошибку - это сбой доставки, а не «не срочно»: повтори "
|
||||
"один раз с тем же текстом; без inject резюме никто не увидит."
|
||||
)
|
||||
|
||||
_last_wake: datetime | None = None
|
||||
_backlog: list[str] = []
|
||||
|
||||
|
||||
async def vibegram(run: JobRun) -> None:
|
||||
global _last_wake # noqa: PLW0603
|
||||
if VIBEGRAM is None:
|
||||
return
|
||||
events = await VIBEGRAM.pending()
|
||||
_backlog.extend(e.line() for e in events)
|
||||
if not _backlog:
|
||||
return
|
||||
now = datetime.now(UTC)
|
||||
addressed = VIBEGRAM.mentioned(events)
|
||||
if (
|
||||
not addressed
|
||||
and _last_wake is not None
|
||||
and now - _last_wake < timedelta(hours=1)
|
||||
):
|
||||
return
|
||||
_last_wake = now
|
||||
items, _backlog[:] = list(_backlog), []
|
||||
await run.spawn_job(
|
||||
agent="beaver-triage",
|
||||
title="вайбграм",
|
||||
text=BRIEF.format(
|
||||
room=VIBEGRAM_REPO.name,
|
||||
nick=VIBEGRAM.nick,
|
||||
items="\n".join(f"- {item}" for item in items[-40:]),
|
||||
),
|
||||
)
|
||||
Reference in New Issue
Block a user