revert: вернуть ui.py к версии ee2c8de (до изменений rate limit)

This commit is contained in:
2026-06-05 17:16:14 +03:00
parent b9b0f58de2
commit 2bfcf5f782
+167 -183
View File
@@ -1,7 +1,7 @@
""" """
Streamlit UI для brojs-agent. Streamlit UI для brojs-agent.
Запуск: python -m streamlit run ui.py Запуск: streamlit run ui.py
""" """
import asyncio import asyncio
import json import json
@@ -68,16 +68,12 @@ st.markdown("""
.thinking { .thinking {
color: #64748b; font-style: italic; font-size: .82em; padding: 4px 0; color: #64748b; font-style: italic; font-size: .82em; padding: 4px 0;
} }
.status-ok { background:#052e16; border:1px solid #22c55e; color:#4ade80; .status-ok { background:#052e16; border:1px solid #22c55e; color:#4ade80; padding:10px 16px; border-radius:8px; }
padding:10px 16px; border-radius:8px; } .status-warn { background:#1c1917; border:1px solid #f59e0b; color:#fbbf24; padding:10px 16px; border-radius:8px; }
.status-warn { background:#1c1917; border:1px solid #f59e0b; color:#fbbf24; .status-err { background:#1c0a0a; border:1px solid #ef4444; color:#f87171; padding:10px 16px; border-radius:8px; }
padding:10px 16px; border-radius:8px; }
.status-err { background:#1c0a0a; border:1px solid #ef4444; color:#f87171; .chat-user { background:#1e3a5f; border-radius:12px 12px 2px 12px; padding:10px 14px; margin:6px 0; }
padding:10px 16px; border-radius:8px; } .chat-agent { background:#1a1a2e; border-radius:12px 12px 12px 2px; padding:10px 14px; margin:6px 0; }
.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> </style>
""", unsafe_allow_html=True) """, unsafe_allow_html=True)
@@ -93,118 +89,76 @@ st.markdown("""
""", unsafe_allow_html=True) """, unsafe_allow_html=True)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Кэш агента — загружается один раз при первом обращении # Кэш агента
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@st.cache_resource(show_spinner="Инициализация агента (подключение к MCP)...") @st.cache_resource(show_spinner="Инициализация агента (~30с)...")
def get_agent(): def get_agent():
from src.agent.agent import homework_direct_agent from src.agent.agent import homework_direct_agent
return homework_direct_agent return homework_direct_agent
@st.cache_resource(show_spinner="Загрузка pipeline...") @st.cache_resource(show_spinner="Загрузка pipeline...")
def get_pipeline(): def get_pipeline():
from src.agent.graph.pipeline import pipeline from src.agent.graph.pipeline import pipeline
return pipeline return pipeline
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Сборщик событий агента (синхронный — не нужна очередь) # Callback — перехватывает события агента и шлёт в очередь
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
_SOLVE_TOOLS = {"validate_teacher_comment", "generate_code_solution"} _SOLVE_TOOLS = {"validate_teacher_comment", "generate_code_solution"}
_JOURNAL_PREFIX = "mcp__journal-bh-professor__" _JOURNAL_PREFIX = "mcp__journal-bh-professor__"
class AgentEventCollector(BaseCallbackHandler): class AgentCallback(BaseCallbackHandler):
"""Накапливает события агента в список во время синхронного вызова.""" def __init__(self, q: queue.Queue):
self.q = q
def __init__(self) -> None:
self.events: list[dict] = []
def _ts(self) -> str: def _ts(self) -> str:
return datetime.now().strftime("%H:%M:%S") return datetime.now().strftime("%H:%M:%S")
def _add(self, ev: dict) -> None:
self.events.append(ev)
def on_tool_start(self, serialized, input_str, **kwargs): def on_tool_start(self, serialized, input_str, **kwargs):
name = serialized.get("name", "?") name = serialized.get("name", "?")
try: try:
args = json.loads(str(input_str)) if isinstance(input_str, str) else input_str args = json.loads(str(input_str)) if isinstance(input_str, str) else input_str
except Exception: except Exception:
args = {} args = {}
self._add({"t": "tool_start", "ts": self._ts(), "name": name, "args": args}) self.q.put({"t": "tool_start", "ts": self._ts(), "name": name, "args": args})
def on_tool_end(self, output, **kwargs): def on_tool_end(self, output, **kwargs):
self._add({"t": "tool_end", "ts": self._ts(), "output": str(output)[:300]}) self.q.put({"t": "tool_end", "ts": self._ts(), "output": str(output)[:300]})
def on_tool_error(self, error, **kwargs): def on_tool_error(self, error, **kwargs):
self._add({"t": "tool_error", "ts": self._ts(), "msg": str(error)[:200]}) self.q.put({"t": "tool_error", "ts": self._ts(), "msg": str(error)[:200]})
def on_llm_start(self, *a, **kw): def on_llm_start(self, *a, **kw):
self._add({"t": "thinking", "ts": self._ts()}) self.q.put({"t": "thinking", "ts": self._ts()})
def on_llm_end(self, response, **kwargs): def on_llm_end(self, response, **kwargs):
try: try:
text = response.generations[0][0].text[:120] text = response.generations[0][0].text[:120]
self._add({"t": "llm_end", "ts": self._ts(), "preview": text}) self.q.put({"t": "llm_end", "ts": self._ts(), "preview": text})
except Exception: except Exception:
pass pass
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Запуск агента — блокирующий вызов в отдельном потоке # Запуск агента в фоне
# (избегаем конфликтов с event loop Streamlit)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
def _invoke_agent( def _run_agent_thread(agent, messages, config, q: queue.Queue, cb: AgentCallback):
agent, messages: dict, config: dict async def _inner():
) -> tuple[dict | None, list[dict], str | None]:
"""
Запускает агента синхронно. Блокирует поток до завершения.
Возвращает (result, events, error_message).
"""
from src.agent.middlewares.retry_on_rate_limit import set_ui_event_queue
collector = AgentEventCollector()
# Очередь для rate-limit событий из middleware
rl_queue: queue.Queue = queue.Queue()
set_ui_event_queue(rl_queue)
result: dict | None = None
error: str | None = None
def _thread_fn() -> None:
nonlocal result, error
async def _inner() -> None:
nonlocal result, error
try: try:
r = await agent.ainvoke(messages, {**config, "callbacks": [collector]}) result = await agent.ainvoke(messages, {**config, "callbacks": [cb]})
result = r q.put({"t": "done", "result": result})
except Exception as e: except Exception as e:
error = str(e) q.put({"t": "fatal", "msg": str(e)})
asyncio.run(_inner()) asyncio.run(_inner())
t = threading.Thread(target=_thread_fn, daemon=True)
t.start()
t.join() # ждём завершения (UI показывает spinner)
set_ui_event_queue(None)
# Переносим rate-limit события в список collector'а
while not rl_queue.empty():
try:
collector.events.append(rl_queue.get_nowait())
except Exception:
break
return result, collector.events, error
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Рендер одного события в HTML # Рендер одного события в лог
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
def _render_event(ev: dict) -> str: def _render_event(ev: dict) -> str:
@@ -218,16 +172,25 @@ def _render_event(ev: dict) -> str:
name = ev["name"] name = ev["name"]
args = ev.get("args", {}) args = ev.get("args", {})
short = name.replace(_JOURNAL_PREFIX, "mcp::") short = name.replace(_JOURNAL_PREFIX, "mcp::")
# Определяем тип инструмента
if name in _SOLVE_TOOLS: if name in _SOLVE_TOOLS:
cls, icon, label = "tool-subagent", "🧠", f"[субагент] {short}" cls = "tool-subagent"
icon = "🧠"
label = f"[субагент] {short}"
elif "gitea" in name: elif "gitea" in name:
cls, icon, label = "tool-call", "📦", short cls = "tool-call"
icon = "📦"
label = short
elif "mcp::" in short or "journal" in name: elif "mcp::" in short or "journal" in name:
cls, icon, label = "tool-call", "📡", short cls = "tool-call"
icon = "📡"
label = short
else: else:
cls, icon, label = "tool-call", "🔧", short cls = "tool-call"
icon = "🔧"
label = short
# Показываем ключевые аргументы
hint = "" hint = ""
for key in ("taskId", "path", "repo", "repo_name", "name"): for key in ("taskId", "path", "repo", "repo_name", "name"):
if key in args: if key in args:
@@ -244,21 +207,6 @@ def _render_event(ev: dict) -> str:
msg = ev["msg"].replace("<", "&lt;") msg = ev["msg"].replace("<", "&lt;")
return f'<div class="tool-result" style="border-color:#ef4444;color:#f87171">⚠ {msg}</div>' return f'<div class="tool-result" style="border-color:#ef4444;color:#f87171">⚠ {msg}</div>'
if kind == "rate_limit_wait":
name, pause, attempt, mx = (
ev.get("name", "?"), ev.get("pause", 30),
ev.get("attempt", 1), ev.get("max", 5),
)
return (f'<div class="thinking" style="color:#f59e0b">'
f'{ts} 429 rate limit — {name} · ждал {pause}с '
f'(попытка {attempt}/{mx})</div>')
if kind == "rate_limit_retry":
name = ev.get("name", "?")
attempt = ev.get("attempt", 2)
return (f'<div class="thinking" style="color:#86efac">'
f'🔄 {ts} повтор {name} (попытка {attempt})...</div>')
if kind == "llm_end": if kind == "llm_end":
preview = ev.get("preview", "").replace("<", "&lt;")[:100] preview = ev.get("preview", "").replace("<", "&lt;")[:100]
return f'<div class="thinking">✏ {ts} {preview}...</div>' return f'<div class="thinking">✏ {ts} {preview}...</div>'
@@ -266,22 +214,11 @@ def _render_event(ev: dict) -> str:
return "" return ""
def _show_events(events: list[dict], expanded: bool = False) -> None:
if not events:
return
with st.expander(f"🔍 Лог инструментов ({len(events)} событий)", expanded=expanded):
html = "".join(_render_event(e) for e in events[-80:])
st.markdown(f'<div style="max-height:320px;overflow-y:auto">{html}</div>',
unsafe_allow_html=True)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# Вкладки # Вкладки
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
tab_chat, tab_pipeline, tab_status = st.tabs( tab_chat, tab_pipeline, tab_status = st.tabs(["💬 Чат с агентом", "⚡ Pipeline", "📊 Статус заданий"])
["💬 Чат с агентом", "⚡ Pipeline", "📊 Статус заданий"]
)
# ══════════════════════════════════════════════════════════════════════════ # ══════════════════════════════════════════════════════════════════════════
@@ -289,8 +226,9 @@ tab_chat, tab_pipeline, tab_status = st.tabs(
# ══════════════════════════════════════════════════════════════════════════ # ══════════════════════════════════════════════════════════════════════════
with tab_chat: with tab_chat:
st.caption("Пиши агенту напрямую. Пока он работает — показывается spinner.") st.caption("Общайся с агентом: задай вопрос, попроси решить задание или разобрать ситуацию.")
# История сообщений
if "chat_history" not in st.session_state: if "chat_history" not in st.session_state:
st.session_state.chat_history = [] st.session_state.chat_history = []
if "chat_events" not in st.session_state: if "chat_events" not in st.session_state:
@@ -298,25 +236,28 @@ with tab_chat:
if "chat_thread_id" not in st.session_state: if "chat_thread_id" not in st.session_state:
st.session_state.chat_thread_id = f"ui-{int(time.time())}" st.session_state.chat_thread_id = f"ui-{int(time.time())}"
# История сообщений # Показываем историю
for msg in st.session_state.chat_history: for msg in st.session_state.chat_history:
role, text = msg["role"], msg["text"] role = msg["role"]
text = msg["text"]
if role == "user": if role == "user":
st.markdown(f'<div class="chat-user">👤 {text}</div>', st.markdown(f'<div class="chat-user">👤 {text}</div>', unsafe_allow_html=True)
unsafe_allow_html=True)
else: else:
st.markdown(f'<div class="chat-agent">🤖 {text}</div>', 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) unsafe_allow_html=True)
# Лог предыдущего запроса # Ввод
_show_events(st.session_state.chat_events)
# Поле ввода
col_input, col_btn = st.columns([5, 1]) col_input, col_btn = st.columns([5, 1])
with col_input: with col_input:
user_input = st.text_input( user_input = st.text_input(
"Сообщение", "Сообщение",
placeholder='"Реши задание 6a1864f7..." или "Какие задания у меня есть?"', placeholder='Например: "Реши задание 6a1864f7..." или "Какие задания у меня есть?"',
label_visibility="collapsed", label_visibility="collapsed",
key="chat_input", key="chat_input",
) )
@@ -328,42 +269,71 @@ with tab_chat:
st.session_state.chat_history.append({"role": "user", "text": msg_text}) st.session_state.chat_history.append({"role": "user", "text": msg_text})
st.session_state.chat_events = [] st.session_state.chat_events = []
lc_messages = [ # Строим историю сообщений для агента
HumanMessage(content=m["text"]) if m["role"] == "user" lc_messages = []
else AIMessage(content=m["text"]) for m in st.session_state.chat_history:
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}} config = {"configurable": {"thread_id": st.session_state.chat_thread_id}}
try:
agent = get_agent() agent = get_agent()
except Exception as e:
st.error(f"Ошибка инициализации агента: {e}")
st.stop()
with st.spinner("🤖 Агент работает... (это может занять несколько минут)"): # Placeholders для обновления в реальном времени
result, events, error = _invoke_agent( events_ph = st.empty()
agent, {"messages": lc_messages}, config 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()
st.session_state.chat_events = events final_result = None
fatal = None
if error: while thread.is_alive() or not evq.empty():
st.session_state.chat_history.append( changed = False
{"role": "agent", "text": f"⚠️ Ошибка: {error}"} 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,
) )
elif result: time.sleep(0.15)
msgs = result.get("messages", [])
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 last = msgs[-1] if msgs else None
reply = last.content if last and hasattr(last, "content") else "Готово." reply = last.content if last and hasattr(last, "content") else "Готово."
st.session_state.chat_history.append({"role": "agent", "text": reply}) st.session_state.chat_history.append({"role": "agent", "text": reply})
else:
st.session_state.chat_history.append(
{"role": "agent", "text": "Агент завершил работу без ответа."}
)
st.rerun() st.rerun()
# Кнопка очистки
if st.session_state.chat_history: if st.session_state.chat_history:
if st.button("🗑 Очистить чат"): if st.button("🗑 Очистить чат"):
st.session_state.chat_history = [] st.session_state.chat_history = []
@@ -377,17 +347,12 @@ with tab_chat:
# ══════════════════════════════════════════════════════════════════════════ # ══════════════════════════════════════════════════════════════════════════
with tab_pipeline: with tab_pipeline:
st.caption("Запусти конкретное задание по ID или все todo-задания сразу.") st.caption("Автоматически решает все todo-задания курса по очереди.")
if "pipe_events" not in st.session_state:
st.session_state.pipe_events = []
if "pipe_result_text" not in st.session_state:
st.session_state.pipe_result_text = ""
col1, col2 = st.columns([3, 1]) col1, col2 = st.columns([3, 1])
with col1: with col1:
task_id_input = st.text_input( task_id_input = st.text_input(
"Task ID", "Task ID (оставь пустым — решить все todo)",
placeholder="6a1864f7fd30e81cf3146d65", placeholder="6a1864f7fd30e81cf3146d65",
label_visibility="visible", label_visibility="visible",
) )
@@ -395,48 +360,68 @@ with tab_pipeline:
st.write("") st.write("")
run_btn = st.button("▶ Запустить", type="primary", use_container_width=True) run_btn = st.button("▶ Запустить", type="primary", use_container_width=True)
# Показываем результат предыдущего запуска
if st.session_state.pipe_result_text:
st.markdown(st.session_state.pipe_result_text, unsafe_allow_html=True)
_show_events(st.session_state.pipe_events)
if run_btn: if run_btn:
if not task_id_input.strip(): result_ph = st.empty()
st.warning("Введи Task ID") events_ph2 = st.empty()
else: agent = get_agent()
if task_id_input.strip():
# Одно задание
task_id = task_id_input.strip() task_id = task_id_input.strip()
repo_url = f"https://git.brojs.ru/{GITEA_OWNER}/task-{task_id}" repo_url = f"https://git.brojs.ru/{GITEA_OWNER}/task-{task_id}"
prompt = f"Реши задание taskId={task_id} курса 698b49da77cb6d4d2e43ce78" prompt = f"Реши задание taskId={task_id} курса 698b49da77cb6d4d2e43ce78"
config = {"configurable": {"thread_id": f"pipe-{task_id}-{int(time.time())}"}} config = {"configurable": {"thread_id": f"pipe-{task_id}-{int(time.time())}"}}
messages = {"messages": [HumanMessage(content=prompt)]} messages = {"messages": [HumanMessage(content=prompt)]}
else:
st.session_state.pipe_events = [] result_ph.info("Pipeline для всех todo-заданий — используй раздел ниже")
st.session_state.pipe_result_text = ""
try:
agent = get_agent()
except Exception as e:
st.error(f"Ошибка инициализации агента: {e}")
st.stop() st.stop()
with st.spinner(f"🤖 Решаю задание {task_id[:8]}..."): evq2: queue.Queue = queue.Queue()
result, events, error = _invoke_agent(agent, messages, config) cb2 = AgentCallback(evq2)
all_events2: list[dict] = []
st.session_state.pipe_events = events thread2 = threading.Thread(
target=_run_agent_thread,
if error: args=(agent, messages, config, evq2, cb2),
st.session_state.pipe_result_text = ( daemon=True,
f'<div class="status-err">❌ Ошибка: {error[:400]}</div>'
) )
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: else:
st.session_state.pipe_result_text = ( all_events2.append(ev)
f'<div class="status-ok">✅ Готово! '
f'<a href="{repo_url}" target="_blank" style="color:#4ade80">'
f'Открыть репозиторий ↗</a></div>'
)
st.rerun() 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.divider()
st.subheader("Запустить все todo-задания") st.subheader("Запустить все todo-задания")
@@ -479,19 +464,18 @@ with tab_status:
async def _fetch(): async def _fetch():
j = load_journal_toolsets() j = load_journal_toolsets()
# Инструменты имеют префикс mcp__journal-bh-professor__
tools = {t.name: t for t in j.tasks_submissions_tools} tools = {t.name: t for t in j.tasks_submissions_tools}
full_name = "mcp__journal-bh-professor__tasks_list" full_name = "mcp__journal-bh-professor__tasks_list"
t = tools.get(full_name) or next( t = tools.get(full_name)
(v for k, v in tools.items() if "tasks_list" in k), None
)
if not t: if not t:
st.warning(f"tasks_list не найден. Доступны: {list(tools.keys())}") # 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 [] return []
raw = await t.ainvoke({"courseId": "698b49da77cb6d4d2e43ce78"}) raw = await t.ainvoke({"courseId": "698b49da77cb6d4d2e43ce78"})
text = ( text = next((x["text"] for x in raw if x.get("type") == "text"), str(raw)) if isinstance(raw, list) else str(raw)
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) data = json.loads(text)
return data.get("tasks", data) if isinstance(data, dict) else data return data.get("tasks", data) if isinstance(data, dict) else data