Give each help file unique sequential extraction codes, show the count, and allow wiping all extractions for a file while testing.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-09-10 13:26:21 +03:00
parent 170c551d6e
commit 3a30959ae5
5 changed files with 761 additions and 39 deletions

View File

@@ -38,6 +38,16 @@ import anthropic
from docx import Document
from bs4 import BeautifulSoup
from help_codes import (
WIPED_HASH,
file_result as _file_result,
make_code,
parse_code,
remove_section_outputs,
source_basename,
source_identity,
)
try:
import pdfplumber
HAS_PDF = True
@@ -204,38 +214,200 @@ class Database:
CREATE INDEX IF NOT EXISTS ix_rip_help_sections_prefix
ON rip_help_sections(prefix)
""")
cur.execute("""
CREATE INDEX IF NOT EXISTS ix_rip_help_sections_source
ON rip_help_sections(prefix, source_file)
""")
cur.execute(
"ALTER TABLE rip_help_files ADD COLUMN IF NOT EXISTS file_index INTEGER"
)
self.conn.commit()
log.info("Схемата е проверена / създадена.")
def matching_source_paths(self, prefix: str, identity: str) -> list[str]:
"""Всички file_path/source_file за същия help файл (basename, вкл. стари temp пътища)."""
name = source_basename(identity).lower()
if not name:
return []
found: list[str] = []
seen: set[str] = set()
for p in self.all_source_files(prefix):
if source_basename(p).lower() == name and p not in seen:
seen.add(p)
found.append(p)
return found
def get_file_hash(self, prefix: str, file_path: str) -> Optional[str]:
paths = self.matching_source_paths(prefix, file_path) or [file_path]
cur = self.conn.cursor()
cur.execute(
"SELECT file_hash FROM rip_help_files WHERE prefix=%s AND file_path=%s",
(prefix, file_path)
"SELECT file_hash FROM rip_help_files "
"WHERE prefix=%s AND file_path = ANY(%s)",
(prefix, list(paths)),
)
row = cur.fetchone()
return row[0] if row else None
for (h,) in cur.fetchall():
if not h:
continue
val = str(h).strip()
if val and val != WIPED_HASH:
return val
return None
def upsert_file(self, prefix: str, file_path: str, file_hash: str, section_count: int):
def upsert_file(
self,
prefix: str,
file_path: str,
file_hash: str,
section_count: int,
file_index: Optional[int] = None,
):
canonical = source_basename(file_path) or file_path
stale = [p for p in self.matching_source_paths(prefix, canonical) if p != canonical]
cur = self.conn.cursor()
if stale:
cur.execute(
"DELETE FROM rip_help_files WHERE prefix=%s AND file_path = ANY(%s)",
(prefix, stale),
)
cur.execute("""
INSERT INTO rip_help_files (prefix, file_path, file_hash, section_count)
VALUES (%s, %s, %s, %s)
INSERT INTO rip_help_files (prefix, file_path, file_hash, section_count, file_index)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (prefix, file_path) DO UPDATE SET
file_hash = EXCLUDED.file_hash,
section_count= EXCLUDED.section_count,
processed_at = NOW()
""", (prefix, file_path, file_hash, section_count))
processed_at = NOW(),
file_index = COALESCE(EXCLUDED.file_index, rip_help_files.file_index)
""", (prefix, canonical, file_hash, section_count, file_index))
self.conn.commit()
def delete_sections_for_file(self, prefix: str, file_path: str):
paths = self.matching_source_paths(prefix, file_path) or [file_path]
cur = self.conn.cursor()
cur.execute(
"DELETE FROM rip_help_sections WHERE prefix=%s AND source_file=%s",
(prefix, file_path)
"DELETE FROM rip_help_sections WHERE prefix=%s AND source_file = ANY(%s)",
(prefix, list(paths)),
)
self.conn.commit()
def sections_for_file(self, prefix: str, identity: str) -> list[tuple[str, str, Optional[str]]]:
"""(code, source_file, output_path) за всички извлечения на файла."""
paths = self.matching_source_paths(prefix, identity)
if not paths:
return []
cur = self.conn.cursor()
cur.execute(
"SELECT code, source_file, output_path FROM rip_help_sections "
"WHERE prefix=%s AND source_file = ANY(%s) ORDER BY code",
(prefix, list(paths)),
)
return [(r[0], r[1], r[2]) for r in cur.fetchall()]
def file_index_for(self, prefix: str, identity: str) -> Optional[int]:
paths = self.matching_source_paths(prefix, identity)
if not paths:
return None
cur = self.conn.cursor()
cur.execute(
"SELECT file_index FROM rip_help_files "
"WHERE prefix=%s AND file_path = ANY(%s) AND file_index IS NOT NULL",
(prefix, list(paths)),
)
from_col = [r[0] for r in cur.fetchall() if r[0]]
if from_col:
return min(from_col)
cur.execute(
"SELECT code FROM rip_help_sections WHERE prefix=%s AND source_file = ANY(%s)",
(prefix, list(paths)),
)
found: list[int] = []
for (code,) in cur.fetchall():
parsed = parse_code(code)
if parsed:
found.append(parsed[1])
if not found:
return None
tally: dict[int, int] = {}
for idx in found:
tally[idx] = tally.get(idx, 0) + 1
return max(tally, key=lambda k: (tally[k], -k))
def max_file_index(self, prefix: str) -> int:
cur = self.conn.cursor()
cur.execute(
"SELECT COALESCE(MAX(file_index), 0) FROM rip_help_files WHERE prefix=%s",
(prefix,),
)
m1 = int(cur.fetchone()[0] or 0)
cur.execute("SELECT code FROM rip_help_sections WHERE prefix=%s", (prefix,))
m2 = 0
for (code,) in cur.fetchall():
parsed = parse_code(code)
if parsed:
m2 = max(m2, parsed[1])
return max(m1, m2)
def file_index_used_by_others(self, prefix: str, identity: str, idx: int) -> bool:
name = source_basename(identity).lower()
cur = self.conn.cursor()
cur.execute(
"SELECT code, source_file FROM rip_help_sections WHERE prefix=%s",
(prefix,),
)
for code, src in cur.fetchall():
parsed = parse_code(code)
if parsed and parsed[1] == idx and source_basename(src).lower() != name:
return True
return False
def allocate_file_index(self, prefix: str, identity: str) -> int:
existing = self.file_index_for(prefix, identity)
if existing and not self.file_index_used_by_others(prefix, identity, existing):
return existing
return self.max_file_index(prefix) + 1
def wipe_extractions_for_file(
self,
prefix: str,
identity: str,
output_dir: Optional[Path] = None,
) -> dict:
"""Изтрива всички извлечения за файла. Запазва file_index; следващият scan почва от SEC_0001."""
rows = self.sections_for_file(prefix, identity)
codes = [r[0] for r in rows]
file_index = self.file_index_for(prefix, identity)
if file_index and self.file_index_used_by_others(prefix, identity, file_index):
file_index = None
paths = self.matching_source_paths(prefix, identity)
if output_dir:
remove_section_outputs(output_dir, codes, [r[2] for r in rows if r[2]])
self.delete_sections_for_file(prefix, identity)
canonical = source_basename(identity) or identity
stale = list(paths) if paths else []
cur = self.conn.cursor()
if stale:
cur.execute(
"DELETE FROM rip_help_files WHERE prefix=%s AND file_path = ANY(%s)",
(prefix, stale),
)
if file_index:
cur.execute("""
INSERT INTO rip_help_files (prefix, file_path, file_hash, section_count, file_index)
VALUES (%s, %s, %s, 0, %s)
ON CONFLICT (prefix, file_path) DO UPDATE SET
file_hash = EXCLUDED.file_hash,
section_count = 0,
processed_at = NOW(),
file_index = EXCLUDED.file_index
""", (prefix, canonical, WIPED_HASH, file_index))
self.conn.commit()
return {
"file": canonical,
"prefix": prefix,
"deleted": len(codes),
"codes": codes,
"file_index": file_index,
}
def all_source_files(self, prefix: str) -> list[str]:
"""Връща всички source_file пътища за даден префикс."""
cur = self.conn.cursor()
@@ -968,12 +1140,9 @@ def classify_section(client: anthropic.Anthropic, title: str, text: str) -> tupl
# ──────────────────────────────────────────────
# Генериране на кодове
# Генериране на кодове — help_codes.py
# ──────────────────────────────────────────────
def make_code(prefix: str, file_index: int, sec_index: int) -> str:
return f"{prefix}_{file_index:04d}_SEC_{sec_index:04d}"
# ──────────────────────────────────────────────
# Основна обработка
@@ -996,46 +1165,56 @@ def process_file(
prefix: str = "HLP",
force: bool = False,
remote_root: Optional[str] = None,
) -> int:
"""Обработва един файл. Връща броя записани секции (0 = пропуснат)."""
rel = str(path)
source_key: Optional[str] = None,
) -> dict:
"""Обработва един файл. Връща статистика (saved, codes, ok)."""
rel = source_key or path.name
fh = file_hash(path)
existing_rows = db.sections_for_file(prefix, rel)
existing_codes = [r[0] for r in existing_rows]
if not force:
stored = db.get_file_hash(prefix, rel)
if stored == fh:
log.info(f" [SKIP] {path.name} (непроменен)")
return 0
return _file_result(rel, file_index, existing_codes, saved=0, skipped=True)
log.info(f" [PROC] {path.name}")
ext = path.suffix.lower()
parser = PARSERS.get(ext)
if not parser:
log.warning(f" Неподдържан формат: {ext}")
return 0
return _file_result(rel, file_index, existing_codes, saved=0)
try:
sections = parser(path)
except Exception as e:
log.error(f" Грешка при парсване: {e}")
return 0
return _file_result(rel, file_index, existing_codes, saved=0)
sections = merge_short_sections(sections)
# Изтриваме старите секции за файла при повторна обработка
remove_section_outputs(
output_dir,
existing_codes,
[r[2] for r in existing_rows if r[2]],
)
db.delete_sections_for_file(prefix, rel)
images_dir = output_dir / "images"
images_dir.mkdir(parents=True, exist_ok=True)
saved = 0
for i, sec in enumerate(sections, 1):
codes: list[str] = []
for sec in sections:
text = clean_text(sec.text)
html_text = sec.html_text or ""
if not text and not sec.images and not html_text:
continue
code = make_code(prefix, file_index, i)
sec_index = saved + 1
code = make_code(prefix, file_index, sec_index)
# Записваме картинките на диск и заменяме placeholder-ите в текста + HTML
image_rel_paths: list[str] = []
@@ -1070,7 +1249,7 @@ def process_file(
title, keywords = classify_section(client, sec.title, text)
except Exception as e:
log.warning(f" AI грешка за {code}: {e}")
title, keywords = sec.title or f"Секция {i}", ""
title, keywords = sec.title or f"Секция {sec_index}", ""
images_json = json.dumps(image_rel_paths, ensure_ascii=False)
ps = ProcessedSection(
@@ -1094,11 +1273,12 @@ def process_file(
db.insert_section(prefix, ps, _db_output_path(out_path, code, remote_root))
saved += 1
codes.append(code)
log.debug(f" {code}: {title[:60]} ({len(image_rel_paths)} img)")
db.upsert_file(prefix, rel, fh, saved)
db.upsert_file(prefix, rel, fh, saved, file_index=file_index)
log.info(f" → {saved} секции записани")
return saved
return _file_result(rel, file_index, codes, saved=saved)
_PREFIX_RE = re.compile(r"^[A-Za-z][A-Za-z0-9_]{0,49}$")
@@ -1141,19 +1321,28 @@ def process_directory(
if remote_root:
log.info(f"DB output_path root: {remote_root.replace(chr(92), '/').rstrip('/')}")
current_paths = {str(p) for p in files}
current_ids = {source_identity(p, input_dir) for p in files}
current_names = {source_basename(i).lower() for i in current_ids}
file_results: list[dict] = []
total_sections = 0
try:
for idx, path in enumerate(sorted(files), 1):
n = process_file(
for path in sorted(files):
identity = source_identity(path, input_dir)
idx = db.allocate_file_index(prefix, identity)
info = process_file(
path, idx, db, client, output_dir,
prefix=prefix, force=force, remote_root=remote_root,
source_key=identity,
)
total_sections += n
file_results.append(info)
total_sections += int(info.get("saved") or 0)
if purge_missing:
existing = set(db.all_source_files(prefix))
orphans = sorted(existing - current_paths)
orphans = sorted(
e for e in existing
if source_basename(e).lower() not in current_names
)
if not orphans:
log.info(f"Purge: няма orphan записи в БД за prefix={prefix}.")
else:
@@ -1191,6 +1380,7 @@ def process_directory(
db.close()
log.info(f"Готово. Prefix={prefix}. Общо нови/обновени секции: {total_sections}")
return {"sections": total_sections, "files": file_results}
# ──────────────────────────────────────────────