v3.0.0 VM: убрать v2/Redis/чанки, чистый FormData, PostgreSQL
This commit is contained in:
@@ -6,6 +6,3 @@ psycopg2-binary
|
||||
python-dotenv
|
||||
pdfplumber
|
||||
camelot-py>=2.0
|
||||
redis
|
||||
httpx
|
||||
h2
|
||||
|
||||
+1
-12
@@ -4,9 +4,6 @@ from dotenv import load_dotenv
|
||||
from flask import Flask, render_template, request, redirect, url_for, Response, jsonify
|
||||
import json
|
||||
|
||||
import sys as _sys, os as _os
|
||||
_sys.path.insert(0, _os.path.join(_os.path.dirname(__file__), '..'))
|
||||
|
||||
import schema
|
||||
import db
|
||||
import parser as parser_mod
|
||||
@@ -17,9 +14,6 @@ import mimeutil
|
||||
from test_routes import test_bp
|
||||
from upload import upload_bp
|
||||
from api import api_bp
|
||||
from v2.chunk_api import bp as v2_bp
|
||||
import redis_client
|
||||
import redis_worker
|
||||
|
||||
load_dotenv()
|
||||
|
||||
@@ -28,8 +22,6 @@ import re as _re
|
||||
|
||||
# Логи в памяти (быстро, не тормозит ответ)
|
||||
_log_memory = []
|
||||
# Хранилище чанков в памяти
|
||||
_chunks_mem = {}
|
||||
|
||||
def _log(step, detail=""):
|
||||
"""Писать лог в память (мгновенно)."""
|
||||
@@ -83,20 +75,17 @@ class ContractsApp:
|
||||
def __init__(self):
|
||||
self.app = Flask(__name__)
|
||||
schema.ensure_schema()
|
||||
redis_worker.start_worker()
|
||||
self._add_routes()
|
||||
|
||||
def _add_routes(self):
|
||||
self.app.add_url_rule("/", "index", self._index, methods=["GET", "POST"])
|
||||
self.app.add_url_rule("/parse/<cid>", "parse", self._parse, methods=["GET"])
|
||||
self.app.add_url_rule("/chunk_json", "chunk_json", self._chunk_json, methods=["POST"])
|
||||
self.app.add_url_rule("/upload", "upload_json", self._upload_json, methods=["POST"])
|
||||
|
||||
self.app.add_url_rule("/health", "health", self._health)
|
||||
self.app.add_url_rule("/logs", "logs", self._logs)
|
||||
self.app.register_blueprint(test_bp)
|
||||
self.app.register_blueprint(upload_bp)
|
||||
self.app.register_blueprint(api_bp)
|
||||
self.app.register_blueprint(v2_bp)
|
||||
|
||||
def _health(self):
|
||||
return "OK", 200, {"Content-Type": "text/plain"}
|
||||
|
||||
@@ -1,33 +0,0 @@
|
||||
"""redis_client.py — Очередь загрузок через Redis."""
|
||||
|
||||
import os, json, redis
|
||||
|
||||
UPLOAD_QUEUE = "contracts:upload_queue"
|
||||
|
||||
_r = None
|
||||
|
||||
def _get():
|
||||
global _r
|
||||
if _r is None:
|
||||
_r = redis.Redis(
|
||||
host=os.getenv("REDIS_HOST", "redisk8s.f6303a9d-4ab1-4b2f-8aa8-08f933d02f82.svc.cluster.local"),
|
||||
port=int(os.getenv("REDIS_PORT", "6379")),
|
||||
password=os.getenv("REDIS_PASS", "aLITloRefJEiCPqUc2xB"),
|
||||
decode_responses=True,
|
||||
socket_connect_timeout=5,
|
||||
socket_timeout=10,
|
||||
)
|
||||
return _r
|
||||
|
||||
|
||||
def push(task: dict):
|
||||
"""Добавить задание в очередь."""
|
||||
_get().rpush(UPLOAD_QUEUE, json.dumps(task, ensure_ascii=False))
|
||||
|
||||
|
||||
def pop(timeout=5):
|
||||
"""Взять задание из очереди (блокирующий)."""
|
||||
result = _get().blpop(UPLOAD_QUEUE, timeout=timeout)
|
||||
if result:
|
||||
return json.loads(result[1])
|
||||
return None
|
||||
@@ -1,69 +0,0 @@
|
||||
"""redis_worker.py — Фоновый обработчик очереди загрузок."""
|
||||
|
||||
import threading, time, base64, json
|
||||
import db, mimeutil, parser as parser_mod, textify
|
||||
|
||||
def start_worker():
|
||||
"""Запустить фоновый поток обработки очереди."""
|
||||
t = threading.Thread(target=_worker_loop, daemon=True)
|
||||
t.start()
|
||||
|
||||
def _worker_loop():
|
||||
from redis_client import pop as queue_pop
|
||||
while True:
|
||||
try:
|
||||
task = queue_pop(timeout=5)
|
||||
if task:
|
||||
_process(task)
|
||||
except Exception as e:
|
||||
time.sleep(1)
|
||||
|
||||
def _process(task):
|
||||
filename = task.get("filename", "")
|
||||
b64 = task.get("b64", "")
|
||||
mime = task.get("mime", "application/octet-stream")
|
||||
cid = task.get("cid", "")
|
||||
|
||||
if not filename or not b64:
|
||||
return
|
||||
|
||||
# Декодировать и сохранить в БД
|
||||
file_bytes = base64.b64decode(b64)
|
||||
|
||||
conn, err = db.connect()
|
||||
if err:
|
||||
return
|
||||
try:
|
||||
cur = conn.cursor()
|
||||
cur.execute(
|
||||
"INSERT INTO documents (filename, mime_type, original_bytes, original_b64, status) VALUES (%s,%s,%s,%s,'uploaded')",
|
||||
(filename, mime, file_bytes, b64),
|
||||
)
|
||||
cur.execute("SELECT id FROM documents WHERE filename=%s AND status='uploaded' ORDER BY created_at DESC LIMIT 1", (filename,))
|
||||
row = cur.fetchone()
|
||||
doc_id = row[0] if row else None
|
||||
if doc_id:
|
||||
cur.execute(
|
||||
"INSERT INTO supplements (contract_id, document_id, type) VALUES (%s,%s,'initial')",
|
||||
(cid, str(doc_id)),
|
||||
)
|
||||
|
||||
# Парсинг
|
||||
pr = parser_mod.parse(file_bytes, mime)
|
||||
elements = pr.get("elements", [])
|
||||
if mime == "application/zip" and "files" in pr:
|
||||
for zf in pr["files"]:
|
||||
elements.extend(zf.get("elements", []))
|
||||
text = textify.to_text(elements) if elements else ""
|
||||
|
||||
cur.execute(
|
||||
"UPDATE documents SET parsed_text=%s, elements_json=%s, status='parsed' WHERE id=%s",
|
||||
(text, json.dumps(elements, ensure_ascii=False), str(doc_id)),
|
||||
)
|
||||
conn.commit()
|
||||
cur.close()
|
||||
except Exception as e:
|
||||
try: conn.rollback()
|
||||
except: pass
|
||||
finally:
|
||||
db.put_conn(conn)
|
||||
+25
-79
@@ -5,8 +5,6 @@
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>Сверка договоров</title>
|
||||
<link rel="icon" type="image/png" href="{{ url_for('static', filename='favicon.png') }}">
|
||||
<script src="https://cdnjs.cloudflare.com/ajax/libs/spark-md5/3.0.2/spark-md5.min.js"></script>
|
||||
<script src="https://unpkg.com/lucide@latest"></script>
|
||||
<script src="https://cdnjs.cloudflare.com/ajax/libs/jszip/3.10.1/jszip.min.js"></script>
|
||||
<style>
|
||||
:root {
|
||||
@@ -60,7 +58,7 @@
|
||||
<body>
|
||||
<div class="topbar">
|
||||
<img src="{{ url_for('static', filename='nubes-logo.svg') }}" alt="Nubes">
|
||||
<span class="title">Сверка договоров <span style="font-weight:400;color:var(--muted);font-size:12px;">v2.0.8</span></span>
|
||||
<span class="title">Сверка договоров <span style="font-weight:400;color:var(--muted);font-size:12px;">v3.0.0</span></span>
|
||||
</div>
|
||||
|
||||
<div class="content">
|
||||
@@ -156,81 +154,35 @@
|
||||
|
||||
window.removeFile = removeFile;
|
||||
|
||||
// ── Чанковая загрузка v2 (30KB + MD5) ────────────────────
|
||||
// ── Простая загрузка FormData (ВМ, без лимитов) ──────────
|
||||
|
||||
function uploadFileChunked(file, cid, onProgress) {
|
||||
function uploadFile(file, onProgress) {
|
||||
return new Promise(function(resolve, reject) {
|
||||
var reader = new FileReader();
|
||||
reader.onload = function() {
|
||||
var buffer = reader.result;
|
||||
var bytes = new Uint8Array(buffer);
|
||||
// Чанковая конвертация в base64 — apply не тянет >125K аргументов
|
||||
var b64Chunks = [];
|
||||
var B64_CHUNK = 0x8000; // 32768 bytes
|
||||
for (var bi = 0; bi < bytes.length; bi += B64_CHUNK) {
|
||||
b64Chunks.push(String.fromCharCode.apply(null, bytes.subarray(bi, bi + B64_CHUNK)));
|
||||
}
|
||||
var b64 = btoa(b64Chunks.join(''));
|
||||
var checksum = SparkMD5.ArrayBuffer.hash(buffer);
|
||||
var CHUNK = 60 * 1024; // ближе к лимиту 64KB → меньше чанков
|
||||
var totalChunks = Math.ceil(b64.length / CHUNK);
|
||||
var uploadId = 'u' + Date.now().toString(36) + Math.random().toString(36).substr(2, 7);
|
||||
var done = 0;
|
||||
var fd = new FormData();
|
||||
fd.append('files', file);
|
||||
if (contractId) fd.append('cid', contractId);
|
||||
|
||||
console.log('uploadFileChunked: ' + file.name + ' b64=' + b64.length + ' chunks=' + totalChunks + ' uploadId=' + uploadId);
|
||||
|
||||
function fetchWithRetry(url, body, retries) {
|
||||
retries = retries || 3;
|
||||
return fetch(url, {
|
||||
method: 'POST',
|
||||
headers: {'Content-Type': 'application/json'},
|
||||
body: JSON.stringify(body)
|
||||
}).catch(function(e) {
|
||||
if (retries > 0) {
|
||||
console.log('retry ' + url + ' (' + retries + ' left): ' + e.message);
|
||||
return new Promise(function(r) { setTimeout(r, 1000); }).then(function() {
|
||||
return fetchWithRetry(url, body, retries - 1);
|
||||
});
|
||||
var xhr = new XMLHttpRequest();
|
||||
xhr.open('POST', '/');
|
||||
xhr.upload.onprogress = function(e) {
|
||||
if (e.lengthComputable && onProgress) {
|
||||
onProgress(Math.round(e.loaded / e.total * 100));
|
||||
}
|
||||
throw e;
|
||||
});
|
||||
}
|
||||
|
||||
function sendChunk(i) {
|
||||
if (i >= totalChunks) {
|
||||
console.log('finalize: ' + uploadId);
|
||||
fetchWithRetry('/v2/finalize', {upload_id: uploadId, filename: file.name, total: totalChunks, checksum: checksum})
|
||||
.then(function(r) {
|
||||
if (!r.ok) throw new Error('HTTP ' + r.status);
|
||||
return r.json();
|
||||
}).then(function(data) {
|
||||
if (data.ok) resolve({contract_id: data.contract_id});
|
||||
else reject(new Error(data.error));
|
||||
}).catch(function(e) { reject(e); });
|
||||
return;
|
||||
}
|
||||
var piece = b64.substring(i * CHUNK, (i + 1) * CHUNK);
|
||||
console.log('chunk ' + i + '/' + totalChunks + ' send ' + piece.length + ' bytes');
|
||||
fetchWithRetry('/v2/chunk', {upload_id: uploadId, index: i, total: totalChunks, data: piece})
|
||||
.then(function(r) {
|
||||
if (!r.ok) throw new Error('HTTP ' + r.status);
|
||||
return r.json();
|
||||
}).then(function(data) {
|
||||
if (!data.ok) { reject(new Error(data.error)); return; }
|
||||
done++;
|
||||
if (onProgress) onProgress({pct: Math.round(done / totalChunks * 100), done: done, total: totalChunks});
|
||||
setTimeout(function() { sendChunk(i + 1); }, 3000);
|
||||
}).catch(function(e) { reject(e); });
|
||||
}
|
||||
sendChunk(0);
|
||||
};
|
||||
reader.onerror = function() { reject(new Error('Read error')); };
|
||||
reader.readAsArrayBuffer(file);
|
||||
xhr.onload = function() {
|
||||
try {
|
||||
var r = JSON.parse(xhr.responseText);
|
||||
if (r.contract_id) resolve(r);
|
||||
else reject(new Error(r.error || 'Неизвестная ошибка'));
|
||||
} catch(e) { reject(new Error('Некорректный ответ')); }
|
||||
};
|
||||
xhr.onerror = function() { reject(new Error('Сеть')); };
|
||||
xhr.ontimeout = function() { reject(new Error('Таймаут')); };
|
||||
xhr.timeout = 120000;
|
||||
xhr.send(fd);
|
||||
});
|
||||
}
|
||||
|
||||
// ── Старая загрузка (одним POST) ─────────────────────────────
|
||||
|
||||
// ── Выбор файла → сразу загрузка ────────────────────────────
|
||||
|
||||
fileInput.addEventListener('change', async function() {
|
||||
@@ -252,21 +204,15 @@
|
||||
var startTime = Date.now();
|
||||
|
||||
try {
|
||||
var resp = await uploadFileChunked(f, contractId, function(info) {
|
||||
fileQueue[rowIdx].status = '<span style="color:var(--muted);">↑ чанк ' + info.done + '/' + info.total + ' (' + info.pct + '%)</span>';
|
||||
var resp = await uploadFile(f, function(pct) {
|
||||
fileQueue[rowIdx].status = '<span style="color:var(--muted);">↑ ' + pct + '%</span>';
|
||||
renderTable();
|
||||
});
|
||||
// uploadFileChunked10k: прогресс по чанкам
|
||||
// но этот же resp приходит только с последнего чанка
|
||||
if (resp && resp.contract_id) {
|
||||
contractId = resp.contract_id;
|
||||
}
|
||||
// Если файл был один и он маленький (один чанк), contract_id уже есть
|
||||
if (resp && resp.contract_id) contractId = resp.contract_id;
|
||||
var elapsed = ((Date.now() - startTime) / 1000).toFixed(1);
|
||||
fileQueue[rowIdx].status = '<span class="status-ok">✓ ' + elapsed + 'с</span>';
|
||||
fileQueue[rowIdx].uploaded = true;
|
||||
renderTable();
|
||||
|
||||
if (f.type === 'application/zip' || f.name.toLowerCase().endsWith('.zip')) {
|
||||
await showZipContents(f);
|
||||
}
|
||||
|
||||
@@ -1,75 +0,0 @@
|
||||
"""v2/chunk_api.py — Приём чанков по 30KB, финализация."""
|
||||
|
||||
import sys, os, base64, hashlib, logging
|
||||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), '..', 'site'))
|
||||
import db, mimeutil
|
||||
sys.path.insert(0, os.path.dirname(__file__))
|
||||
import init as v2
|
||||
|
||||
from flask import Blueprint, request, jsonify
|
||||
|
||||
bp = Blueprint("v2", __name__, url_prefix="/v2")
|
||||
|
||||
@bp.route("/chunk", methods=["POST"])
|
||||
def receive_chunk():
|
||||
body = request.get_json(silent=True) or {}
|
||||
if not body.get("upload_id") or not body.get("data"):
|
||||
return jsonify({"ok": False, "error": "missing fields"}), 400
|
||||
try:
|
||||
v2.store_chunk(body["upload_id"], int(body.get("index", 0)), body["data"])
|
||||
return jsonify({"ok": True, "received": int(body.get("index", 0))})
|
||||
except Exception as e:
|
||||
return jsonify({"ok": False, "error": str(e)}), 500
|
||||
|
||||
@bp.route("/finalize", methods=["POST"])
|
||||
def finalize_upload():
|
||||
body = request.get_json(silent=True) or {}
|
||||
upload_id = body.get("upload_id", "")
|
||||
filename = body.get("filename", "")
|
||||
total = int(body.get("total", 0))
|
||||
checksum = body.get("checksum", "")
|
||||
|
||||
if not all([upload_id, filename, total]):
|
||||
return jsonify({"ok": False, "error": "missing fields"}), 400
|
||||
|
||||
chunks = v2.get_chunks(upload_id)
|
||||
if len(chunks) != total:
|
||||
return jsonify({"ok": False, "error": f"expected {total}, got {len(chunks)}"}), 409
|
||||
|
||||
try:
|
||||
file_bytes = b"".join(base64.b64decode(c) for c in chunks)
|
||||
except Exception as e:
|
||||
return jsonify({"ok": False, "error": f"decode: {e}"}), 422
|
||||
|
||||
if checksum and hashlib.md5(file_bytes).hexdigest() != checksum:
|
||||
return jsonify({"ok": False, "error": "checksum mismatch"}), 422
|
||||
|
||||
mime = mimeutil.guess_mime(filename) or "application/octet-stream"
|
||||
|
||||
# Сохранить в БД
|
||||
conn, err = db.connect()
|
||||
if err:
|
||||
return jsonify({"ok": False, "error": f"db: {err}"}), 500
|
||||
try:
|
||||
cur = conn.cursor()
|
||||
cur.execute(
|
||||
"INSERT INTO contracts (number, client) VALUES (%s,%s) RETURNING id",
|
||||
("б/н " + __import__("datetime").datetime.now().strftime("%Y%m%d-%H%M"), ""),
|
||||
)
|
||||
cid = cur.fetchone()[0]
|
||||
cur.execute(
|
||||
"INSERT INTO documents (filename, mime_type, original_bytes, status) VALUES (%s,%s,%s,'uploaded')",
|
||||
(filename, mime, file_bytes),
|
||||
)
|
||||
cur.execute("SELECT id FROM documents WHERE filename=%s ORDER BY created_at DESC LIMIT 1", (filename,))
|
||||
doc_id = cur.fetchone()[0]
|
||||
cur.execute("INSERT INTO supplements (contract_id, document_id, type) VALUES (%s,%s,'initial')", (str(cid), str(doc_id)))
|
||||
conn.commit()
|
||||
cur.close()
|
||||
finally:
|
||||
db.put_conn(conn)
|
||||
|
||||
v2.delete_chunks(upload_id)
|
||||
v2.push_finalize({"doc_id": str(doc_id), "filename": filename, "mime_type": mime})
|
||||
|
||||
return jsonify({"ok": True, "contract_id": str(cid), "doc_id": str(doc_id)})
|
||||
@@ -1,213 +0,0 @@
|
||||
/**
|
||||
* v2/chunk_upload.js — Клиентская загрузка файла чанками по 32KB.
|
||||
*
|
||||
* Использование:
|
||||
* <input type="file" id="file-input">
|
||||
* <div id="upload-ui"></div>
|
||||
* <script src="/static/v2/chunk_upload.js"></script>
|
||||
*
|
||||
* После успешной загрузки появляется кнопка «Парсинг».
|
||||
*/
|
||||
|
||||
(function () {
|
||||
"use strict";
|
||||
|
||||
const CHUNK_SIZE = 32 * 1024; // 32 KB
|
||||
|
||||
// ── MD5 (встроенная реализация без зависимостей) ──────────────────────────
|
||||
// RFC 1321 — используется только для checksum, не для безопасности.
|
||||
|
||||
function md5(buffer) {
|
||||
// SparkMD5 должен быть подключён отдельно, либо используем SubtleCrypto
|
||||
// Здесь используем встроенный Web Crypto API (SHA-256 недостаточен — нужен MD5)
|
||||
// Простая MD5 реализация:
|
||||
return SparkMD5.ArrayBuffer.hash(buffer);
|
||||
}
|
||||
|
||||
// ── UI ────────────────────────────────────────────────────────────────────
|
||||
|
||||
function getUI() {
|
||||
return document.getElementById("upload-ui");
|
||||
}
|
||||
|
||||
function setStatus(html) {
|
||||
const ui = getUI();
|
||||
if (ui) ui.innerHTML = html;
|
||||
}
|
||||
|
||||
function renderProgress(filename, pct, elapsed) {
|
||||
return `
|
||||
<div class="upload-progress">
|
||||
<strong>${escHtml(filename)}</strong><br>
|
||||
<progress value="${pct}" max="100" style="width:100%"></progress>
|
||||
<span>${pct}% | ${elapsed}с</span>
|
||||
</div>`;
|
||||
}
|
||||
|
||||
function escHtml(s) {
|
||||
return String(s)
|
||||
.replace(/&/g, "&")
|
||||
.replace(/</g, "<")
|
||||
.replace(/>/g, ">");
|
||||
}
|
||||
|
||||
// ── XHR helper ────────────────────────────────────────────────────────────
|
||||
|
||||
function post(url, body) {
|
||||
return new Promise(function (resolve, reject) {
|
||||
const xhr = new XMLHttpRequest();
|
||||
xhr.open("POST", url, true);
|
||||
xhr.setRequestHeader("Content-Type", "application/json");
|
||||
xhr.timeout = 15000;
|
||||
xhr.onload = function () {
|
||||
if (xhr.status >= 200 && xhr.status < 300) {
|
||||
try { resolve(JSON.parse(xhr.responseText)); }
|
||||
catch (e) { reject(new Error("bad json: " + xhr.responseText)); }
|
||||
} else {
|
||||
reject(new Error("HTTP " + xhr.status + ": " + xhr.responseText));
|
||||
}
|
||||
};
|
||||
xhr.onerror = function () { reject(new Error("network error")); };
|
||||
xhr.ontimeout = function () { reject(new Error("timeout")); };
|
||||
xhr.send(JSON.stringify(body));
|
||||
});
|
||||
}
|
||||
|
||||
// ── Base64 encode ArrayBuffer ─────────────────────────────────────────────
|
||||
|
||||
function toBase64(buffer) {
|
||||
let binary = "";
|
||||
const bytes = new Uint8Array(buffer);
|
||||
for (let i = 0; i < bytes.byteLength; i++) {
|
||||
binary += String.fromCharCode(bytes[i]);
|
||||
}
|
||||
return btoa(binary);
|
||||
}
|
||||
|
||||
// ── Генерация upload_id ───────────────────────────────────────────────────
|
||||
|
||||
function genId() {
|
||||
return "u" + Date.now().toString(36) + Math.random().toString(36).slice(2, 7);
|
||||
}
|
||||
|
||||
// ── Основная функция загрузки ─────────────────────────────────────────────
|
||||
|
||||
async function uploadFile(file) {
|
||||
const uploadId = genId();
|
||||
const total = Math.ceil(file.size / CHUNK_SIZE);
|
||||
const startTs = Date.now();
|
||||
|
||||
// Считаем MD5 всего файла через SparkMD5
|
||||
const fullBuffer = await file.arrayBuffer();
|
||||
const checksum = md5(fullBuffer);
|
||||
|
||||
setStatus(renderProgress(file.name, 0, 0));
|
||||
|
||||
for (let i = 0; i < total; i++) {
|
||||
const start = i * CHUNK_SIZE;
|
||||
const end = Math.min(start + CHUNK_SIZE, file.size);
|
||||
const slice = fullBuffer.slice(start, end);
|
||||
const b64 = toBase64(slice);
|
||||
const elapsed = Math.round((Date.now() - startTs) / 1000);
|
||||
const pct = Math.round(((i + 1) / total) * 90); // до 90%, финал = 100%
|
||||
|
||||
setStatus(renderProgress(file.name, pct, elapsed));
|
||||
|
||||
let res;
|
||||
try {
|
||||
res = await post("/v2/chunk", {
|
||||
upload_id: uploadId,
|
||||
index: i,
|
||||
total: total,
|
||||
data: b64,
|
||||
});
|
||||
} catch (e) {
|
||||
setStatus(`<p class="error">Ошибка чанка ${i}: ${escHtml(e.message)}</p>`);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!res.ok) {
|
||||
setStatus(`<p class="error">Сервер отклонил чанк ${i}: ${escHtml(res.error)}</p>`);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
// Финализация
|
||||
const elapsed = Math.round((Date.now() - startTs) / 1000);
|
||||
setStatus(renderProgress(file.name, 95, elapsed));
|
||||
|
||||
let fin;
|
||||
try {
|
||||
fin = await post("/v2/finalize", {
|
||||
upload_id: uploadId,
|
||||
filename: file.name,
|
||||
mime_type: file.type || "application/octet-stream",
|
||||
total: total,
|
||||
checksum: checksum,
|
||||
});
|
||||
} catch (e) {
|
||||
setStatus(`<p class="error">Ошибка финализации: ${escHtml(e.message)}</p>`);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!fin.ok) {
|
||||
setStatus(`<p class="error">Финализация не удалась: ${escHtml(fin.error)}</p>`);
|
||||
return;
|
||||
}
|
||||
|
||||
const docId = fin.doc_id;
|
||||
const done = Math.round((Date.now() - startTs) / 1000);
|
||||
setStatus(`
|
||||
<div class="upload-done">
|
||||
<strong>${escHtml(file.name)}</strong> загружен за ${done}с.<br>
|
||||
<button id="btn-parse" data-docid="${escHtml(docId)}">Парсинг</button>
|
||||
<div id="parse-status"></div>
|
||||
</div>`);
|
||||
|
||||
document.getElementById("btn-parse").addEventListener("click", function () {
|
||||
pollStatus(docId);
|
||||
});
|
||||
}
|
||||
|
||||
// ── Опрос статуса парсинга ────────────────────────────────────────────────
|
||||
|
||||
function pollStatus(docId) {
|
||||
const ps = document.getElementById("parse-status");
|
||||
if (ps) ps.textContent = "Парсинг запущен, ожидаем…";
|
||||
|
||||
let attempts = 0;
|
||||
const timer = setInterval(async function () {
|
||||
attempts++;
|
||||
try {
|
||||
const r = await fetch("/v2/status/" + docId);
|
||||
const d = await r.json();
|
||||
if (d.status === "parsed") {
|
||||
clearInterval(timer);
|
||||
if (ps) ps.innerHTML = `<span style="color:green">Готово: ${escHtml(d.filename)}</span>`;
|
||||
} else if (d.status === "error") {
|
||||
clearInterval(timer);
|
||||
if (ps) ps.innerHTML = `<span style="color:red">Ошибка: ${escHtml(d.error || "unknown")}</span>`;
|
||||
} else {
|
||||
if (ps) ps.textContent = `Статус: ${d.status} (${attempts * 2}с)`;
|
||||
}
|
||||
} catch (e) {
|
||||
if (ps) ps.textContent = `Ошибка опроса: ${e.message}`;
|
||||
}
|
||||
if (attempts > 60) { clearInterval(timer); }
|
||||
}, 2000);
|
||||
}
|
||||
|
||||
// ── Инициализация ─────────────────────────────────────────────────────────
|
||||
|
||||
document.addEventListener("DOMContentLoaded", function () {
|
||||
const input = document.getElementById("file-input");
|
||||
if (!input) return;
|
||||
input.addEventListener("change", function () {
|
||||
const file = input.files[0];
|
||||
if (!file) return;
|
||||
uploadFile(file).catch(function (e) {
|
||||
setStatus(`<p class="error">Критическая ошибка: ${escHtml(e.message)}</p>`);
|
||||
});
|
||||
});
|
||||
});
|
||||
})();
|
||||
-45
@@ -1,45 +0,0 @@
|
||||
"""v2/init.py — Redis и хранилище чанков. БД — через site/db.py."""
|
||||
|
||||
import os, json, logging, redis
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
CHUNK_STORE_KEY = "v2:chunks:"
|
||||
FINALIZE_QUEUE = "v2:finalize"
|
||||
|
||||
_r = None
|
||||
|
||||
def get_redis():
|
||||
global _r
|
||||
if _r is None:
|
||||
_r = redis.Redis(
|
||||
host=os.getenv("REDIS_HOST", "redisk8s.f6303a9d-4ab1-4b2f-8aa8-08f933d02f82.svc.cluster.local"),
|
||||
port=int(os.getenv("REDIS_PORT", "6379")),
|
||||
password=os.getenv("REDIS_PASS", "aLITloRefJEiCPqUc2xB"),
|
||||
decode_responses=True,
|
||||
socket_connect_timeout=5,
|
||||
socket_timeout=10,
|
||||
)
|
||||
return _r
|
||||
|
||||
def store_chunk(upload_id, index, data_b64):
|
||||
get_redis().rpush(f"{CHUNK_STORE_KEY}{upload_id}", json.dumps([index, data_b64]))
|
||||
|
||||
def get_chunks(upload_id):
|
||||
raw = get_redis().lrange(f"{CHUNK_STORE_KEY}{upload_id}", 0, -1)
|
||||
chunks = []
|
||||
for r in raw:
|
||||
idx, data = json.loads(r)
|
||||
chunks.append((idx, data))
|
||||
chunks.sort()
|
||||
return [c[1] for c in chunks]
|
||||
|
||||
def delete_chunks(upload_id):
|
||||
get_redis().delete(f"{CHUNK_STORE_KEY}{upload_id}")
|
||||
|
||||
def push_finalize(task):
|
||||
get_redis().rpush(FINALIZE_QUEUE, json.dumps(task, ensure_ascii=False))
|
||||
|
||||
def pop_finalize(timeout=5):
|
||||
result = get_redis().blpop(FINALIZE_QUEUE, timeout=timeout)
|
||||
return json.loads(result[1]) if result else None
|
||||
@@ -1,51 +0,0 @@
|
||||
"""v2/worker.py — Фоновый парсинг через site/parser.py."""
|
||||
|
||||
import sys, os, json, logging, time
|
||||
sys.path.insert(0, os.path.join(os.path.dirname(__file__), '..', 'site'))
|
||||
import db, parser as parser_mod, textify
|
||||
sys.path.insert(0, os.path.dirname(__file__))
|
||||
import init as v2
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="[worker] %(message)s")
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
def run():
|
||||
logger.info("started")
|
||||
while True:
|
||||
try:
|
||||
task = v2.pop_finalize(timeout=5)
|
||||
if task:
|
||||
_process(task)
|
||||
except Exception as e:
|
||||
logger.error("loop: %s", e)
|
||||
time.sleep(1)
|
||||
|
||||
def _process(task):
|
||||
doc_id = task.get("doc_id")
|
||||
filename = task.get("filename", "")
|
||||
mime = task.get("mime_type", "")
|
||||
|
||||
doc, err = db.query_one("SELECT original_bytes, mime_type FROM documents WHERE id=%s", (doc_id,))
|
||||
if not doc:
|
||||
logger.info("%s: not found", doc_id)
|
||||
return
|
||||
|
||||
raw = doc["original_bytes"]
|
||||
file_bytes = bytes(raw) if isinstance(raw, memoryview) else raw
|
||||
mime = doc["mime_type"] or mime
|
||||
|
||||
pr = parser_mod.parse(file_bytes, mime)
|
||||
elements = pr.get("elements", [])
|
||||
if mime == "application/zip" and "files" in pr:
|
||||
for zf in pr["files"]:
|
||||
elements.extend(zf.get("elements", []))
|
||||
|
||||
text = textify.to_text(elements) if elements else ""
|
||||
db.execute(
|
||||
"UPDATE documents SET parsed_text=%s, elements_json=%s, status='parsed' WHERE id=%s",
|
||||
(text, json.dumps(elements, ensure_ascii=False), doc_id),
|
||||
)
|
||||
logger.info("%s: parsed, %d elements", filename, len(elements))
|
||||
|
||||
if __name__ == "__main__":
|
||||
run()
|
||||
Reference in New Issue
Block a user