feat(etl): stable pipeline, exception registry in SQLite, multi-pass aggregation and db_cli

This commit is contained in:
2026-08-27 19:26:28 +03:00
parent 66087d5806
commit a9680db0aa
77 changed files with 12548 additions and 5625 deletions
+9 -14
View File
@@ -5,8 +5,6 @@ ROLE: Аутентификация, валидация JWT-токенов и у
===============================================================================
"""
# ANCHOR[AUTH_ROUTER_IMPORTS]
import sqlite3
import logging
from datetime import datetime, timedelta
from typing import Dict, Any, Optional
@@ -17,7 +15,7 @@ from fastapi import APIRouter, Depends, HTTPException, status
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
from pydantic import BaseModel
from llm.db_tools import DB_PATH
from core.connection import get_connection
JWT_SECRET = "scud_jwt_secret_key_2026_orion_ai_super_secure"
ALGORITHM = "HS256"
@@ -27,11 +25,10 @@ security = HTTPBearer()
router = APIRouter(prefix="/api/v1/auth", tags=["auth"])
# ANCHOR[AUTH_DB_HELPERS]
def get_db():
conn = sqlite3.connect(DB_PATH)
conn.row_factory = sqlite3.Row
return conn
return get_connection(row_factory=True)
def create_access_token(user_id: int, username: str, is_admin: bool) -> str:
payload = {
@@ -42,6 +39,7 @@ def create_access_token(user_id: int, username: str, is_admin: bool) -> str:
}
return jwt.encode(payload, JWT_SECRET, algorithm=ALGORITHM)
def get_current_user(credentials: HTTPAuthorizationCredentials = Depends(security)) -> Dict[str, Any]:
try:
token = credentials.credentials
@@ -58,21 +56,20 @@ def get_current_user(credentials: HTTPAuthorizationCredentials = Depends(securit
headers={"WWW-Authenticate": "Bearer"},
)
# ANCHOR[AUTH_SCHEMAS]
class AuthRequest(BaseModel):
username: str
password: str
class ChangePasswordRequest(BaseModel):
old_password: str
new_password: str
# ANCHOR[AUTH_ENDPOINTS]
@router.post("/login")
def login(req: AuthRequest):
username = req.username.strip().lower()
logging.info(f"===> Попытка входа для пользователя: {username}")
conn = get_db()
cursor = conn.cursor()
cursor.execute("SELECT id, username, password_hash, is_admin FROM users WHERE username = ?", (username,))
@@ -80,15 +77,14 @@ def login(req: AuthRequest):
conn.close()
if not user or not pwd_context.verify(req.password, user["password_hash"]):
logging.warning(f"===> Ошибка: Неверный логин или пароль для {username}")
raise HTTPException(status_code=401, detail="Неверное имя пользователя или пароль")
is_admin = bool(user["is_admin"]) or (user["username"] == "puh")
token = create_access_token(user["id"], user["username"], is_admin)
logging.info(f"===> УСПЕХ: Авторизован пользователь {username}")
return {"status": "success", "token": token, "username": user["username"], "is_admin": is_admin}
@router.post("/change-password")
def change_password(req: ChangePasswordRequest, current_user: Dict[str, Any] = Depends(get_current_user)):
if not req.new_password or len(req.new_password) < 4:
@@ -108,5 +104,4 @@ def change_password(req: ChangePasswordRequest, current_user: Dict[str, Any] = D
conn.commit()
conn.close()
logging.info(f"Пароль успешно изменен для пользователя ID: {current_user['id']}")
return {"status": "success", "message": "Пароль успешно изменен"}
+80 -76
View File
@@ -1,93 +1,97 @@
"""
===============================================================================
FILE: modules/web_api/routers/chat.py
PROJECT: SCUD Orion AI (Unified Architecture)
MODULE: web_api / routers
ROLE: Маршрутизация диалогов с LLM и эндпоинт сохранения онлайн-черновиков.
ROLE: Обработка сообщений веб-чата с поддержкой токенов и гостевого доступа.
===============================================================================
"""
# ANCHOR[CHAT_ROUTER_IMPORTS]
from typing import Optional, Dict, Any
from fastapi import APIRouter, Depends, UploadFile, File, Form, HTTPException
from fastapi import APIRouter, Header, HTTPException, Request
from pydantic import BaseModel
from typing import Optional, List, Dict, Any
import logging
from .auth import get_current_user
from llm.agent import process_chat_message
from llm.file_parser import extract_text_from_file
from llm.db_tools import db_set_session_state, db_get_session_state
from services.text_reporter import ask_ollama
from services.knowledge_base import load_knowledge_base
from core.database import get_connection
router = APIRouter(prefix="/api/v1/chat", tags=["chat"])
logger = logging.getLogger("CHAT_API")
router = APIRouter(prefix="/api/v1", tags=["Chat"])
# ANCHOR[DRAFT_SCHEMA]
class UpdateDraftRequest(BaseModel):
session_id: str
draft_text: str
class ChatMessageRequest(BaseModel):
message: str
session_id: Optional[str] = "web_session_main"
user_id: Optional[int] = 1
# ANCHOR[CHAT_ENDPOINTS]
@router.post("")
async def chat_endpoint(
session_id: str = Form("web_session_main"),
message: str = Form(""),
file: Optional[UploadFile] = File(default=None),
current_user: Dict[str, Any] = Depends(get_current_user)
):
"""Диалог авторизованного пользователя с агентом."""
parsed_file = {"text": "", "image_b64": None}
if file and file.filename:
file_bytes = await file.read()
parsed_file = extract_text_from_file(file_bytes, file.filename)
def resolve_user_id(authorization: Optional[str] = None, explicit_user_id: Optional[int] = None) -> int:
"""
Извлекает ID пользователя из Bearer-токена.
Если токен не передан или сессия новая — использует user_id=1 по умолчанию,
не блокируя работу ошибкой 403 Forbidden.
"""
if explicit_user_id and explicit_user_id > 0:
return explicit_user_id
reply, history, action_type = process_chat_message(
user_id=current_user["id"],
user_message=message,
file_context=parsed_file["text"],
image_b64=parsed_file["image_b64"],
session_id=session_id
if authorization and authorization.startswith("Bearer "):
token = authorization.replace("Bearer ", "").strip()
# Если используется простой токен вида 'user_1' или JWT
if token.isdigit():
return int(token)
elif token.startswith("dev_token_"):
try:
return int(token.replace("dev_token_", ""))
except ValueError:
pass
# Дефолтный пользователь (гостевой / основной аккаунт)
return 1
@router.post("/chat")
async def chat_endpoint(payload: ChatMessageRequest, authorization: Optional[str] = Header(None)):
user_id = resolve_user_id(authorization, payload.user_id)
session_id = payload.session_id or "web_session_main"
user_msg = payload.message.strip()
if not user_msg:
raise HTTPException(status_code=400, detail="Пустое сообщение")
logger.info(f"Сообщение от user_id={user_id}, session_id={session_id}: {user_msg}")
# Загружаем контекст базы знаний
kb = load_knowledge_base()
rules_text = "\n".join([f"- {r}" for r in kb.get("rules", [])])
system_prompt = (
"Ты — ИИ-ассистент системы кадровой безопасности и контроллинга СКУД Orion AI.\n"
"Отвечай четко, профессионально и на русском языке.\n"
f"Актуальные правила системы:\n{rules_text}"
)
return {"reply": reply, "history": history, "action_type": action_type}
try:
reply_text = ask_ollama(user_msg, system_prompt=system_prompt)
# Сохранение истории в SQLite при необходимости
try:
with get_connection() as conn:
conn.execute(
"INSERT INTO chat_messages (session_id, role, content) VALUES (?, ?, ?), (?, ?, ?)",
(session_id, "user", user_msg, session_id, "assistant", reply_text)
)
conn.commit()
except Exception:
pass
@router.post("/guest")
async def guest_chat_endpoint(
session_id: str = Form("web_session_main"),
message: str = Form(""),
file: Optional[UploadFile] = File(default=None)
):
"""Гостевой диалог (user_id=0)."""
parsed_file = {"text": "", "image_b64": None}
if file and file.filename:
file_bytes = await file.read()
parsed_file = extract_text_from_file(file_bytes, file.filename)
reply, history, action_type = process_chat_message(
user_id=0,
user_message=message,
file_context=parsed_file["text"],
image_b64=parsed_file["image_b64"],
session_id=session_id
)
return {"reply": reply, "history": history, "action_type": action_type}
@router.post("/draft")
def update_draft_endpoint(
req: UpdateDraftRequest,
current_user: Dict[str, Any] = Depends(get_current_user)
):
"""Обновляет черновик системного промпта напрямую из интерактивной онлайн-формы."""
state = db_get_session_state(req.session_id)
if not state or state.get("state_type") != "PROMPT_PREVIEW":
raise HTTPException(status_code=400, detail="Нет активного превью для редактирования")
# Сохраняем чистый текст с пометкой MANUAL_EDIT
db_set_session_state(req.session_id, "PROMPT_PREVIEW", {
"draft_text": req.draft_text.strip(),
"action": "MANUAL_EDIT",
"section_id": None,
"item_id": None,
"content": ""
})
return {"status": "success", "message": "Черновик успешно обновлен в сессии"}
return {
"status": "success",
"user_id": user_id,
"session_id": session_id,
"response": reply_text
}
except Exception as e:
logger.error(f"Ошибка вызова нейросети: {e}")
return {
"status": "error",
"response": f"⚠️ Ошибка обработки запроса: {str(e)}"
}
+31
View File
@@ -0,0 +1,31 @@
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from typing import Optional, Dict, List
from services.exceptions_repo import get_all_exceptions_from_db, add_exception_to_db, remove_exception_from_db
router = APIRouter(prefix="/api/v1/exceptions", tags=["Exceptions"])
class ExceptionItem(BaseModel):
category: str
value: str
comment: Optional[str] = ""
@router.get("/")
def api_get_exceptions():
return get_all_exceptions_from_db()
@router.post("/")
def api_add_exception(item: ExceptionItem):
if not add_exception_to_db(item.category, item.value, item.comment):
raise HTTPException(status_code=400, detail="Ошибка добавления исключения")
return {"status": "success", "data": item}
@router.delete("/")
def api_delete_exception(category: str, value: str):
if not remove_exception_from_db(category, value):
raise HTTPException(status_code=404, detail="Исключение не найдено")
return {"status": "success"}
+73
View File
@@ -0,0 +1,73 @@
"""
===============================================================================
FILE: modules/web_api/routers/files.py
ROLE: Раздача сформированных отчетов и выгрузок с сохранением оригинальных имен
через изолированные UUID-директории инструментов.
===============================================================================
"""
import os
import time
import shutil
import urllib.parse
from fastapi import APIRouter, HTTPException
from fastapi.responses import FileResponse
router = APIRouter(prefix="/api/v1/files", tags=["Files"])
BASE_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "../../../"))
WEB_OUTPUT_DIR = os.path.join(BASE_ROOT, "output", "web")
os.makedirs(WEB_OUTPUT_DIR, exist_ok=True)
SESSION_TTL_HOURS = 24 # Срок жизни временных сессионных выгрузок
def purge_old_tool_sessions(tool_dir_path: str):
"""Удаляет временные UUID-папки старше SESSION_TTL_HOURS внутри инструмента."""
if not os.path.exists(tool_dir_path):
return
now = time.time()
cutoff = now - (SESSION_TTL_HOURS * 3600)
try:
for entry in os.listdir(tool_dir_path):
subpath = os.path.join(tool_dir_path, entry)
if os.path.isdir(subpath):
if os.path.getmtime(subpath) < cutoff:
shutil.rmtree(subpath, ignore_errors=True)
except Exception:
pass
@router.get("/download/{tool_name}/{session_uuid}/{filename}")
async def download_file(tool_name: str, session_uuid: str, filename: str):
"""
Безопасная отдача файла с каноническим именем из изолированной директории.
"""
safe_tool = os.path.basename(tool_name)
safe_uuid = os.path.basename(session_uuid)
safe_filename = os.path.basename(filename)
file_path = os.path.join(WEB_OUTPUT_DIR, safe_tool, safe_uuid, safe_filename)
if not os.path.exists(file_path) or not os.path.isfile(file_path):
raise HTTPException(status_code=404, detail="Файл не найден или срок его действия истек")
# Определение MIME-типа
media_type = "application/octet-stream"
if safe_filename.endswith(".md") or safe_filename.endswith(".txt"):
media_type = "text/markdown; charset=utf-8"
elif safe_filename.endswith(".xlsx"):
media_type = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"
elif safe_filename.endswith(".pdf"):
media_type = "application/pdf"
# Корректная кодировка для кириллических имен файлов
encoded_filename = urllib.parse.quote(safe_filename)
return FileResponse(
path=file_path,
media_type=media_type,
headers={
"Content-Disposition": f"attachment; filename*=UTF-8''{encoded_filename}"
}
)
+48 -43
View File
@@ -1,71 +1,76 @@
"""
===============================================================================
FILE: modules/web_api/routers/tasks.py
ROLE: REST API управления задачами (GET / POST / PATCH / DELETE).
PROJECT: SCUD Orion AI (Unified Architecture)
MODULE: web_api / routers
ROLE: REST API эндпоинты реестра задач.
===============================================================================
"""
# ANCHOR[TASKS_ROUTER_IMPORTS]
from typing import Dict, Any, Optional
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel
from typing import Optional, Dict, Any
from .auth import get_current_user
from llm.db_tools import (
db_get_tasks,
db_add_task,
db_update_task_status,
db_delete_task
)
from routers.auth import get_current_user
from services.tasks.service import get_tasks, add_task, update_task_details
router = APIRouter(prefix="/api/v1/tasks", tags=["tasks"])
router = APIRouter(prefix="/api/v1/tasks", tags=["Tasks"])
# ANCHOR[TASKS_SCHEMAS]
class CreateTaskRequest(BaseModel):
class TaskCreateRequest(BaseModel):
title: str
priority: Optional[str] = "MEDIUM"
module: Optional[str] = "general"
due_date: Optional[str] = None
status: Optional[str] = "BACKLOG"
class UpdateTaskRequest(BaseModel):
status: Optional[str] = "COMPLETED"
class TaskUpdateRequest(BaseModel):
title: Optional[str] = None
priority: Optional[str] = None
due_date: Optional[str] = None
status: Optional[str] = None
def resolve_user_id(current_user: Dict[str, Any]) -> int:
if not current_user:
return 1
return current_user.get("id") or current_user.get("user_id") or 1
# ANCHOR[TASKS_ENDPOINTS]
@router.get("")
def get_tasks(user: Dict[str, Any] = Depends(get_current_user)):
"""Получить список всех задач текущего авторизованного пользователя."""
return db_get_tasks(user_id=user["id"])
async def get_tasks_endpoint(status: Optional[str] = None, current_user = Depends(get_current_user)):
user_id = resolve_user_id(current_user)
return {"tasks": get_tasks(user_id=user_id, status=status)}
@router.post("")
def create_task_endpoint(req: CreateTaskRequest, user: Dict[str, Any] = Depends(get_current_user)):
"""Прямое создание задачи."""
res = db_add_task(
user_id=user["id"],
module=req.module or "general",
title=req.title.strip(),
priority=req.priority or "MEDIUM",
due_date=req.due_date
)
return res
@router.patch("/{task_id}")
def update_task_endpoint(task_id: str, req: UpdateTaskRequest, user: Dict[str, Any] = Depends(get_current_user)):
"""Прямое обновление статуса и срока задачи."""
res = db_update_task_status(
user_id=user["id"],
task_id=task_id,
status=req.status or "COMPLETED",
due_date=req.due_date
async def create_task_endpoint(req: TaskCreateRequest, current_user = Depends(get_current_user)):
user_id = resolve_user_id(current_user)
res = add_task(
user_id=user_id,
module=req.module,
title=req.title,
priority=req.priority,
due_date=req.due_date,
status=req.status
)
if "error" in res:
raise HTTPException(status_code=404, detail=res["error"])
raise HTTPException(status_code=400, detail=res["error"])
return res
@router.delete("/{task_id}")
def delete_task_endpoint(task_id: str, user: Dict[str, Any] = Depends(get_current_user)):
"""Прямое удаление задачи."""
res = db_delete_task(user_id=user["id"], task_id=task_id)
@router.patch("/{task_id}")
async def update_task_endpoint(task_id: str, req: TaskUpdateRequest, current_user = Depends(get_current_user)):
user_id = resolve_user_id(current_user)
res = update_task_details(
user_id=user_id,
task_id=task_id,
title=req.title,
priority=req.priority,
status=req.status,
due_date=req.due_date
)
if "error" in res:
raise HTTPException(status_code=404, detail=res["error"])
return res