/
man4j
/
agent-server
Обзор
Документация
Войти
/
man4j
/
agent-server
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/agent_server/chat/runtime.py
253 строки
9 KB
Vladimir
fixes
21 май 2026, 11:09
21 май 2026, 11:09
ae5cb53
Код
Авторство
О чём код?
from chainlit.data.sql_alchemy import SQLAlchemyDataLayer from chainlit.types import ThreadDict import chainlit as cl from agent_server.chat.client import build_openai_client from agent_server.chat.policy import build_chat_behavior_policy from agent_server.chat.presenter import create_answer_message, finalize_answer_message from agent_server.chat.session import ( apply_runtime_history_state, clear_tool_session_state, get_tool_session_state, get_char_per_token, get_runtime_history_state, get_session_lock, RuntimeHistoryState, set_char_per_token, ) from agent_server.chat.types import ChatMessage from agent_server.chat.loop import run_chat_loop from agent_server.chat.runtime_state import ( apply_runtime_state, build_new_chat_runtime_state, build_resumed_chat_runtime_state, cache_runtime_state_for_thread, ) from agent_server.chat.profile_selection import ( ensure_profile_selected, filter_profiles_for_user, get_current_bot_name, get_profiles_cached, get_selected_agent, get_selected_profile, normalize_profiles, restore_profile_for_thread, ) from agent_server.chat.summary import prepare_messages_for_next_llm_round from agent_server.chat.tools import ( connect_selected_agent_tools, ensure_mcp_ready, send_loaded_context_step, send_mcp_errors, shutdown_selected_agent_tools, ) from agent_server.chat.threads import with_timestamp from agent_server.llm.agent import one_llm_round from agent_server.message_window import clear_reasoning_content from agent_server.message_window import estimated_prompt_tokens from agent_server.pg_storage import PostgresStorageClient from agent_server.config.agents import load_agent_definition from agent_server.config.env import DATABASE_URL, DEFAULT_CHAR_PER_TOKEN from agent_server.db.urls import build_asyncpg_database_url from agent_server.thread_store import persist_current_thread_metadata from agent_server.ui_meta import UsageTotals class ChatRuntime: def __init__(self, storage_client: PostgresStorageClient): self.storage_client = storage_client def get_data_layer(self): return SQLAlchemyDataLayer( conninfo=build_asyncpg_database_url(DATABASE_URL), storage_provider=self.storage_client, ) async def get_chat_profiles(self, current_user: cl.User) -> list[cl.ChatProfile]: profiles = normalize_profiles(get_profiles_cached()) profiles = filter_profiles_for_user(profiles, current_user) return [ cl.ChatProfile( name=p.id, display_name=p.display_name, markdown_description=p.markdown_description, icon=p.icon, ) for p in profiles ] async def on_chat_start(self) -> None: async with get_session_lock(): selected_profile = ensure_profile_selected() selected_agent = load_agent_definition(selected_profile.agent_id) self._ensure_runtime_defaults() mcp_state = await connect_selected_agent_tools() runtime_state = build_new_chat_runtime_state( selected_profile=selected_profile, selected_agent=selected_agent, mcp_state=mcp_state, ) apply_runtime_state(runtime_state) await persist_current_thread_metadata() if mcp_state.get("errors"): await send_mcp_errors(mcp_state) else: await self._send_context_loaded_probe( system_prompt=runtime_state.system_prompt, selected_profile=selected_profile, policy=build_chat_behavior_policy(selected_profile), ) async def on_chat_resume(self, thread: ThreadDict) -> None: async with get_session_lock(): selected_profile = restore_profile_for_thread(thread) selected_agent = load_agent_definition(selected_profile.agent_id) policy = build_chat_behavior_policy(selected_profile) self._ensure_runtime_defaults() await ensure_mcp_ready() runtime_state = build_resumed_chat_runtime_state( thread=thread, selected_profile=selected_profile, selected_agent=selected_agent, policy=policy, ) apply_runtime_state(runtime_state) cache_runtime_state_for_thread(thread, runtime_state) await persist_current_thread_metadata() async def on_chat_end(self) -> None: async with get_session_lock(): await shutdown_selected_agent_tools() clear_tool_session_state() async def on_message(self, input_msg: cl.Message) -> None: async with get_session_lock(): await ensure_mcp_ready() history_state = get_runtime_history_state(DEFAULT_CHAR_PER_TOKEN) messages: list[ChatMessage] = history_state.messages char_per_token: float = history_state.char_per_token or DEFAULT_CHAR_PER_TOKEN selected_profile = get_selected_profile() selected_agent = get_selected_agent() policy = build_chat_behavior_policy(selected_profile) client = build_openai_client(selected_profile) messages.append({"role": "user", "content": with_timestamp(input_msg.content)}) await self._prepare_history_for_next_user_turn( messages, client=client, profile=selected_profile, agent=selected_agent, policy=policy, char_per_token=char_per_token, ) loop_result = await run_chat_loop( client=client, selected_profile=selected_profile, selected_agent=selected_agent, policy=policy, messages=messages, char_per_token=char_per_token, response_headroom_tokens=policy.response_headroom_tokens, ) await self._prepare_history_for_next_user_turn( messages, client=client, profile=selected_profile, agent=selected_agent, policy=policy, char_per_token=loop_result.calibrated_char_per_token, ) apply_runtime_history_state( RuntimeHistoryState( messages=messages, char_per_token=loop_result.calibrated_char_per_token, ) ) await persist_current_thread_metadata() async def _send_context_loaded_probe( self, *, system_prompt: str, selected_profile, policy, ) -> None: messages: list[ChatMessage] = [ {"role": "system", "content": system_prompt}, { "role": "user", "content": "Ответь ровно одной фразой: Контекст успешно загружен", }, ] char_per_token = get_char_per_token(DEFAULT_CHAR_PER_TOKEN) usage_totals = UsageTotals() client = build_openai_client(selected_profile) async with cl.Step(name="Thinking", type="llm") as thinking_step: await send_loaded_context_step(system_prompt) answer_msg = await create_answer_message(get_current_bot_name()) round_result = await one_llm_round( client=client, profile=selected_profile, messages=messages, tools=[], thinking_step=thinking_step, answer_msg=answer_msg, ) usage_totals.add(round_result.usage) if answer_msg.content.strip() != "Контекст успешно загружен": answer_msg.content = "Контекст успешно загружен" messages.append( { "role": "assistant", "content": answer_msg.content, "reasoning_content": round_result.reasoning_content, } ) mcp_state = get_tool_session_state().mcp_state or {} await finalize_answer_message( answer_msg, usage_totals=usage_totals, messages=messages, char_per_token=char_per_token, connected_servers=list((mcp_state.get("servers") or {}).keys()), context_size=policy.context_size, response_headroom_tokens=policy.response_headroom_tokens, estimated_prompt_tokens_fn=estimated_prompt_tokens, ) def _ensure_runtime_defaults(self) -> None: if get_char_per_token() is None: set_char_per_token(DEFAULT_CHAR_PER_TOKEN) async def _prepare_history_for_next_user_turn( self, messages: list[ChatMessage], *, client, profile, agent, policy, char_per_token: float, ) -> None: if policy.clear_reasoning_before_user_turn: clear_reasoning_content(messages) await prepare_messages_for_next_llm_round( messages, client=client, profile=profile, agent=agent, policy=policy, char_per_token=char_per_token, )