From ba83d0cfe458cddac22190a55e0bb5dc9bf29848 Mon Sep 17 00:00:00 2001 From: gleb Date: Thu, 11 Jun 2026 18:43:40 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20LLM-first=20architecture=20=E2=80=94=20?= =?UTF-8?q?agent=20orchestrates=20everything=20via=20solve=5Ftask=20tool?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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. --- cli.py | 32 ++++++---- src/agent/agent.py | 5 +- src/agent/prompts.py | 58 ++++++++++-------- src/agent/runner_tools.py | 108 +++++++++++++++++++++++++++++++++ ui.py | 124 ++++++++++++++++++++++---------------- 5 files changed, 237 insertions(+), 90 deletions(-) create mode 100644 src/agent/runner_tools.py diff --git a/cli.py b/cli.py index 2f5097e..6a2e10b 100644 --- a/cli.py +++ b/cli.py @@ -50,20 +50,30 @@ async def cmd_solve(task_id: str) -> None: async def cmd_run() -> None: - """Решить все todo-задания курса.""" - from src.agent.graph.pipeline import pipeline + """Решить все todo-задания курса (LLM-оркестратор управляет всем).""" + import time + from langchain_core.messages import HumanMessage + from src.agent.agent import agent - print("[cli] Запускаю pipeline для всех todo-заданий...") - result = await pipeline.ainvoke( - {"tasks": [], "current_index": 0, "results": [], "errors": []} + 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}") - for r in result.get("results", []): - icon = "✅" if r.get("status") == "ok" else "❌" - url = r.get("url", r.get("error", "")) - print(f" {icon} {r.get('task_id','')[:8]}... → {url}") - for e in result.get("errors", []): - print(f" ⚠️ {e}") + print(getattr(final, "content", str(final))) print('='*60) diff --git a/src/agent/agent.py b/src/agent/agent.py index 882e77d..3a6b2ed 100644 --- a/src/agent/agent.py +++ b/src/agent/agent.py @@ -19,6 +19,7 @@ from src.agent.prompts import ( rework_instructions, ) 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 @@ -75,7 +76,7 @@ _git_names = {t.name for t in GIT_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]] = { "web_search": _BUILTIN | _web_names, @@ -113,7 +114,7 @@ subagents = [ agent = create_deep_agent( model=llm, - tools=list(GITEA_TOOLS), + tools=[*GITEA_TOOLS, *_journal_tools, solve_task], system_prompt=main_agent_instructions, backend=_composite_backend, memory=[AGENTS_MD_VFS_PATH], diff --git a/src/agent/prompts.py b/src/agent/prompts.py index 4993a69..dbf43cd 100644 --- a/src/agent/prompts.py +++ b/src/agent/prompts.py @@ -163,43 +163,49 @@ Gitea owner = "glevelll" main_agent_instructions = """\ Ты — главный агент-исполнитель домашних заданий курса KFU-26-1 на platform.brojs.ru. -Твоя роль — получать задания из журнала и выполнять их качественно. +Ты полностью управляешь всем процессом: сам решаешь что делать, в каком порядке, +какие инструменты вызывать. Python-инструменты — это только твои руки. ## Известные курсы - KFU-26-1 = courseId `698b49da77cb6d4d2e43ce78` -## Доступные субагенты (вызывай через инструмент `task`) -- `journal_bh_tasks_submissions`: читает задания, проверяет статусы, отправляет ответы -- `homework_doing`: ВЫПОЛНЯЕТ задание (пишет код, создаёт репо, сдаёт) -- `web_search`: ищет информацию в интернете (только если нужно) +## Прямые инструменты (вызывай напрямую) -## Прямые Gitea-инструменты (доступны напрямую без субагента) -- `gitea_list_repos` — список репозиториев, также возвращает username -- `gitea_create_repo` — создать репозиторий -- `gitea_write_file` — создать/обновить файл (автокоммит) -- `gitea_get_file` — получить файл +### Journal (MCP) +- `mcp__journal-bh-professor__tasks_list` — список всех заданий курса со статусами +- `mcp__journal-bh-professor__task_get` — детали задания (статус, ответ, комментарии) +- `mcp__journal-bh-professor__task_text` — полный текст задания +- `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="")` +4. Дождись "OK ..." перед следующим заданием +5. Доложи итоги -### Получить список заданий -1. Делегируй `journal_bh_tasks_submissions`: получить tasks_list для courseId +## Как выполнить одно задание +1. `solve_task(task_id="")` — сделает всё сам -### Выполнить задание -1. Делегируй `homework_doing`: выполни задание с taskId= - (он сам прочитает текст, создаст репо, напишет код и сдаст) -2. Верни пользователю ссылку на репозиторий - -### Проверить статусы -1. Делегируй `journal_bh_tasks_submissions`: получить статусы всех заданий курса +## Как проверить статусы +1. `mcp__journal-bh-professor__tasks_list(courseId="698b49da77cb6d4d2e43ce78")` ## Жёсткие ограничения -- Не вызывай больше одного инструмента за шаг -- Не делегируй субагенту несколько независимых задач сразу -- Не говори что задание выполнено, если оно не было реально выполнено +- Не говори что задание выполнено, если `solve_task` не вернул "OK" +- Выполняй задания строго последовательно, не параллельно +- Не вызывай `solve_task` для заданий со статусом `done` или `ready_for_review` - Не подменяй требования задания своими догадками """ diff --git a/src/agent/runner_tools.py b/src/agent/runner_tools.py new file mode 100644 index 0000000..b6cfd9e --- /dev/null +++ b/src/agent/runner_tools.py @@ -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}" diff --git a/ui.py b/ui.py index 0eccbd8..5b31446 100644 --- a/ui.py +++ b/ui.py @@ -94,14 +94,14 @@ st.markdown(""" @st.cache_resource(show_spinner="Инициализация агента (~30с)...") def get_agent(): - from src.agent.agent import homework_direct_agent - return homework_direct_agent + from src.agent.agent import agent + return agent -@st.cache_resource(show_spinner="Загрузка pipeline...") -def get_pipeline(): - from src.agent.graph.pipeline import pipeline - return pipeline +@st.cache_resource(show_spinner="Загрузка главного агента...") +def get_main_agent(): + from src.agent.agent import agent + return agent # --------------------------------------------------------------------------- # Callback — перехватывает события агента и шлёт в очередь @@ -425,29 +425,45 @@ with tab_pipeline: st.divider() st.subheader("Запустить все todo-задания") - if st.button("⚡ Запустить pipeline для всех заданий", use_container_width=True): - pl = get_pipeline() + 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} - def _run_pipeline(): + 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: - _pl_state["result"] = await pl.ainvoke( - {"tasks": [], "current_index": 0, "results": [], "errors": []}, - {"callbacks": [cb_pl]}, + 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_pipeline, daemon=True) + t_pl = threading.Thread(target=_run_all, daemon=True) t_pl.start() events_pl_ph = st.empty() - with st.spinner("Pipeline работает... (может занять несколько минут)"): + with st.spinner("Агент-оркестратор работает... (LLM управляет всем)"): while t_pl.is_alive() or not evq_pl.empty(): while not evq_pl.empty(): ev = evq_pl.get_nowait() @@ -465,7 +481,7 @@ with tab_pipeline: events_pl_ph.empty() if all_pl_events: - with st.expander(f"🔍 Лог pipeline ({len(all_pl_events)} событий)", expanded=False): + 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'
{html}
', unsafe_allow_html=True) @@ -473,19 +489,13 @@ with tab_pipeline: if _pl_state["error"]: st.error(_pl_state["error"]) elif _pl_state["result"]: - results = _pl_state["result"].get("results", []) - errors = _pl_state["result"].get("errors", []) - md = [f"### Результат: {len(results)} заданий\n"] - for r in results: - tid = r.get("task_id", "") - url = f"https://git.brojs.ru/{GITEA_OWNER}/task-{tid}" - icon = "✅" if r.get("status") == "ok" else "❌" - md.append(f"- {icon} `{tid[:8]}...` — [{r.get('status','')}]({url})") - if errors: - md.append(f"\n**Ошибки ({len(errors)}):**") - for e in errors: - md.append(f"- {e}") - st.markdown("\n".join(md)) + 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'
✅ Агент завершил работу:
{reply[:600]}
', + unsafe_allow_html=True, + ) # ══════════════════════════════════════════════════════════════════════════ @@ -497,31 +507,43 @@ with tab_status: if st.button("🔄 Обновить статусы", type="primary"): with st.spinner("Загружаю статусы..."): - try: - from src.agent.mcp_client import load_journal_toolsets + from src.agent.mcp_client import load_journal_toolsets - async def _fetch(): - j = load_journal_toolsets() - # Инструменты имеют префикс mcp__journal-bh-professor__ - tools = {t.name: t for t in j.tasks_submissions_tools} - full_name = "mcp__journal-bh-professor__tasks_list" - t = tools.get(full_name) - if not t: - # fallback: ищем по любому имени содержащему tasks_list - t = next((v for k, v in tools.items() if "tasks_list" in k), None) - if not t: - st.warning(f"Инструмент tasks_list не найден. Доступны: {list(tools.keys())}") - return [] - 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 + 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 - items = asyncio.run(_fetch()) - st.session_state["task_statuses"] = items - except Exception as e: - st.error(str(e)) - items = [] + _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", [])