206 lines
8.7 KiB
Python
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)]
|