tender / worker.py
Levin-Aleksey's picture
fix
83eb096
Raw
History Blame Contribute Delete
49.4 kB
import os
import json
import asyncio
import tempfile
from pathlib import Path
import mimetypes
import urllib.parse
import re
import httpx
import aio_pika
import pandas as pd
import openpyxl
from openpyxl.styles import Alignment, Font, PatternFill
from supabase import create_client, Client
# Настройки
SUPABASE_URL = os.getenv("SUPABASE_URL", "")
SUPABASE_KEY = os.getenv("SUPABASE_KEY", "")
CLOUDAMQP_URL = os.getenv("CLOUDAMQP_URL", "")
OPENROUTER_API_KEY = os.getenv("OPENROUTER_API_KEY", "")
UNSTRUCTURED_API_KEY = os.getenv("UNSTRUCTURED_API_KEY", "")
UNSTRUCTURED_API_URL = os.getenv(
"UNSTRUCTURED_API_URL", "https://api.unstructuredapp.io/general/v0/general"
).strip()
if UNSTRUCTURED_API_URL.startswith("https://api.unstructured.io"):
UNSTRUCTURED_API_URL = UNSTRUCTURED_API_URL.replace(
"https://api.unstructured.io", "https://api.unstructuredapp.io", 1
)
supabase: Client = create_client(SUPABASE_URL, SUPABASE_KEY)
def _parse_quantity(value) -> float:
if value is None:
return 0.0
if isinstance(value, (int, float)):
return float(value)
raw = str(value).strip().replace(",", ".")
if not raw:
return 0.0
match = re.search(r"-?\d+(?:\.\d+)?", raw)
if not match:
return 0.0
try:
return float(match.group(0))
except ValueError:
return 0.0
def _parse_llm_json(raw_text: str) -> dict:
try:
data = json.loads(raw_text)
if isinstance(data, dict):
return data
except json.JSONDecodeError:
pass
match = re.search(r"\{[\s\S]*\}", raw_text)
if match:
try:
data = json.loads(match.group(0))
if isinstance(data, dict):
return data
except json.JSONDecodeError:
pass
raise RuntimeError("LLM вернул невалидный JSON")
def _normalize_catalog_key(value: str) -> str:
# убирает регистр/кавычки/тире и лишние пробелы, чтобы LLM-варианты имени находили совпадения в справочнике
text = (value or "").strip().lower()
text = text.replace("\u2013", "-").replace("\u2014", "-")
text = re.sub(r"[\"'\u00ab\u00bb]", "", text)
text = re.sub(r"\s+", " ", text)
return text
def _normalize_unrecognized_entry(entry, source: str) -> dict:
if isinstance(entry, dict):
return {
"source": entry.get("source", source),
"description": entry.get("description") or entry.get("name") or str(entry),
"quantity": entry.get("quantity", ""),
"raw_context": entry.get("raw_context", ""),
}
return {
"source": source,
"description": str(entry),
"quantity": "",
"raw_context": "",
}
TENDER_DATA_SECTIONS = (
("НМЦК И ДОГОВОР", "nmck_contract"),
("ОБЪЕКТЫ И ПЛОЩАДИ", "sites_areas"),
("ПЕРСОНАЛ", "personnel"),
("ОБОРУДОВАНИЕ И ТЕХНИКА", "equipment"),
("РАСХОДНИКИ И МАТЕРИАЛЫ", "consumables"),
("ДОПОЛНИТЕЛЬНЫЕ РАБОТЫ", "additional_services"),
("РИСКИ И ТРЕБОВАНИЯ", "risks_requirements"),
("ПРОЧЕЕ", "other"),
)
FACT_STATUSES = {
"found": "Найдено",
"missing": "Не найдено",
"insufficient": "Недостаточно данных",
"conflict": "Конфликт документов",
"manual_review": "Требуется ручная проверка",
}
TENDER_DATA_PARAMETERS = {
"nmck_contract": (
"НМЦК и валюта",
"Разбивка начальной цены по видам работ",
"Срок оказания услуг",
"Лицензии",
"Допуски и разрешения",
"Страхование",
"Арендные платежи исполнителя заказчику",
),
"sites_areas": (
"Вид деятельности и назначение объекта",
"Адреса объектов",
"Удаленность объектов друг от друга",
"Площади внутренних помещений",
"Площади внешней территории",
"Удаленность от остановок общественного транспорта",
),
"personnel": (
"Вид и должности персонала",
"Количество персонала",
"Гражданство персонала",
"Санитарные и медицинские книжки",
"Работа по субботам и воскресеньям",
"График и часы работы",
"Дополнительные требования к персоналу",
"Доставка персонала исполнителем",
),
"equipment": (
"Оборудование для уборки и поставщик оборудования",
"Механизированная уборка и спецтехника",
),
"consumables": (
"Санузловые расходные материалы",
"Характеристики санузловых расходников",
"Количество санузловых расходников",
"Количество сотрудников или посетителей для расчета расходников",
"Прочие расходники и их количество",
"Диспенсеры, дозаторы и электросушилки",
"Материалы для ремонта и сторона предоставления",
),
"additional_services": (
"Химчистка мебели, штор, тюля и ковров",
"Мытье окон и витражей",
"Высотные работы",
"Вывоз снега",
"Замена грязезащитных ковров",
"Чистка кровли, сосулек и наледи",
"Благоустройство и озеленение",
"Дезинсекция и дератизация",
"Мытье бассейна или фонтана",
),
"risks_requirements": (
"Критичные условия и риски договора",
"Требования, влияющие на стоимость",
),
"other": ("Прочие значимые данные",),
}
# --- БЛОК 1: Парсеры ---
def _safe_filename_from_url_or_response(url: str, response: httpx.Response) -> str:
parsed = urllib.parse.urlparse(url)
query = urllib.parse.parse_qs(parsed.query)
filename_values = query.get("filename")
if filename_values and filename_values[0].strip():
return Path(filename_values[0].strip()).name
content_disposition = response.headers.get("content-disposition", "")
match = re.search(r'filename\*?=(?:UTF-8\'\')?"?([^";]+)"?', content_disposition)
if match:
return Path(urllib.parse.unquote(match.group(1))).name
suffix = Path(parsed.path).suffix
if suffix:
return "document" + suffix
content_type = response.headers.get("content-type", "")
ext = mimetypes.guess_extension(content_type.split(";")[0].strip()) or ""
return "document" + ext
def extract_text_from_excel(file_path: str) -> str:
sheets = pd.read_excel(file_path, sheet_name=None, header=None, dtype=str)
parts = []
for sheet_name, df in sheets.items():
df = df.fillna("")
parts.append(f"\n\n# Лист: {sheet_name}")
for row_num, row in enumerate(df.itertuples(index=False), start=1):
values = [str(cell).strip() for cell in row if str(cell).strip()]
if values:
parts.append(f"Строка {row_num}: " + " | ".join(values))
text = "\n".join(parts).strip()
if not text:
raise RuntimeError("Excel-файл пустой или не удалось извлечь текст")
return text
def _elements_to_text(elements) -> str:
parts = []
for element in elements:
if not isinstance(element, dict):
continue
category = element.get("type") or element.get("category") or "Element"
text = (element.get("text") or "").strip()
metadata = element.get("metadata") or {}
table_html = (metadata.get("text_as_html") or "").strip()
if table_html:
parts.append(f"\n\n## {category}\n{table_html}")
elif text:
parts.append(f"\n\n## {category}\n{text}")
result = "".join(parts).strip()
if not result:
raise RuntimeError("Unstructured вернул пустой текст")
return result
async def call_smart_parser(file_path: str) -> str:
if not UNSTRUCTURED_API_KEY:
raise RuntimeError("UNSTRUCTURED_API_KEY is not set")
filename = os.path.basename(file_path) or "document"
mime_type = mimetypes.guess_type(filename)[0] or "application/octet-stream"
headers = {
"accept": "application/json",
"unstructured-api-key": UNSTRUCTURED_API_KEY,
}
data = {
"strategy": "hi_res",
"languages": "rus",
"skip_infer_table_types": "[]",
"include_page_breaks": "true",
}
timeout_config = httpx.Timeout(900.0, read=900.0, connect=60.0)
async with httpx.AsyncClient(
timeout=timeout_config, follow_redirects=True
) as http_client:
with open(file_path, "rb") as f:
response = await http_client.post(
UNSTRUCTURED_API_URL,
headers=headers,
data=data,
files={"files": (filename, f, mime_type)},
)
try:
response.raise_for_status()
except httpx.HTTPStatusError as e:
raise RuntimeError(
f"Ошибка Unstructured API {response.status_code}: {response.text[:2000]}"
) from e
elements = response.json()
if not isinstance(elements, list):
raise RuntimeError(f"Unstructured вернул неожиданный ответ: {elements}")
return _elements_to_text(elements)
# --- БЛОК 2: LLM (DeepSeek) + Интеграция с БД ---
def get_pricebook_catalog_sync() -> list:
catalog = []
try:
p_res = (
supabase.table("personnel_rates")
.select("raw_value")
.eq("is_active", True)
.execute()
)
for r in p_res.data:
if r.get("raw_value"):
catalog.append(r.get("raw_value").split("|")[0].strip())
except:
pass
try:
i_res = (
supabase.table("price_items").select("name").eq("is_active", True).execute()
)
for r in i_res.data:
if r.get("name"):
catalog.append(r.get("name").strip())
except:
pass
result = sorted(list(set(catalog)))
if not result:
raise RuntimeError("Справочник пуст: не удалось выгрузить прайс-лист из БД")
return result
async def analyze_text_with_deepseek_analytics(text: str) -> dict:
if not OPENROUTER_API_KEY:
raise RuntimeError("OPENROUTER_API_KEY is not set")
api_url = "https://openrouter.ai/api/v1/chat/completions"
headers = {
"Authorization": f"Bearer {OPENROUTER_API_KEY}",
"Content-Type": "application/json",
}
system_prompt = """Ты — Главный аналитик тендерной документации и старший сметчик клининговой компании.
Твоя задача — подготовить КРАТКИЙ УПРАВЛЕНЧЕСКИЙ АУДИТ тендера: понять, можно ли переходить к ручному расчету, какие условия влияют на стоимость и какие риски нужно снять до участия.
НЕ рассчитывай цены, ФОТ, прибыль, налоги и итоговую стоимость. Не придумывай данные.
Верни только строгий JSON по указанной схеме:
{
"tender_card": {
"nmck": {"value": "", "status": "found|missing|insufficient|conflict|manual_review", "source_document": "", "raw_context": ""},
"service_term": {"value": "", "status": "", "source_document": "", "raw_context": ""},
"object_type": {"value": "", "status": "", "source_document": "", "raw_context": ""},
"addresses_and_areas": {"value": "", "status": "", "source_document": "", "raw_context": ""},
"personnel_summary": {"value": "", "status": "", "source_document": "", "raw_context": ""},
"work_schedule": {"value": "", "status": "", "source_document": "", "raw_context": ""}
},
"cost_drivers": [
{"factor": "", "impact": "", "status": "", "source_document": "", "raw_context": ""}
],
"risks": [
{"risk": "", "severity": "high|medium|low", "status": "", "source_document": "", "raw_context": "", "recommendation": ""}
],
"conflicts": [
{"topic": "", "details": "", "source_document": "", "raw_context": ""}
],
"critical_questions": [
{"question": "", "reason": "", "severity": "high|medium|low"}
],
"overall_assessment": {
"readiness": "ready|needs_clarification|high_risk",
"summary": ""
}
}
ПРАВИЛА:
1. Статусы: found — данные достаточны; missing — данных нет; insufficient — данные есть, но для решения недостаточны; conflict — документы противоречат; manual_review — нужна ручная трактовка.
2. У каждого найденного факта указывай source_document и дословную raw_context. Для missing источник и цитата могут быть пустыми.
3. При конфликте НЕ выбирай версию самостоятельно. Зафиксируй оба значения и источники в details/raw_context, а также добавь запись в conflicts.
4. К критичным вопросам относись строго: включай только то, что блокирует расчет, меняет цену или создает договорной риск.
5. В cost_drivers отрази персонал и графики, площади, работу по выходным, расходники, оборудование, спецтехнику, дополнительные работы, удаленность, лицензии и другие условия, влияющие на стоимость.
6. В risks отрази лицензии, допуски, медкнижки, гражданство, страхование, аренду, сторону предоставления оборудования/материалов, неясные объемы и рискованные договорные требования.
7. Если по карточке тендера данных нет — ставь status="missing", value="". НИКОГДА НЕ ДОДУМЫВАЙ.
8. В overall_assessment.readiness: ready — критичных пробелов нет; needs_clarification — есть вопросы до расчета; high_risk — есть существенные риски или конфликты.
"""
payload = {
"model": "deepseek/deepseek-v4-pro",
"messages": [
{"role": "system", "content": system_prompt},
{"role": "user", "content": f"Проанализируй текст и верни JSON:\n\n{text}"},
],
"reasoning": {"enabled": True},
"temperature": 0.1,
"response_format": {"type": "json_object"},
}
async with httpx.AsyncClient(timeout=240.0) as client:
response = await client.post(api_url, headers=headers, json=payload)
response.raise_for_status()
return _parse_llm_json(
response.json()["choices"][0]["message"]
.get("content", "")
.replace("```json", "")
.replace("```", "")
.strip()
)
async def extract_tender_data_for_estimator(text: str) -> dict:
if not OPENROUTER_API_KEY:
raise RuntimeError("OPENROUTER_API_KEY is not set")
api_url = "https://openrouter.ai/api/v1/chat/completions"
headers = {
"Authorization": f"Bearer {OPENROUTER_API_KEY}",
"Content-Type": "application/json",
}
sections = ", ".join(section_key for _, section_key in TENDER_DATA_SECTIONS)
system_prompt = f"""Ты — старший сметчик клининговой компании. Прочитай ОБЪЕДИНЕННЫЙ текст документов тендера и подготовь КАРТУ ИСХОДНЫХ ДАННЫХ для человека-сметчика.
Твоя задача — только извлечь факты. НЕ рассчитывай цены, ФОТ, прибыль, налоги и не сопоставляй с прайс-листом.
РАЗДЕЛЫ JSON: {sections}.
Для каждого найденного факта верни объект:
{{
"parameter": "точное название проверяемого параметра",
"status": "found | missing | insufficient | conflict | manual_review",
"value": "значение только из документов",
"source_document": "имя документа после заголовка ДОКУМЕНТ",
"source_location": "страница, лист, строка или раздел, если есть",
"raw_context": "точная цитата из документа",
"comment": "что должен проверить сметчик"
}}
ПРАВИЛА:
1. Не придумывай данные. Если значение отсутствует — status="missing", value="", а в comment сформулируй вопрос заказчику.
2. Если данные есть, но не позволяют считать (например, расходники без количества) — status="insufficient" и укажи чего не хватает.
3. Если документы противоречат друг другу — status="conflict"; сохрани оба значения и оба источника в value/raw_context. Не выбирай одно значение сам.
4. Если формулировка двусмысленна — status="manual_review" и объясни причину.
5. Для каждого статуса кроме missing обязательно заполняй source_document и raw_context дословной цитатой.
6. Все количественные значения сохраняй вместе с единицами измерения, периодичностью, площадью, графиком или стороной предоставления, когда они указаны.
7. Обязательно ищи: НМЦК и разбивку цены; срок, лицензии, допуски, страхование и аренду; типы объектов, адреса, площади; персонал, гражданство, медкнижки и графики; оборудование и кто его предоставляет; расходники и характеристики; все дополнительные работы.
8. Не создавай разделов вне списка. Возвращай только строгий JSON.
Формат ответа:
{{
"sections": {{
"personnel": [{{"parameter": "Количество персонала", "status": "found", "value": "8 уборщиков", "source_document": "ТЗ.pdf", "source_location": "стр. 12", "raw_context": "...", "comment": ""}}]
}},
"critical_questions": ["вопрос заказчику, без которого нельзя рассчитывать"],
"documents_summary": ["краткое замечание по документу"]
}}"""
payload = {
"model": "deepseek/deepseek-v4-pro",
"messages": [
{"role": "system", "content": system_prompt},
{
"role": "user",
"content": f"Вот объединенный текст всех документов тендера:\n\n{text}",
},
],
"reasoning": {"enabled": True},
"temperature": 0.1,
"response_format": {"type": "json_object"},
}
async with httpx.AsyncClient(timeout=240.0) as client:
response = await client.post(api_url, headers=headers, json=payload)
response.raise_for_status()
return _parse_llm_json(
response.json()["choices"][0]["message"]
.get("content", "")
.replace("```json", "")
.replace("```", "")
.strip()
)
# --- БЛОК 3: Обработка очередей ---
async def download_file(url: str, temp_dir: str) -> str:
if "yandex.ru" in url or "ya.disk" in url:
async with httpx.AsyncClient(follow_redirects=True) as client:
resp = await client.get(
f"https://cloud-api.yandex.net/v1/disk/public/resources/download?public_key={urllib.parse.quote(url)}"
)
resp.raise_for_status()
url = resp.json().get("href", url)
async with httpx.AsyncClient(follow_redirects=True) as client:
resp = await client.get(url)
resp.raise_for_status()
filename = _safe_filename_from_url_or_response(url, resp)
path = os.path.join(temp_dir, filename)
with open(path, "wb") as f:
f.write(resp.content)
return path
# --- БЛОК 2.5: МАГИЧЕСКИЙ КАЛЬКУЛЯТОР СМЕТЫ ---
class TenderCalculator:
def __init__(self, supabase_client):
self.supabase = supabase_client
self.price_book = {}
self.normalized_lookup = {}
self.cleaning_group_id = None
self._load_reference_data()
def _load_reference_data(self):
groups = self.supabase.table("payroll_groups").select("id, code").execute()
for g in groups.data:
if g.get("code") == "CLEANING_STAFF":
self.cleaning_group_id = g.get("id")
if self.cleaning_group_id is None:
raise RuntimeError("В payroll_groups отсутствует код CLEANING_STAFF.")
personnel = (
self.supabase.table("personnel_rates")
.select("raw_value, rate_min, payroll_group_id")
.eq("is_active", True)
.execute()
)
for p in personnel.data:
val = p.get("raw_value", "")
if val:
key = val.split("|")[0].strip()
self.price_book[key] = {
"type": "personnel",
"price": float(p.get("rate_min") or 0),
"group_id": p.get("payroll_group_id"),
}
self.normalized_lookup[_normalize_catalog_key(key)] = key
items = (
self.supabase.table("price_items")
.select("name, price_min, category_id")
.eq("is_active", True)
.execute()
)
for i in items.data:
name = i.get("name")
if name:
key = name.strip()
self.price_book[key] = {
"type": "item",
"price": float(i.get("price_min") or 0),
"category_id": i.get("category_id"),
}
self.normalized_lookup[_normalize_catalog_key(key)] = key
def _find_price_entry(self, name: str):
# сначала точное совпадение, затем поиск без учёта регистра/пробелов/тире
if not name:
return None, None
if name in self.price_book:
return self.price_book[name], name
canonical_name = self.normalized_lookup.get(_normalize_catalog_key(name))
if canonical_name:
return self.price_book[canonical_name], canonical_name
return None, None
def calculate(self, extracted_data: dict) -> dict:
estimate = {
"details": [],
"totals": {
"cleaning_payroll": 0,
"other_payroll": 0,
"materials_and_items": 0,
"calculated_chemistry": 0,
"grand_total": 0,
},
"unrecognized": [
_normalize_unrecognized_entry(entry, source="llm_unrecognized")
for entry in extracted_data.get("unrecognized_items", [])
],
}
for item in extracted_data.get("extracted_items", []):
name, quantity = item.get("exact_name"), _parse_quantity(
item.get("quantity", 0)
)
if not name or quantity <= 0:
continue
ref_data, canonical_name = self._find_price_entry(name)
if ref_data:
price = ref_data["price"]
total_cost = price * quantity
estimate["details"].append(
{
"name": canonical_name or name,
"quantity": quantity,
"unit_price": price,
"total_cost": total_cost,
"type": ref_data["type"],
"raw_context": item.get("raw_context", ""),
}
)
if ref_data["type"] == "personnel":
if ref_data["group_id"] == self.cleaning_group_id:
estimate["totals"]["cleaning_payroll"] += total_cost
else:
estimate["totals"]["other_payroll"] += total_cost
else:
estimate["totals"]["materials_and_items"] += total_cost
else:
estimate["unrecognized"].append(
{
"source": "no_price_match",
"description": name,
"quantity": quantity,
"raw_context": item.get("raw_context", ""),
}
)
chemistry_cost = estimate["totals"]["cleaning_payroll"] * 0.10
estimate["totals"]["calculated_chemistry"] = chemistry_cost
if chemistry_cost > 0:
estimate["details"].append(
{
"name": "Бюджет химии и расходников (10% от ФОТ уборщиц)",
"quantity": 1,
"unit_price": chemistry_cost,
"total_cost": chemistry_cost,
"type": "calculated_rule",
"raw_context": "Расчётное правило: не найдено в тексте документа",
}
)
estimate["totals"]["grand_total"] = sum(
[
estimate["totals"]["cleaning_payroll"],
estimate["totals"]["other_payroll"],
estimate["totals"]["materials_and_items"],
estimate["totals"]["calculated_chemistry"],
]
)
return estimate
# --- БЛОК 2.7: ГЕНЕРАТОР EXCEL ---
def create_excel_report(estimate_data: dict, temp_dir: str, task_id: str) -> str:
wb = openpyxl.Workbook()
ws_summary = wb.active
ws_summary.title = "СВОДНАЯ"
headers_font, header_fill, bold_font = (
Font(bold=True, color="FFFFFF"),
PatternFill("solid", fgColor="4F81BD"),
Font(bold=True),
)
ws_summary.append(["Статья расходов", "Сумма (руб.)"])
ws_summary.cell(row=1, column=1).font, ws_summary.cell(row=1, column=1).fill = (
headers_font,
header_fill,
)
ws_summary.cell(row=1, column=2).font, ws_summary.cell(row=1, column=2).fill = (
headers_font,
header_fill,
)
totals = estimate_data.get("totals", {})
ws_summary.append(["ФОТ (Уборочный персонал)", totals.get("cleaning_payroll", 0)])
ws_summary.append(["ФОТ (Прочий персонал)", totals.get("other_payroll", 0)])
ws_summary.append(
["Материалы, инвентарь, техника", totals.get("materials_and_items", 0)]
)
ws_summary.append(
["Химия и расходники (по правилу 10%)", totals.get("calculated_chemistry", 0)]
)
ws_summary.append(["", ""])
ws_summary.append(["ИТОГО РАСХОДЫ (Себестоимость):", "=SUM(B2:B5)"])
ws_summary.cell(row=7, column=1).font = bold_font
ws_summary.cell(row=7, column=2).font = bold_font
ws_summary.append(["Налоги УСН (6%):", "=B7*0.06"])
ws_summary.append(["Желаемая прибыль (10%):", "=B7*0.10"])
ws_summary.append(["", ""])
ws_summary.append(["ИТОГО СТОИМОСТЬ КЛИЕНТУ:", "=B7+B8+B9"])
ws_summary.cell(row=11, column=1).font, ws_summary.cell(row=11, column=2).font = (
Font(bold=True, color="FF0000", size=12),
Font(bold=True, color="FF0000", size=12),
)
ws_summary.column_dimensions["A"].width, ws_summary.column_dimensions["B"].width = (
45,
20,
)
ws_details = wb.create_sheet(title="ДЕТАЛИЗАЦИЯ")
ws_details.append(
[
"Категория",
"Наименование",
"Кол-во",
"Цена за ед.",
"Итоговая сумма",
"Цитата из ТЗ",
]
)
for col in range(1, 7):
(
ws_details.cell(row=1, column=col).font,
ws_details.cell(row=1, column=col).fill,
) = (headers_font, header_fill)
type_labels = {
"personnel": "Персонал",
"item": "Материал/инвентарь",
"calculated_rule": "Расчётное правило (не из ТЗ)",
}
row_idx = 2
for item in estimate_data.get("details", []):
category = type_labels.get(
item.get("type", "unknown"), item.get("type", "unknown")
)
ws_details.append(
[
category,
item.get("name", ""),
item.get("quantity", 0),
item.get("unit_price", 0),
f"=C{row_idx}*D{row_idx}",
item.get("raw_context", ""),
]
)
row_idx += 1
for idx, w in enumerate([20, 55, 10, 15, 15, 50], 1):
ws_details.column_dimensions[openpyxl.utils.get_column_letter(idx)].width = w
ws_unrecognized = wb.create_sheet(title="НЕРАСПОЗНАНО")
ws_unrecognized.append(["Источник", "Описание", "Кол-во", "Цитата из ТЗ"])
for col in range(1, 5):
(
ws_unrecognized.cell(row=1, column=col).font,
ws_unrecognized.cell(row=1, column=col).fill,
) = (headers_font, PatternFill("solid", fgColor="C0504D"))
source_labels = {
"llm_unrecognized": "Не найдено в справочнике (по мнению ИИ)",
"no_price_match": "Не сопоставлено с прайс-листом",
}
unrecognized_items = estimate_data.get("unrecognized", [])
if not unrecognized_items:
ws_unrecognized.append(["-", "Нераспознанных позиций нет", "-", "-"])
else:
for entry in unrecognized_items:
normalized = _normalize_unrecognized_entry(entry, source="unknown")
ws_unrecognized.append(
[
source_labels.get(normalized["source"], normalized["source"]),
normalized["description"],
normalized["quantity"],
normalized["raw_context"],
]
)
for idx, w in enumerate([35, 45, 10, 60], 1):
ws_unrecognized.column_dimensions[
openpyxl.utils.get_column_letter(idx)
].width = w
file_path = os.path.join(temp_dir, f"Smeta_Tender_{task_id[:8]}.xlsx")
wb.save(file_path)
return file_path
def upload_excel_to_supabase(file_path: str, filename: str) -> str:
with open(file_path, "rb") as f:
supabase.storage.from_("tenders").upload(
filename,
f.read(),
file_options={
"content-type": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet"
},
)
public_url = supabase.storage.from_("tenders").get_public_url(filename)
return (
public_url.get("publicURL") or public_url.get("public_url") or ""
if isinstance(public_url, dict)
else str(public_url)
)
def _normalize_tender_fact(fact, parameter: str) -> dict:
fact = fact if isinstance(fact, dict) else {}
status = fact.get("status", "missing")
if status not in FACT_STATUSES:
status = "manual_review"
return {
"parameter": fact.get("parameter") or parameter,
"status": status,
"value": fact.get("value", ""),
"source_document": fact.get("source_document", ""),
"source_location": fact.get("source_location", ""),
"raw_context": fact.get("raw_context", ""),
"comment": fact.get("comment", ""),
}
def normalize_tender_data(raw_data: dict) -> dict:
raw_sections = raw_data.get("sections", {})
normalized_sections = {}
clarification_questions = list(raw_data.get("critical_questions", []))
for _, section_key in TENDER_DATA_SECTIONS:
facts_by_parameter = {}
for fact in raw_sections.get(section_key, []):
normalized = _normalize_tender_fact(fact, "Прочие значимые данные")
facts_by_parameter.setdefault(normalized["parameter"], []).append(
normalized
)
section_facts = []
for parameter in TENDER_DATA_PARAMETERS[section_key]:
matched_facts = facts_by_parameter.pop(parameter, [])
if matched_facts:
section_facts.extend(matched_facts)
else:
section_facts.append(
_normalize_tender_fact(
{
"status": "missing",
"comment": f"Уточнить у заказчика: {parameter}.",
},
parameter,
)
)
for extra_facts in facts_by_parameter.values():
section_facts.extend(extra_facts)
normalized_sections[section_key] = section_facts
for facts in normalized_sections.values():
for fact in facts:
if fact["status"] in {"missing", "insufficient", "conflict"}:
question = fact["comment"] or fact["parameter"]
if question not in clarification_questions:
clarification_questions.append(question)
return {
"sections": normalized_sections,
"critical_questions": clarification_questions,
"documents_summary": raw_data.get("documents_summary", []),
}
def create_tender_data_report(data: dict, temp_dir: str, task_id: str) -> str:
wb = openpyxl.Workbook()
header_font = Font(bold=True, color="FFFFFF")
header_fill = PatternFill("solid", fgColor="4F81BD")
status_fills = {
"found": PatternFill("solid", fgColor="C6EFCE"),
"missing": PatternFill("solid", fgColor="FFC7CE"),
"insufficient": PatternFill("solid", fgColor="FFEB9C"),
"conflict": PatternFill("solid", fgColor="F4B183"),
"manual_review": PatternFill("solid", fgColor="D9EAD3"),
}
columns = [
"Параметр",
"Статус",
"Найденные данные",
"Документ-источник",
"Страница / лист / строка",
"Цитата из документа",
"Комментарий для сметчика",
"Решение сметчика",
]
ws_summary = wb.active
ws_summary.title = "СВОДКА"
ws_summary.append(["Карта исходных данных тендера", ""])
ws_summary.cell(row=1, column=1).font = Font(bold=True, size=14)
ws_summary.append(["Раздел", "Найдено", "Пробелы/уточнения", "Конфликты"])
for column in range(1, 5):
ws_summary.cell(row=2, column=column).font = header_font
ws_summary.cell(row=2, column=column).fill = header_fill
row = 3
total_found = total_gaps = total_conflicts = 0
for section_name, section_key in TENDER_DATA_SECTIONS:
facts = data["sections"][section_key]
found = sum(fact["status"] == "found" for fact in facts)
gaps = sum(fact["status"] in {"missing", "insufficient"} for fact in facts)
conflicts = sum(fact["status"] == "conflict" for fact in facts)
ws_summary.append([section_name, found, gaps, conflicts])
total_found += found
total_gaps += gaps
total_conflicts += conflicts
row += 1
ws_summary.append(["ИТОГО", total_found, total_gaps, total_conflicts])
for column in range(1, 5):
ws_summary.cell(row=row, column=column).font = Font(bold=True)
ws_summary.append([])
ws_summary.append(["Документы", ""])
for document_summary in data.get("documents_summary", []):
ws_summary.append([str(document_summary), ""])
for column, width in enumerate([42, 15, 20, 15], start=1):
ws_summary.column_dimensions[openpyxl.utils.get_column_letter(column)].width = (
width
)
ws_questions = wb.create_sheet(title="КРИТИЧНЫЕ УТОЧНЕНИЯ")
ws_questions.append(["Вопрос или риск", "Действие сметчика"])
for column in range(1, 3):
ws_questions.cell(row=1, column=column).font = header_font
ws_questions.cell(row=1, column=column).fill = PatternFill(
"solid", fgColor="C0504D"
)
questions = data.get("critical_questions", [])
if questions:
for question in questions:
ws_questions.append([str(question), "Запросить уточнение у заказчика"])
else:
ws_questions.append(["Критичных уточнений не выявлено", ""])
ws_questions.column_dimensions["A"].width = 90
ws_questions.column_dimensions["B"].width = 38
ws_request = wb.create_sheet(title="ЗАПРОС ЗАКАЗЧИКУ")
ws_request.append(["№", "Вопрос заказчику", "Причина"])
for column in range(1, 4):
ws_request.cell(row=1, column=column).font = header_font
ws_request.cell(row=1, column=column).fill = header_fill
for index, question in enumerate(questions, start=1):
ws_request.append([index, str(question), "Данные нужны для расчета сметы"])
if not questions:
ws_request.append([1, "Уточнения не требуются", ""])
ws_request.column_dimensions["A"].width = 8
ws_request.column_dimensions["B"].width = 90
ws_request.column_dimensions["C"].width = 38
ws_sources = wb.create_sheet(title="ИСТОЧНИКИ")
ws_sources.append(["Раздел", *columns[:-1]])
for column in range(1, 8):
ws_sources.cell(row=1, column=column).font = header_font
ws_sources.cell(row=1, column=column).fill = header_fill
for section_name, section_key in TENDER_DATA_SECTIONS:
ws_section = wb.create_sheet(title=section_name[:31])
ws_section.append(columns)
for column in range(1, len(columns) + 1):
ws_section.cell(row=1, column=column).font = header_font
ws_section.cell(row=1, column=column).fill = header_fill
for fact in data["sections"][section_key]:
values = [
fact["parameter"],
FACT_STATUSES[fact["status"]],
fact["value"],
fact["source_document"],
fact["source_location"],
fact["raw_context"],
fact["comment"],
"",
]
ws_section.append(values)
current_row = ws_section.max_row
ws_section.cell(row=current_row, column=2).fill = status_fills[
fact["status"]
]
ws_sources.append([section_name, *values[:-1]])
for column, width in enumerate([34, 24, 48, 28, 24, 65, 50, 30], start=1):
ws_section.column_dimensions[
openpyxl.utils.get_column_letter(column)
].width = width
ws_section.freeze_panes = "A2"
ws_section.auto_filter.ref = ws_section.dimensions
for worksheet_row in ws_section.iter_rows(min_row=2):
for cell in worksheet_row:
cell.alignment = Alignment(wrap_text=True, vertical="top")
for column, width in enumerate([28, 34, 24, 48, 28, 24, 65], start=1):
ws_sources.column_dimensions[openpyxl.utils.get_column_letter(column)].width = (
width
)
ws_sources.freeze_panes = "A2"
ws_sources.auto_filter.ref = ws_sources.dimensions
for worksheet in (ws_summary, ws_questions, ws_request, ws_sources):
for worksheet_row in worksheet.iter_rows():
for cell in worksheet_row:
cell.alignment = Alignment(wrap_text=True, vertical="top")
file_path = os.path.join(temp_dir, f"Karta_Dannyh_Tendera_{task_id[:8]}.xlsx")
wb.save(file_path)
return file_path
# --- БЛОК 2.8: Финализация пакета после успешной обработки всех документов ---
def _normalize_extracted_data(raw):
if isinstance(raw, dict):
return raw
if isinstance(raw, str):
try:
parsed = json.loads(raw)
if isinstance(parsed, dict):
return parsed
except json.JSONDecodeError:
return {}
return {}
async def finalize_task_if_ready(task_id: str, mode: str):
docs_resp = await asyncio.to_thread(
lambda: supabase.table("documents")
.select("id, status, document_name, extracted_data")
.eq("task_id", task_id)
.execute()
)
docs = docs_resp.data or []
if not docs:
raise RuntimeError(f"По task_id={task_id} не найдено документов")
statuses = [doc.get("status", "pending") for doc in docs]
if any(status in ("pending", "processing") for status in statuses):
return
if any(status == "error" for status in statuses):
return
extracted_by_doc = []
for doc in docs:
doc_name = doc.get("document_name", "Unknown_Doc")
parsed_data = _normalize_extracted_data(doc.get("extracted_data"))
text = parsed_data.get("parsed_text", "")
if not text:
raise RuntimeError(f"У документа '{doc_name}' отсутствует parsed_text")
extracted_by_doc.append((doc_name, text))
combined_text = ""
for doc_name, text in extracted_by_doc:
combined_text += f"\n\n{'='*30}\nДОКУМЕНТ: {doc_name}\n{'='*30}\n\n{text}"
if mode == "analytics":
final_result = await analyze_text_with_deepseek_analytics(combined_text)
elif mode == "estimate":
raw_tender_data = await extract_tender_data_for_estimator(combined_text)
final_result = normalize_tender_data(raw_tender_data)
final_result["report_type"] = "tender_data_map"
final_result["report_note"] = (
"Карта исходных данных для ручного расчета сметчиком. "
"Финансовый расчет не выполнялся."
)
with tempfile.TemporaryDirectory() as temp_dir:
excel_path = await asyncio.to_thread(
create_tender_data_report, final_result, temp_dir, task_id
)
final_result["excel_download_url"] = await asyncio.to_thread(
upload_excel_to_supabase,
excel_path,
f"Karta_Dannyh_Tendera_{task_id[:8]}.xlsx",
)
else:
raise RuntimeError(f"Неизвестный режим: {mode}")
payload = {
"final_result": final_result,
"stage": "finalized",
}
await asyncio.to_thread(
lambda: supabase.table("documents")
.update({"status": "completed", "extracted_data": payload})
.eq("task_id", task_id)
.execute()
)
# --- ИЗМЕНЕНИЕ: Воркер теперь обрабатывает документы поштучно и ждёт все документы перед расчетом ---
async def process_message(message: aio_pika.IncomingMessage):
async with message.process(requeue=True, reject_on_redelivered=True):
body = json.loads(message.body.decode())
task_id = body.get("task_id")
mode = body.get("mode", "analytics")
if not task_id:
raise RuntimeError("Некорректное сообщение: нет task_id")
try:
docs_resp = await asyncio.to_thread(
lambda: supabase.table("documents")
.select("id, document_url, document_name, status")
.eq("task_id", task_id)
.execute()
)
documents = docs_resp.data or []
if not documents:
raise RuntimeError(f"По task_id={task_id} не найдено документов")
with tempfile.TemporaryDirectory() as temp_dir:
for doc in documents:
doc_id = doc.get("id")
url = doc.get("document_url")
doc_name = doc.get("document_name") or "Unknown_Doc"
if not doc_id or not url:
continue
await asyncio.to_thread(
lambda current_doc_id=doc_id: supabase.table("documents")
.update(
{
"status": "processing",
"extracted_data": {"stage": "downloading"},
}
)
.eq("id", current_doc_id)
.execute()
)
file_path = await download_file(url, temp_dir)
ext = Path(file_path).suffix.lower()
if ext in [".xlsx", ".xls"]:
text = extract_text_from_excel(file_path)
else:
text = await call_smart_parser(file_path)
payload = {
"stage": "parsed",
"parsed_text": text,
"source_file": os.path.basename(file_path),
"document_name": doc_name,
}
await asyncio.to_thread(
lambda current_doc_id=doc_id, current_payload=payload: supabase.table(
"documents"
)
.update(
{
"status": "completed",
"extracted_data": current_payload,
}
)
.eq("id", current_doc_id)
.execute()
)
await finalize_task_if_ready(task_id, mode)
except Exception as e:
error_msg = str(e)
if "ReadTimeout" in error_msg or "timeout" in error_msg.lower():
error_msg = "Пакет документов слишком тяжелый. Сервер парсинга не смог прочитать документы за 15 минут. Вырежи лишние страницы из ТЗ."
await asyncio.to_thread(
lambda: supabase.table("documents")
.update({"status": "error", "extracted_data": {"error": error_msg}})
.eq("task_id", task_id)
.execute()
)
async def main():
conn = await aio_pika.connect_robust(CLOUDAMQP_URL)
async with conn:
channel = await conn.channel()
await channel.set_qos(prefetch_count=1)
queue = await channel.declare_queue("tender_tasks", durable=True)
await queue.consume(process_message)
await asyncio.Future()
if __name__ == "__main__":
asyncio.run(main())