Files
creator3/backend/laden.py
2026-07-24 11:23:18 +02:00

206 lines
8.7 KiB
Python

"""Ingestion (R6): Bytes → trafilatura(markdown) → ftfy NFC → Snapshot.
Normalisierung passiert GENAU EINMAL, danach ist der Snapshot unveränderlich —
alle Verbatim-Zitate stammen ausschließlich aus dem Snapshot."""
import asyncio
import hashlib
import os
import re
import time
from urllib.parse import urlsplit
import httpx
from . import config, db, engine, graph, llm, textkit
_letzter_fetch: dict[str, float] = {}
def _text_reparieren(text: str) -> str:
text = re.sub(r"[\x00-\x08\x0b\x0c\x0e-\x1f]", " ", text) # VOR ftfy
try:
import ftfy
return ftfy.fix_text(text, normalization="NFC")
except ImportError:
import unicodedata
return unicodedata.normalize("NFC", text)
async def _domain_bremse(url: str) -> None:
domain = urlsplit(url).netloc
delta = time.monotonic() - _letzter_fetch.get(domain, 0)
if delta < config.DOMAIN_RATE_S:
await asyncio.sleep(config.DOMAIN_RATE_S - delta)
_letzter_fetch[domain] = time.monotonic()
_WIKI_RE = re.compile(r"https?://([a-z]{2})\.wikipedia\.org/wiki/(.+)$")
async def _wikipedia_api(url: str) -> str:
"""Wikipedia über die Action-API laden (bot-freundlich; die HTML-Seiten
blockierten per robots.txt die Kernquelle des ersten Echtlaufs)."""
m = _WIKI_RE.match(url)
if not m:
return ""
sprache, titel = m.group(1), m.group(2).split("#")[0]
async with httpx.AsyncClient(
timeout=config.FETCH_READ_S,
headers={"User-Agent": config.user_agent()}) as client:
r = await client.get(
f"https://{sprache}.wikipedia.org/w/api.php",
params={"action": "query", "prop": "extracts", "explaintext": 1,
"format": "json", "redirects": 1, "titles": titel})
r.raise_for_status()
seiten = r.json().get("query", {}).get("pages", {})
return "\n\n".join(p.get("extract", "") for p in seiten.values())
async def _fetch(url: str) -> tuple[bytes, str]:
async with httpx.AsyncClient(
follow_redirects=True,
timeout=httpx.Timeout(config.FETCH_READ_S, connect=config.FETCH_CONNECT_S),
headers={"User-Agent": config.BROWSER_UA},
transport=httpx.AsyncHTTPTransport(retries=2)) as client:
r = await client.get(url)
r.raise_for_status()
return r.content, r.headers.get("content-type", "")
def _extrahieren(inhalt: bytes) -> str:
import trafilatura
return trafilatura.extract(inhalt, output_format="markdown",
include_tables=True, include_links=False,
include_images=False) or ""
def snapshot_pfad(topic: str, hash_: str) -> str:
return str(config.KORPUS_DIR / topic / f"q{hash_}.md")
def snapshot_lesen(quelle: dict) -> str:
from pathlib import Path
return Path(quelle["snapshot"]).read_text(encoding="utf-8")
def _atomar_schreiben(pfad: str, daten: bytes) -> None:
os.makedirs(os.path.dirname(pfad), exist_ok=True)
tmp = pfad + ".tmp"
with open(tmp, "wb") as f:
f.write(daten)
os.replace(tmp, pfad)
@engine.worker("laden.laden")
async def laden(task: dict) -> engine.Ergebnis:
quelle = db.one("SELECT * FROM quellen WHERE id=?",
db.uj(task["payload"], {}).get("quelle_id"))
if quelle is None:
raise RuntimeError("Quelle fehlt")
if quelle["status"] == "geladen": # Resume
return engine.Ergebnis(daten={"skip": True},
neue_tasks=_folge_tasks(task, quelle))
if config.FAKE:
from . import fakes
text = fakes.laden(quelle["url"])
roh = b""
else:
if _WIKI_RE.match(quelle["url"]):
try:
text = await _wikipedia_api(quelle["url"])
except Exception:
text = ""
if text:
text = _text_reparieren(text)
return await _snapshot_ablegen(task, quelle, text, b"")
# API-Fehler/Titel unbekannt → normaler Seiten-Weg als Fallback
await _domain_bremse(quelle["url"])
try:
roh, ctype = await _fetch(quelle["url"])
except Exception as e:
db.update("quellen", "id=?", (quelle["id"],), status="fehler",
grund=f"{type(e).__name__}"[:80])
return engine.Ergebnis(daten={"fetch_fehler": str(e)[:120]})
if "pdf" in ctype.lower(): # Content-Type, nicht URL-Suffix (Lektion 106)
db.update("quellen", "id=?", (quelle["id"],), status="fehler",
grund="pdf_uebersprungen")
db.insert("befunde", run_id=task["run_id"], stufe="korpus",
knoten="laden", art="pdf", item=quelle["url"][:120],
detail="PDF in v1 übersprungen")
return engine.Ergebnis(daten={"pdf": True})
text = _extrahieren(roh)
text = _text_reparieren(text)
return await _snapshot_ablegen(task, quelle, text, roh)
async def _inhalt_nuetzlich(task: dict, quelle: dict, text: str) -> tuple[bool, str]:
"""Inhalts-Urteil statt Titel-Filter (Lektion 102): erst laden, dann auf
dem echten Text entscheiden. AUSNAHMSLOS für jede Quelle — auch Seed/
Wikipedia (Lektion 106: privilegierte Quellen luden „Kosovo" ungeprüft)."""
thema = db.one("SELECT COALESCE(NULLIF(titel,''), name) t FROM topics "
"WHERE name=?", task["topic"])["t"]
antwort = await llm.call(
run_id=task["run_id"], stufe="korpus", knoten="laden",
item=f"quelle:{quelle['id']}:urteil", role="judge",
skill_namen=graph.skills_von(task["knoten"], extra="urteil"),
werte={"thema": thema, "titel": quelle["titel"],
"auszug": text[:config.URTEIL_AUSZUG_ZEICHEN]})
bloecke = textkit.bloecke(antwort or "", "URTEIL")
urteil = bloecke[0].get("nuetzlich", "").strip().lower() if bloecke else ""
if urteil not in ("ja", "nein"): # Ausfall/kaputt → Retry, kein Silent-Keep
raise RuntimeError("quellen-urteil ohne verwertbare Antwort")
return urteil == "ja", bloecke[0].get("grund", "")[:80]
async def _snapshot_ablegen(task: dict, quelle: dict, text: str,
roh: bytes) -> engine.Ergebnis:
if len(text) < config.SNAPSHOT_MIN_ZEICHEN:
# zu wenig extrahierbarer Text (auch JS-Walls landen hier) — keine
# Sprach-/Phrasen-Heuristik (Lektion 106)
db.update("quellen", "id=?", (quelle["id"],), status="leer")
db.insert("befunde", run_id=task["run_id"], stufe="korpus", knoten="laden",
art="leer", item=quelle["url"][:120])
return engine.Ergebnis(daten={"leer": True})
hash_ = hashlib.sha256(text.encode()).hexdigest()[:16]
dublette = db.one("SELECT id FROM quellen WHERE topic=? AND hash=? AND id!=?",
task["topic"], hash_, quelle["id"])
if dublette:
db.update("quellen", "id=?", (quelle["id"],), status="fehler",
grund="inhalt_dublette")
return engine.Ergebnis(daten={"dublette": True})
nuetzlich, grund = await _inhalt_nuetzlich(task, quelle, text)
if not nuetzlich:
db.update("quellen", "id=?", (quelle["id"],), status="aussortiert",
grund=grund)
return engine.Ergebnis(daten={"aussortiert": True})
pfad = snapshot_pfad(task["topic"], hash_)
_atomar_schreiben(pfad, text.encode("utf-8"))
roh_pfad = ""
if roh:
roh_pfad = pfad.replace(".md", ".html")
_atomar_schreiben(roh_pfad, roh) # Re-Extraktion ohne Re-Fetch (R6)
db.update("quellen", "id=?", (quelle["id"],), status="geladen", hash=hash_,
snapshot=pfad, roh=roh_pfad)
quelle = {**quelle, "snapshot": pfad, "status": "geladen"}
return engine.Ergebnis(daten={"zeichen": len(text)},
neue_tasks=_folge_tasks(task, quelle))
def _folge_tasks(task: dict, quelle: dict) -> list:
"""korpus-/seed-Quellen → Soll-Extraktion je 40k-Chunk; Lücken-Quellen
(soll_id gesetzt) → direkte Atom-Extraktion (Soll steht schon fest)."""
if quelle["zweck"] == "luecke" and quelle.get("soll_id"):
text = snapshot_lesen(quelle)
return [{"knoten": "atom_extraktion",
"item": f"quelle:{quelle['id']}:a{i}:r1",
"payload": {"quelle_id": quelle["id"], "offset": o,
"chars": len(c), "soll_id": quelle["soll_id"]}}
for i, (o, c) in enumerate(textkit.abschnitte(text))]
text = snapshot_lesen(quelle)
chunks = textkit.abschnitte(text, config.SOLL_CHUNK_CHARS)
return [{"knoten": "soll_extraktion", "item": f"quelle:{quelle['id']}:c{i}",
"payload": {"quelle_id": quelle["id"], "offset": o, "chars": len(c)}}
for i, (o, c) in enumerate(chunks)]