Files
scud_ai/modules/web_api/llm/agent.py
T

401 lines
24 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
===============================================================================
FILE: modules/web_api/llm/agent.py
PROJECT: SCUD Orion AI (Unified Architecture)
MODULE: web_api / llm (Core Agent & Function Calling Dispatcher)
ROLE: Главный оркестратор взаимодействия с Ollama LLM (Qwen 2.5), разбор вызовов
инструментов (Function Calling), интеграция с декларативным реестром
действий SQLite (tool_action_registry), детерминированный Fast-Path
для подтверждений, динамические кнопки и автоочистка эфемерных сообщений.
===============================================================================
"""
# --- [SECTION 1: SYSTEM PATHS & IMPORTS] --- # ANCHOR[SYS_IMPORTS]
import os
import sys
import json
import urllib.request
import urllib.error
import logging
import re
from typing import List, Dict, Any, Tuple, Optional
# Импорт фасада базы данных
from .db_tools import (
db_get_active_system_prompt,
db_add_system_prompt,
db_get_tool_action,
db_get_tasks,
db_update_task_status,
db_delete_task,
db_add_task,
db_get_rules,
db_set_session_state,
db_get_session_state,
db_get_snapshots,
db_delete_snapshots,
db_clear_session_state,
db_get_current_server_time,
db_save_chat_message,
db_get_chat_history,
db_purge_ephemeral_messages,
db_get_stats,
db_get_anomalies,
db_get_session_states,
db_get_reference
)
from .schemas import TOOLS_SCHEMA
from .core.calendar_utils import get_dynamic_calendar_context, parse_relative_date_ru
from .core.tool_injector import clean_raw_tool_tags, clean_output, inject_tools_if_needed
# --- [SECTION 2: LOGGING CONFIGURATION] --- # ANCHOR[LOGGING_CONFIG]
logger = logging.getLogger("SCUD_AGENT")
logger.setLevel(logging.INFO)
logger.propagate = False
if not logger.handlers:
handler = logging.StreamHandler(sys.stdout)
formatter = logging.Formatter("%(asctime)s [%(levelname)s] [%(name)s] %(message)s")
handler.setFormatter(formatter)
logger.addHandler(handler)
OLLAMA_URL = "http://192.168.11.3:11434/api/chat"
TEXT_MODEL = "qwen2.5:14b"
VISION_MODEL = "qwen2.5vl:7b-q8_0"
# --- [SECTION 3: MAIN CHAT PROCESSING PIPELINE] --- # ANCHOR[CHAT_PROCESSOR]
def process_chat_message(
user_id: int,
user_message: str,
file_context: str = "",
image_b64: Optional[str] = None,
chat_history: List[Dict[str, Any]] = None,
session_id: str = "web_session_main"
) -> Tuple[str, List[Dict[str, Any]], Optional[Dict[str, Any]]]:
"""
Главный конвейер обработки входящего сообщения:
1. Fast-Path перехват подтверждений/отмен при активном session_state.
2. Перехват завершения работы ('нет, закончить настройку') с автоочисткой эфемерных сообщений.
3. Формирование системного контекста и вызов Ollama.
4. Выполнение вызванного Tool и опрос Data-Driven реестра действий.
5. Возврат кортежа: (reply_text, chat_history, action_metadata).
"""
logger.info(f"Получено сообщение от user_id={user_id}, session_id={session_id}: {user_message}")
# 3.1. Обогащение текста вложением (при наличии)
full_user_content = user_message
if file_context:
full_user_content = f"{user_message}\n\n[СОДЕРЖИМОЕ ПРИКРЕПЛЕННОГО ФАЙЛА]:\n{file_context}"
session_state = db_get_session_state(session_id)
user_msg_clean = user_message.lower().strip(" .!?:;")
# --- [FAST-PATH 1: ПЕРЕХВАТ ЗАВЕРШЕНИЯ НАСТРОЙКИ С ОЧИСТКОЙ ЭФЕМЕРНОЙ ПАМЯТИ] ---
if user_msg_clean in ["нет, спасибо", "нет, закончить настройку", "закончить настройку", "завершить", "нет"]:
db_clear_session_state(session_id)
deleted_count = db_purge_ephemeral_messages(session_id)
logger.info(f"🧹 Завершена работа с инструментом. Удалено эфемерных сообщений: {deleted_count}")
reply_text = "Хорошо. Настройка завершена, контекст диалога чист. Чем я могу помочь дальше?"
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=0)
return reply_text, db_get_chat_history(session_id), None
# --- [FAST-PATH 2: ПЕРЕХВАТ ПОДТВЕРЖДЕНИЯ / ОТМЕНЫ ПРЕВЬЮ ПРОМПТА] ---
if session_state and session_state.get("state_type") == "PROMPT_PREVIEW":
if user_msg_clean in ["подтверждаю", "подтвердить", "да", "сохранить", "применить", "ок", "хорошо"]:
db_add_system_prompt("main_agent", session_state.get("pending_data", ""))
db_clear_session_state(session_id)
action_cfg = db_get_tool_action("db_confirm_prompt_preview")
reply_text = action_cfg["success_template"] if action_cfg else "✅ Системный промпт успешно сохранен и применен в базе данных."
if action_cfg and action_cfg.get("follow_up_question"):
reply_text += f"\n\n{action_cfg['follow_up_question']}"
db_save_chat_message(session_id, "user", full_user_content, is_ephemeral=1)
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=1)
return reply_text, db_get_chat_history(session_id), {
"type": action_cfg.get("action_type") if action_cfg else "FOLLOW_UP_ACTION",
"buttons": action_cfg.get("buttons", []) if action_cfg else []
}
elif user_msg_clean in ["отмена", "отменить", "отклонить", "назад", "стоп"]:
db_clear_session_state(session_id)
action_cfg = db_get_tool_action("db_cancel_prompt_preview")
reply_text = action_cfg["success_template"] if action_cfg else "❌ Изменения системного промпта отменены."
if action_cfg and action_cfg.get("follow_up_question"):
reply_text += f"\n\n{action_cfg['follow_up_question']}"
db_save_chat_message(session_id, "user", full_user_content, is_ephemeral=1)
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=1)
return reply_text, db_get_chat_history(session_id), {
"type": action_cfg.get("action_type") if action_cfg else "FOLLOW_UP_ACTION",
"buttons": action_cfg.get("buttons", []) if action_cfg else []
}
# 3.2. Сохраняем входящее сообщение в историю диалога
is_user_ephemeral = 1 if session_state else 0
db_save_chat_message(session_id, "user", full_user_content, is_ephemeral=is_user_ephemeral)
# 3.3. Извлекаем полную актуальную историю для передачи в LLM
db_history = db_get_chat_history(session_id, limit=20)
dynamic_prompt_text = db_get_active_system_prompt()
calendar_context = get_dynamic_calendar_context()
user_info = f"Пользователь ID={user_id}" if user_id != 0 else "Гость"
system_prompt_content = (
f"Ты — интеллектуальный ассистент SCUD Orion AI. "
f"Ты управляешь задачами, аналитикой СКУД и системными настройками с помощью инструментов (tools).\n\n"
f"[ОКРУЖЕНИЕ]\n"
f"- Пользователь: {user_info}\n"
f"- {calendar_context}\n\n"
f"[ПРАВИЛА ИСПОЛЬЗОВАНИЯ ИНСТРУМЕНТОВ]\n"
f"1. Для любых изменений системного промпта (добавить, удалить, изменить пункт) ВСЕГДА вызывай функцию db_preview_prompt_merge(prompt_text=...).\n"
f"2. Для просмотра системного промпта ВСЕГДА вызывай db_get_system_prompt().\n"
f"3. Для просмотра задач ВСЕГДА вызывай db_get_tasks().\n"
f"4. Никогда не симулируй выполнение функций в виде обычного текста. Если требуется действие — сразу вызывай соответствующий инструмент.\n\n"
f"[ПРИМЕРЫ ВЫЗОВА ИНСТРУМЕНТОВ]:\n"
f"- Пользователь: 'добавь пункт 3.4. Работать от сюда и до заката.' -> Вызов: db_preview_prompt_merge(prompt_text='3.4. Работать от сюда и до заката.')\n"
f"- Пользователь: 'удали пункт 3.4' -> Вызов: db_preview_prompt_merge(prompt_text='3.4')\n"
f"- Пользователь: 'покажи системный промпт' -> Вызов: db_get_system_prompt()\n"
f"- Пользователь: 'покажи мои задачи' -> Вызов: db_get_tasks()\n\n"
f"[ТЕКУЩИЙ АКТИВНЫЙ СИСТЕМНЫЙ ПРОМПТ]:\n{dynamic_prompt_text}"
)
llm_options = {
"num_predict": 8192,
"num_ctx": 8192,
"temperature": 0.1,
"repeat_penalty": 1.1,
"presence_penalty": 0.5,
"top_p": 0.9
}
user_msg_object = {"role": "user", "content": full_user_content}
# --- [SECTION 4: ROUTING & OLLAMA PAYLOAD] --- # ANCHOR[PAYLOAD_BUILD]
if image_b64:
user_msg_object["images"] = [image_b64]
messages = [
{"role": "system", "content": "Ты — строгий модуль OCR. Перепиши весь текст с изображения буква в букву без отсебятины."},
user_msg_object
]
payload = {"model": VISION_MODEL, "messages": messages, "stream": False, "options": llm_options}
else:
clean_db_history = [dict(m) for m in db_history]
for m in clean_db_history:
m.pop("images", None)
messages = [{"role": "system", "content": system_prompt_content}] + clean_db_history + [user_msg_object]
payload = {"model": TEXT_MODEL, "messages": messages, "tools": TOOLS_SCHEMA, "stream": False, "options": llm_options}
# --- [SECTION 5: EXECUTION & TOOL ROUTING] --- # ANCHOR[TOOL_ROUTER]
try:
req = urllib.request.Request(
OLLAMA_URL,
data=json.dumps(payload).encode("utf-8"),
headers={"Content-Type": "application/json"}
)
with urllib.request.urlopen(req) as response:
res_data = json.loads(response.read().decode("utf-8"))
msg = res_data.get("message", {})
tool_calls = msg.get("tool_calls", [])
raw_text_content = msg.get("content", "")
#tool_calls = inject_tools_if_needed(user_message, raw_text_content, tool_calls)
if tool_calls:
logger.info(f"Ответ от Ollama получен. Tool calls: True (кол-во: {len(tool_calls)})")
messages.append(msg)
for tool in tool_calls:
fn_name = tool["function"]["name"]
fn_args = tool["function"].get("arguments", {})
logger.info(f"🚀 Вызов функции (Tool): {fn_name} с аргументами: {fn_args}")
tool_result_content = ""
action_cfg = db_get_tool_action(fn_name)
if fn_name == "db_confirm_prompt_preview":
if session_state and session_state.get("state_type") == "PROMPT_PREVIEW":
db_add_system_prompt("main_agent", session_state.get("pending_data", ""))
db_clear_session_state(session_id)
if action_cfg and action_cfg.get("bypass_llm"):
reply_text = action_cfg["success_template"]
if action_cfg.get("follow_up_question"):
reply_text += f"\n\n{action_cfg['follow_up_question']}"
is_eph = action_cfg.get("is_ephemeral", 1)
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=is_eph)
return reply_text, db_get_chat_history(session_id), {
"type": action_cfg.get("action_type"),
"buttons": action_cfg.get("buttons", [])
}
tool_result_content = json.dumps({"status": "success"}, ensure_ascii=False)
else:
err_reply = "⚠️ Нет активного превью для подтверждения."
db_save_chat_message(session_id, "assistant", err_reply, is_ephemeral=1)
return err_reply, db_get_chat_history(session_id), None
elif fn_name == "db_cancel_prompt_preview":
db_clear_session_state(session_id)
if action_cfg and action_cfg.get("bypass_llm"):
reply_text = action_cfg["success_template"]
if action_cfg.get("follow_up_question"):
reply_text += f"\n\n{action_cfg['follow_up_question']}"
is_eph = action_cfg.get("is_ephemeral", 1)
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=is_eph)
return reply_text, db_get_chat_history(session_id), {
"type": action_cfg.get("action_type"),
"buttons": action_cfg.get("buttons", [])
}
tool_result_content = json.dumps({"status": "cancelled"}, ensure_ascii=False)
elif fn_name == "db_get_snapshots":
snapshots_res = db_get_snapshots(session_id=session_id, date_str=fn_args.get("date_str"), original_user_message=user_message)
tool_result_content = json.dumps(snapshots_res, ensure_ascii=False)
elif fn_name == "db_get_current_server_time":
tool_result_content = json.dumps(db_get_current_server_time(), ensure_ascii=False)
elif fn_name == "db_get_tasks":
tool_result_content = json.dumps(db_get_tasks(user_id), ensure_ascii=False)
elif fn_name in ["db_get_system_prompt", "db_get_system_prompts"]:
tool_result_content = json.dumps({"system_prompt": db_get_active_system_prompt()}, ensure_ascii=False)
elif fn_name == "db_get_stats":
tool_result_content = json.dumps(db_get_stats(), ensure_ascii=False)
elif fn_name == "db_get_anomalies":
tool_result_content = json.dumps(db_get_anomalies(limit=fn_args.get("limit", 100), date_str=fn_args.get("date_str")), ensure_ascii=False)
elif fn_name == "db_get_session_states":
tool_result_content = json.dumps(db_get_session_states(), ensure_ascii=False)
elif fn_name == "db_delete_snapshots":
res = db_delete_snapshots(snapshot_id=fn_args.get("snapshot_id"), day_str=fn_args.get("day_str"))
if action_cfg and action_cfg.get("bypass_llm"):
reply_text = action_cfg["success_template"]
if action_cfg.get("follow_up_question"):
reply_text += f"\n\n{action_cfg['follow_up_question']}"
is_eph = action_cfg.get("is_ephemeral", 1)
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=is_eph)
return reply_text, db_get_chat_history(session_id), {
"type": action_cfg.get("action_type"),
"buttons": action_cfg.get("buttons", [])
}
tool_result_content = json.dumps(res, ensure_ascii=False)
elif fn_name == "db_get_reference":
tool_result_content = json.dumps(db_get_reference(category=fn_args.get("category")), ensure_ascii=False)
# --- [SECTION 6: PROMPT MERGE & PREVIEW ENGINE] --- # ANCHOR[PROMPT_MERGE_LOGIC]
elif fn_name == "db_preview_prompt_merge":
proposed_text = fn_args.get("prompt_text") or fn_args.get("proposed_prompt") or ""
if isinstance(fn_args, str):
proposed_text = fn_args
current_prompt = db_get_active_system_prompt()
user_msg_lower = user_message.lower()
# 1. ОБРАБОТКА УДАЛЕНИЯ ПУНКТА
if any(w in user_msg_lower for w in ["удали", "стереть", "убрать", "вырежи", "удалить"]):
target_num_match = re.search(r'\d+(\.\d+)*', user_message)
target_num = target_num_match.group(0) if target_num_match else ""
lines = current_prompt.splitlines()
if target_num:
new_lines = [line for line in lines if not line.strip().startswith(f"{target_num}.")]
else:
new_lines = lines
proposed_text = "\n".join(new_lines)
# 2. ОБРАБОТКА ДОБАВЛЕНИЯ / ИЗМЕНЕНИЯ ПУНКТА
elif proposed_text:
if len(proposed_text) < 500:
clean_item = proposed_text.strip()
for prefix in ["добавь пункт", "добавить пункт", "вставь пункт", "добавь"]:
if prefix in clean_item.lower():
clean_item = re.sub(prefix, "", clean_item, flags=re.IGNORECASE).strip(" .:")
lines = current_prompt.splitlines()
new_lines = []
inserted = False
for line in lines:
new_lines.append(line)
if "3.3." in line and not inserted:
item_str = clean_item if re.match(r'^\d+\.\d+\.', clean_item) else f"3.4. {clean_item}"
new_lines.append(f" {item_str}")
inserted = True
if not inserted:
new_lines.append(f" {clean_item}")
proposed_text = "\n".join(new_lines)
db_set_session_state(session_id, "PROMPT_PREVIEW", proposed_text)
preview_reply = (
f"Ваше изменение успешно предпросмотрено. Полный обновленный системный промпт теперь выглядит так:\n\n"
f"{proposed_text}\n\n"
f"Для применения изменений подтвердите действие («подтверждаю») или отмените («отмена»)."
)
db_save_chat_message(session_id, "assistant", preview_reply, is_ephemeral=1)
return clean_raw_tool_tags(preview_reply), db_get_chat_history(session_id), {
"type": "PROMPT_PREVIEW",
"buttons": [
{"label": "Подтвердить", "value": "подтверждаю", "style": "primary"},
{"label": "Отменить", "value": "отмена", "style": "danger"}
]
}
elif fn_name == "db_get_rules":
tool_result_content = json.dumps(db_get_rules(), ensure_ascii=False)
elif fn_name == "db_add_task":
res = db_add_task(user_id=user_id, module=fn_args.get("module", "general"), title=fn_args.get("title"), priority=fn_args.get("priority", "MEDIUM"), due_date=fn_args.get("due_date"))
tool_result_content = json.dumps(res, ensure_ascii=False)
elif fn_name == "db_update_task_status":
res = db_update_task_status(user_id=user_id, task_id=str(fn_args.get("task_id")), status=fn_args.get("status", "COMPLETED"), due_date=fn_args.get("due_date"))
tool_result_content = json.dumps(res, ensure_ascii=False)
elif fn_name == "db_delete_task":
res = db_delete_task(user_id=user_id, task_id=str(fn_args.get("task_id", "")).upper())
if action_cfg and action_cfg.get("bypass_llm"):
reply_text = action_cfg["success_template"]
if action_cfg.get("follow_up_question"):
reply_text += f"\n\n{action_cfg['follow_up_question']}"
is_eph = action_cfg.get("is_ephemeral", 1)
db_save_chat_message(session_id, "assistant", reply_text, is_ephemeral=is_eph)
return reply_text, db_get_chat_history(session_id), {
"type": action_cfg.get("action_type"),
"buttons": action_cfg.get("buttons", [])
}
tool_result_content = json.dumps(res, ensure_ascii=False)
messages.append({"role": "tool", "content": tool_result_content})
# --- [SECTION 7: SECONDARY LLM PASS] --- # ANCHOR[SECONDARY_PASS]
second_payload = {"model": TEXT_MODEL, "messages": messages, "stream": False, "options": llm_options}
sec_req = urllib.request.Request(OLLAMA_URL, data=json.dumps(second_payload).encode("utf-8"), headers={"Content-Type": "application/json"})
with urllib.request.urlopen(sec_req) as sec_response:
sec_res_data = json.loads(sec_response.read().decode("utf-8"))
raw_content = sec_res_data.get("message", {}).get("content", "").strip().replace("**", "")
final_content = clean_raw_tool_tags(clean_output(raw_content))
db_save_chat_message(session_id, "assistant", final_content, is_ephemeral=0)
return final_content, db_get_chat_history(session_id), None
# Если вызовов функций не было
raw_str = msg.get("content", "").strip().replace("**", "")
content_str = clean_raw_tool_tags(clean_output(raw_str))
final_reply = content_str or "Запрос обработан."
db_save_chat_message(session_id, "assistant", final_reply, is_ephemeral=0)
return final_reply, db_get_chat_history(session_id), None
except Exception as ex:
logger.exception(f"Непредвиденная ошибка: {ex}")
error_reply = f"Внутренняя ошибка сервера: {ex}"
return error_reply, db_get_chat_history(session_id), None