Step 1: copy db/ modules + llm_prompt.py to app/

This commit is contained in:
2026-07-14 22:25:38 +04:00
parent 976e59b496
commit fc08db0c76
9 changed files with 943 additions and 0 deletions
View File
+91
View File
@@ -0,0 +1,91 @@
"""
Database connection — psycopg2 connection pool.
Архитектурное решение (Opus):
- ThreadedConnectionPool(minconn=1, maxconn=10) — оптимально для ThreadingMixIn HTTP-сервера.
Каждый HTTP-запрос в отдельном потоке получает своё соединение из пула.
- RealDictCursor для query() — возвращает dict с lowercase ключами (консистентно с Python API).
- execute()/execute_returning() — для INSERT/UPDATE/DELETE с автокоммитом и откатом при ошибке.
- Пул создаётся лениво (при первом запросе) через get_pool().
"""
import os
import psycopg2
import psycopg2.pool
import psycopg2.extras
# Глобальный пул соединений (singleton)
_pool = None
# Параметры подключения из переменных окружения (systemd Environment)
DB_CONFIG = {
"host": os.getenv("DB_HOST", "127.0.0.1"),
"port": int(os.getenv("DB_PORT", "5432")),
"dbname": os.getenv("DB_NAME", "baza"),
"user": os.getenv("DB_USER", "super"),
"password": os.getenv("DB_PASS", ""),
}
def get_pool():
"""Возвращает глобальный пул соединений. Создаёт при первом вызове."""
global _pool
if _pool is None:
_pool = psycopg2.pool.ThreadedConnectionPool(
minconn=1, maxconn=10, **DB_CONFIG
)
return _pool
def query(sql, params=None):
"""
SELECT → list[dict] с lowercase ключами.
Используется всеми db/*.py модулями для чтения данных.
"""
pool = get_pool()
conn = pool.getconn()
try:
with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
cur.execute(sql, params)
rows = cur.fetchall()
return [{k.lower(): v for k, v in r.items()} for r in rows]
finally:
pool.putconn(conn)
def execute(sql, params=None):
"""
INSERT/UPDATE/DELETE → количество затронутых строк.
Автокоммит. При ошибке — rollback и проброс исключения.
"""
pool = get_pool()
conn = pool.getconn()
try:
with conn.cursor() as cur:
cur.execute(sql, params)
conn.commit()
return cur.rowcount
except Exception:
conn.rollback()
raise
finally:
pool.putconn(conn)
def execute_returning(sql, params=None):
"""
INSERT/UPDATE/DELETE с RETURNING → dict первой строки.
Используется для insert с автогенерацией UUID (gen_random_uuid()).
"""
pool = get_pool()
conn = pool.getconn()
try:
with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
cur.execute(sql, params)
conn.commit()
row = cur.fetchone()
return {k.lower(): v for k, v in row.items()} if row else None
except Exception:
conn.rollback()
raise
finally:
pool.putconn(conn)
+25
View File
@@ -0,0 +1,25 @@
"""Contracts CRUD."""
from .connection import query, execute, execute_returning
def insert(number, client=""):
return execute_returning(
"INSERT INTO contracts (number, client) VALUES (%s, %s) RETURNING *",
(number, client),
)
def get(contract_id):
rows = query("SELECT * FROM contracts WHERE id = %s", (contract_id,))
return rows[0] if rows else None
def delete(contract_id):
return execute("DELETE FROM contracts WHERE id = %s", (contract_id,))
def delete_orphaned():
"""Remove contracts with no supplements."""
execute(
"DELETE FROM contracts WHERE id NOT IN (SELECT DISTINCT contract_id FROM supplements)"
)
+115
View File
@@ -0,0 +1,115 @@
"""Documents CRUD."""
from .connection import query, execute, execute_returning
def insert(filename, mime_type, original_bytes, status="uploaded", batch_id=None, zip_source=None, content_hash=None):
"""Insert document, return row dict."""
return execute_returning(
"""INSERT INTO documents (filename, mime_type, original_bytes, status, batch_id, zip_source, content_hash)
VALUES (%s, %s, %s, %s, %s, %s, %s) RETURNING *""",
(filename, mime_type, original_bytes, status, batch_id, zip_source, content_hash),
)
def get(doc_id):
rows = query("SELECT * FROM documents WHERE id = %s", (doc_id,))
return rows[0] if rows else None
def get_by_hash(batch_id, content_hash):
"""Find document by content hash within batch."""
rows = query(
"SELECT id FROM documents WHERE batch_id = %s AND content_hash = %s LIMIT 1",
(batch_id, content_hash),
)
return rows[0] if rows else None
def set_parsed(doc_id, elements_json):
"""Update elements_json + status='parsed'."""
import json
return execute(
"UPDATE documents SET elements_json = %s::jsonb, status = 'parsed' WHERE id = %s",
(json.dumps(elements_json, ensure_ascii=False), doc_id),
)
def set_error(doc_id, error_message):
return execute(
"UPDATE documents SET status = 'error', error_message = %s WHERE id = %s",
(error_message, doc_id),
)
def delete(doc_id):
return execute("DELETE FROM documents WHERE id = %s", (doc_id,))
def set_classification(doc_id, doc_type, own_number, parent_number, doc_date, counterparty, classify_raw=None, classify_input=None):
"""Store LLM classification results + raw response + input text."""
return execute(
"""UPDATE documents SET doc_type=%s, own_number=%s, parent_number=%s,
doc_date=%s, counterparty=%s, classify_status='classified', classify_raw=%s, classify_input=%s
WHERE id=%s""",
(doc_type, own_number, parent_number, doc_date, counterparty, classify_raw, classify_input, doc_id),
)
def set_classify_failed(doc_id, error):
return execute(
"UPDATE documents SET classify_status='failed', error_message=%s WHERE id=%s",
(error, doc_id),
)
def set_classify_garbage(doc_id, reason=""):
"""Mark document as garbage (Stage 1-2 filter, no LLM call)."""
return execute(
"UPDATE documents SET doc_type='garbage', classify_status='garbage', error_message=%s WHERE id=%s",
(f"garbage: {reason}", doc_id),
)
def list_pending(batch_id):
"""Documents waiting for classification."""
return query(
"SELECT id, filename, elements_json FROM documents WHERE batch_id=%s AND classify_status='pending'",
(batch_id,),
)
def reset_classify_status(batch_id):
"""Сбросить classify_status на 'pending' только для 'processing' (crash recovery).
Уже классифицированные ('classified', 'garbage', 'failed') НЕ трогаем."""
return execute(
"UPDATE documents SET classify_status='pending', error_message=NULL WHERE batch_id=%s AND classify_status='processing'",
(batch_id,),
)
def set_classify_processing(doc_id):
"""Mark document as being processed (for crash recovery)."""
return execute(
"UPDATE documents SET classify_status='processing' WHERE id=%s",
(doc_id,),
)
def list_by_batch(batch_id):
"""All documents in a batch with classification fields."""
return query(
"""SELECT id, filename, status, doc_type, own_number, parent_number,
doc_date, counterparty, classify_status, error_message, classify_raw, classify_input,
zip_source
FROM documents WHERE batch_id=%s ORDER BY created_at""",
(batch_id,),
)
def count_by_status(batch_id):
"""Count documents by classify_status."""
rows = query(
"SELECT classify_status, count(*) as cnt FROM documents WHERE batch_id=%s GROUP BY classify_status",
(batch_id,),
)
return {r["classify_status"]: r["cnt"] for r in rows}
+151
View File
@@ -0,0 +1,151 @@
"""Prompts CRUD."""
from .connection import query, execute, execute_returning
def _serialize(row):
"""Convert datetime fields to strings for JSON serialization (non-mutating)."""
if row and row.get("created_at"):
row = dict(row)
row["created_at"] = str(row["created_at"])
return row
def get_active(role):
rows = query(
"SELECT * FROM prompts WHERE role = %s AND is_active = true ORDER BY created_at DESC LIMIT 1",
(role,),
)
return _serialize(rows[0]) if rows else None
def get(prompt_id):
rows = query("SELECT * FROM prompts WHERE id = %s", (prompt_id,))
return rows[0] if rows else None
def seed_defaults():
"""Auto-seed default prompts if table is empty."""
# Check each role separately — don't skip if one is missing
existing_roles = set(r["role"] for r in query("SELECT DISTINCT role FROM prompts"))
if "extract" not in existing_roles:
extract_body = (
"Ты — анализатор договоров облачного провайдера.\n\n"
"Ниже текст спецификации услуг из ПЕРВОГО документа (базовый договор).\n"
"Извлеки ВСЕ строки спецификации в JSON-массив.\n\n"
"Верни СТРОГО JSON без пояснений. Не используй markdown-блоки.\n\n"
"ФОРМАТ:\n"
'{\n "mode": "partial",\n "ops": [\n'
' {"action": "ADD", "new_row": {"name": "полное наименование", "price": число, "qty": число, "sum": число, "date_start": "YYYY-MM-DD"}, "comment": ""}\n'
" ]\n}\n\n"
"ПРАВИЛА:\n"
"1. Извлеки КАЖДУЮ строку таблицы спецификации как отдельную ADD-операцию.\n"
'2. Пропускай итоговые строки ("Итого...") и строки с подписями.\n'
"3. Если ячейка пустая — ставь null (не пиши 0).\n"
"4. price, qty, sum — ЧИСЛА, не строки.\n\n"
"ТЕКСТ ДОКУМЕНТА:\n---\n{doc_text}\n---"
)
execute(
"INSERT INTO prompts (role, name, body, is_active, notes) VALUES (%s, %s, %s, %s, %s)",
("extract", "default-v1", extract_body, True, "Авто-создан из llm_prompt.py _build_initial"),
)
if "diff" not in existing_roles:
diff_body = (
"Ты — анализатор допсоглашений к договорам облачного провайдера.\n\n"
"У тебя есть текущая спецификация услуг и текст нового допсоглашения (ДС).\n"
"Твоя задача — определить, какие изменения вносит ДС в текущую спецификацию.\n\n"
"Верни СТРОГО JSON без пояснений. Не используй markdown-блоки.\n\n"
"ФОРМАТ:\n"
'{\n "mode": "partial" | "full_replace",\n "ops": [\n'
' {"action": "ADD", "new_row": {"name": "...", "price": число, "qty": число, "sum": число, "date_start": "YYYY-MM-DD"}, "comment": "..."},\n'
' {"action": "UPDATE", "target_id": "rN", "new_values": {"price": число}, "comment": "..."},\n'
' {"action": "DELETE", "target_id": "rN", "comment": "..."},\n'
' {"action": "UNRESOLVED", "new_values": {"name": "...", "price": число, ...}, "reason": "почему не смог сопоставить"}\n'
" ]\n}\n\n"
"ПРАВИЛА:\n"
'1. mode = "full_replace" — если в тексте есть фразы: «в следующей редакции», «заменить приложение», «излагается в следующей редакции». При full_replace — опиши ВСЕ новые строки как ADD.\n'
'2. mode = "partial" — если ДС меняет только отдельные строки.\n'
"3. Для UPDATE/DELETE — укажи target_id (r1, r2...) из списка текущей спецификации. НЕ придумывай новые id.\n"
"4. new_values в UPDATE — только ИЗМЕНЁННЫЕ поля (не все).\n"
"5. Если не можешь однозначно сопоставить строку — action: UNRESOLVED с reason.\n"
"6. Пропускай итоговые строки и подписи.\n"
"7. price, qty, sum — ЧИСЛА (не строки).\n\n"
"ТЕКУЩАЯ СПЕЦИФИКАЦИЯ:\n{spec_current}\n\n"
"ТЕКСТ ДОПСОГЛАШЕНИЯ:\n---\n{doc_text}\n---"
)
execute(
"INSERT INTO prompts (role, name, body, is_active, notes) VALUES (%s, %s, %s, %s, %s)",
("diff", "default-v1", diff_body, True, "Авто-создан из llm_prompt.py _build_diff"),
)
_ensure_classify_prompt()
def list_by_role(role):
"""List all versions for a role, newest first."""
rows = query(
"SELECT id, role, name, is_active, created_by, notes, created_at FROM prompts WHERE role=%s ORDER BY created_at DESC",
(role,),
)
for r in rows:
if r.get("created_at"):
r["created_at"] = str(r["created_at"])
return rows
def save_new_version(role, name, body, notes="", is_active=True):
"""Save new prompt version. Deactivates all others for this role, inserts new one."""
if is_active:
execute("UPDATE prompts SET is_active=false WHERE role=%s", (role,))
return execute_returning(
"""INSERT INTO prompts (role, name, body, is_active, notes)
VALUES (%s, %s, %s, %s, %s) RETURNING *""",
(role, name, body, is_active, notes),
)
def activate(prompt_id):
"""Activate a prompt version (deactivates others for same role)."""
row = query("SELECT role FROM prompts WHERE id=%s", (prompt_id,))
if not row:
return False
role = row[0]["role"]
execute("UPDATE prompts SET is_active=false WHERE role=%s", (role,))
execute("UPDATE prompts SET is_active=true WHERE id=%s", (prompt_id,))
return True
def delete_prompt(prompt_id):
"""Delete a prompt version (cannot delete active one)."""
row = query("SELECT is_active FROM prompts WHERE id=%s", (prompt_id,))
if not row:
return False
if row[0]["is_active"]:
return False # cannot delete active
execute("DELETE FROM prompts WHERE id=%s", (prompt_id,))
return True
def _ensure_classify_prompt():
"""Ensure classify prompt exists (idempotent)."""
existing = query("SELECT id FROM prompts WHERE role = 'classify' AND is_active = true LIMIT 1")
if existing:
return
body = (
"Ты — система классификации договорных документов. "
"Проанализируй текст и верни СТРОГО ВАЛИДНЫЙ JSON ОДНОЙ СТРОКОЙ "
"(без переносов строк, без markdown, без лишних пробелов в начале/конце).\n\n"
"Поля:\n"
'- doc_type: "contract" (договор) / "supplement" (допсоглашение) / "specification" (спецификация/приложение) / "other"\n'
"- own_number: номер ЭТОГО документа, строка без лишних пробелов (или null)\n"
'- parent_number: номер родительского договора из фразы «к Договору №...» (или null)\n'
'- doc_date: дата в YYYY-MM-DD. «27 февраля 2026» → 2026-02-27 (или null)\n'
"- counterparty: полное название контрагента (Заказчик/Арендатор), без сокращений (или null)\n\n"
"Пример вывода:\n"
'{"doc_type":"supplement","own_number":"1","parent_number":"01300_2","doc_date":"2026-02-27","counterparty":"АО XXX003"}\n\n'
"ДОКУМЕНТ:\n---\n{header_text}\n---"
)
execute(
"INSERT INTO prompts (role, name, body, is_active, notes) VALUES (%s, %s, %s, %s, %s)",
("classify", "default-v1", body, True, "Авто-создан: classify prompt v2 — LLM сам предложил"),
)
+19
View File
@@ -0,0 +1,19 @@
"""spec_current — текущее состояние спецификации."""
from .connection import query
def list_by_contract(contract_id):
"""Return list of dicts with name_hash, name, price, qty, sum, date_start."""
return query(
"""SELECT name_hash, name, price, qty, sum, date_start
FROM spec_current WHERE contract_id = %s ORDER BY name""",
(contract_id,),
)
def get_elements_json(document_id):
"""Get elements_json for a document."""
rows = query(
"SELECT elements_json FROM documents WHERE id = %s", (document_id,)
)
return rows[0]["elements_json"] if rows and rows[0]["elements_json"] else None
+202
View File
@@ -0,0 +1,202 @@
"""spec_events — event sourcing: apply ops, reset contract."""
import json, uuid
from .connection import query, execute, get_pool
def reset(contract_id):
"""Clear spec_events + spec_current for contract."""
execute("DELETE FROM spec_current WHERE contract_id = %s", (contract_id,))
execute("DELETE FROM spec_events WHERE contract_id = %s", (contract_id,))
def get_next_seq(contract_id):
"""Get next sequence number with row lock to prevent race conditions."""
pool = get_pool()
conn = pool.getconn()
try:
conn.autocommit = False
with conn.cursor() as cur:
cur.execute(
"SELECT seq FROM spec_events WHERE contract_id = %s ORDER BY seq DESC LIMIT 1 FOR UPDATE",
(contract_id,),
)
row = cur.fetchone()
seq = (row[0] + 1) if row else 1
conn.commit()
return seq
except Exception:
conn.rollback()
raise
finally:
conn.autocommit = True
pool.putconn(conn)
def apply_ops(contract_id, supplement_id, document_id, ops, prompt_id, raw_llm_response):
"""Apply ADD/UPDATE/DELETE ops. Returns summary dict."""
added = 0
updated = 0
deleted = 0
seq = get_next_seq(contract_id)
for op in ops:
action = op.get("action", "").upper()
if action == "ADD":
nr = op.get("new_row", {})
name = nr.get("name", "")
if not name:
# ADD without name → UNRESOLVED
_log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, "ADD with empty name")
seq += 1
continue
name_hash = _hash(name, nr.get("date_start"))
execute(
"""INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash,
new_values, comment, status, prompt_version, source_document_id, raw_llm_response)
VALUES (%s, %s, %s, 'ADD', %s, %s, %s, 'applied', %s, %s, %s)""",
(
contract_id, supplement_id, seq, name_hash,
json.dumps(nr, ensure_ascii=False),
op.get("comment", ""), prompt_id, document_id,
json.dumps(raw_llm_response, ensure_ascii=False),
),
)
_upsert_spec_current(contract_id, name_hash, nr)
seq += 1
added += 1
elif action == "UPDATE":
nv = op.get("new_values", {})
th = op.get("target_hash", "")
if not th:
_log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, "UPDATE with empty target_hash")
seq += 1
continue
execute(
"""INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash,
new_values, comment, status, prompt_version, source_document_id, raw_llm_response)
VALUES (%s, %s, %s, 'UPDATE', %s, %s, %s, 'applied', %s, %s, %s)""",
(
contract_id, supplement_id, seq, th,
json.dumps(nv, ensure_ascii=False),
op.get("comment", ""), prompt_id, document_id,
json.dumps(raw_llm_response, ensure_ascii=False),
),
)
_update_spec_current(contract_id, th, nv)
seq += 1
updated += 1
elif action == "DELETE":
th = op.get("target_hash", "")
if not th:
_log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, "DELETE with empty target_hash")
seq += 1
continue
execute(
"""INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash,
new_values, comment, status, prompt_version, source_document_id, raw_llm_response)
VALUES (%s, %s, %s, 'DELETE', %s, %s, %s, 'applied', %s, %s, %s)""",
(
contract_id, supplement_id, seq, th,
json.dumps({}), op.get("comment", ""),
prompt_id, document_id,
json.dumps(raw_llm_response, ensure_ascii=False),
),
)
execute("DELETE FROM spec_current WHERE contract_id = %s AND name_hash = %s", (contract_id, th))
seq += 1
deleted += 1
elif action == "UNRESOLVED":
# Log but don't apply
execute(
"""INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash,
new_values, comment, status, prompt_version, source_document_id, raw_llm_response)
VALUES (%s, %s, %s, 'UNRESOLVED', %s, %s, %s, 'unresolved', %s, %s, %s)""",
(
contract_id, supplement_id, seq,
op.get("target_hash", ""),
json.dumps(op.get("new_values", {}), ensure_ascii=False),
op.get("reason", op.get("comment", "")),
prompt_id, document_id,
json.dumps(raw_llm_response, ensure_ascii=False),
),
)
seq += 1
else:
# Unknown action — log as UNRESOLVED
_log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response,
f"unknown action: {action}")
return {"added": added, "updated": updated, "deleted": deleted}
def _log_unresolved(contract_id, supplement_id, seq, op, prompt_id, document_id, raw_llm_response, reason):
"""Log an op as UNRESOLVED instead of silently ignoring it."""
execute(
"""INSERT INTO spec_events (contract_id, supplement_id, seq, action, target_hash,
new_values, comment, status, prompt_version, source_document_id, raw_llm_response)
VALUES (%s, %s, %s, 'UNRESOLVED', %s, %s, %s, 'unresolved', %s, %s, %s)""",
(
contract_id, supplement_id, seq,
op.get("target_hash", ""),
json.dumps(op.get("new_values", op.get("new_row", {})) or {}, ensure_ascii=False),
reason,
prompt_id, document_id,
json.dumps(raw_llm_response, ensure_ascii=False),
),
)
def _hash(name, date_start=None):
"""Нормализованный хеш услуги. Включает нормализованный date_start чтобы различать периоды."""
import hashlib
key = name.strip().lower()
if date_start:
from compare.metrics import normalize_date
nd = normalize_date(str(date_start))
if nd:
key += "|" + nd
return hashlib.sha256(key.encode()).hexdigest()[:16]
def _upsert_spec_current(contract_id, name_hash, row):
"""INSERT or UPDATE spec_current."""
existing = query(
"SELECT id FROM spec_current WHERE contract_id = %s AND name_hash = %s",
(contract_id, name_hash),
)
if existing:
execute(
"""UPDATE spec_current SET name=%s, price=%s, qty=%s, sum=%s, date_start=%s, updated_at=now()
WHERE contract_id=%s AND name_hash=%s""",
(row.get("name"), row.get("price"), row.get("qty"), row.get("sum"),
row.get("date_start"), contract_id, name_hash),
)
else:
execute(
"""INSERT INTO spec_current (contract_id, name_hash, name, price, qty, sum, date_start)
VALUES (%s, %s, %s, %s, %s, %s, %s)""",
(contract_id, name_hash, row.get("name"), row.get("price"),
row.get("qty"), row.get("sum"), row.get("date_start")),
)
def _update_spec_current(contract_id, name_hash, new_values):
"""Update specific fields in spec_current."""
sets = []
params = []
for field in ("name", "price", "qty", "sum", "date_start"):
if field in new_values:
sets.append(f"{field} = %s")
params.append(new_values[field])
if sets:
sets.append("updated_at = now()")
params.extend([contract_id, name_hash])
execute(
f"UPDATE spec_current SET {', '.join(sets)} WHERE contract_id = %s AND name_hash = %s",
params,
)
+59
View File
@@ -0,0 +1,59 @@
"""Supplements CRUD."""
from .connection import query, execute, execute_returning
def insert(contract_id, document_id, supp_type="additional"):
return execute_returning(
"""INSERT INTO supplements (contract_id, document_id, type)
VALUES (%s, %s, %s) RETURNING *""",
(contract_id, document_id, supp_type),
)
def list_by_contract(contract_id):
"""Supplements with parsed documents, ordered by created_at."""
return query(
"""SELECT s.id, s.type, s.document_id, d.filename
FROM supplements s
JOIN documents d ON s.document_id = d.id
WHERE s.contract_id = %s AND d.elements_json IS NOT NULL
ORDER BY s.created_at""",
(contract_id,),
)
def get(supp_id):
rows = query("SELECT * FROM supplements WHERE id = %s", (supp_id,))
return rows[0] if rows else None
def delete_by_document(contract_id, filename):
"""Delete ALL supplements+documents by contract+filename (cascade: spec_events first).
Handles duplicates from previously failed uploads."""
rows = query(
"""SELECT s.id as sid, s.document_id FROM supplements s
JOIN documents d ON d.id = s.document_id
WHERE s.contract_id = %s AND d.filename = %s""",
(contract_id, filename),
)
if not rows:
return False
for r in rows:
# 1. Delete spec_current rows referencing this supplement's events
execute(
"""DELETE FROM spec_current WHERE contract_id = %s
AND last_event_id IN (SELECT id FROM spec_events WHERE supplement_id = %s)""",
(contract_id, r["sid"]),
)
# 2. Delete spec_events referencing this supplement
execute("DELETE FROM spec_events WHERE supplement_id = %s", (r["sid"],))
# 3. Delete supplement
execute("DELETE FROM supplements WHERE id = %s", (r["sid"],))
# 4. Delete document
execute("DELETE FROM documents WHERE id = %s", (r["document_id"],))
return True
def delete(supp_id):
return execute("DELETE FROM supplements WHERE id = %s", (supp_id,))