import asyncio
import ast
import json
import queue as q_mod
import threading
import time
from datetime import datetime
import streamlit as st
from langchain_core.callbacks.base import BaseCallbackHandler
from langchain_core.messages import HumanMessage
OWNER = "KirillKutlakhmetov"
# ── Page config ───────────────────────────────────────────────────────────────
st.set_page_config(
page_title="BroJS Agent",
layout="wide",
initial_sidebar_state="collapsed",
)
# ── CSS ───────────────────────────────────────────────────────────────────────
st.markdown("""
""", unsafe_allow_html=True)
# ── Header ────────────────────────────────────────────────────────────────────
st.markdown("""
""", unsafe_allow_html=True)
# ── Cache resources ───────────────────────────────────────────────────────────
@st.cache_resource(show_spinner="Инициализация агента и MCP-подключения (~30с)...")
def get_agent():
from src.agent.agent import homework_direct_agent
from src.agent.graph.pipeline import _invoke_with_retry
return homework_direct_agent, _invoke_with_retry
@st.cache_resource(show_spinner="Загрузка pipeline...")
def get_pipeline():
from src.agent.graph.pipeline import pipeline
return pipeline
# ── Callback (realtime events → queue) ───────────────────────────────────────
class UICallback(BaseCallbackHandler):
def __init__(self, q: q_mod.Queue):
self.q = q
def _parse_inputs(self, input_str: str, kwargs: dict) -> dict:
"""Пробуем достать структурированные inputs из kwargs или распарсить строку."""
inputs = kwargs.get("inputs") or {}
if inputs:
return inputs
try:
return json.loads(input_str)
except Exception:
pass
try:
return ast.literal_eval(input_str)
except Exception:
pass
return {}
def on_tool_start(self, serialized, input_str, **kwargs):
name = serialized.get("name", "?")
inputs = self._parse_inputs(str(input_str), kwargs)
self.q.put({"t": "tool_start", "name": name, "inputs": inputs})
def on_tool_end(self, output, **kwargs):
self.q.put({"t": "tool_end", "output": str(output)[:500]})
def on_tool_error(self, error, **kwargs):
self.q.put({"t": "tool_error", "msg": str(error)[:300]})
def on_llm_start(self, *a, **kw):
self.q.put({"t": "thinking"})
# ── Steps ─────────────────────────────────────────────────────────────────────
STEPS = [
("task_get", "Читаю состояние задания"),
("task_text", "Читаю текст задания"),
("gitea_create_repo", "Создаю репозиторий"),
("gitea_write_file", "Загружаю файлы"),
("task_update_answer", "Устанавливаю ответ"),
("task_submit", "Сдаю задание"),
]
_STEP_CLS = {"done": "step-done", "active": "step-active", "pending": "step-pending"}
_STEP_ICON = {"done": "✓", "active": "⟳", "pending": "○"}
def render_steps(states: dict, file_count: int = 0) -> str:
parts = []
for key, label in STEPS:
sv = states.get(key, "pending")
cls = _STEP_CLS[sv]
ico = _STEP_ICON[sv]
extra = f" ({file_count} файлов)" if key == "gitea_write_file" and file_count > 0 else ""
parts.append(f'{ico} {label}{extra}
')
return "".join(parts)
def step_key_for(tool_name: str) -> str | None:
for key, _ in STEPS:
if key in tool_name:
return key
return None
# ── Background runner ─────────────────────────────────────────────────────────
def run_agent_background(agent, messages, config, q: q_mod.Queue, cb):
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())
# ══════════════════════════════════════════════════════════════════════════════
# Tabs
# ══════════════════════════════════════════════════════════════════════════════
tab1, tab2 = st.tabs(["Одно задание", "Все задания (Pipeline)"])
# ── Tab 1 — одно задание ──────────────────────────────────────────────────────
with tab1:
task_id_raw = st.text_input(
"Task ID",
placeholder="6a1867fa8a94f887e50d52bd",
help="Скопируй из URL на platform.brojs.ru",
label_visibility="visible",
)
go = st.button("Выполнить задание", type="primary", use_container_width=True)
if go:
task_id = task_id_raw.strip()
if not task_id:
st.warning("Введи Task ID")
else:
repo = f"task-{task_id}"
url = f"https://git.brojs.ru/{OWNER}/{repo}"
prompt = f"""Выполни задание из курса BroJS.
taskId: {task_id}
owner: {OWNER}
ШАГ 1. mcp__journal-bh-professor__task_get(taskId='{task_id}')
→ Смотри поле answer.content:
• Пусто → ПЕРВАЯ СДАЧА (шаг 2а)
• Ссылка на репо → ПЕРЕСДАЧА (шаг 2б)
→ Читай комментарии преподавателя — они в приоритете
ШАГ 2а (ПЕРВАЯ СДАЧА):
a. mcp__journal-bh-professor__task_text(taskId='{task_id}') — прочитай полный текст задания
b. gitea_create_repo(name='{repo}') — создай репозиторий
c. gitea_write_file — загрузи ВСЕ файлы (main.py, requirements.txt, README.md + модули). Полный код, без заглушек.
ШАГ 2б (ПЕРЕСДАЧА):
a. mcp__journal-bh-professor__task_text(taskId='{task_id}') — прочитай оригинальное задание
b. Сравни комментарий с требованиями:
• Противоречит заданию → task_comment с объяснением + цитата, НЕ меняй код
• Реальная ошибка → исправь через gitea_write_file
ШАГ 3. mcp__journal-bh-professor__task_update_answer(taskId='{task_id}', answerType='link', content='{url}')
ШАГ 4. mcp__journal-bh-professor__task_submit(taskId='{task_id}', confirmSubmit=True)"""
config = {"configurable": {"thread_id": f"ui-{task_id}-{int(time.time())}"}}
# ── UI placeholders ────────────────────────────────────────────
col_left, col_right = st.columns([1, 2])
with col_left:
st.markdown("**Pipeline**")
steps_ph = st.empty()
with col_right:
st.markdown("**Лог событий**")
log_ph = st.empty()
st.markdown("**Код (последний записанный файл)**")
code_header_ph = st.empty()
code_ph = st.empty()
status_ph = st.empty()
# ── Init ───────────────────────────────────────────────────────
step_states = {k: "pending" for k, _ in STEPS}
logs: list[str] = []
files: dict[str, str] = {}
file_count = 0
active_key = None
thinking_shown = False
steps_ph.markdown(render_steps(step_states), unsafe_allow_html=True)
# ── Start thread ───────────────────────────────────────────────
update_q: q_mod.Queue = q_mod.Queue()
agent, _ = get_agent()
cb = UICallback(update_q)
thread = threading.Thread(
target=run_agent_background,
args=(agent, {"messages": [HumanMessage(content=prompt)]}, config, update_q, cb),
daemon=True,
)
thread.start()
# ── Poll loop ──────────────────────────────────────────────────
final_result = None
fatal_error = None
while thread.is_alive() or not update_q.empty():
dirty = False
while not update_q.empty():
ev = update_q.get_nowait()
ts = datetime.now().strftime("%H:%M:%S")
if ev["t"] == "thinking":
if not thinking_shown:
logs.append(
f''
f'{ts} '
f'модель думает...'
f'
'
)
thinking_shown = True
dirty = True
elif ev["t"] == "tool_start":
thinking_shown = False
name = ev["name"]
inputs = ev.get("inputs", {})
key = step_key_for(name)
if key:
if active_key and active_key != key:
step_states[active_key] = "done"
step_states[key] = "active"
active_key = key
if key == "gitea_write_file":
file_count += 1
path = inputs.get("path", "")
content = inputs.get("content", "")
if path and content:
files[path] = content
short = name.replace("mcp__journal-bh-professor__", "mcp::")
path = inputs.get("path", "")
finfo = f" {path}" if path else ""
logs.append(
f''
f'{ts} '
f'[tool] {short}{finfo}'
f'
'
)
dirty = True
elif ev["t"] == "tool_end":
out = ev["output"][:140].replace("<", "<").replace(">", ">")
logs.append(
f''
f'{ts} '
f'-> {out}'
f'
'
)
dirty = True
elif ev["t"] == "tool_error":
if active_key:
step_states[active_key] = "pending"
msg = ev["msg"][:140].replace("<", "<").replace(">", ">")
logs.append(
f''
f'{ts} '
f'[error] {msg}'
f'
'
)
dirty = True
elif ev["t"] == "done":
final_result = ev["result"]
if active_key:
step_states[active_key] = "done"
dirty = True
elif ev["t"] == "fatal":
fatal_error = ev["msg"]
dirty = True
if dirty:
# Steps
steps_ph.markdown(
render_steps(step_states, file_count),
unsafe_allow_html=True,
)
# Log
log_ph.markdown(
'' + "".join(logs[-60:]) + '
',
unsafe_allow_html=True,
)
# Code panel — показываем последний файл
if files:
last_path = list(files)[-1]
lang = "python" if last_path.endswith(".py") else (
"text" if last_path.endswith(".txt") else "markdown"
)
code_header_ph.markdown(
f'',
unsafe_allow_html=True,
)
code_ph.code(files[last_path], language=lang)
time.sleep(0.15)
# ── Final ──────────────────────────────────────────────────────
if fatal_error:
status_ph.markdown(
f'Ошибка: {fatal_error[:200]}
',
unsafe_allow_html=True,
)
elif final_result is not None:
msgs = final_result.get("messages", [])
submitted = any(
"task_submit" in tc.get("name", "")
for m in msgs
for tc in getattr(m, "tool_calls", [])
)
# финально помечаем все шаги
for k, _ in STEPS:
if step_states[k] in ("active", "done"):
step_states[k] = "done"
steps_ph.markdown(render_steps(step_states, file_count), unsafe_allow_html=True)
if submitted:
status_ph.markdown(
f'',
unsafe_allow_html=True,
)
else:
status_ph.markdown(
'Агент не вызвал task_submit — проверь лог
',
unsafe_allow_html=True,
)
else:
status_ph.markdown(
'Агент завершился без результата
',
unsafe_allow_html=True,
)
# ── Tab 2 — pipeline ──────────────────────────────────────────────────────────
with tab2:
st.info(
"Загружает **все незакрытые задания** курса KFU-26-1 и выполняет по очереди. "
"Может работать долго — не закрывай вкладку."
)
if st.button("Запустить Pipeline", type="primary", use_container_width=True):
pl = get_pipeline()
result_ph = st.empty()
with st.spinner("Pipeline работает..."):
try:
res = asyncio.run(pl.ainvoke({
"tasks": [], "current_index": 0, "results": [], "errors": [],
}))
results = res.get("results", [])
errors = res.get("errors", [])
lines = [f"### Выполнено заданий: {len(results)}\n"]
for r in results:
tid = r.get("task_id", "")
repo_link = f"https://git.brojs.ru/{OWNER}/task-{tid}"
lines.append(
f"- `{tid[:8]}...` [{r.get('status','?')}] "
f"mode=**{r.get('mode','?')}** "
f"retries={r.get('retries', 0)} "
f"[→ репо]({repo_link})"
)
if errors:
lines.append(f"\n### Ошибки ({len(errors)})")
for e in errors:
lines.append(f"- {e}")
result_ph.markdown("\n".join(lines))
except Exception as e:
result_ph.error(str(e))