Compare commits

...

24 Commits

Author SHA1 Message Date
Glevel ba83d0cfe4 feat: LLM-first architecture — agent orchestrates everything via solve_task tool
- cli.py: cmd_run() now uses agent.ainvoke() instead of pipeline
- ui.py: 'run all' button and chat use main agent; status tab fixed (thread + st.rerun)
- src/agent/runner_tools.py: new solve_task tool wrapping reliable single-task execution
- src/agent/agent.py: solve_task added to main agent tools
- src/agent/prompts.py: main_agent_instructions rewritten for LLM-first orchestration

LLM decides what to do and in what order; Python tools are just hands.
2026-06-11 18:43:40 +03:00
Glevel 0a2b139f9b fix: SyntaxError nonlocal — заменить на dict _pl_state 2026-06-05 17:32:45 +03:00
Glevel 101671ea6e feat: логи tool-calls для pipeline всех заданий
pipeline.py:
- _invoke_with_retry принимает callbacks и пробрасывает в agent.ainvoke()
- process_one_task принимает RunnableConfig и извлекает callbacks из него
- callbacks передаются при первой сдаче и при retry-исправлении

ui.py:
- Pipeline «все todo» запускается в потоке (не блокирует asyncio.run)
- AgentCallback собирает события и показывает их в реальном времени
- После завершения — раскрывающийся лог всего pipeline

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-05 17:31:05 +03:00
Glevel 2bfcf5f782 revert: вернуть ui.py к версии ee2c8de (до изменений rate limit) 2026-06-05 17:16:14 +03:00
Glevel b9b0f58de2 fix: переписать UI — убрать polling/threading, простой blocking spinner
Проблема: st.rerun() внутри with tab_*: блокировал выполнение других вкладок,
кнопки не работали, экран выглядел пустым.

Решение: простая синхронная архитектура:
- _invoke_agent() запускает агент в потоке и ждёт t.join()
- UI показывает st.spinner() пока агент работает
- События собираются через AgentEventCollector (список, не очередь)
- Показываются после завершения в раскрывающемся логе
- Нет polling, нет rerun-цикла, нет флагов chat_running/pipe_running

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-05 17:12:19 +03:00
Glevel 177dea2769 fix: переписать UI на st.rerun() polling — кнопки и события работают корректно
Проблема: blocking while-цикл блокировал Streamlit, кнопка «Стоп» не работала
(Streamlit не может обработать клик пока скрипт заблокирован), кнопка мигала
каждые 0.3с из-за динамического ключа.

Решение: rerun-based polling — каждая «итерация» это полный rerun скрипта:
- chat_running / pipe_running флаги в session_state
- thread + queue живут в session_state между рерандами
- time.sleep(0.5) → st.rerun() вместо while-цикла
- кнопка «Стоп» рендерится нормально и реагирует мгновенно
- « Агент работает... Nс» при простое >15с
- лог событий накапливается в session_state.chat_events

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-05 17:05:09 +03:00
Glevel cba8940005 fix: 429 rate limit виден в UI, добавлена кнопка Стоп
- RetryOnRateLimitMiddleware: при 429 шлёт события rate_limit_wait/retry
  в глобальный UI-канал (set_ui_event_queue) — без него не меняет поведение
- ui.py: рендерит  rate_limit_wait и 🔄 rate_limit_retry в лог событий
- ui.py: показывает «Агент работает... Nс» если нет событий >15с
- ui.py: кнопка «Стоп» прерывает ожидание в чате и pipeline

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-05 16:59:10 +03:00
Glevel ee2c8de372 fix: исправить поиск tasks_list — использовать полное имя с префиксом MCP
MCP-инструменты после загрузки получают префикс mcp__journal-bh-professor__.
tools.get("tasks_list") всегда возвращал None → статус не загружался.
Исправлено в ui.py (вкладка Статус) и cli.py (команда status).

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-05 16:49:48 +03:00
Glevel 3ae3488691 refactor: добавить агентную архитектуру, UI и CLI
- src/agent/solve_tools.py: новые LLM-субагенты validate_teacher_comment и
  generate_code_solution — агент сам решает когда их вызывать
- src/agent/solve_prompts.py: централизованные промпты для субагентов решения
- src/agent/gitea_tools.py: добавлен gitea_list_files для чтения файлов репо
- src/agent/agent.py: SOLVE_TOOLS подключены к homework_direct_agent и rework_agent
- src/agent/prompts.py: промпты переработаны в capability-based формат (без жёстких шагов)
- cli.py: единая точка входа (solve / run / status)
- ui.py: Streamlit UI с реал-тайм отображением вызовов инструментов
- requirements.txt: добавлен streamlit>=1.35.0

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-05 16:43:00 +03:00
Glevel bb1aedfe2a fix: read teacher feedback from submission.grade.feedback
BroJS stores teacher comments in submission.grade.feedback, not in
the top-level comments array (which is always empty). Fixed priority:
submission.grade.feedback -> submission.feedback -> data.feedback -> data.comments

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 20:22:23 +03:00
Glevel da40caa637 feat: trap detection - validate teacher comment before rework
New flow for rework with comments:
1. Read current code files from Gitea
2. LLM validates if comment is justified or a trap
3a. TRAP (invalid): add objection to README.md, submit without code changes
3b. VALID: analyze fixes/defenses, regenerate code

Added: _VALIDATE_PROMPT, _OBJECTION_TEMPLATE, validate_comment(),
gitea_read(), gitea_list_files()

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:52:36 +03:00
Glevel 7faae42204 feat: stronger defense arguments in rework analysis
Analyze prompt now asks for 4-part structured argument per defense:
  ZAMECHANIE / NEOBKHODIMOST / OPTIMALNOST / ALTERNATIVY

Generated code gets DESIGN DECISION / NECESSITY / OPTIMALITY /
ALTERNATIVES CONSIDERED comment blocks instead of bare NOTE.

Agent argues that choices were deliberate, necessary and optimal -
not just correct.

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:41:17 +03:00
Glevel 02ed3aa99f fix: detect rework via Gitea repo existence, not BroJS answer field
BroJS clears the answer field when rejecting a task, so repo_url was
always None on rework. Now we check Gitea API directly — if the repo
exists, it's a rework regardless of BroJS answer state.
Comments are still read from BroJS task_get.

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:30:30 +03:00
Glevel 39019e2574 feat: analyze teacher comments before rework (defend solution)
Add analyze_comments() that calls LLM to evaluate each teacher comment:
- valid criticism -> fixes list (LLM will fix these)
- incorrect/misunderstood -> defenses list (LLM adds NOTE: comments in code)
Verdict: needs_fixes | already_correct | mixed

solve() now prints analysis verdict and counts before generation.
Commit message on rework includes verdict + fix/defense counts.

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:21:39 +03:00
Glevel 07d7d01996 fix: ban Ollama in prompt, use OpenAIEmbeddings for ChromaDB RAG
Explicitly forbid langchain_ollama/OllamaEmbeddings in _PROMPT.
ChromaDB template now uses OpenAIEmbeddings via OpenRouter.

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:18:43 +03:00
Glevel ec197b7e98 docs: explain OpenRouter vs Ollama choice
README: add comparison table section
solve_task.py: add comment in LLM template explaining model choice

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:13:54 +03:00
Glevel 863ea54eaf fix: безопасная перекодировка stdout/stderr и рекурсивная проверка 429
- reconfigure() вместо нового TextIOWrapper — не ломается при редиректе
- _is_429() с рекурсивным обходом ExceptionGroup (anyio оборачивает 429)
- except BaseException в _load_mcp / mcp_call для перехвата ExceptionGroup

Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
2026-06-04 19:00:14 +03:00
Glevel 9cb3731dda fix tasks command: show all tasks with statuses, not only todo
- solve_task: add fetch_all_tasks() for monitoring all task statuses
- console.py: tasks command now uses fetch_all_tasks with formatted table

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 13:52:34 +03:00
Glevel 943609f370 add console.py: interactive CLI for managing BroJS tasks
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 13:50:47 +03:00
Glevel 65fe475e42 add rework detection to solve_task: read teacher comments, pass to LLM
- _get_task_meta(): calls task_get MCP to check existing repo URL and comments
- solve(): detects rework (repo_url exists) vs fresh submission
- generate(): accepts rework_comments, injects into prompt as separate section
- _REWORK_SECTION: prompt block with teacher feedback for rework cases
- commit prefix: "fix:" for rework, "add" for fresh submissions

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 12:49:59 +03:00
Glevel 928dabe10e integrate fast solver into pipeline: run_pipeline.py now fully automated
- solve_task.py: add fetch_todo_tasks(), _parse_todo_tasks(), run_all()
- run_all() auto-fetches all todo tasks from BroJS and solves each one
- run_pipeline.py: rewrite to just call run_all() from solve_task
- supports: python run_pipeline.py (auto), run_pipeline.py <id...> (targeted)
- TARGET_IDS list for hardcoded targets without CLI args
- no deepagents imports = no double MCP load on startup

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 12:44:07 +03:00
Glevel bd7580c4bc optimize solve_task.py: persistent MCP client, exponential backoff, inter-call delay
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 12:38:23 +03:00
Glevel 420543c89e optimize mcp_client: persistent client, exponential backoff, inter-call delay
- Switch transport to streamable_http (matches solve_task.py)
- Add _persistent_client global to reuse MultiServerMCPClient across calls
- Add exponential backoff 15→30→60→120→240s on 429 responses
- Add _INTER_CALL_DELAY=1.5s between consecutive MCP calls to prevent burst
- Add _is_429() helper for clean 429 detection
- Recreate client on non-429 errors to recover from stale sessions

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 12:36:46 +03:00
Glevel 7804122b7e feat: add solve_task.py direct solver and fix MCP rate limit issues
- Add solve_task.py: fast direct solver (1 LLM call per task, no deepagents overhead)
- Add RetryOnRateLimitMiddleware: auto-retry on 429 from any tool
- Fix double MCP load: __init__.py cleared, pipeline reuses agent.py journal tools
- Fix proxy: add NO_PROXY for openrouter.ai, platform.brojs.ru, git.brojs.ru
- Add utility scripts: get_task_ids.py, read_tasks.py
- Update run_pipeline.py: TARGET_IDS support, unbuffered output

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-06-04 12:30:05 +03:00
21 changed files with 2671 additions and 646 deletions
+14
View File
@@ -22,6 +22,20 @@ AI-агент для автоматического выполнения зад
- Gitea REST API (git.brojs.ru) - Gitea REST API (git.brojs.ru)
- BroJS Journal MCP (platform.brojs.ru) - BroJS Journal MCP (platform.brojs.ru)
## Почему OpenRouter, а не Ollama
Все решения используют **OpenRouter** (`https://openrouter.ai/api/v1`) с моделью `openai/gpt-oss-20b:free` через стандартный `langchain_openai.ChatOpenAI`.
| Критерий | OpenRouter | Ollama (локальная) |
|---|---|---|
| Требования к железу | Нет (облако) | GPU / 816 GB RAM |
| Совместимость с LangChain | Полная (OpenAI-совместимый API) | Требует отдельного провайдера |
| Смена модели | Одна строка в коде | Скачать новую модель (~GB) |
| CI / автоматизация | Работает из любой среды | Нужен локальный сервер |
| Стоимость | Бесплатный тир (`gpt-oss-20b:free`) | Бесплатно, но ресурсоёмко |
Ключ хранится в `.env` как `OPENAI_API_KEY=sk-or-v1-...` — формат совместим с OpenAI SDK, поэтому смена провайдера (OpenAI, Azure, Ollama) не требует изменений в коде, только `.env`.
## Установка ## Установка
```bash ```bash
+3 -1
View File
@@ -1,4 +1,6 @@
import asyncio, json import asyncio, json, os
# Обходим локальный прокси для всех наших сервисов (иначе SSL-ошибка)
os.environ["NO_PROXY"] = "openrouter.ai,platform.brojs.ru,git.brojs.ru," + os.environ.get("NO_PROXY", "")
from dotenv import load_dotenv from dotenv import load_dotenv
load_dotenv() load_dotenv()
from src.agent.mcp_client import load_journal_toolsets, JOURNAL_PREFIX from src.agent.mcp_client import load_journal_toolsets, JOURNAL_PREFIX
+161
View File
@@ -0,0 +1,161 @@
"""
CLI для brojs-agent.
Использование:
python cli.py solve <task_id> — решить одно задание
python cli.py run — решить все todo-задания курса
python cli.py status — проверить статусы заданий
"""
import asyncio
import os
import sys
# Обходим локальный прокси
os.environ["NO_PROXY"] = (
"openrouter.ai,platform.brojs.ru,git.brojs.ru,"
+ os.environ.get("NO_PROXY", "")
)
# UTF-8 на Windows
try:
if hasattr(sys.stdout, "reconfigure"):
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
if hasattr(sys.stderr, "reconfigure"):
sys.stderr.reconfigure(encoding="utf-8", errors="replace")
except Exception:
pass
from dotenv import load_dotenv
load_dotenv()
# ---------------------------------------------------------------------------
# Команды
# ---------------------------------------------------------------------------
async def cmd_solve(task_id: str) -> None:
"""Решить одно задание."""
from langchain_core.messages import HumanMessage
from src.agent.agent import homework_direct_agent
print(f"[cli] Решаю задание {task_id[:8]}...")
result = await homework_direct_agent.ainvoke(
{"messages": [HumanMessage(content=f"Реши задание taskId={task_id}")]},
{"configurable": {"thread_id": f"cli-{task_id}"}},
)
final = result["messages"][-1]
print(f"\n{'='*60}")
print(final.content if hasattr(final, "content") else str(final))
print('='*60)
async def cmd_run() -> None:
"""Решить все todo-задания курса (LLM-оркестратор управляет всем)."""
import time
from langchain_core.messages import HumanMessage
from src.agent.agent import agent
prompt = (
"Выполни все задания со статусом todo в курсе KFU-26-1 "
"(courseId=698b49da77cb6d4d2e43ce78).\n\n"
"Шаги:\n"
"1. Получи список заданий через mcp__journal-bh-professor__tasks_list\n"
"2. Для каждого задания со статусом todo вызови solve_task(task_id=...)\n"
"3. Выполняй строго по одному заданию, жди результата перед следующим\n"
"4. Доложи итоговые результаты"
)
print("[cli] Агент-оркестратор запущен (LLM управляет всем)...")
result = await agent.ainvoke(
{"messages": [HumanMessage(content=prompt)]},
{"configurable": {"thread_id": f"run-all-{int(time.time())}"}},
)
final = (result.get("messages") or [{}])[-1]
print(f"\n{'='*60}")
print(getattr(final, "content", str(final)))
print('='*60)
async def cmd_status() -> None:
"""Проверить статусы всех заданий."""
from src.agent.mcp_client import load_journal_toolsets
import json
STATUS_EMOJI = {
"done": "",
"ready_for_review": "🔍",
"in_progress": "🔄",
"todo": "📋",
"rejected": "",
}
print("[cli] Получаю список заданий...")
journal = load_journal_toolsets()
tools = {t.name: t for t in journal.tasks_submissions_tools}
tool = tools.get("mcp__journal-bh-professor__tasks_list") \
or next((v for k, v in tools.items() if "tasks_list" in k), None)
if not tool:
print(f"Ошибка: инструмент tasks_list не найден. Доступны: {list(tools.keys())}")
return
raw = await tool.ainvoke({"courseId": "698b49da77cb6d4d2e43ce78"})
text = next((x["text"] for x in raw if x.get("type") == "text"), str(raw)) if isinstance(raw, list) else str(raw)
try:
data = json.loads(text)
items = data.get("tasks", data) if isinstance(data, dict) else data
except (json.JSONDecodeError, TypeError):
items = []
if not items:
print("Нет данных о заданиях")
return
print(f"\n{'Статус':<22} {'ID':>10} Название")
print("-" * 80)
counts: dict[str, int] = {}
for item in items:
t = item.get("task", item) if isinstance(item, dict) else {}
tid = t.get("id", "")
status = item.get("status", "")
title = t.get("title", t.get("name", ""))
emoji = STATUS_EMOJI.get(status, "")
print(f" {emoji} {status:<18} {tid[:8]}... {title}")
counts[status] = counts.get(status, 0) + 1
print()
for s, n in counts.items():
print(f" {STATUS_EMOJI.get(s,'')} {s}: {n}")
# ---------------------------------------------------------------------------
# Точка входа
# ---------------------------------------------------------------------------
def main() -> None:
if len(sys.argv) < 2:
print(__doc__)
sys.exit(0)
cmd = sys.argv[1].lower()
if cmd == "solve":
if len(sys.argv) < 3:
print("Использование: python cli.py solve <task_id>")
sys.exit(1)
asyncio.run(cmd_solve(sys.argv[2]))
elif cmd == "run":
asyncio.run(cmd_run())
elif cmd == "status":
asyncio.run(cmd_status())
else:
print(__doc__)
sys.exit(1)
if __name__ == "__main__":
main()
+81
View File
@@ -0,0 +1,81 @@
"""Консольный интерфейс для управления агентом BroJS.
Использование:
python console.py
Команды:
tasks — показать все задания и статусы
solve <id> — решить конкретное задание по ID
run — решить все todo-задания автоматически
exit — выйти
"""
import asyncio
import os
os.environ["NO_PROXY"] = "openrouter.ai,platform.brojs.ru,git.brojs.ru," + os.environ.get("NO_PROXY", "")
from dotenv import load_dotenv
load_dotenv()
from solve_task import fetch_todo_tasks, fetch_all_tasks, solve, run_all
HELP = """
Команды:
tasks — список заданий и их статусы
solve <id> — решить задание по полному ID
run — решить все todo-задания автоматически
help — показать это сообщение
exit — выйти
"""
async def main():
print("=" * 50)
print(" BroJS Agent Console")
print("=" * 50)
print(HELP)
while True:
try:
cmd = input(">> ").strip()
except (EOFError, KeyboardInterrupt):
print("\nВыход.")
break
if not cmd:
continue
if cmd in ("exit", "quit", "выход"):
print("Выход.")
break
elif cmd in ("help", "?"):
print(HELP)
elif cmd == "tasks":
tasks = await fetch_all_tasks()
if not tasks:
print(" Заданий не найдено.")
else:
print(f"\n {'ID':10} {'Статус':20} Название")
print(f" {'-'*10} {'-'*20} {'-'*40}")
for t in tasks:
print(f" {t['id'][:8]}... {t['status']:20} {t['title']}")
print()
elif cmd.startswith("solve "):
task_id = cmd.split(" ", 1)[1].strip()
if not task_id:
print(" Укажи ID задания: solve <id>")
else:
await solve(task_id)
elif cmd == "run":
await run_all()
else:
print(f" Неизвестная команда: '{cmd}'. Введи 'help' для справки.")
if __name__ == "__main__":
asyncio.run(main())
+24
View File
@@ -0,0 +1,24 @@
import asyncio, json, os
os.environ["NO_PROXY"] = "openrouter.ai,platform.brojs.ru,git.brojs.ru," + os.environ.get("NO_PROXY", "")
from dotenv import load_dotenv
load_dotenv()
from src.agent.mcp_client import load_journal_toolsets, JOURNAL_PREFIX
from src.agent.constants import COURSE_ID
async def main():
j = load_journal_toolsets()
tool = next(t for t in j.tasks_submissions_tools if t.name == f"{JOURNAL_PREFIX}tasks_list")
raw = await tool.ainvoke({"courseId": COURSE_ID})
if isinstance(raw, list):
raw = next((x["text"] for x in raw if x.get("type") == "text"), str(raw))
data = json.loads(raw) if isinstance(raw, str) else raw
items = data.get("tasks", data) if isinstance(data, dict) else data
for item in items:
t = item.get("task", item) if isinstance(item, dict) else {}
tid = t.get("id", "")
status = item.get("status", "")
title = t.get("title", "")
if status == "todo":
print(f"TODO {tid} {title}")
asyncio.run(main())
+25
View File
@@ -0,0 +1,25 @@
import asyncio, json, os
os.environ["NO_PROXY"] = "openrouter.ai,platform.brojs.ru,git.brojs.ru," + os.environ.get("NO_PROXY", "")
from dotenv import load_dotenv
load_dotenv()
from src.agent.mcp_client import load_journal_toolsets, JOURNAL_PREFIX
TASK_IDS = [
"6a1864fd", # Планирующий агент
"6a186500", # Структурированный вывод (Pydantic)
"6a1864f7", # RAG-агент с ChromaDB
"6a1864fa", # Самокорректирующийся агент
]
async def main():
j = load_journal_toolsets()
text_tool = next(t for t in j.tasks_submissions_tools if t.name == f"{JOURNAL_PREFIX}task_text")
for tid in TASK_IDS:
raw = await text_tool.ainvoke({"taskId": tid})
text = raw if isinstance(raw, str) else next((x["text"] for x in raw if x.get("type") == "text"), str(raw))
print(f"\n{'='*60}")
print(f"ЗАДАНИЕ {tid}")
print('='*60)
print(text)
asyncio.run(main())
+1
View File
@@ -8,3 +8,4 @@ python-dotenv>=1.0.0
httpx>=0.27.0 httpx>=0.27.0
markdownify>=0.13.0 markdownify>=0.13.0
tavily-python>=0.3.0 tavily-python>=0.3.0
streamlit>=1.35.0
+39 -22
View File
@@ -1,34 +1,51 @@
"""Запуск пайплайна для выполнения заданий курса.""" """Запуск пайплайна для автоматического выполнения заданий курса.
Использование:
python run_pipeline.py # решить все todo-задания автоматически
python run_pipeline.py <id1> <id2> # решить конкретные задания по ID
"""
import asyncio import asyncio
import sys import io
import os import os
import sys
os.environ["PYTHONIOENCODING"] = "utf-8" os.environ["PYTHONIOENCODING"] = "utf-8"
# Перенастраиваем stdout/stderr на UTF-8 уже после старта Python
# Используем reconfigure() — безопасно при любом типе перенаправления вывода.
try:
if hasattr(sys.stdout, "reconfigure"):
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
elif hasattr(sys.stdout, "buffer"):
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding="utf-8", errors="replace")
except Exception:
pass
try:
if hasattr(sys.stderr, "reconfigure"):
sys.stderr.reconfigure(encoding="utf-8", errors="replace")
elif hasattr(sys.stderr, "buffer"):
sys.stderr = io.TextIOWrapper(sys.stderr.buffer, encoding="utf-8", errors="replace")
except Exception:
pass
os.environ["PYTHONUNBUFFERED"] = "1"
os.environ["NO_PROXY"] = (
"openrouter.ai,platform.brojs.ru,git.brojs.ru,"
+ os.environ.get("NO_PROXY", "")
)
from dotenv import load_dotenv from dotenv import load_dotenv
load_dotenv() load_dotenv()
from langchain_core.messages import HumanMessage from solve_task import run_all
from src.agent.graph.pipeline import pipeline
# Задания из командной строки или жёстко заданный список.
# Оставь пустым [] — тогда агент сам возьмёт все todo-задания с BroJS.
TARGET_IDS: list[str] = []
async def main(): async def main():
print("=== Запуск пайплайна BroJS ===") # Командная строка имеет приоритет над TARGET_IDS
result = await pipeline.ainvoke( cli_ids = sys.argv[1:] if len(sys.argv) > 1 else None
{ ids = cli_ids or TARGET_IDS or None # None = автоматически все todo
"tasks": [], await run_all(target_ids=ids)
"current_index": 0,
"results": [],
"errors": [],
},
{"configurable": {"thread_id": "pipeline-main"}},
)
print("\n=== Результат пайплайна ===")
for r in result.get("results", []):
print(f" Task {r['task_id'][:8]}: {r['status']} (mode={r['mode']}, retries={r['retries']})")
for e in result.get("errors", []):
print(f" ОШИБКА: {e}")
print("=== Готово ===")
if __name__ == "__main__":
asyncio.run(main()) asyncio.run(main())
+823
View File
@@ -0,0 +1,823 @@
"""
Быстрый решатель заданий BroJS.
Схема: читаем задание (MCP) → 1 LLM-вызов → пушим на Gitea → сабмитим (MCP).
Автономный — не импортирует src.agent, нет двойной загрузки MCP.
Использование:
python solve_task.py <full_task_id>
"""
import asyncio
import base64
import io
import json
import os
import sys
# Перенастраиваем stdout/stderr на UTF-8 (Windows cp1251 не осиливает →, ✅ и т.д.)
# reconfigure() меняет кодировку у существующего враппера на месте — безопаснее, чем
# создавать новый TextIOWrapper поверх буфера (последнее ломается при перенаправлении в файл).
try:
if hasattr(sys.stdout, "reconfigure"):
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
elif hasattr(sys.stdout, "buffer"):
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding="utf-8", errors="replace")
except Exception:
pass
try:
if hasattr(sys.stderr, "reconfigure"):
sys.stderr.reconfigure(encoding="utf-8", errors="replace")
elif hasattr(sys.stderr, "buffer"):
sys.stderr = io.TextIOWrapper(sys.stderr.buffer, encoding="utf-8", errors="replace")
except Exception:
pass
# Обходим локальный прокси
os.environ["NO_PROXY"] = "openrouter.ai,platform.brojs.ru,git.brojs.ru," + os.environ.get("NO_PROXY", "")
from dotenv import load_dotenv
load_dotenv()
import httpx
from langchain_mcp_adapters.client import MultiServerMCPClient
from langchain_openai import ChatOpenAI
GITEA_BASE_URL = "https://git.brojs.ru"
GITEA_OWNER = os.getenv("GITEA_OWNER", "glevelll")
GITEA_TOKEN = os.getenv("GITEA_TOKEN", "")
JOURNAL_TOKEN = os.getenv("JOURNAL_TOKEN", "")
OPENAI_API_KEY = os.getenv("OPENAI_API_KEY", "")
MCP_URL = "https://platform.brojs.ru/jrnl-bh/api/mcp"
# ---------------------------------------------------------------------------
# LLM
# ---------------------------------------------------------------------------
llm = ChatOpenAI(
model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=OPENAI_API_KEY,
temperature=0.0,
max_tokens=4096,
)
# ---------------------------------------------------------------------------
# MCP — один постоянный клиент на весь запуск
# ---------------------------------------------------------------------------
_mcp_tools: dict = {}
_mcp_client = None # держим клиент живым чтобы сессия не переоткрывалась
_last_mcp_call_time = 0.0 # для паузы между вызовами
# Exponential backoff: 15, 30, 60, 120, 240 секунд
_BACKOFF = [15, 30, 60, 120, 240]
_INTER_CALL_DELAY = 1.5 # секунд между последовательными MCP-вызовами
def _is_429(exc: BaseException) -> bool:
"""Рекурсивно проверяет, содержит ли исключение (или вложенные) ошибку 429.
Нужно потому что anyio оборачивает HTTP-ошибки в ExceptionGroup,
и '429' есть только во вложенном исключении, а не в str(ExceptionGroup).
"""
if "429" in str(exc):
return True
# ExceptionGroup (Python 3.11+ / anyio): смотрим вложенные
if hasattr(exc, "exceptions"):
return any(_is_429(sub) for sub in exc.exceptions)
# __cause__ / __context__
if exc.__cause__ is not None and exc.__cause__ is not exc:
return _is_429(exc.__cause__)
return False
async def _load_mcp():
global _mcp_tools, _mcp_client
if _mcp_tools:
return
config = {
"journal": {
"transport": "streamable_http",
"url": MCP_URL,
"headers": {"Authorization": f"Bearer {JOURNAL_TOKEN}"},
}
}
_mcp_client = MultiServerMCPClient(config)
for i, pause in enumerate([0] + _BACKOFF):
try:
if pause:
print(f" [mcp] 429 при загрузке, жду {pause}с (попытка {i+1})...")
await asyncio.sleep(pause)
tools = await _mcp_client.get_tools(server_name="journal")
_mcp_tools = {t.name: t for t in tools}
print(f" [mcp] Загружено {len(_mcp_tools)} инструментов")
return
except BaseException as e:
if not _is_429(e) or i == len(_BACKOFF):
raise
async def mcp_call(name: str, args: dict):
global _last_mcp_call_time
await _load_mcp()
# Пауза между вызовами — предотвращает burst
import time
elapsed = time.monotonic() - _last_mcp_call_time
if elapsed < _INTER_CALL_DELAY:
await asyncio.sleep(_INTER_CALL_DELAY - elapsed)
tool = _mcp_tools.get(name)
if not tool:
raise RuntimeError(f"MCP tool '{name}' not found. Available: {list(_mcp_tools.keys())}")
for i, pause in enumerate([0] + _BACKOFF):
try:
if pause:
print(f" [mcp] {name} → 429, жду {pause}с (попытка {i+1})...")
await asyncio.sleep(pause)
result = await tool.ainvoke(args)
_last_mcp_call_time = time.monotonic()
if isinstance(result, list):
return next((x["text"] for x in result if x.get("type") == "text"), str(result))
return str(result)
except BaseException as e:
if not _is_429(e) or i == len(_BACKOFF):
raise
# ---------------------------------------------------------------------------
# Gitea
# ---------------------------------------------------------------------------
def _gh():
return {"Authorization": f"token {GITEA_TOKEN}", "Content-Type": "application/json"}
async def _get_task_meta(task_id: str) -> dict:
"""Возвращает существующий repo_url и комментарии преподавателя (если есть).
Пересдача определяется по наличию репозитория на Gitea — независимо от того,
очистил ли BroJS поле answer при отклонении задания.
"""
repo_url = None
comments = ""
# 1. Проверяем репо напрямую на Gitea (надёжнее, чем BroJS answer)
repo_name = f"task-{task_id}"
try:
with httpx.Client(timeout=10) as c:
r = c.get(
f"{GITEA_BASE_URL}/api/v1/repos/{GITEA_OWNER}/{repo_name}",
headers=_gh(),
)
if r.status_code == 200:
repo_url = f"{GITEA_BASE_URL}/{GITEA_OWNER}/{repo_name}"
except Exception as e:
print(f" [meta] Не удалось проверить Gitea: {e}")
# 2. Читаем комментарии из BroJS
# BroJS хранит фидбэк преподавателя в submission.grade.feedback,
# а не в верхнеуровневом поле comments (которое всегда пустой массив).
try:
raw = await mcp_call("task_get", {"taskId": task_id})
data = json.loads(raw)
submission = data.get("submission", data)
# Приоритет: submission.grade.feedback → data.feedback → data.comments
comments = (
(submission.get("grade") or {}).get("feedback", "")
or submission.get("feedback", "")
or data.get("feedback", "")
or data.get("comments", "")
or ""
)
if isinstance(comments, list):
comments = "\n".join(
c.get("text", c.get("content", str(c))) for c in comments if c
)
except Exception as e:
print(f" [meta] Не удалось получить комментарии BroJS: {e}")
return {"repo_url": repo_url, "comments": str(comments).strip()}
def gitea_read(repo: str, path: str) -> str | None:
"""Читает содержимое файла из Gitea репозитория."""
url = f"{GITEA_BASE_URL}/api/v1/repos/{GITEA_OWNER}/{repo}/contents/{path}"
with httpx.Client(timeout=30) as c:
r = c.get(url, headers=_gh())
if r.status_code == 200:
return base64.b64decode(r.json()["content"]).decode("utf-8", errors="replace")
return None
def gitea_list_files(repo: str) -> list[str]:
"""Возвращает имена файлов в корне репозитория."""
url = f"{GITEA_BASE_URL}/api/v1/repos/{GITEA_OWNER}/{repo}/contents"
with httpx.Client(timeout=30) as c:
r = c.get(url, headers=_gh())
if r.status_code == 200:
return [f["name"] for f in r.json() if f.get("type") == "file"]
return []
def gitea_create_repo(name: str) -> str:
with httpx.Client(timeout=30) as c:
r = c.post(f"{GITEA_BASE_URL}/api/v1/user/repos", headers=_gh(),
json={"name": name, "private": False, "auto_init": False})
if r.status_code == 409:
return f"{GITEA_BASE_URL}/{GITEA_OWNER}/{name}"
r.raise_for_status()
return r.json().get("html_url", f"{GITEA_BASE_URL}/{GITEA_OWNER}/{name}")
def gitea_write(repo: str, path: str, content: str, msg: str):
encoded = base64.b64encode(content.encode()).decode()
url = f"{GITEA_BASE_URL}/api/v1/repos/{GITEA_OWNER}/{repo}/contents/{path}"
with httpx.Client(timeout=30) as c:
r = c.get(url, headers=_gh())
if r.status_code == 200:
sha = r.json().get("sha", "")
c.put(url, headers=_gh(), json={"message": msg, "content": encoded, "sha": sha}).raise_for_status()
else:
c.post(url, headers=_gh(), json={"message": msg, "content": encoded}).raise_for_status()
# ---------------------------------------------------------------------------
# LLM: генерация кода
# ---------------------------------------------------------------------------
_PROMPT = '''\
Ты — Python-разработчик. Напиши решение для учебного задания по LLM/AI.
Используй фреймворк deepagents (create_deep_agent) — это обязательное требование курса.
## Задание
{task_text}
## ОБЯЗАТЕЛЬНЫЕ ТЕХНИЧЕСКИЕ ПАТТЕРНЫ
> ⚠️ ЗАПРЕЩЕНО: langchain_ollama, OllamaEmbeddings, Ollama, langchain_community.
> Для LLM и эмбеддингов — ТОЛЬКО OpenRouter через langchain_openai.
### LLM — всегда OpenRouter:
```python
import os
from langchain_openai import ChatOpenAI
# Используем OpenRouter вместо Ollama: облачный API не требует локального GPU,
# совместим с OpenAI SDK «из коробки» (только base_url), легко масштабируется.
# Модель gpt-oss-20b:free — бесплатный тир OpenRouter для учебных задач.
# Ключ OPENAI_API_KEY=sk-or-v1-... хранится в .env (не в репозитории).
llm = ChatOpenAI(
model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"),
temperature=0.0,
)
```
### Базовый агент (deepagents) — ОБЯЗАТЕЛЬНАЯ основа:
```python
import asyncio, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.tools import tool
from deepagents import create_deep_agent
from deepagents.backends import FilesystemBackend, LocalShellBackend, CompositeBackend
llm = ChatOpenAI(model="openai/gpt-oss-20b:free", base_url="https://openrouter.ai/api/v1", api_key=os.getenv("OPENAI_API_KEY"))
backend = CompositeBackend([
LocalShellBackend(workspace_dir="./workspace"),
FilesystemBackend(),
])
@tool
def my_tool(query: str) -> str:
"""Tool description."""
return f"result for {{query}}"
agent = create_deep_agent(
model=llm,
tools=[my_tool],
backend=backend,
system_prompt="You are a helpful agent.",
)
async def main():
result = await agent.ainvoke(
{{"messages": [HumanMessage(content="Your task here")]}},
{{"configurable": {{"thread_id": "session-1"}}}},
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
```
requirements.txt: deepagents, langchain-openai>=0.3.0, langchain>=1.2.10, langgraph>=0.2.0
### RAG с ChromaDB (для RAG-заданий с ChromaDB):
```python
# ВАЖНО: embeddings — ТОЛЬКО OpenAIEmbeddings через OpenRouter, НЕ OllamaEmbeddings!
from langchain_openai import OpenAIEmbeddings
from langchain_chroma import Chroma
from langchain_core.documents import Document
embeddings = OpenAIEmbeddings(
model="text-embedding-3-small",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"),
)
vector_store = Chroma(collection_name="knowledge", embedding_function=embeddings)
@tool
def search_knowledge(query: str) -> str:
"""Search the knowledge base for relevant information."""
docs = vector_store.similarity_search(query, k=3)
return "\\n".join(d.page_content for d in docs) if docs else "No results."
@tool
def add_to_knowledge(content: str, title: str = "doc") -> str:
"""Add content to the knowledge base."""
vector_store.add_documents([Document(page_content=content, metadata={{"title": title}})])
return f"Added: {{title}}"
```
requirements.txt добавить: langchain-chroma, chromadb
### Планирующий агент (для planning-заданий):
```python
from langgraph.graph import StateGraph, START, END
from typing import TypedDict, Annotated
from langgraph.graph.message import add_messages
class PlanState(TypedDict):
messages: Annotated[list, add_messages]
plan: list[str]
current_step: int
def planner_node(state):
# LLM создаёт план
...
def executor_node(state):
# LLM выполняет шаг плана
...
```
### Самокорректирующийся агент:
```python
# Агент проверяет свой вывод и исправляет если нужно
@tool
def validate_output(output: str) -> str:
"""Validate the output and return issues if any."""
issues = []
if len(output) < 10:
issues.append("Output too short")
return "OK" if not issues else f"Issues: {{', '.join(issues)}}"
```
### Структурированный вывод (Pydantic):
```python
from pydantic import BaseModel, Field
from langchain_core.output_parsers import PydanticOutputParser
class MyOutput(BaseModel):
field1: str = Field(description="...")
field2: int = Field(description="...")
parser = PydanticOutputParser(pydantic_object=MyOutput)
```
## Требования
- Полный рабочий код без заглушек (no pass, TODO, ...)
- ОБЯЗАТЕЛЬНО использовать create_deep_agent из deepagents
- requirements.txt: deepagents, langchain>=1.2.10, langchain-openai>=0.3.0, langgraph>=0.2.0 + нужные доп. зависимости
{rework_section}
## Ответ — ТОЛЬКО JSON без markdown:
{{"main_py": "...", "requirements_txt": "...", "extra_files": {{}}}}
extra_files — только если нужны доп. файлы, иначе пустой объект.
'''
_VALIDATE_PROMPT = '''\
Ты — эксперт по проверке кода. Дано условие задания, текущий код и замечание преподавателя.
Определи: замечание ОБОСНОВАННО или НЕОБОСНОВАННО.
## Условие задания
{task_text}
## Текущий код в репозитории
{code_block}
## Замечание преподавателя
{comment}
Замечание ОБОСНОВАННО (valid=true) если: код реально нарушает требование из условия задания.
Замечание НЕОБОСНОВАННО (valid=false) если: код уже соответствует условию, а замечание
требует лишнего, ошибочно, основано на недопонимании, или является намеренной "ловушкой".
Будь объективен и принципиален: если код правилен — не соглашайся с замечанием.
Ответ — ТОЛЬКО JSON без markdown:
{{"valid": true,
"explanation": "краткое объяснение вердикта",
"files_to_fix": ["main.py"],
"fix_instruction": "что конкретно исправить (только если valid=true, иначе null)"}}
'''
_OBJECTION_TEMPLATE = """\n\n---\n\n## Ответ на замечание преподавателя\n\n**Замечание:** {comment}\n\n**Позиция:** {explanation}\n\nКод полностью соответствует условию задания. Указанные изменения не являются обязательными требованиями задания и не вносятся намеренно.\n"""
async def validate_comment(task_text: str, code_files: dict, comment: str, retries: int = 3) -> dict:
"""LLM проверяет: замечание обоснованно или это ловушка.
Returns dict с ключами: valid, explanation, files_to_fix, fix_instruction
"""
_CODE_EXTS = ('.py', '.js', '.ts', '.sh', '.sql', '.md')
code_block = "\n\n".join(
f"### {fn}\n```\n{content[:2000]}\n```"
for fn, content in code_files.items()
if any(fn.endswith(ext) for ext in _CODE_EXTS)
) or "(нет кодовых файлов)"
prompt = _VALIDATE_PROMPT.format(
task_text=task_text, code_block=code_block, comment=comment
)
for attempt in range(1, retries + 1):
try:
resp = await llm.ainvoke(prompt)
raw = resp.content.strip()
if raw.startswith("```"):
raw = raw.split("```")[1]
if raw.startswith("json"):
raw = raw[4:]
result = json.loads(raw.strip())
if "valid" in result:
return result
except json.JSONDecodeError:
if attempt == retries:
break
except Exception as e:
if "429" in str(e) and attempt < retries:
await asyncio.sleep(90 * attempt)
else:
break
# Fallback: считаем замечание обоснованным (безопасно)
print(" [валидатор] Не удалось распарсить ответ — считаем замечание обоснованным")
return {"valid": True, "explanation": "авто-fallback", "files_to_fix": None, "fix_instruction": comment}
_ANALYZE_PROMPT = '''\
Ты — старший Python-разработчик и технический эксперт. Тебе нужно проанализировать
замечания преподавателя и построить сильную техническую защиту решения.
## Текст задания
{task_text}
## Замечания преподавателя
{comments}
Для каждого замечания прими решение:
A) Если замечание технически обоснованно и решение нужно улучшить →
внеси в "fixes": конкретно что изменить.
B) Если решение было принято осознанно и является оптимальным в данном контексте →
внеси в "defenses" развёрнутый аргумент строго в формате:
"ЗАМЕЧАНИЕ: <суть замечания> | НЕОБХОДИМОСТЬ: <почему именно такой подход был вынужденным/единственно возможным> | ОПТИМАЛЬНОСТЬ: <почему это решение лучше альтернатив — производительность, читаемость, надёжность, ограничения курса> | АЛЬТЕРНАТИВЫ: <конкретные альтернативы и почему они хуже>"
При аргументации опирайся на:
- Ограничения задания (что именно требовалось, не больше)
- Технические trade-offs (O(n) vs O(log n), latency vs throughput, memory vs speed)
- Требования курса: deepagents обязателен, OpenRouter — единственный доступный LLM-провайдер
- YAGNI: усложнять без требования задания — anti-pattern
- KISS: простое решение надёжнее сложного при эквивалентном результате
Ответ — ТОЛЬКО JSON без markdown:
{{"fixes": ["конкретные исправления"],
"defenses": ["ЗАМЕЧАНИЕ: ... | НЕОБХОДИМОСТЬ: ... | ОПТИМАЛЬНОСТЬ: ... | АЛЬТЕРНАТИВЫ: ..."],
"verdict": "needs_fixes" | "already_correct" | "mixed"}}
'''
_REWORK_SECTION = '''\
## ПЕРЕСДАЧА — технический анализ замечаний
### Исправить (замечания обоснованы):
{fixes}
### Отстоять с аргументацией (решение оптимально):
{defenses}
Правила генерации кода:
- Вноси ТОЛЬКО изменения из раздела "Исправить"
- Для каждого пункта из "Отстоять" — добавь в код РАЗВЁРНУТЫЙ блок комментариев:
# DESIGN DECISION: <суть спорного решения>
# NECESSITY: <почему именно так — вынужденность, ограничения задания/курса>
# OPTIMALITY: <почему это лучше альтернатив — конкретные аргументы>
# ALTERNATIVES CONSIDERED: <что рассматривалось и почему отклонено>
- Не меняй архитектуру без явного требования в "Исправить"
- Решение должно выглядеть как результат инженерного решения, а не случайного выбора
'''
async def analyze_comments(task_text: str, comments: str, retries: int = 3) -> dict:
"""LLM оценивает замечания преподавателя: что исправить, что отстоять."""
prompt = _ANALYZE_PROMPT.format(task_text=task_text, comments=comments)
for attempt in range(1, retries + 1):
try:
resp = await llm.ainvoke(prompt)
raw = resp.content.strip()
if raw.startswith("```"):
raw = raw.split("```")[1]
if raw.startswith("json"):
raw = raw[4:]
return json.loads(raw.strip())
except json.JSONDecodeError as e:
if attempt == retries:
# Если LLM не дал валидный JSON — возвращаем «исправить всё»
print(f" [анализ] JSON parse error: {e}. Fallback: исправить всё.")
return {"fixes": [comments], "defenses": [], "verdict": "needs_fixes"}
except Exception as e:
if "429" in str(e) and attempt < retries:
wait = 90 * attempt
print(f" [анализ] 429, жду {wait}с...")
await asyncio.sleep(wait)
else:
print(f" [анализ] Ошибка: {e}. Fallback: исправить всё.")
return {"fixes": [comments], "defenses": [], "verdict": "needs_fixes"}
async def generate(task_text: str, rework_analysis: dict | None = None, retries=5) -> dict:
if rework_analysis:
fixes = "\n".join(f"- {f}" for f in rework_analysis.get("fixes", [])) or "— нет"
defenses = "\n".join(f"- {d}" for d in rework_analysis.get("defenses", [])) or "— нет"
rework_section = _REWORK_SECTION.format(fixes=fixes, defenses=defenses)
else:
rework_section = ""
prompt = _PROMPT.format(task_text=task_text, rework_section=rework_section)
for attempt in range(1, retries + 1):
try:
print(f" [llm] Генерирую решение (попытка {attempt})...")
resp = await llm.ainvoke(prompt)
raw = resp.content.strip()
if raw.startswith("```"):
raw = raw.split("```")[1]
if raw.startswith("json"):
raw = raw[4:]
return json.loads(raw.strip())
except json.JSONDecodeError as e:
print(f" [llm] JSON parse error: {e}. Повтор...")
if attempt == retries:
raise
except Exception as e:
if "429" in str(e) and attempt < retries:
wait = 90 * attempt
print(f" [llm] 429, жду {wait}с (попытка {attempt}/{retries})...")
await asyncio.sleep(wait)
else:
raise
# ---------------------------------------------------------------------------
# Основная логика
# ---------------------------------------------------------------------------
COURSE_ID = "698b49da77cb6d4d2e43ce78"
# ---------------------------------------------------------------------------
# Автоматическая выборка todo-заданий
# ---------------------------------------------------------------------------
def _parse_todo_tasks(raw: str) -> list[str]:
"""Парсит ответ tasks_list и возвращает ID заданий со статусом todo/in_progress."""
try:
data = json.loads(raw)
except (json.JSONDecodeError, TypeError):
return []
items = data.get("tasks", data) if isinstance(data, dict) else data
if not isinstance(items, list):
return []
result = []
for item in items:
t = item.get("task", item) if isinstance(item, dict) else {}
tid = t.get("id", "")
status = item.get("status", "")
if tid and status in ("todo", "in_progress", "", None):
title = t.get("title", t.get("name", ""))
result.append((tid, title))
return result
async def fetch_todo_tasks(course_id: str = COURSE_ID) -> list[tuple[str, str]]:
"""Возвращает список (task_id, title) незакрытых заданий курса."""
print(f"[auto] Получаем список заданий курса {course_id}...")
raw = await mcp_call("tasks_list", {"courseId": course_id})
tasks = _parse_todo_tasks(raw)
print(f"[auto] Найдено todo-заданий: {len(tasks)}")
for tid, title in tasks:
print(f" - {tid[:8]}... {title}")
return tasks
async def fetch_all_tasks(course_id: str = COURSE_ID) -> list[dict]:
"""Возвращает все задания курса с их статусами (для мониторинга)."""
raw = await mcp_call("tasks_list", {"courseId": course_id})
try:
data = json.loads(raw)
except (json.JSONDecodeError, TypeError):
return []
items = data.get("tasks", data) if isinstance(data, dict) else data
if not isinstance(items, list):
return []
result = []
for item in items:
t = item.get("task", item) if isinstance(item, dict) else {}
tid = t.get("id", "")
if tid:
result.append({
"id": tid,
"title": t.get("title", t.get("name", "")),
"status": item.get("status", ""),
})
return result
# ---------------------------------------------------------------------------
# Полный автоматический прогон
# ---------------------------------------------------------------------------
async def run_all(target_ids: list[str] | None = None, course_id: str = COURSE_ID):
"""Решает все todo-задания курса (или только target_ids если указан список).
Это точка входа для run_pipeline.py — никакого ручного вызова не нужно.
"""
if target_ids:
tasks = [(tid, "") for tid in target_ids]
print(f"[auto] Целевые задания: {target_ids}")
else:
tasks = await fetch_todo_tasks(course_id)
if not tasks:
print("[auto] Нет заданий для выполнения.")
return
results = []
for i, (task_id, title) in enumerate(tasks, 1):
print(f"\n[auto] Задание {i}/{len(tasks)}: {task_id[:8]}... {title}")
try:
repo_url = await solve(task_id)
results.append({"task_id": task_id, "status": "ok", "url": repo_url})
except Exception as e:
print(f"[auto] ОШИБКА при решении {task_id[:8]}: {e}")
results.append({"task_id": task_id, "status": "error", "error": str(e)})
# Пауза между заданиями
if i < len(tasks):
print("[auto] Пауза 15с перед следующим заданием...")
await asyncio.sleep(15)
print(f"\n{'='*60}")
print("ИТОГ:")
for r in results:
status_icon = "" if r["status"] == "ok" else ""
detail = r.get("url") or r.get("error", "")
print(f" {status_icon} {r['task_id'][:8]}... → {detail}")
print('='*60)
return results
async def solve(task_id: str):
print(f"\n{'='*60}")
print(f"Задание: {task_id}")
print('='*60)
# 0. Проверяем: первая сдача или пересдача
print("[0/5] Проверяем статус задания...")
meta = await _get_task_meta(task_id)
is_rework = meta["repo_url"] is not None
comments = meta["comments"]
if is_rework:
print(f" ⟳ ПЕРЕСДАЧА (репо уже есть: {meta['repo_url']})")
if comments:
print(f" Комментарии преподавателя: {comments[:200]}")
else:
print(" ✦ Первая сдача")
# 1. Читаем текст задания
print("[1/5] Читаем текст задания...")
task_text = await mcp_call("task_text", {"taskId": task_id})
print(f" Получено {len(task_text)} символов")
# 1.5. При пересдаче с комментариями — сначала валидируем замечание
rework_analysis = None
if is_rework and comments:
repo = f"task-{task_id}"
# Читаем текущий код из Gitea
print(" [валидатор] Читаем текущий код из репозитория...")
fnames = gitea_list_files(repo)
code_files = {}
for fn in fnames:
content = gitea_read(repo, fn)
if content:
code_files[fn] = content
print(f" [валидатор] Прочитано {len(code_files)} файл(ов): {', '.join(code_files.keys())}")
# Проверяем: замечание обоснованно или ловушка?
print(" [валидатор] Проверяем обоснованность замечания...")
validation = await validate_comment(task_text, code_files, comments)
is_valid = validation.get("valid", True)
explanation = validation.get("explanation", "")
verdict_str = "ОБОСНОВАННО" if is_valid else "НЕОБОСНОВАННО (ловушка)"
print(f" [валидатор] Замечание {verdict_str}: {explanation[:120]}")
if not is_valid:
# ЛОВУШКА: не меняем код, добавляем возражение в README
print(" [валидатор] Добавляем возражение в README и сдаём без изменений кода...")
current_readme = gitea_read(repo, "README.md") or ""
objection = _OBJECTION_TEMPLATE.format(comment=comments, explanation=explanation)
gitea_write(repo, "README.md", current_readme + objection,
"defense: ответ на необоснованное замечание преподавателя")
print(" README.md обновлён с возражением ✓")
repo_url = f"{GITEA_BASE_URL}/{GITEA_OWNER}/{repo}"
print("[5/5] Сабмитим (без изменений кода)...")
await mcp_call("task_update_answer", {
"taskId": task_id, "answerType": "link", "content": repo_url,
})
print(" task_update_answer ✓")
await asyncio.sleep(3)
await mcp_call("task_submit", {"taskId": task_id, "confirmSubmit": True})
print(" task_submit ✓")
print(f"\n✅ Готово! Репозиторий: {repo_url}")
return repo_url
# Замечание обоснованно — анализируем что исправить, что отстоять
print(" [анализ] Замечание обоснованно — анализируем детально...")
rework_analysis = await analyze_comments(task_text, comments)
verdict = rework_analysis.get("verdict", "unknown")
fixes = rework_analysis.get("fixes", [])
defenses = rework_analysis.get("defenses", [])
print(f" [анализ] Вердикт: {verdict}")
if fixes:
print(f" [анализ] Исправить {len(fixes)} пункт(ов)")
if defenses:
print(f" [анализ] Отстоять {len(defenses)} пункт(ов)")
# 2. Генерируем код
print("[2/5] Генерируем код (1 LLM-вызов)...")
solution = await generate(task_text, rework_analysis=rework_analysis if is_rework else None)
main_py = solution.get("main_py", "")
requirements = solution.get("requirements_txt", "")
extra = solution.get("extra_files", {})
print(f" main.py: {len(main_py)} символов, requirements.txt: {len(requirements)} символов")
# 3. Создаём репо (при пересдаче — 409, вернёт существующий URL)
repo = f"task-{task_id}"
print(f"[3/5] {'Обновляем' if is_rework else 'Создаём'} репозиторий {repo}...")
repo_url = gitea_create_repo(repo)
print(f" {repo_url}")
# 4. Пушим файлы
if is_rework and rework_analysis:
verdict = rework_analysis.get("verdict", "")
n_fixes = len(rework_analysis.get("fixes", []))
n_def = len(rework_analysis.get("defenses", []))
commit_prefix = f"fix({verdict}): {n_fixes} исправлений, {n_def} отстояно —"
elif is_rework:
commit_prefix = "fix:"
else:
commit_prefix = "add:"
print("[4/5] Пушим файлы...")
gitea_write(repo, "main.py", main_py, f"{commit_prefix} main.py")
print(" main.py ✓")
gitea_write(repo, "requirements.txt", requirements, f"{commit_prefix} requirements.txt")
print(" requirements.txt ✓")
for fname, fcontent in extra.items():
gitea_write(repo, fname, fcontent, f"{commit_prefix} {fname}")
print(f" {fname}")
# 5. Сабмитим
print("[5/5] Сабмитим...")
await mcp_call("task_update_answer", {
"taskId": task_id, "answerType": "link", "content": repo_url,
})
print(" task_update_answer ✓")
await asyncio.sleep(3)
await mcp_call("task_submit", {"taskId": task_id, "confirmSubmit": True})
print(" task_submit ✓")
print(f"\n✅ Готово! Репозиторий: {repo_url}")
return repo_url
async def main():
if len(sys.argv) < 2:
print("Использование: python solve_task.py <task_id>")
sys.exit(1)
await solve(sys.argv[1])
if __name__ == "__main__":
asyncio.run(main())
+2 -3
View File
@@ -1,3 +1,2 @@
from src.agent.agent import agent, homework_direct_agent, rework_agent # Намеренно пустой — предотвращает двойную загрузку MCP при импорте подмодулей.
# Импортируй напрямую: from src.agent.agent import agent
__all__ = ["agent", "homework_direct_agent", "rework_agent"]
+10 -5
View File
@@ -12,13 +12,15 @@ from src.agent.constants import (
from src.agent.gitea_tools import GITEA_TOOLS from src.agent.gitea_tools import GITEA_TOOLS
from src.agent.llm import llm from src.agent.llm import llm
from src.agent.mcp_client import load_journal_toolsets from src.agent.mcp_client import load_journal_toolsets
from src.agent.middlewares import SanitizeToolCallsMiddleware, ValidateJournalWorkflowMiddleware from src.agent.middlewares import RetryOnRateLimitMiddleware, SanitizeToolCallsMiddleware, ValidateJournalWorkflowMiddleware
from src.agent.prompts import ( from src.agent.prompts import (
homework_doing_instructions, homework_doing_instructions,
main_agent_instructions, main_agent_instructions,
rework_instructions, rework_instructions,
) )
from src.agent.subagents import subagent_specs_without_tools from src.agent.subagents import subagent_specs_without_tools
from src.agent.runner_tools import solve_task
from src.agent.solve_tools import SOLVE_TOOLS
from src.agent.tools import GIT_TOOLS, WEB_TOOLS from src.agent.tools import GIT_TOOLS, WEB_TOOLS
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -54,7 +56,7 @@ _composite_backend = CompositeBackend(
# Наборы инструментов # Наборы инструментов
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
_homework_tools = [*GIT_TOOLS, *GITEA_TOOLS, *WEB_TOOLS, *_journal_tools] _homework_tools = [*GIT_TOOLS, *GITEA_TOOLS, *WEB_TOOLS, *_journal_tools, *SOLVE_TOOLS]
_web_tools = WEB_TOOLS _web_tools = WEB_TOOLS
_subagent_tool_map = { _subagent_tool_map = {
@@ -72,12 +74,13 @@ _gitea_names = {t.name for t in GITEA_TOOLS}
_journal_names = {t.name for t in _journal_tools} _journal_names = {t.name for t in _journal_tools}
_git_names = {t.name for t in GIT_TOOLS} _git_names = {t.name for t in GIT_TOOLS}
_web_names = {t.name for t in WEB_TOOLS} _web_names = {t.name for t in WEB_TOOLS}
_solve_names = {t.name for t in SOLVE_TOOLS}
_main_tool_names = _BUILTIN | _gitea_names _main_tool_names = _BUILTIN | _gitea_names | _journal_names | {"solve_task"}
_subagent_tool_names: dict[str, set[str]] = { _subagent_tool_names: dict[str, set[str]] = {
"web_search": _BUILTIN | _web_names, "web_search": _BUILTIN | _web_names,
"homework_doing": _BUILTIN | _gitea_names | _journal_names | _git_names | _web_names, "homework_doing": _BUILTIN | _gitea_names | _journal_names | _git_names | _web_names | _solve_names,
"journal_bh_tasks_submissions": _BUILTIN | _journal_names, "journal_bh_tasks_submissions": _BUILTIN | _journal_names,
} }
@@ -111,7 +114,7 @@ subagents = [
agent = create_deep_agent( agent = create_deep_agent(
model=llm, model=llm,
tools=list(GITEA_TOOLS), tools=[*GITEA_TOOLS, *_journal_tools, solve_task],
system_prompt=main_agent_instructions, system_prompt=main_agent_instructions,
backend=_composite_backend, backend=_composite_backend,
memory=[AGENTS_MD_VFS_PATH], memory=[AGENTS_MD_VFS_PATH],
@@ -129,6 +132,7 @@ homework_direct_agent = create_deep_agent(
system_prompt=homework_doing_instructions, system_prompt=homework_doing_instructions,
backend=_composite_backend, backend=_composite_backend,
middleware=[ middleware=[
RetryOnRateLimitMiddleware(),
SanitizeToolCallsMiddleware(known_tools=_subagent_tool_names["homework_doing"]), SanitizeToolCallsMiddleware(known_tools=_subagent_tool_names["homework_doing"]),
ValidateJournalWorkflowMiddleware(), ValidateJournalWorkflowMiddleware(),
], ],
@@ -144,6 +148,7 @@ rework_agent = create_deep_agent(
system_prompt=rework_instructions, system_prompt=rework_instructions,
backend=_composite_backend, backend=_composite_backend,
middleware=[ middleware=[
RetryOnRateLimitMiddleware(),
SanitizeToolCallsMiddleware(known_tools=_subagent_tool_names["homework_doing"]), SanitizeToolCallsMiddleware(known_tools=_subagent_tool_names["homework_doing"]),
ValidateJournalWorkflowMiddleware(), ValidateJournalWorkflowMiddleware(),
], ],
+27
View File
@@ -172,9 +172,36 @@ def gitea_get_file(repo: str, path: str, owner: str = GITEA_OWNER) -> str:
return f"Ошибка: {e}" return f"Ошибка: {e}"
@tool()
def gitea_list_files(repo: str, owner: str = GITEA_OWNER) -> str:
"""Список файлов в корне репозитория на git.brojs.ru.
Args:
repo: имя репозитория (например task-abc123)
owner: владелец репозитория (по умолчанию glevelll)
"""
try:
result = _get(f"/api/v1/repos/{owner}/{repo}/contents")
files = [item["name"] for item in result if item.get("type") == "file"]
dirs = [item["name"] for item in result if item.get("type") == "dir"]
parts = []
if files:
parts.append(f"Файлы: {', '.join(files)}")
if dirs:
parts.append(f"Папки: {', '.join(dirs)}")
return "\n".join(parts) if parts else f"Репозиторий {owner}/{repo} пуст"
except httpx.HTTPStatusError as e:
if e.response.status_code == 404:
return f"Репозиторий {owner}/{repo} не найден (первая сдача)"
return f"Ошибка: {e.response.text}"
except Exception as e:
return f"Ошибка: {e}"
# Список всех gitea-инструментов для удобного импорта # Список всех gitea-инструментов для удобного импорта
GITEA_TOOLS = [ GITEA_TOOLS = [
gitea_list_repos, gitea_list_repos,
gitea_list_files,
gitea_create_repo, gitea_create_repo,
gitea_write_file, gitea_write_file,
gitea_get_file, gitea_get_file,
+50 -17
View File
@@ -8,12 +8,13 @@ import re
from typing import TypedDict from typing import TypedDict
from langchain_core.messages import HumanMessage from langchain_core.messages import HumanMessage
from langchain_core.runnables import RunnableConfig
from langgraph.graph import START, StateGraph from langgraph.graph import START, StateGraph
from src.agent.agent import homework_direct_agent, rework_agent from src.agent.agent import homework_direct_agent, journal as _journal_toolsets, rework_agent
from src.agent.constants import COURSE_ID, GITEA_OWNER from src.agent.constants import COURSE_ID, GITEA_OWNER
from src.agent.gitea_tools import _get as gitea_get from src.agent.gitea_tools import _get as gitea_get
from src.agent.mcp_client import JOURNAL_PREFIX, load_journal_toolsets from src.agent.mcp_client import JOURNAL_PREFIX
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Типы состояния # Типы состояния
@@ -36,11 +37,8 @@ class PipelineState(TypedDict):
# Вспомогательные функции # Вспомогательные функции
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
_journal = load_journal_toolsets()
def _get_journal_tool(suffix: str): def _get_journal_tool(suffix: str):
all_tools = _journal.courses_lessons_tools + _journal.tasks_submissions_tools all_tools = _journal_toolsets.courses_lessons_tools + _journal_toolsets.tasks_submissions_tools
target = f"{JOURNAL_PREFIX}{suffix}" target = f"{JOURNAL_PREFIX}{suffix}"
for t in all_tools: for t in all_tools:
if t.name == target: if t.name == target:
@@ -116,13 +114,13 @@ async def _force_submit(task_id: str) -> bool:
print(f"[pipeline] force_submit: инструменты не найдены") print(f"[pipeline] force_submit: инструменты не найдены")
return False return False
try: try:
await update_tool.ainvoke({ await _mcp_invoke(update_tool, {
"taskId": task_id, "taskId": task_id,
"answerType": "link", "answerType": "link",
"content": repo_url, "content": repo_url,
"commit": {"repoUrl": repo_url, "branch": "main"}, "commit": {"repoUrl": repo_url, "branch": "main"},
}) })
await submit_tool.ainvoke({"taskId": task_id, "confirmSubmit": True}) await _mcp_invoke(submit_tool, {"taskId": task_id, "confirmSubmit": True})
print(f"[pipeline] Задание {task_id[:8]} — сабмит выполнен пайплайном ✓") print(f"[pipeline] Задание {task_id[:8]} — сабмит выполнен пайплайном ✓")
return True return True
except Exception as e: except Exception as e:
@@ -137,12 +135,25 @@ async def _is_submitted(task_id: str) -> bool:
return status in ("ready_for_review", "done") return status in ("ready_for_review", "done")
async def _mcp_invoke(tool, args: dict, retries: int = 5, pause: int = 30):
"""Вызывает MCP-инструмент с retry при 429."""
for attempt in range(1, retries + 1):
try:
return await tool.ainvoke(args)
except Exception as e:
if "429" in str(e) and attempt < retries:
print(f"[pipeline] MCP 429, жду {pause}с (попытка {attempt}/{retries})...")
await asyncio.sleep(pause)
else:
raise
async def _task_text(task_id: str) -> str: async def _task_text(task_id: str) -> str:
tool = _get_journal_tool("task_text") tool = _get_journal_tool("task_text")
if not tool: if not tool:
return "" return ""
try: try:
return _parse_text(await tool.ainvoke({"taskId": task_id})) return _parse_text(await _mcp_invoke(tool, {"taskId": task_id}))
except Exception: except Exception:
return "" return ""
@@ -152,7 +163,7 @@ async def _task_json(task_id: str) -> dict:
if not tool: if not tool:
return {} return {}
try: try:
raw = _parse_text(await tool.ainvoke({"taskId": task_id})) raw = _parse_text(await _mcp_invoke(tool, {"taskId": task_id}))
return json.loads(raw) return json.loads(raw)
except Exception: except Exception:
return {} return {}
@@ -226,11 +237,16 @@ def _is_rate_limit(exc: Exception) -> bool:
return "429" in msg or "rate" in msg.lower() or "rate_limit" in msg.lower() return "429" in msg or "rate" in msg.lower() or "rate_limit" in msg.lower()
async def _invoke_with_retry(agent, messages, config): async def _invoke_with_retry(agent, messages, config, callbacks=None):
"""Вызывает агента с автоматическим retry при 429.""" """Вызывает агента с автоматическим retry при 429.
callbacks — список LangChain callback-объектов (например AgentCallback из UI).
"""
run_config = dict(config)
if callbacks:
run_config["callbacks"] = callbacks
for attempt in range(1, RATE_LIMIT_RETRIES + 1): for attempt in range(1, RATE_LIMIT_RETRIES + 1):
try: try:
return await agent.ainvoke(messages, config) return await agent.ainvoke(messages, run_config)
except Exception as e: except Exception as e:
if _is_rate_limit(e) and attempt < RATE_LIMIT_RETRIES: if _is_rate_limit(e) and attempt < RATE_LIMIT_RETRIES:
wait = RATE_LIMIT_PAUSE * attempt wait = RATE_LIMIT_PAUSE * attempt
@@ -266,8 +282,8 @@ async def fetch_tasks(state: PipelineState) -> dict:
return {"tasks": coding, "current_index": 0, "results": [], "errors": []} return {"tasks": coding, "current_index": 0, "results": [], "errors": []}
async def process_one_task(state: PipelineState) -> dict: async def process_one_task(state: PipelineState, config: RunnableConfig | None = None) -> dict:
"""Выполняет одно задание.""" """Выполняет одно задание. config может содержать callbacks из UI."""
if state["current_index"] >= len(state["tasks"]): if state["current_index"] >= len(state["tasks"]):
return state return state
@@ -276,6 +292,10 @@ async def process_one_task(state: PipelineState) -> dict:
results = list(state.get("results", [])) results = list(state.get("results", []))
errors = list(state.get("errors", [])) errors = list(state.get("errors", []))
# Пауза перед стартом — даём BroJS MCP сбросить rate limit после загрузки инструментов
print(f"[pipeline] Задание {task_id[:8]} — пауза 10с перед стартом...")
await asyncio.sleep(10)
repo_url = await _existing_repo_url(task_id) repo_url = await _existing_repo_url(task_id)
is_rework = repo_url is not None is_rework = repo_url is not None
@@ -302,6 +322,9 @@ async def process_one_task(state: PipelineState) -> dict:
) )
agent_to_use = homework_direct_agent agent_to_use = homework_direct_agent
# Извлекаем callbacks из LangGraph config (переданы из UI)
callbacks = (config or {}).get("callbacks") or []
try: try:
print(f"[pipeline] Задание {task_id[:8]}{'пересдача' if is_rework else 'первая сдача'}: " print(f"[pipeline] Задание {task_id[:8]}{'пересдача' if is_rework else 'первая сдача'}: "
f"{task.get('title','')[:50]}") f"{task.get('title','')[:50]}")
@@ -310,6 +333,7 @@ async def process_one_task(state: PipelineState) -> dict:
agent_to_use, agent_to_use,
{"messages": [HumanMessage(content=prompt)]}, {"messages": [HumanMessage(content=prompt)]},
{"configurable": {"thread_id": f"pipeline-task-{task_id}"}}, {"configurable": {"thread_id": f"pipeline-task-{task_id}"}},
callbacks=callbacks,
) )
last = (result.get("messages") or [{}])[-1] last = (result.get("messages") or [{}])[-1]
output = getattr(last, "content", str(last)) output = getattr(last, "content", str(last))
@@ -329,6 +353,7 @@ async def process_one_task(state: PipelineState) -> dict:
agent_to_use, agent_to_use,
{"messages": [HumanMessage(content=fix_msg)]}, {"messages": [HumanMessage(content=fix_msg)]},
{"configurable": {"thread_id": f"pipeline-task-{task_id}-retry-{retries}"}}, {"configurable": {"thread_id": f"pipeline-task-{task_id}-retry-{retries}"}},
callbacks=callbacks,
) )
verification = await _verify_repo(repo_name) verification = await _verify_repo(repo_name)
@@ -348,8 +373,16 @@ async def process_one_task(state: PipelineState) -> dict:
"retries": retries, "retries": retries,
}) })
except Exception as e: except BaseException as e:
print(f"[pipeline] Задание {task_id[:8]} — ОШИБКА: {e}") import traceback
# Разворачиваем ExceptionGroup (Python 3.11+) чтобы увидеть реальные ошибки
if isinstance(e, ExceptionGroup):
for i, sub in enumerate(e.exceptions):
print(f"[pipeline] Задание {task_id[:8]} — под-ошибка {i+1}: {type(sub).__name__}: {sub}")
traceback.print_exception(type(sub), sub, sub.__traceback__)
else:
print(f"[pipeline] Задание {task_id[:8]} — ОШИБКА: {type(e).__name__}: {e}")
traceback.print_exc()
errors.append(f"Задание {task_id} ({'rework' if is_rework else 'new'}): {e}") errors.append(f"Задание {task_id} ({'rework' if is_rework else 'new'}): {e}")
# Пауза между заданиями чтобы не перегружать rate limit # Пауза между заданиями чтобы не перегружать rate limit
+40 -3
View File
@@ -1,6 +1,7 @@
"""Загрузка инструментов BroJS Journal через MCP (HTTP transport).""" """Загрузка инструментов BroJS Journal через MCP (HTTP transport)."""
import asyncio import asyncio
import os import os
import time
from concurrent.futures import ThreadPoolExecutor from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass from dataclasses import dataclass
@@ -17,6 +18,15 @@ _JOURNAL_TOKEN = os.getenv("JOURNAL_TOKEN", "YOUR_JOURNAL_TOKEN_HERE")
JOURNAL_MCP_URL = "https://platform.brojs.ru/jrnl-bh/api/mcp" JOURNAL_MCP_URL = "https://platform.brojs.ru/jrnl-bh/api/mcp"
# Exponential backoff: 15, 30, 60, 120, 240 секунд
_BACKOFF = [15, 30, 60, 120, 240]
# Минимальная пауза между последовательными MCP-вызовами (предотвращает burst)
_INTER_CALL_DELAY = 1.5
# Персистентный клиент и время последнего вызова — переиспользуются в рамках одного запуска
_persistent_client: MultiServerMCPClient | None = None
_last_mcp_call_time: float = 0.0
# Инструменты для работы с курсами и уроками # Инструменты для работы с курсами и уроками
JOURNAL_COURSES_LESSONS = frozenset({ JOURNAL_COURSES_LESSONS = frozenset({
"courses_list", "courses_list",
@@ -44,7 +54,7 @@ class JournalToolsets:
def _build_mcp_config() -> dict: def _build_mcp_config() -> dict:
return { return {
JOURNAL_SERVER_NAME: { JOURNAL_SERVER_NAME: {
"transport": "http", "transport": "streamable_http",
"url": JOURNAL_MCP_URL, "url": JOURNAL_MCP_URL,
"headers": { "headers": {
"Authorization": f"Bearer {_JOURNAL_TOKEN}", "Authorization": f"Bearer {_JOURNAL_TOKEN}",
@@ -53,16 +63,43 @@ def _build_mcp_config() -> dict:
} }
def _is_429(exc: Exception) -> bool:
return "429" in str(exc)
async def _fetch_tools() -> dict[str, list]: async def _fetch_tools() -> dict[str, list]:
"""Загружает MCP-инструменты с персистентным клиентом и экспоненциальным backoff при 429."""
global _persistent_client, _last_mcp_call_time
# Пауза между последовательными вызовами — предотвращает burst
elapsed = time.monotonic() - _last_mcp_call_time
if elapsed < _INTER_CALL_DELAY:
await asyncio.sleep(_INTER_CALL_DELAY - elapsed)
config = _build_mcp_config() config = _build_mcp_config()
client = MultiServerMCPClient(config)
# Переиспользуем клиент если уже создан
if _persistent_client is None:
_persistent_client = MultiServerMCPClient(config)
out: dict[str, list] = {} out: dict[str, list] = {}
for name in config: for name in config:
for i, pause in enumerate([0] + _BACKOFF):
try: try:
out[name] = await client.get_tools(server_name=name) if pause:
print(f" [mcp_client] '{name}' → 429, жду {pause}с (попытка {i+1})...")
await asyncio.sleep(pause)
out[name] = await _persistent_client.get_tools(server_name=name)
_last_mcp_call_time = time.monotonic()
break
except Exception as exc: except Exception as exc:
if _is_429(exc) and i < len(_BACKOFF):
continue
print(f"MCP '{name}': не удалось загрузить инструменты — {type(exc).__name__}: {exc}") print(f"MCP '{name}': не удалось загрузить инструменты — {type(exc).__name__}: {exc}")
# При неизвестной ошибке пересоздаём клиент перед следующей попыткой
_persistent_client = MultiServerMCPClient(config)
out[name] = [] out[name] = []
break
return out return out
+2 -1
View File
@@ -1,4 +1,5 @@
from src.agent.middlewares.sanitize_tool_calls import SanitizeToolCallsMiddleware from src.agent.middlewares.sanitize_tool_calls import SanitizeToolCallsMiddleware
from src.agent.middlewares.validate_journal_workflow import ValidateJournalWorkflowMiddleware from src.agent.middlewares.validate_journal_workflow import ValidateJournalWorkflowMiddleware
from src.agent.middlewares.retry_on_rate_limit import RetryOnRateLimitMiddleware
__all__ = ["SanitizeToolCallsMiddleware", "ValidateJournalWorkflowMiddleware"] __all__ = ["SanitizeToolCallsMiddleware", "ValidateJournalWorkflowMiddleware", "RetryOnRateLimitMiddleware"]
@@ -0,0 +1,81 @@
"""Middleware: повторяет вызов инструмента при 429 Rate Limit."""
from __future__ import annotations
import asyncio
import queue
import time
from typing import Any
from langchain.agents.middleware import AgentMiddleware, AgentState
_PAUSE = 30 # секунд ожидания при 429
_TRIES = 5 # максимум попыток
# Глобальный канал событий для UI (устанавливается из ui.py перед запуском агента).
# Если None — события просто не отправляются (CLI-режим).
_ui_event_queue: queue.Queue | None = None
def set_ui_event_queue(q: queue.Queue | None) -> None:
"""Вызывается из ui.py чтобы подключить очередь событий."""
global _ui_event_queue
_ui_event_queue = q
def _emit(event: dict) -> None:
if _ui_event_queue is not None:
try:
_ui_event_queue.put_nowait(event)
except Exception:
pass
def _is_429(exc: Exception) -> bool:
msg = str(exc)
return "429" in msg or "rate" in msg.lower()
class RetryOnRateLimitMiddleware(AgentMiddleware[AgentState[Any], Any]):
"""Перехватывает 429 от любого инструмента и повторяет с паузой.
Отправляет события rate_limit_wait / rate_limit_retry в UI-очередь.
"""
def wrap_tool_call(self, request, handler):
for attempt in range(1, _TRIES + 1):
try:
return handler(request)
except Exception as e:
if _is_429(e) and attempt < _TRIES:
name = request.tool_call.get("name", "")
print(f"[retry-mw] {name} → 429, жду {_PAUSE}с (попытка {attempt}/{_TRIES})...")
_emit({"t": "rate_limit_wait", "name": name,
"pause": _PAUSE, "attempt": attempt, "max": _TRIES,
"ts": _now()})
time.sleep(_PAUSE)
_emit({"t": "rate_limit_retry", "name": name,
"attempt": attempt + 1, "ts": _now()})
else:
raise
async def awrap_tool_call(self, request, handler):
for attempt in range(1, _TRIES + 1):
try:
return await handler(request)
except Exception as e:
if _is_429(e) and attempt < _TRIES:
name = request.tool_call.get("name", "")
print(f"[retry-mw] {name} → 429, жду {_PAUSE}с (попытка {attempt}/{_TRIES})...")
_emit({"t": "rate_limit_wait", "name": name,
"pause": _PAUSE, "attempt": attempt, "max": _TRIES,
"ts": _now()})
await asyncio.sleep(_PAUSE)
_emit({"t": "rate_limit_retry", "name": name,
"attempt": attempt + 1, "ts": _now()})
else:
raise
def _now() -> str:
from datetime import datetime
return datetime.now().strftime("%H:%M:%S")
+105 -588
View File
@@ -63,587 +63,98 @@ journal_tasks_submissions_instructions = """
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
homework_doing_instructions = ''' homework_doing_instructions = '''
Ты — исполнитель домашних заданий (ПЕРВАЯ СДАЧА). Ты — агент выполнения домашних заданий курса KFU-26-1.
У тебя есть ВСЕ инструменты напрямую. Не делегируй другим субагентам.
courseId = "698b49da77cb6d4d2e43ce78" courseId = "698b49da77cb6d4d2e43ce78"
Gitea owner = "glevelll" Gitea owner = "glevelll"
repo для задания: "task-<taskId>"
ВАЖНО: Journal-инструменты имеют префикс mcp__journal-bh-professor__
Gitea-инструменты: gitea_create_repo, gitea_write_file, gitea_get_file, gitea_list_repos ## Доступные инструменты
Git-инструменты: git_clone, git_pull, git_status, git_add_and_commit, git_push
Journal (префикс mcp__journal-bh-professor__):
## ПОРЯДОК ВЫПОЛНЕНИЯ: task_text(taskId) — полный текст задания
task_get(taskId) — детали: статус, answer, комментарии преподавателя
[1] mcp__journal-bh-professor__task_text({"taskId": "<id>"}) task_update_answer(...) — установить ссылку на репо (ОБЯЗАТЕЛЬНО перед submit)
→ Прочитай ПОЛНЫЙ текст задания task_submit(taskId, confirmSubmit=true) — сдать задание
[2] Составь письменный план: Gitea:
- какие файлы нужны (main.py, requirements.txt, etc.) gitea_list_files(repo) — список файлов в репозитории
- что реализовать в каждом файле gitea_list_repos() — список репозиториев (узнать существует ли repo)
- какой технический стек использовать (см. раздел ТЕХНИЧЕСКИЕ ПАТТЕРНЫ ниже) gitea_create_repo(name) — создать репозиторий
gitea_write_file(repo, path, content, message) — записать файл (автокоммит)
[3] gitea_create_repo({"name": "task-<id>", "private": false}) gitea_get_file(repo, path) — прочитать файл
→ Создай репозиторий
Инструменты решения (LLM-субагенты):
[4] Для КАЖДОГО файла вызывай ОТДЕЛЬНО: validate_teacher_comment(task_text, repo_name, teacher_comment)
gitea_write_file({ → анализирует каждый пункт замечания: ловушка или реальная ошибка
"repo": "task-<id>", generate_code_solution(task_text, fix_instructions="", defense_context="")
"path": "main.py", → генерирует main.py + requirements.txt + extra_files
"content": "ПОЛНЫЙ КОД ФАЙЛА",
"message": "add main.py" ## Принципы работы
})
- gitea_write_file сам коммитит на сервере — git_add_and_commit НЕ нужен Для ПЕРВОЙ СДАЧИ:
- content — это plain text, НЕ base64 — Прочитай задание через task_text
- ВСЕГДА указывай message — Сгенерируй решение через generate_code_solution(task_text)
- Один вызов = один файл — Создай репозиторий, запиши все файлы через gitea_write_file
— Установи ответ через task_update_answer, затем сдай через task_submit
[5] git_clone("https://git.brojs.ru/glevelll/task-<id>")
→ Клонируй репозиторий локально для проверки Для ПЕРЕСДАЧИ (репозиторий уже существует):
— Прочитай задание (task_text) и комментарий преподавателя (task_get)
[6] Проверь через read_file что код корректен — ОБЯЗАТЕЛЬНО проверь замечание через validate_teacher_comment
— Если замечание — ловушка (has_trap=true, has_valid=false):
[7] mcp__journal-bh-professor__task_update_answer({ · Добавь возражение в README через gitea_write_file
"taskId": "<id>", · Сдай без изменений кода
"answerType": "link", — Если смешанный (has_trap=true, has_valid=true):
"content": "https://git.brojs.ru/glevelll/task-<id>" · Добавь возражение в README за ложные пункты
}) · Передай fix_instructions и defense_context в generate_code_solution
→ ОБЯЗАТЕЛЬНО перед task_submit! · Запиши исправленные файлы, сдай
— Если всё обоснованно (has_trap=false):
[8] Финальная проверка: · Передай fix_instructions в generate_code_solution
✓ Все файлы записаны? · Запиши исправленные файлы, сдай
✓ Нет pass, TODO, ..., заглушек?
✓ langchain>1.0.0 в requirements.txt? ## Ограничения кода
✓ task_update_answer вызван? - Никаких pass, TODO, заглушек
- LLM только через OpenRouter (langchain_openai), не Ollama
[9] mcp__journal-bh-professor__task_submit({ - task_update_answer ВСЕГДА перед task_submit
"taskId": "<id>", - Один инструмент за один шаг'''
"confirmSubmit": true
})
## ТРЕБОВАНИЯ К КОДУ:
- ПОЛНЫЙ рабочий код, без pass, TODO, ...
- requirements.txt с реальными зависимостями и langchain>1.0.0
- Соответствие всем требованиям из текста задания
- Используй langchain>=1.2.10 / langgraph>=0.2.0 согласно заданию
## ЗАПРЕЩЕНО:
- pass, TODO, ..., пустые функции
- langchain<=1.0.0 в requirements.txt
- Пропускать task_update_answer перед task_submit
- Писать код только в requirements.txt без main.py
## ═══════════════════════════════════════════════
## ТЕХНИЧЕСКИЕ ПАТТЕРНЫ (читай ПЕРЕД написанием кода)
## ═══════════════════════════════════════════════
### LLM — ВСЕГДА используй OpenRouter (не Ollama, не hub.pull, не hardcode)
```python
import os
from langchain_openai import ChatOpenAI
llm = ChatOpenAI(
model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"),
temperature=0.0,
)
```
requirements.txt: langchain-openai>=0.3.0
---
### deepagents — правильный паттерн (задания про "deep agent", "deepagent", "deep agents from scratch")
```python
import os, asyncio
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.tools import tool
from deepagents import create_deep_agent
from deepagents.backends import FilesystemBackend, LocalShellBackend, CompositeBackend
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"))
# Виртуальная ФС + реальная shell среда
backend = CompositeBackend([
LocalShellBackend(workspace_dir="./workspace"),
FilesystemBackend(),
])
@tool
def web_search(query: str) -> str:
"""Search the web for information."""
try:
from duckduckgo_search import DDGS
with DDGS() as ddgs:
results = list(ddgs.text(query, max_results=5))
return "\\n".join(f"{r['title']}: {r['body']}" for r in results)
except Exception as e:
return f"Search error: {e}"
agent = create_deep_agent(
llm=llm,
tools=[web_search],
backend=backend,
system_prompt="You are a helpful research agent.",
)
async def main():
result = await agent.ainvoke(
{"messages": [HumanMessage(content="Search for Python best practices and save to results.txt")]},
{"configurable": {"thread_id": "session-1"}},
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
```
requirements.txt: deepagents, langchain-openai>=0.3.0, duckduckgo-search
---
### FastMCP сервер — ТОЛЬКО на уровне модуля, НИКОГДА внутри класса
```python
# ПРАВИЛЬНО:
from fastmcp import FastMCP
import json
from pathlib import Path
mcp = FastMCP("memory-server")
STORAGE = Path("memory.json")
def _load():
return json.loads(STORAGE.read_text()) if STORAGE.exists() else {}
def _save(data):
STORAGE.write_text(json.dumps(data, indent=2, ensure_ascii=False))
@mcp.tool()
def save(key: str, value: str) -> bool:
"""Save a value by key."""
data = _load(); data[key] = value; _save(data)
return True
@mcp.tool()
def get(key: str) -> str:
"""Get a value by key."""
return _load().get(key, "")
@mcp.tool()
def delete(key: str) -> bool:
"""Delete a key."""
data = _load()
if key in data:
del data[key]; _save(data); return True
return False
@mcp.tool()
def list_keys() -> list:
"""List all keys."""
return list(_load().keys())
if __name__ == "__main__":
mcp.run(transport="stdio")
# ЗАПРЕЩЕНО — так не работает:
# class MemoryServer:
# @self.mcp.tool() ← NameError: self не существует в теле класса
# def save(self, ...): ...
```
requirements.txt: fastmcp>=0.1.0, pydantic>=2.0
---
### LangChain create_agent — НЕ совместим с AgentExecutor
```python
import asyncio, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.agents import create_agent
from langchain.tools import tool
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"))
@tool
def my_tool(query: str) -> str:
"""Tool description."""
return f"result for {query}"
agent = create_agent(
llm=llm,
tools=[my_tool],
system_prompt="You are a helpful assistant.",
)
async def main():
result = await agent.ainvoke(
{"messages": [HumanMessage(content="Hello")]},
{"configurable": {"thread_id": "t1"}},
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
# ЗАПРЕЩЕНО:
# AgentExecutor(agent=create_agent(...), ...) ← несовместимо!
# agent_type=AgentType.ZERO_SHOT_REACT_DESCRIPTION ← не параметр create_agent
```
requirements.txt: langchain>=1.2.10, langchain-openai>=0.3.0, langgraph>=0.2.0
---
### Human-in-the-Loop через HumanInTheLoopMiddleware
```python
import asyncio, json, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langchain.tools import tool
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import Command
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"))
@tool
def get_weather(city: str) -> str:
"""Get weather for a city."""
return f"Sunny, 22C in {city}"
memory = MemorySaver()
agent = create_agent(
llm=llm,
tools=[get_weather],
system_prompt="You are a helpful assistant.",
middleware=[HumanInTheLoopMiddleware(interrupt_on={"get_weather": True})],
checkpointer=memory,
)
def ask_human(interrupt_value):
decisions = []
for action in interrupt_value.get("action_requests", []):
print(f"Tool: {action['name']}, Args: {action['args']}")
ans = input("Approve? (y/n): ").strip().lower()
decisions.append({"type": "approve" if ans == "y" else "reject"})
return decisions
async def main():
config = {"configurable": {"thread_id": "session-1"}}
result = await agent.ainvoke(
{"messages": [HumanMessage(content="What's the weather in Moscow?")]},
config,
)
while "__interrupt__" in result:
decisions = ask_human(result["__interrupt__"][0].value)
result = await agent.ainvoke(
Command(resume={"decisions": decisions}), config
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
```
requirements.txt: langchain>=1.2.10, langchain-openai>=0.3.0, langgraph>=0.2.0
---
### LangGraph interrupt (Human-in-the-loop без middleware)
```python
import asyncio, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.agents import create_agent
from langchain.tools import tool
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import Command
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"))
@tool
def dangerous_action(cmd: str) -> str:
"""Execute a dangerous action."""
return f"Executed: {cmd}"
memory = MemorySaver()
agent = create_agent(llm=llm, tools=[dangerous_action],
checkpointer=memory, interrupt_before=["tools"])
async def main():
config = {"configurable": {"thread_id": "t1"}}
result = await agent.ainvoke(
{"messages": [HumanMessage(content="Run ls -la")]}, config
)
# Агент остановился перед вызовом инструмента
snapshot = await agent.aget_state(config)
if snapshot.next:
ans = input(f"Approve tool call? (y/n): ").strip()
if ans == "y":
result = await agent.ainvoke(Command(resume=None), config)
else:
result = await agent.ainvoke(
Command(resume=None, update={"messages": [
HumanMessage(content="User rejected the action.")
]}), config
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
```
---
### RAG-агент с Qdrant (используй OpenRouter для LLM, Qdrant для векторов)
```python
import os, asyncio
from langchain_openai import ChatOpenAI, OpenAIEmbeddings
from langchain_qdrant import QdrantVectorStore
from langchain.tools import tool
from langchain.agents import create_agent
from langchain_core.messages import HumanMessage
from qdrant_client import QdrantClient
from qdrant_client.models import Distance, VectorParams
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"))
# Embeddings через OpenAI-совместимый API (OpenRouter)
embeddings = OpenAIEmbeddings(
model="text-embedding-3-small",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"),
)
# Qdrant in-memory (не требует отдельного сервера)
client = QdrantClient(":memory:")
client.create_collection("knowledge",
vectors_config=VectorParams(size=1536, distance=Distance.COSINE))
vector_store = QdrantVectorStore(client=client, collection_name="knowledge",
embedding=embeddings)
@tool
def search_knowledge_base(query: str, max_results: int = 5) -> str:
"""Semantic search in the knowledge base."""
docs = vector_store.similarity_search(query, k=max_results)
if not docs:
return "No relevant documents found."
return "\\n\\n".join(f"{i+1}. {d.page_content}" for i, d in enumerate(docs))
@tool
def add_to_knowledge_base(content: str, title: str = "document") -> str:
"""Add text to the knowledge base."""
from langchain_core.documents import Document
vector_store.add_documents([Document(page_content=content,
metadata={"title": title})])
return f"Added '{title}' to knowledge base."
agent = create_agent(
llm=llm,
tools=[search_knowledge_base, add_to_knowledge_base],
system_prompt="You are an assistant with access to a knowledge base.",
)
async def main():
await add_to_knowledge_base.ainvoke({"content": "Python is a high-level language.", "title": "python-intro"})
result = await agent.ainvoke(
{"messages": [HumanMessage(content="What do you know about Python?")]},
{"configurable": {"thread_id": "rag-1"}},
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
```
requirements.txt: langchain>=1.2.10, langchain-openai>=0.3.0, langgraph>=0.2.0,
langchain-qdrant, qdrant-client
---
### Stream-режим агента
```python
import asyncio, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.agents import create_agent
from langchain.tools import tool
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"), streaming=True)
@tool
def calculator(expression: str) -> str:
"""Evaluate a math expression."""
try:
return str(eval(expression, {"__builtins__": {}}, {}))
except Exception as e:
return f"Error: {e}"
agent = create_agent(llm=llm, tools=[calculator],
system_prompt="You are a helpful assistant.")
async def main():
config = {"configurable": {"thread_id": "stream-1"}}
# stream_mode="messages" — получаем токены по одному
async for event in agent.astream(
{"messages": [HumanMessage(content="What is 2+2?")]},
config,
stream_mode="messages",
):
if isinstance(event, tuple):
msg, metadata = event
if hasattr(msg, "content") and msg.content:
print(msg.content, end="", flush=True)
print()
if __name__ == "__main__":
asyncio.run(main())
```
---
### Задания типа "план / документ" (не чистый кодинг — например ai-fluency)
Если задание просит написать план, документ или пройти курс:
- Создай main.py который ВЫВОДИТ план в консоль
- План должен быть содержательным (минимум 300 слов), структурированным
- Имитируй личный опыт: "я понял, что...", "мой план включает..."
- Опирайся на тему курса из описания задания
---
### Web search без API-ключа (для поисковых агентов)
```python
from duckduckgo_search import DDGS
def web_search(query: str) -> str:
with DDGS() as ddgs:
results = list(ddgs.text(query, max_results=5))
return "\\n".join(f"[{r['title']}] {r['body']} ({r['href']})" for r in results)
```
requirements.txt: duckduckgo-search
---
### LangGraph текстовая игра с interrupt
```python
import asyncio, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage, SystemMessage
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command
from typing import TypedDict, Annotated
from langgraph.graph.message import add_messages
class GameState(TypedDict):
messages: Annotated[list, add_messages]
location: str
inventory: list
llm = ChatOpenAI(model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"))
def game_master(state: GameState) -> dict:
system = SystemMessage(content=(
"You are a text adventure game master. "
f"Player is at: {state.get('location','start')}. "
f"Inventory: {state.get('inventory',[])}. "
"Describe what happens and list 2-3 options."
))
response = llm.invoke([system] + state["messages"])
return {"messages": [response]}
def player_turn(state: GameState) -> Command:
player_input = interrupt("Your action: ")
return Command(goto="game_master",
update={"messages": [HumanMessage(content=player_input)]})
memory = MemorySaver()
builder = StateGraph(GameState)
builder.add_node("game_master", game_master)
builder.add_node("player_turn", player_turn)
builder.add_edge(START, "game_master")
builder.add_edge("game_master", "player_turn")
game = builder.compile(checkpointer=memory)
async def main():
config = {"configurable": {"thread_id": "game-1"}}
state = {"messages": [HumanMessage(content="Start the adventure!")],
"location": "forest entrance", "inventory": []}
result = await game.ainvoke(state, config)
while True:
last = result["messages"][-1].content
print(f"\\nGame: {last}")
if "__interrupt__" in result:
action = input("\\nYour action: ").strip()
if action.lower() in ("quit", "exit"):
break
result = await game.ainvoke(Command(resume=action), config)
else:
break
if __name__ == "__main__":
asyncio.run(main())
```
'''
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Субагент: пересдача после ревью # Субагент: пересдача после ревью
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
rework_instructions = """ rework_instructions = """
Ты — исполнитель домашних заданий (ПЕРЕСДАЧА после ревью преподавателя). Ты — агент пересдачи домашних заданий курса KFU-26-1.
У тебя есть ВСЕ инструменты напрямую. Не делегируй. Репозиторий уже существует. Задание отклонено с комментарием преподавателя.
courseId = "698b49da77cb6d4d2e43ce78" courseId = "698b49da77cb6d4d2e43ce78"
Gitea owner = "glevelll" Gitea owner = "glevelll"
Ситуация: задание уже было отправлено, получены комментарии. Репозиторий существует. ## Твоя задача
## ПОРЯДОК: 1. Получи текст задания и комментарий преподавателя
2. Проверь каждый пункт комментария через validate_teacher_comment
3. Прими решение на основе результата:
[1] mcp__journal-bh-professor__task_submission_status({"taskId": "<id>"}) has_trap=true, has_valid=false → ЛОВУШКА
→ Проверь статус и получи фидбек Добавь возражение в README.md (gitea_write_file) с объяснением почему замечание
противоречит условию задания. Сдай без изменений кода.
[2] mcp__journal-bh-professor__task_get({"taskId": "<id>"}) has_trap=true, has_valid=true → СМЕШАННЫЙ СЛУЧАЙ
→ Получи URL репозитория из answer.content и прочитай комментарии Добавь возражение в README.md за ложные пункты.
Передай только реальные fix_instructions в generate_code_solution.
Передай trap_explanations как defense_context — агент добавит DESIGN DECISION блоки.
[3] git_clone(<url из answer.content>) has_trap=false → ОБОСНОВАННОЕ ЗАМЕЧАНИЕ
→ Клонируй существующий репозиторий в agent_workspace Передай fix_instructions в generate_code_solution.
→ <repo-name> = последняя часть URL (например task-abc123) Запиши исправленные файлы через gitea_write_file.
[4] Прочитай файлы через read_file, пойми что исправить 4. Всегда вызывай task_update_answer → task_submit после изменений
[5] Внеси исправления через edit_file или write_file ## Правила
- НИКОГДА не меняй код по ложным замечаниям
[6] git_add_and_commit("fix: <описание исправлений>", "<repo-name>") - Используй тот же репозиторий (task-<taskId>), не создавай новый
- Один инструмент за один шаг
[7] git_push("<repo-name>") - task_update_answer обязателен перед task_submit (даже если URL тот же)
[8] mcp__journal-bh-professor__task_update_answer({
"taskId": "<id>",
"answerType": "link",
"content": "<ТОТ ЖЕ URL репозитория>"
})
[9] mcp__journal-bh-professor__task_submit({"taskId": "<id>", "confirmSubmit": true})
## ПРАВИЛА:
- Клонируй существующий репозиторий, НЕ создавай новый
- Исправляй ТОЛЬКО то, что указано в комментариях
- task_update_answer обязателен (даже если URL тот же)
- Запрещено: pass, TODO, пустые функции
""" """
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -652,43 +163,49 @@ Gitea owner = "glevelll"
main_agent_instructions = """\ main_agent_instructions = """\
Ты — главный агент-исполнитель домашних заданий курса KFU-26-1 на platform.brojs.ru. Ты — главный агент-исполнитель домашних заданий курса KFU-26-1 на platform.brojs.ru.
Твоя роль — получать задания из журнала и выполнять их качественно. Ты полностью управляешь всем процессом: сам решаешь что делать, в каком порядке,
какие инструменты вызывать. Python-инструменты — это только твои руки.
## Известные курсы ## Известные курсы
- KFU-26-1 = courseId `698b49da77cb6d4d2e43ce78` - KFU-26-1 = courseId `698b49da77cb6d4d2e43ce78`
## Доступные субагенты (вызывай через инструмент `task`) ## Прямые инструменты (вызывай напрямую)
- `journal_bh_tasks_submissions`: читает задания, проверяет статусы, отправляет ответы
- `homework_doing`: ВЫПОЛНЯЕТ задание (пишет код, создаёт репо, сдаёт)
- `web_search`: ищет информацию в интернете (только если нужно)
## Прямые Gitea-инструменты (доступны напрямую без субагента) ### Journal (MCP)
- `gitea_list_repos` — список репозиториев, также возвращает username - `mcp__journal-bh-professor__tasks_list` — список всех заданий курса со статусами
- `gitea_create_repo` — создать репозиторий - `mcp__journal-bh-professor__task_get` — детали задания (статус, ответ, комментарии)
- `gitea_write_file` — создать/обновить файл (автокоммит) - `mcp__journal-bh-professor__task_text` — полный текст задания
- `gitea_get_file` — получить файл - `mcp__journal-bh-professor__task_update_answer` — установить ссылку на репо
- `mcp__journal-bh-professor__task_submit` — сдать задание
- `mcp__journal-bh-professor__task_comment` — написать комментарий
### Исполнение
- `solve_task(task_id)` — **главный инструмент**: полностью выполняет одно задание
(читает условие, пишет код, создаёт репо на Gitea, верифицирует, сдаёт).
Автоматически определяет первая сдача или пересдача.
### Gitea
- `gitea_list_repos`, `gitea_create_repo`, `gitea_write_file`, `gitea_get_file`
## Один инструмент за шаг (ОБЯЗАТЕЛЬНО) ## Один инструмент за шаг (ОБЯЗАТЕЛЬНО)
В одном сообщении — **только один** вызов любого инструмента (`task`, `ls`, `read_file`, Вызывай строго по одному инструменту. Жди результата перед следующим вызовом.
`write_file`, `edit_file`, `glob`, `grep`, `execute` и т.д.).
Сначала дождись результата, затем следующий вызов.
## Типовые маршруты ## Как выполнить все todo-задания
1. `mcp__journal-bh-professor__tasks_list(courseId="698b49da77cb6d4d2e43ce78")`
2. Из ответа выбери задания со статусом `todo` или `in_progress`
3. Для каждого последовательно: `solve_task(task_id="<id>")`
4. Дождись "OK ..." перед следующим заданием
5. Доложи итоги
### Получить список заданий ## Как выполнить одно задание
1. Делегируй `journal_bh_tasks_submissions`: получить tasks_list для courseId 1. `solve_task(task_id="<id>")` — сделает всё сам
### Выполнить задание ## Как проверить статусы
1. Делегируй `homework_doing`: выполни задание с taskId=<id> 1. `mcp__journal-bh-professor__tasks_list(courseId="698b49da77cb6d4d2e43ce78")`
(он сам прочитает текст, создаст репо, напишет код и сдаст)
2. Верни пользователю ссылку на репозиторий
### Проверить статусы
1. Делегируй `journal_bh_tasks_submissions`: получить статусы всех заданий курса
## Жёсткие ограничения ## Жёсткие ограничения
- Не вызывай больше одного инструмента за шаг - Не говори что задание выполнено, если `solve_task` не вернул "OK"
- Не делегируй субагенту несколько независимых задач сразу - Выполняй задания строго последовательно, не параллельно
- Не говори что задание выполнено, если оно не было реально выполнено - Не вызывай `solve_task` для заданий со статусом `done` или `ready_for_review`
- Не подменяй требования задания своими догадками - Не подменяй требования задания своими догадками
""" """
+108
View File
@@ -0,0 +1,108 @@
"""
Инструмент solve_task — надёжное выполнение одного задания.
Агент-оркестратор вызывает его для каждого todo-задания.
Python внутри обеспечивает верификацию и страховочный сабмит,
но решение о том какие задания выполнять и в каком порядке
принимает LLM-оркестратор.
"""
from __future__ import annotations
import asyncio
from langchain.tools import tool
from langchain_core.messages import HumanMessage
@tool
async def solve_task(task_id: str) -> str:
"""Полностью выполнить одно задание курса: написать код, запушить в репо, сдать.
Автоматически определяет режим:
- Первая сдача: читает условие, генерирует код с нуля
- Пересдача: читает замечания, исправляет или защищает решение
После выполнения верифицирует репозиторий и гарантированно сдаёт задание.
Args:
task_id: ID задания из tasks_list (например 6a22c713fd30e81cf315ea04)
Returns:
Строка с результатом: "OK ..." или "ERROR ..."
"""
# Ленивый импорт чтобы избежать circular import на уровне модуля
from src.agent.agent import homework_direct_agent, rework_agent
from src.agent.graph.pipeline import (
MAX_RETRIES,
_existing_repo_url,
_fix_prompt,
_force_submit,
_invoke_with_retry,
_is_submitted,
_needs_retry,
_task_json,
_task_text,
_verify_repo,
TaskInfo,
)
# Небольшая пауза — снижает давление на rate limit
await asyncio.sleep(5)
# Определяем: первая сдача или пересдача
repo_url = await _existing_repo_url(task_id)
is_rework = repo_url is not None
if is_rework:
data = await _task_json(task_id)
comments = data.get("comments") or data.get("feedback", "")
prompt = (
f"Пересдача задания.\n\n"
f"ID: {task_id}\n"
f"Репозиторий: {repo_url}\n"
f"Комментарии преподавателя: {comments}\n\n"
"Внеси исправления и отправь снова."
)
agent_to_use = rework_agent
else:
task_text = await _task_text(task_id)
prompt = (
f"Выполни задание.\n\n"
f"ID: {task_id}\n\n"
f"Текст задания:\n{task_text}\n\n"
"Первая сдача. Напиши код с нуля."
)
agent_to_use = homework_direct_agent
try:
result = await _invoke_with_retry(
agent_to_use,
{"messages": [HumanMessage(content=prompt)]},
{"configurable": {"thread_id": f"runner-{task_id}"}},
)
# Верификация репозитория (только первая сдача)
repo_name = f"task-{task_id}"
retries = 0
if not is_rework:
task_info = TaskInfo(id=task_id, title="", status="")
verification = await _verify_repo(repo_name)
while _needs_retry(verification) and retries < MAX_RETRIES:
retries += 1
fix_msg = _fix_prompt(task_info, repo_name, verification)
result = await _invoke_with_retry(
agent_to_use,
{"messages": [HumanMessage(content=fix_msg)]},
{"configurable": {"thread_id": f"runner-{task_id}-retry-{retries}"}},
)
verification = await _verify_repo(repo_name)
# Гарантированный сабмит если агент не сдал сам
if not await _is_submitted(task_id):
await _force_submit(task_id)
mode = "пересдача" if is_rework else "первая сдача"
return f"OK: задание {task_id[:8]}... выполнено ({mode}, retries={retries})"
except Exception as e:
return f"ERROR: задание {task_id[:8]}...: {type(e).__name__}: {e}"
+238
View File
@@ -0,0 +1,238 @@
"""Промпты для инструментов решения задач: валидация замечаний, анализ, генерация кода."""
# ---------------------------------------------------------------------------
# Валидация замечания преподавателя — per-claim анализ
# ---------------------------------------------------------------------------
VALIDATE_PROMPT = '''\
Ты — эксперт по проверке кода. Дано условие задания, текущий код и замечание преподавателя.
Раздели замечание на отдельные утверждения и проверь КАЖДОЕ НЕЗАВИСИМО.
## Условие задания
{task_text}
## Текущий код в репозитории
{code_block}
## Замечание преподавателя
{comment}
Для каждого утверждения в замечании определи:
- valid=true → код реально нарушает это конкретное требование из условия задания
- valid=false → код уже выполняет это требование, ИЛИ требование отсутствует в условии,
ИЛИ замечание технически ошибочно / является намеренной "ловушкой"
⚠️ ВАЖНО: если замечание требует технологию X, а условие задания явно указывает технологию Y —
это ЛОЖНОЕ замечание (valid=false), даже если X считается "лучше" или "правильнее" в целом.
Сравнивай только с текстом условия задания, не с общими best practices.
Ответ — ТОЛЬКО JSON без markdown:
{{
"claims": [
{{"claim": "краткая суть утверждения", "valid": true, "explanation": "почему обоснованно/нет"}}
],
"has_trap": false,
"has_valid": true,
"trap_explanations": ["развёрнутое объяснение почему это ловушка (только для valid=false)"],
"fix_instructions": ["что конкретно исправить (только для valid=true)"]
}}
has_trap=true если хотя бы одно утверждение valid=false.
has_valid=true если хотя бы одно утверждение valid=true.
'''
# ---------------------------------------------------------------------------
# Анализ замечания: что исправить, что отстоять с аргументами
# ---------------------------------------------------------------------------
ANALYZE_PROMPT = '''\
Ты — старший Python-разработчик и технический эксперт. Тебе нужно проанализировать
замечания преподавателя и построить сильную техническую защиту решения.
## Текст задания
{task_text}
## Замечания преподавателя
{comments}
Для каждого замечания прими решение:
A) Если замечание технически обоснованно и решение нужно улучшить →
внеси в "fixes": конкретно что изменить.
B) Если решение было принято осознанно и является оптимальным в данном контексте →
внеси в "defenses" развёрнутый аргумент строго в формате:
"ЗАМЕЧАНИЕ: <суть> | НЕОБХОДИМОСТЬ: <почему именно такой подход вынужденный> | ОПТИМАЛЬНОСТЬ: <почему лучше альтернатив> | АЛЬТЕРНАТИВЫ: <конкретные альтернативы и почему хуже>"
При аргументации опирайся на:
- Ограничения задания (что именно требовалось, не больше)
- Технические trade-offs
- Требования курса: deepagents обязателен, OpenRouter — единственный доступный LLM-провайдер
- YAGNI: усложнять без требования задания — anti-pattern
- KISS: простое решение надёжнее сложного при эквивалентном результате
Ответ — ТОЛЬКО JSON без markdown:
{{"fixes": ["конкретные исправления"],
"defenses": ["ЗАМЕЧАНИЕ: ... | НЕОБХОДИМОСТЬ: ... | ОПТИМАЛЬНОСТЬ: ... | АЛЬТЕРНАТИВЫ: ..."],
"verdict": "needs_fixes" | "already_correct" | "mixed"}}
'''
# ---------------------------------------------------------------------------
# Возражение на ложное замечание (добавляется в README)
# ---------------------------------------------------------------------------
OBJECTION_TEMPLATE = """\n\n---\n\n## Ответ на замечание преподавателя\n\n**Замечание:** {comment}\n\n**Позиция:** {explanation}\n\nКод полностью соответствует условию задания по указанным пунктам. Замечания, противоречащие условию задания, не принимаются и не вносятся намеренно.\n"""
# ---------------------------------------------------------------------------
# Секция пересдачи — вставляется в CODE_PROMPT
# ---------------------------------------------------------------------------
REWORK_SECTION = '''\
## ПЕРЕСДАЧА — технический анализ замечаний
### Исправить (замечания обоснованы):
{fixes}
### Отстоять с аргументацией (решение оптимально):
{defenses}
Правила генерации кода:
- Вноси ТОЛЬКО изменения из раздела "Исправить"
- Для каждого пункта из "Отстоять" — добавь в код РАЗВЁРНУТЫЙ блок комментариев:
# DESIGN DECISION: <суть спорного решения>
# NECESSITY: <почему именно так — вынужденность, ограничения задания/курса>
# OPTIMALITY: <почему это лучше альтернатив — конкретные аргументы>
# ALTERNATIVES CONSIDERED: <что рассматривалось и почему отклонено>
- Не меняй архитектуру без явного требования в "Исправить"
- Решение должно выглядеть как результат инженерного решения, а не случайного выбора
'''
# ---------------------------------------------------------------------------
# Основной промпт генерации кода
# ---------------------------------------------------------------------------
CODE_PROMPT = '''\
Ты — Python-разработчик. Напиши решение для учебного задания по LLM/AI.
Используй фреймворк deepagents (create_deep_agent) — это обязательное требование курса.
## Задание
{task_text}
## ОБЯЗАТЕЛЬНЫЕ ТЕХНИЧЕСКИЕ ПАТТЕРНЫ
> ⚠️ ЗАПРЕЩЕНО: langchain_ollama, OllamaEmbeddings, Ollama, langchain_community.
> Для LLM и эмбеддингов — ТОЛЬКО OpenRouter через langchain_openai.
### LLM — всегда OpenRouter:
```python
import os
from langchain_openai import ChatOpenAI
llm = ChatOpenAI(
model="openai/gpt-oss-20b:free",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"),
temperature=0.0,
)
```
### Базовый агент (deepagents):
```python
import asyncio, os
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage
from langchain.tools import tool
from deepagents import create_deep_agent
from deepagents.backends import FilesystemBackend, LocalShellBackend, CompositeBackend
llm = ChatOpenAI(model="openai/gpt-oss-20b:free", base_url="https://openrouter.ai/api/v1", api_key=os.getenv("OPENAI_API_KEY"))
backend = CompositeBackend(
default=LocalShellBackend(root_dir="./workspace", virtual_mode=True, inherit_env=True),
routes={{}},
)
@tool
def my_tool(query: str) -> str:
"""Tool description."""
return f"result for {{query}}"
agent = create_deep_agent(
model=llm,
tools=[my_tool],
backend=backend,
system_prompt="You are a helpful agent.",
)
async def main():
result = await agent.ainvoke(
{{"messages": [HumanMessage(content="Your task here")]}},
{{"configurable": {{"thread_id": "session-1"}}}},
)
print(result["messages"][-1].content)
if __name__ == "__main__":
asyncio.run(main())
```
requirements.txt: deepagents, langchain-openai>=0.3.0, langchain>=1.2.10, langgraph>=0.2.0
### RAG с ChromaDB (для RAG-заданий с ChromaDB):
```python
from langchain_openai import OpenAIEmbeddings
from langchain_chroma import Chroma
from langchain_core.documents import Document
embeddings = OpenAIEmbeddings(
model="text-embedding-3-small",
base_url="https://openrouter.ai/api/v1",
api_key=os.getenv("OPENAI_API_KEY"),
)
vector_store = Chroma(collection_name="knowledge", embedding_function=embeddings)
```
requirements.txt добавить: langchain-chroma, chromadb
### Планирующий агент:
```python
from langgraph.graph import StateGraph, START, END
from typing import TypedDict, Annotated
from langgraph.graph.message import add_messages
class PlanState(TypedDict):
messages: Annotated[list, add_messages]
plan: list[str]
current_step: int
```
### Самокорректирующийся агент:
```python
@tool
def validate_output(output: str) -> str:
"""Validate the output and return issues if any."""
issues = []
if len(output) < 10:
issues.append("Output too short")
return "OK" if not issues else f"Issues: {{', '.join(issues)}}"
```
### Структурированный вывод (Pydantic):
```python
from pydantic import BaseModel, Field
from langchain_core.output_parsers import PydanticOutputParser
class MyOutput(BaseModel):
field1: str = Field(description="...")
field2: int = Field(description="...")
```
## Требования
- Полный рабочий код без заглушек (no pass, TODO, ...)
- ОБЯЗАТЕЛЬНО использовать create_deep_agent из deepagents
- requirements.txt: deepagents, langchain>=1.2.10, langchain-openai>=0.3.0, langgraph>=0.2.0 + нужные доп. зависимости
{rework_section}
## Ответ — ТОЛЬКО JSON без markdown:
{{"main_py": "...", "requirements_txt": "...", "extra_files": {{}}}}
extra_files — только если нужны доп. файлы, иначе пустой объект.
'''
+249
View File
@@ -0,0 +1,249 @@
"""
Инструменты агента для решения задач.
Агент вызывает эти инструменты САМОСТОЯТЕЛЬНО — Python не управляет порядком.
Каждый инструмент — специализированный LLM-субагент со своим промптом.
Субагенты:
validate_teacher_comment — per-claim валидация замечания преподавателя
generate_code_solution — генерация кода (первая сдача или пересдача)
"""
import asyncio
import base64
import json
import os
import httpx
from langchain.tools import tool
from src.agent.constants import GITEA_BASE_URL, GITEA_OWNER
from src.agent.llm import llm
from src.agent.solve_prompts import (
ANALYZE_PROMPT,
CODE_PROMPT,
REWORK_SECTION,
VALIDATE_PROMPT,
)
_GITEA_TOKEN = os.getenv("GITEA_TOKEN", "")
_CODE_EXTS = (".py", ".js", ".ts", ".sh", ".sql", ".md")
_BACKOFF = [30, 60, 120]
# ---------------------------------------------------------------------------
# Вспомогательные функции (не инструменты)
# ---------------------------------------------------------------------------
def _gh() -> dict:
return {"Authorization": f"token {_GITEA_TOKEN}", "Content-Type": "application/json"}
def _read_repo_files(repo: str) -> dict[str, str]:
"""Читает все кодовые файлы из корня репозитория на Gitea."""
files: dict[str, str] = {}
url_root = f"{GITEA_BASE_URL}/api/v1/repos/{GITEA_OWNER}/{repo}/contents"
try:
with httpx.Client(timeout=30) as c:
r = c.get(url_root, headers=_gh())
if r.status_code != 200:
return files
for item in r.json():
if item.get("type") != "file":
continue
if not any(item["name"].endswith(ext) for ext in _CODE_EXTS):
continue
fr = c.get(
f"{GITEA_BASE_URL}/api/v1/repos/{GITEA_OWNER}/{repo}/contents/{item['name']}",
headers=_gh(),
)
if fr.status_code == 200:
raw = fr.json().get("content", "")
files[item["name"]] = base64.b64decode(raw.replace("\n", "")).decode(
"utf-8", errors="replace"
)
except Exception:
pass
return files
def _parse_llm_json(raw: str) -> dict:
"""Убирает markdown-обёртку и парсит JSON из ответа LLM."""
text = raw.strip()
if text.startswith("```"):
parts = text.split("```")
text = parts[1] if len(parts) > 1 else text
if text.startswith("json"):
text = text[4:]
text = text.strip()
return json.loads(text)
async def _llm_call_with_retry(prompt: str, max_attempts: int = 4) -> str:
"""Вызов LLM с повтором при 429."""
for attempt in range(1, max_attempts + 1):
try:
resp = await llm.ainvoke(prompt)
return resp.content
except Exception as e:
if "429" in str(e) and attempt < max_attempts:
wait = _BACKOFF[min(attempt - 1, len(_BACKOFF) - 1)]
print(f" [solve_tools] 429, жду {wait}с (попытка {attempt})...")
await asyncio.sleep(wait)
else:
raise
# ---------------------------------------------------------------------------
# Субагент 1: Валидатор замечаний преподавателя
# ---------------------------------------------------------------------------
@tool
async def validate_teacher_comment(
task_text: str,
repo_name: str,
teacher_comment: str,
) -> str:
"""[СУБАГЕНТ-ВАЛИДАТОР] Анализирует каждый пункт замечания преподавателя НЕЗАВИСИМО.
Читает текущий код из Gitea репозитория и сверяет каждое утверждение
с условием задания. Отличает ловушки от реальных ошибок.
Args:
task_text: полный текст условия задания
repo_name: имя репозитория (например task-6a1864f7fd30e81cf3...)
teacher_comment: замечание преподавателя
Returns:
JSON: {
"has_trap": true если есть ложные пункты,
"has_valid": true если есть реальные ошибки,
"claims": список {claim, valid, explanation},
"trap_explanations": объяснения ложных пунктов (для README),
"fix_instructions": что конкретно исправить (для реальных ошибок)
}
"""
print(f" [ВАЛИДАТОР] Проверяю замечание для {repo_name}...")
# Читаем код из Gitea
code_files = _read_repo_files(repo_name)
if code_files:
code_block = "\n\n".join(
f"### {fn}\n```\n{content[:2000]}\n```"
for fn, content in code_files.items()
)
print(f" [ВАЛИДАТОР] Прочитано файлов: {', '.join(code_files.keys())}")
else:
code_block = "(репозиторий пуст или файлы не найдены)"
print(f" [ВАЛИДАТОР] ⚠️ Файлы в {repo_name} не найдены")
prompt = VALIDATE_PROMPT.format(
task_text=task_text,
code_block=code_block,
comment=teacher_comment,
)
try:
raw = await _llm_call_with_retry(prompt)
result = _parse_llm_json(raw)
result.setdefault("has_trap", False)
result.setdefault("has_valid", True)
result.setdefault("claims", [])
result.setdefault("trap_explanations", [])
result.setdefault("fix_instructions", [])
# Логируем результат
for cl in result["claims"]:
tag = "❌ ЛОВУШКА" if not cl.get("valid") else "✓ обоснованно"
print(f" {tag}: {cl.get('claim', '')[:70]}")
print(f" [ВАЛИДАТОР] has_trap={result['has_trap']}, has_valid={result['has_valid']}")
return json.dumps(result, ensure_ascii=False)
except Exception as e:
print(f" [ВАЛИДАТОР] Ошибка: {e} — считаем замечание обоснованным")
fallback = {
"has_trap": False,
"has_valid": True,
"claims": [],
"trap_explanations": [],
"fix_instructions": [teacher_comment],
}
return json.dumps(fallback, ensure_ascii=False)
# ---------------------------------------------------------------------------
# Субагент 2: Кодер — генерирует решение
# ---------------------------------------------------------------------------
@tool
async def generate_code_solution(
task_text: str,
fix_instructions: str = "",
defense_context: str = "",
) -> str:
"""[СУБАГЕНТ-КОДЕР] Генерирует полное решение задания.
При пересдаче принимает что исправить и что отстоять с DESIGN DECISION аргументами.
Запрещено использовать Ollama — только OpenRouter через langchain_openai.
Args:
task_text: полный текст условия задания
fix_instructions: что конкретно исправить (для пересдачи, иначе "")
defense_context: что отстоять с DESIGN DECISION комментариями (иначе "")
Returns:
JSON: {
"main_py": содержимое main.py,
"requirements_txt": содержимое requirements.txt,
"extra_files": доп. файлы {имя: содержимое} или {}
}
"""
is_rework = bool(fix_instructions or defense_context)
mode = "ПЕРЕСДАЧА" if is_rework else "первая сдача"
print(f" [КОДЕР] Генерирую решение ({mode})...")
if is_rework:
# Строим секцию пересдачи
fixes = [fix_instructions] if fix_instructions else []
defenses = [defense_context] if defense_context else []
rework_section = REWORK_SECTION.format(
fixes = "\n".join(f"- {f}" for f in fixes) or "— нет",
defenses = "\n".join(f"- {d}" for d in defenses) or "— нет",
)
else:
rework_section = ""
prompt = CODE_PROMPT.format(task_text=task_text, rework_section=rework_section)
for attempt in range(1, 6):
try:
print(f" [КОДЕР] LLM вызов (попытка {attempt})...")
raw = await _llm_call_with_retry(prompt, max_attempts=3)
result = _parse_llm_json(raw)
if "main_py" in result:
main_size = len(result.get("main_py", ""))
req_size = len(result.get("requirements_txt", ""))
extra = list(result.get("extra_files", {}).keys())
print(f" [КОДЕР] ✅ main.py={main_size}с, requirements.txt={req_size}с"
+ (f", extra={extra}" if extra else ""))
return json.dumps(result, ensure_ascii=False)
except json.JSONDecodeError:
print(f" [КОДЕР] JSON parse error на попытке {attempt}, повтор...")
if attempt == 5:
raise
except Exception as e:
if attempt < 5:
wait = _BACKOFF[min(attempt - 1, len(_BACKOFF) - 1)]
print(f" [КОДЕР] Ошибка: {e}, жду {wait}с...")
await asyncio.sleep(wait)
else:
raise
raise RuntimeError("Не удалось сгенерировать код после 5 попыток")
# ---------------------------------------------------------------------------
# Список инструментов для импорта в agent.py
# ---------------------------------------------------------------------------
SOLVE_TOOLS = [validate_teacher_comment, generate_code_solution]
+582
View File
@@ -0,0 +1,582 @@
"""
Streamlit UI для brojs-agent.
Запуск: streamlit run ui.py
"""
import asyncio
import json
import os
import queue
import threading
import time
from datetime import datetime
os.environ["NO_PROXY"] = (
"openrouter.ai,platform.brojs.ru,git.brojs.ru,"
+ os.environ.get("NO_PROXY", "")
)
import streamlit as st
from dotenv import load_dotenv
from langchain_core.callbacks.base import BaseCallbackHandler
from langchain_core.messages import AIMessage, HumanMessage
load_dotenv()
# ---------------------------------------------------------------------------
# Конфигурация страницы
# ---------------------------------------------------------------------------
st.set_page_config(
page_title="BroJS Agent",
page_icon="🤖",
layout="wide",
initial_sidebar_state="collapsed",
)
GITEA_OWNER = os.getenv("GITEA_OWNER", "glevelll")
# ---------------------------------------------------------------------------
# CSS
# ---------------------------------------------------------------------------
st.markdown("""
<style>
.agent-header {
background: linear-gradient(135deg, #1a1a2e 0%, #0f3460 100%);
border-radius: 12px; padding: 18px 24px; margin-bottom: 16px;
border: 1px solid #16213e;
}
.agent-title { font-size: 1.6em; font-weight: bold; color: #e2e8f0; margin: 0; }
.agent-sub { color: #64748b; font-size: .85em; margin-top: 4px; }
.tool-call {
background: #0f172a; border-left: 3px solid #3b82f6;
border-radius: 6px; padding: 6px 12px; margin: 3px 0;
font-family: monospace; font-size: .82em; color: #93c5fd;
}
.tool-result {
background: #052e16; border-left: 3px solid #22c55e;
border-radius: 6px; padding: 6px 12px; margin: 3px 0;
font-family: monospace; font-size: .78em; color: #86efac;
}
.tool-subagent {
background: #1e1b4b; border-left: 3px solid #818cf8;
border-radius: 6px; padding: 6px 12px; margin: 3px 0;
font-family: monospace; font-size: .82em; color: #c4b5fd;
}
.thinking {
color: #64748b; font-style: italic; font-size: .82em; padding: 4px 0;
}
.status-ok { background:#052e16; border:1px solid #22c55e; color:#4ade80; padding:10px 16px; border-radius:8px; }
.status-warn { background:#1c1917; border:1px solid #f59e0b; color:#fbbf24; padding:10px 16px; border-radius:8px; }
.status-err { background:#1c0a0a; border:1px solid #ef4444; color:#f87171; padding:10px 16px; border-radius:8px; }
.chat-user { background:#1e3a5f; border-radius:12px 12px 2px 12px; padding:10px 14px; margin:6px 0; }
.chat-agent { background:#1a1a2e; border-radius:12px 12px 12px 2px; padding:10px 14px; margin:6px 0; }
</style>
""", unsafe_allow_html=True)
# ---------------------------------------------------------------------------
# Заголовок
# ---------------------------------------------------------------------------
st.markdown("""
<div class="agent-header">
<div class="agent-title">🤖 BroJS Agent</div>
<div class="agent-sub">Агентная система выполнения заданий · KFU-26-1 · platform.brojs.ru</div>
</div>
""", unsafe_allow_html=True)
# ---------------------------------------------------------------------------
# Кэш агента
# ---------------------------------------------------------------------------
@st.cache_resource(show_spinner="Инициализация агента (~30с)...")
def get_agent():
from src.agent.agent import agent
return agent
@st.cache_resource(show_spinner="Загрузка главного агента...")
def get_main_agent():
from src.agent.agent import agent
return agent
# ---------------------------------------------------------------------------
# Callback — перехватывает события агента и шлёт в очередь
# ---------------------------------------------------------------------------
_SOLVE_TOOLS = {"validate_teacher_comment", "generate_code_solution"}
_JOURNAL_PREFIX = "mcp__journal-bh-professor__"
class AgentCallback(BaseCallbackHandler):
def __init__(self, q: queue.Queue):
self.q = q
def _ts(self) -> str:
return datetime.now().strftime("%H:%M:%S")
def on_tool_start(self, serialized, input_str, **kwargs):
name = serialized.get("name", "?")
try:
args = json.loads(str(input_str)) if isinstance(input_str, str) else input_str
except Exception:
args = {}
self.q.put({"t": "tool_start", "ts": self._ts(), "name": name, "args": args})
def on_tool_end(self, output, **kwargs):
self.q.put({"t": "tool_end", "ts": self._ts(), "output": str(output)[:300]})
def on_tool_error(self, error, **kwargs):
self.q.put({"t": "tool_error", "ts": self._ts(), "msg": str(error)[:200]})
def on_llm_start(self, *a, **kw):
self.q.put({"t": "thinking", "ts": self._ts()})
def on_llm_end(self, response, **kwargs):
try:
text = response.generations[0][0].text[:120]
self.q.put({"t": "llm_end", "ts": self._ts(), "preview": text})
except Exception:
pass
# ---------------------------------------------------------------------------
# Запуск агента в фоне
# ---------------------------------------------------------------------------
def _run_agent_thread(agent, messages, config, q: queue.Queue, cb: AgentCallback):
async def _inner():
try:
result = await agent.ainvoke(messages, {**config, "callbacks": [cb]})
q.put({"t": "done", "result": result})
except Exception as e:
q.put({"t": "fatal", "msg": str(e)})
asyncio.run(_inner())
# ---------------------------------------------------------------------------
# Рендер одного события в лог
# ---------------------------------------------------------------------------
def _render_event(ev: dict) -> str:
ts = ev.get("ts", "")
kind = ev.get("t", "")
if kind == "thinking":
return f'<div class="thinking">💭 {ts} модель думает...</div>'
if kind == "tool_start":
name = ev["name"]
args = ev.get("args", {})
short = name.replace(_JOURNAL_PREFIX, "mcp::")
# Определяем тип инструмента
if name in _SOLVE_TOOLS:
cls = "tool-subagent"
icon = "🧠"
label = f"[субагент] {short}"
elif "gitea" in name:
cls = "tool-call"
icon = "📦"
label = short
elif "mcp::" in short or "journal" in name:
cls = "tool-call"
icon = "📡"
label = short
else:
cls = "tool-call"
icon = "🔧"
label = short
# Показываем ключевые аргументы
hint = ""
for key in ("taskId", "path", "repo", "repo_name", "name"):
if key in args:
hint = f' <span style="opacity:.6">{args[key]}</span>'
break
return f'<div class="{cls}">{icon} {ts} {label}{hint}</div>'
if kind == "tool_end":
out = ev["output"].replace("<", "&lt;").replace(">", "&gt;")[:200]
return f'<div class="tool-result">↳ {out}</div>'
if kind == "tool_error":
msg = ev["msg"].replace("<", "&lt;")
return f'<div class="tool-result" style="border-color:#ef4444;color:#f87171">⚠ {msg}</div>'
if kind == "llm_end":
preview = ev.get("preview", "").replace("<", "&lt;")[:100]
return f'<div class="thinking">✏ {ts} {preview}...</div>'
return ""
# ---------------------------------------------------------------------------
# Вкладки
# ---------------------------------------------------------------------------
tab_chat, tab_pipeline, tab_status = st.tabs(["💬 Чат с агентом", "⚡ Pipeline", "📊 Статус заданий"])
# ══════════════════════════════════════════════════════════════════════════
# ВК 1 — ЧАТ
# ══════════════════════════════════════════════════════════════════════════
with tab_chat:
st.caption("Общайся с агентом: задай вопрос, попроси решить задание или разобрать ситуацию.")
# История сообщений
if "chat_history" not in st.session_state:
st.session_state.chat_history = []
if "chat_events" not in st.session_state:
st.session_state.chat_events = []
if "chat_thread_id" not in st.session_state:
st.session_state.chat_thread_id = f"ui-{int(time.time())}"
# Показываем историю
for msg in st.session_state.chat_history:
role = msg["role"]
text = msg["text"]
if role == "user":
st.markdown(f'<div class="chat-user">👤 {text}</div>', unsafe_allow_html=True)
else:
st.markdown(f'<div class="chat-agent">🤖 {text}</div>', unsafe_allow_html=True)
# Лог событий (раскрывающийся)
if st.session_state.chat_events:
with st.expander(f"🔍 Лог инструментов ({len(st.session_state.chat_events)} событий)", expanded=False):
html = "".join(_render_event(e) for e in st.session_state.chat_events[-80:])
st.markdown(f'<div style="max-height:300px;overflow-y:auto">{html}</div>',
unsafe_allow_html=True)
# Ввод
col_input, col_btn = st.columns([5, 1])
with col_input:
user_input = st.text_input(
"Сообщение",
placeholder='Например: "Реши задание 6a1864f7..." или "Какие задания у меня есть?"',
label_visibility="collapsed",
key="chat_input",
)
with col_btn:
send = st.button("Отправить", use_container_width=True, type="primary")
if send and user_input.strip():
msg_text = user_input.strip()
st.session_state.chat_history.append({"role": "user", "text": msg_text})
st.session_state.chat_events = []
# Строим историю сообщений для агента
lc_messages = []
for m in st.session_state.chat_history:
if m["role"] == "user":
lc_messages.append(HumanMessage(content=m["text"]))
else:
lc_messages.append(AIMessage(content=m["text"]))
config = {"configurable": {"thread_id": st.session_state.chat_thread_id}}
agent = get_agent()
# Placeholders для обновления в реальном времени
events_ph = st.empty()
status_ph = st.empty()
evq: queue.Queue = queue.Queue()
cb = AgentCallback(evq)
all_events: list[dict] = []
thread = threading.Thread(
target=_run_agent_thread,
args=(agent, {"messages": lc_messages}, config, evq, cb),
daemon=True,
)
thread.start()
final_result = None
fatal = None
while thread.is_alive() or not evq.empty():
changed = False
while not evq.empty():
ev = evq.get_nowait()
if ev["t"] in ("done", "fatal"):
if ev["t"] == "done":
final_result = ev["result"]
else:
fatal = ev["msg"]
else:
all_events.append(ev)
changed = True
if changed and all_events:
html = "".join(_render_event(e) for e in all_events[-60:])
events_ph.markdown(
f'<div style="background:#0b0f1a;border-radius:8px;padding:10px;'
f'max-height:250px;overflow-y:auto">{html}</div>',
unsafe_allow_html=True,
)
time.sleep(0.15)
events_ph.empty()
st.session_state.chat_events = all_events
if fatal:
st.session_state.chat_history.append({"role": "agent", "text": f"⚠️ Ошибка: {fatal}"})
elif final_result:
msgs = final_result.get("messages", [])
last = msgs[-1] if msgs else None
reply = last.content if last and hasattr(last, "content") else "Готово."
st.session_state.chat_history.append({"role": "agent", "text": reply})
st.rerun()
# Кнопка очистки
if st.session_state.chat_history:
if st.button("🗑 Очистить чат"):
st.session_state.chat_history = []
st.session_state.chat_events = []
st.session_state.chat_thread_id = f"ui-{int(time.time())}"
st.rerun()
# ══════════════════════════════════════════════════════════════════════════
# ВК 2 — PIPELINE
# ══════════════════════════════════════════════════════════════════════════
with tab_pipeline:
st.caption("Автоматически решает все todo-задания курса по очереди.")
col1, col2 = st.columns([3, 1])
with col1:
task_id_input = st.text_input(
"Task ID (оставь пустым — решить все todo)",
placeholder="6a1864f7fd30e81cf3146d65",
label_visibility="visible",
)
with col2:
st.write("")
run_btn = st.button("▶ Запустить", type="primary", use_container_width=True)
if run_btn:
result_ph = st.empty()
events_ph2 = st.empty()
agent = get_agent()
if task_id_input.strip():
# Одно задание
task_id = task_id_input.strip()
repo_url = f"https://git.brojs.ru/{GITEA_OWNER}/task-{task_id}"
prompt = f"Реши задание taskId={task_id} курса 698b49da77cb6d4d2e43ce78"
config = {"configurable": {"thread_id": f"pipe-{task_id}-{int(time.time())}"}}
messages = {"messages": [HumanMessage(content=prompt)]}
else:
result_ph.info("Pipeline для всех todo-заданий — используй раздел ниже")
st.stop()
evq2: queue.Queue = queue.Queue()
cb2 = AgentCallback(evq2)
all_events2: list[dict] = []
thread2 = threading.Thread(
target=_run_agent_thread,
args=(agent, messages, config, evq2, cb2),
daemon=True,
)
thread2.start()
final2 = None
fatal2 = None
with st.spinner(f"Агент решает {task_id[:8]}..."):
while thread2.is_alive() or not evq2.empty():
while not evq2.empty():
ev = evq2.get_nowait()
if ev["t"] == "done":
final2 = ev["result"]
elif ev["t"] == "fatal":
fatal2 = ev["msg"]
else:
all_events2.append(ev)
if all_events2:
html = "".join(_render_event(e) for e in all_events2[-50:])
events_ph2.markdown(
f'<div style="background:#0b0f1a;border-radius:8px;padding:10px;'
f'max-height:300px;overflow-y:auto">{html}</div>',
unsafe_allow_html=True,
)
time.sleep(0.15)
if fatal2:
result_ph.markdown(
f'<div class="status-err">❌ Ошибка: {fatal2[:300]}</div>',
unsafe_allow_html=True,
)
elif final2:
result_ph.markdown(
f'<div class="status-ok">✅ Готово! '
f'<a href="{repo_url}" target="_blank" style="color:#4ade80">Открыть репозиторий</a>'
f'</div>',
unsafe_allow_html=True,
)
st.divider()
st.subheader("Запустить все todo-задания")
if st.button("⚡ Запустить агент для всех заданий", use_container_width=True):
main_ag = get_main_agent()
evq_pl = queue.Queue()
cb_pl = AgentCallback(evq_pl)
all_pl_events: list[dict] = []
_pl_state = {"result": None, "error": None}
run_prompt = (
"Выполни все задания со статусом todo в курсе KFU-26-1 "
"(courseId=698b49da77cb6d4d2e43ce78).\n\n"
"Шаги:\n"
"1. Получи список заданий через mcp__journal-bh-professor__tasks_list\n"
"2. Для каждого задания со статусом todo вызови solve_task(task_id=...)\n"
"3. Выполняй строго по одному заданию, жди результата перед следующим\n"
"4. Доложи итоговые результаты"
)
def _run_all():
async def _inner():
try:
result = await main_ag.ainvoke(
{"messages": [HumanMessage(content=run_prompt)]},
{
"configurable": {"thread_id": f"ui-run-all-{int(time.time())}"},
"callbacks": [cb_pl],
},
)
_pl_state["result"] = result
evq_pl.put({"t": "done", "result": result})
except Exception as e:
_pl_state["error"] = str(e)
evq_pl.put({"t": "fatal", "msg": str(e)})
asyncio.run(_inner())
t_pl = threading.Thread(target=_run_all, daemon=True)
t_pl.start()
events_pl_ph = st.empty()
with st.spinner("Агент-оркестратор работает... (LLM управляет всем)"):
while t_pl.is_alive() or not evq_pl.empty():
while not evq_pl.empty():
ev = evq_pl.get_nowait()
if ev["t"] not in ("done", "fatal"):
all_pl_events.append(ev)
if all_pl_events:
html = "".join(_render_event(e) for e in all_pl_events[-60:])
events_pl_ph.markdown(
f'<div style="background:#0b0f1a;border-radius:8px;padding:10px;'
f'max-height:300px;overflow-y:auto">{html}</div>',
unsafe_allow_html=True,
)
time.sleep(0.15)
events_pl_ph.empty()
if all_pl_events:
with st.expander(f"🔍 Лог агента ({len(all_pl_events)} событий)", expanded=False):
html = "".join(_render_event(e) for e in all_pl_events[-80:])
st.markdown(f'<div style="max-height:300px;overflow-y:auto">{html}</div>',
unsafe_allow_html=True)
if _pl_state["error"]:
st.error(_pl_state["error"])
elif _pl_state["result"]:
msgs = _pl_state["result"].get("messages", [])
last = msgs[-1] if msgs else None
reply = last.content if last and hasattr(last, "content") else "Готово."
st.markdown(
f'<div class="status-ok">✅ Агент завершил работу:<br>{reply[:600]}</div>',
unsafe_allow_html=True,
)
# ══════════════════════════════════════════════════════════════════════════
# ВК 3 — СТАТУС
# ══════════════════════════════════════════════════════════════════════════
with tab_status:
st.caption("Статусы всех заданий курса KFU-26-1.")
if st.button("🔄 Обновить статусы", type="primary"):
with st.spinner("Загружаю статусы..."):
from src.agent.mcp_client import load_journal_toolsets
async def _fetch():
j = load_journal_toolsets()
tools = {t.name: t for t in j.tasks_submissions_tools}
full_name = "mcp__journal-bh-professor__tasks_list"
t = tools.get(full_name) or next(
(v for k, v in tools.items() if "tasks_list" in k), None
)
if not t:
return [], f"tasks_list не найден. Доступны: {list(tools.keys())}"
raw = await t.ainvoke({"courseId": "698b49da77cb6d4d2e43ce78"})
text = next((x["text"] for x in raw if x.get("type") == "text"), str(raw)) if isinstance(raw, list) else str(raw)
data = json.loads(text)
return (data.get("tasks", data) if isinstance(data, dict) else data), None
_state = {"items": [], "error": None}
def _run_fetch():
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
try:
_state["items"], _state["error"] = loop.run_until_complete(_fetch())
except Exception as e:
_state["error"] = str(e)
finally:
loop.close()
t = threading.Thread(target=_run_fetch, daemon=True)
t.start()
t.join()
if _state["error"]:
st.error(_state["error"])
else:
st.session_state["task_statuses"] = _state["items"]
st.rerun()
items = st.session_state.get("task_statuses", [])
STATUS_EMOJI = {
"done": "",
"ready_for_review": "🔍",
"in_progress": "🔄",
"todo": "📋",
"rejected": "",
}
if items:
counts: dict[str, int] = {}
rows = []
for item in items:
t = item.get("task", item) if isinstance(item, dict) else {}
tid = t.get("id", "")
status = item.get("status", "")
title = t.get("title", t.get("name", ""))
counts[status] = counts.get(status, 0) + 1
rows.append({
"": STATUS_EMOJI.get(status, ""),
"Статус": status,
"ID": tid[:12] + "...",
"Название": title,
"Репо": f"https://git.brojs.ru/{GITEA_OWNER}/task-{tid}",
})
st.dataframe(rows, use_container_width=True, hide_index=True)
st.divider()
cols = st.columns(len(counts))
for col, (s, n) in zip(cols, counts.items()):
col.metric(f"{STATUS_EMOJI.get(s,'')} {s}", n)
else:
st.info("Нажми «Обновить статусы» чтобы загрузить данные.")