@@ -20,7 +20,7 @@ from pathlib import Path
import database as db
from agents import kill_process , cancel_scope , clear_scope
from config import KONSENS_GRACE , RECHERCHE_GRACE , KONSENS_MAX_RUNDEN , DEFAULT_PROVIDER
from config import KONSENS_GRACE , RECHERCHE_GRACE , KONSENS_MAX_RUNDEN , DEFAULT_PROVIDER , CRAWL_KEEP_PATTERNS , CRAWL_NOISE_PATTERNS , CRAWL_MIN_CHARS
from fsutil import atomic_write_text , atomic_write_json
from jsonio import read_json_file as _json_datei
from lernen import FRAGETYPEN
@@ -46,10 +46,11 @@ RECHERCHE_BATCH = 20 # Crawl-Seiten je Batch
RECHERCHE_READERS = 2 # Reader-Agenten je Batch (Konsens ≥2 innerhalb des Batches)
RECHERCHE_THEMA_AGENTEN = 5 # Web-Modus (Quelle „thema", kein Crawl-Ordner)
RECHERCHE_KAPPE = 1800 # Sicherheits-Deckel je Batch-Agent
SICHTUNG_CHUNK = 80 # Crawl-Seiten je Sichtungs-Chunk (Content/Noise-Triage)
SICHTUNG_READERS = 3 # Sichter je Chunk; Merge per Mehrheit, Tie-Break = Content
# Sichtung (Content/Noise) ist jetzt ein deterministischer Regel-Filter (config.CRAWL_*).
SUBBAUSTEIN_KAPPE = 900 # Subbaustein-Finde-Loop je Chunk (15 min)
KONSOLIDIERUNG_CHUNK = 600 # bis hierher EIN globaler Judge (dedupt alles); darüber chunked + Merge-Pass
# Frage-Muster-Chunks per LPT nach Sub-Last balancieren (Makespan), statt nach Baustein-Anzahl.
FRAGE_CHUNK_SUBS = 50 # Ziel-Summe relevanter Subs je Chunk
log = logging . getLogger ( " creator.bausteine " )
@@ -324,16 +325,17 @@ def active_bausteine() -> list[dict]:
def reset_bausteine ( topic : str ) - > None :
""" „Lösch en " : räumt nur die generierten Bausteine weg (Dateien ab Inventar). BEHÄLT
Quelle (quelle.json), Crawl-Ordner und Sichtung. Voller Wipe inkl. Quelle/C rawl nur
über „Thema löschen " (DELETE /topics → rmtree(topic_dir)). """
_reset_ab_phase ( topic , " Inventar " )
""" „Entfern en " : löscht den GESAMTEN Bausteine-Bereich — Crawl, Sichtung, Inventar … Fragen.
BEHÄLT nur die Themen-Config `quelle.json` (Typ/Link/Spec). Re-Generieren c rawlt neu.
(Crawl/Sichtung gehören zu den Bausteinen; nur die Config ist „Thema " .) """
files = _bausteine_files ( topic )
files [ " final " ] . unlink ( missing_ok = True )
files [ " sidecar " ] . unlink ( missing_ok = True )
files [ " frage_muster " ] . unlink ( missing_ok = True )
shutil . rmtree ( quelle_crawl_dir ( topic ) , ignore_errors = True ) # Crawl gehört zu Bausteinen
shutil . rmtree ( files [ " arbeit " ] , ignore_errors = True )
_bausteine_errors . pop ( topic , None )
async def reset_bausteine_db ( topic : str ) - > None :
""" DB-Pendant zu reset_bausteine: Inventar…Fragen aus der DB; Coverage/Sichtung + Quelle bleiben. """
await _reset_db_ab_phase ( topic , " Inventar " )
# quelle.json bleibt bewusst stehen — das ist die Themen-Config.
def _phase_idx ( label : str ) - > int :
@@ -550,6 +552,21 @@ def _n_chunks(count: int, size: int = SUBBAUSTEIN_CHUNK) -> int:
return min ( SUBBAUSTEIN_MAX , max ( 1 , math . ceil ( count / size ) ) )
def _lpt_chunks ( gewichte : list [ int ] , target : int ) - > list [ list [ int ] ] :
""" Indizes lastbalanciert auf Chunks verteilen (LPT, Makespan-minimal). Gewicht = Kosten je Index.
K = ceil(Gesamtgewicht/target); schwerste zuerst in den jeweils leichtesten Bin. → Index-Listen. """
if not gewichte :
return [ ]
K = max ( 1 , math . ceil ( sum ( gewichte ) / max ( 1 , target ) ) )
bins : list [ list [ int ] ] = [ [ ] for _ in range ( K ) ]
last = [ 0 ] * K
for i in sorted ( range ( len ( gewichte ) ) , key = lambda x : gewichte [ x ] , reverse = True ) :
j = min ( range ( K ) , key = lambda b : last [ b ] )
bins [ j ] . append ( i )
last [ j ] + = gewichte [ i ]
return [ b for b in bins if b ]
def _merge_finder ( num : int , idx : dict , finder : list [ dict ] ) - > tuple [ list [ str ] , list [ str ] ] :
""" Subbausteine eines Bausteins über die Finder mergen → (konsens ≥2, rest ==1). """
counts : dict [ str , int ] = { }
@@ -753,7 +770,18 @@ async def _stufen_block(ctx: GenContext, set_p, files: dict, roh: dict, instruct
items = [ ( titel , sub ) for titel , subs in roh . items ( ) for sub in subs ] # globale id = index+1
if not items :
return { titel : [ ] for titel in roh }
c hunks = _chunk_nums ( list ( range ( len ( items ) ) ) , _n_chunks ( len ( items ) , STUFE_CHUNK ) )
# C hunks aus GANZEN Bausteinen packen (keinen Baustein splitten) → Rater sieht je Baustein
# alle Subs und kann relativ einstufen. Item-Indizes je Baustein in roh-Reihenfolge.
chunks , cur , i = [ ] , [ ] , 0
for _titel_b , subs in roh . items ( ) :
g = list ( range ( i , i + len ( subs ) ) )
i + = len ( subs )
if cur and len ( cur ) + len ( g ) > STUFE_CHUNK :
chunks . append ( cur )
cur = [ ]
cur . extend ( g )
if cur :
chunks . append ( cur )
n = len ( chunks )
def rater_paths ( c ) :
@@ -769,7 +797,14 @@ async def _stufen_block(ctx: GenContext, set_p, files: dict, roh: dict, instruct
vorhanden = sum ( 1 for p in paths if _stufen_schema ( _json_datei ( p ) , local_set ) )
if vorhanden > = 2 :
return True
enum = " \n " . join ( f " { k } . [ { items [ j ] [ 0 ] } ] { items [ j ] [ 1 ] } " for k , j in enumerate ( item_idxs , 1 ) )
enum_zeilen , cur_b = [ ] , None
for k , j in enumerate ( item_idxs , 1 ) :
b , sub = items [ j ]
if b != cur_b :
enum_zeilen . append ( f " \n BAUSTEIN: { b } " )
cur_b = b
enum_zeilen . append ( f " { k } . { sub } " )
enum = " \n " . join ( enum_zeilen ) . strip ( )
offen = [ ( i , p ) for i , p in enumerate ( paths , 1 ) if not _stufen_schema ( _json_datei ( p ) , local_set ) ]
slots = [ {
" key " : f " bausteine- { topic } -stufe-c { c } - { i } " ,
@@ -1002,7 +1037,7 @@ async def _frage_muster_block(ctx: GenContext, set_p, files: dict, sidecar: dict
if not bausteine :
return { }
typen_block = " \n " . join ( f " - { k } : { v } " for k , v in FRAGETYPEN . items ( ) )
chunks = _chunk_num s ( list ( range ( len ( bausteine ) ) ) , _n_chunks ( len ( bausteine ) ) ) # ~10 Bausteine/Chunk
chunks = _lpt_ chunks ( [ len ( rel ) for _ , rel in bausteine ] , FRAGE_CHUNK_SUBS ) # lastbalanciert nach Sub-Zahl
def roh_path ( ci ) :
return arbeit / f " frage-muster-c { ci } .json "
@@ -1223,65 +1258,25 @@ async def _set_inventar(topic: str, eintrag: str, status: str) -> None:
await db . set_baustein_status ( topic , norm , status )
def _sichtung_schema ( data ) - > list [ str ] | None :
""" { " noise " : [ " datei.txt " , …]} → Liste Dateinamen (leer erlaubt) · sonst None. """
if not isinstance ( data , dict ) or not isinstance ( data . get ( " noise") , list ) :
return None
return [ s for x in data [ " noise " ] if ( s := str ( x ) . strip ( ) ) ]
async def _sichte_inhalt ( ctx : GenContext , set_p , files : dict , ordner , pages : list [ str ] ) - > set [ str ] :
""" Content/Noise-Triage nach dem Crawl. Je Chunk SICHTUNG_READERS Sichter; Merge je Seite
per Mehrheit, Tie-Break = Content (konservativ). Fehler/Ausfall → Content (fail-open).
→ Menge der Content-Dateinamen (nie leer, solange pages nicht leer). """
topic , provider , is_cancelled = ctx . topic , ctx . provider , ctx . is_cancelled
pages = list ( pages )
if not pages :
return set ( )
arbeit = files [ " arbeit " ]
def _sichte_regeln ( ordner , pages : list [ str ] ) - > tuple [ list[ str ] , list [ str ] ] :
""" Deterministischer Content/Noise-Filter (config.CRAWL_*). Substring-Match (klein) gegen
URL + Dateiname. Reihenfolge: keep > noise > min_chars > behalten. → (content, noise). """
ordner = Path ( ordner )
def _snippet ( fn : str ) - > str :
content , noise = [ ] , [ ]
for fn in pages :
zeilen = _read ( ordner / fn ) . splitlines ( )
url = zeilen [ 0 ] [ len ( " QUELLE: " ) : ] . strip ( ) if zeilen and zeilen [ 0 ] . startswith ( " QUELLE: " ) else " "
body = " " . join ( z . strip ( ) for z in zeilen[ 1 : ] if z . strip ( ) ) [ : 200 ]
return f " - { fn } · { url } · { body } "
chunks = _chunk_nums ( sorted ( pages ) , max ( 1 , math . ceil ( l en( pages ) / SICHTUNG_CHUNK ) ) )
async def _chunk ( ci : int , chunk : list [ str ] ) - > set [ str ] :
seiten = " \n " . join ( _snippet ( fn ) for fn in chunk )
paths = [ arbeit / f " sichtung-c { ci } - { i } .json " for i in range ( 1 , SICHTUNG_READERS + 1 ) ]
for p in paths :
p . unlink ( missing_ok = True )
slots = [ {
" key " : f " bausteine- { topic } -sichtung-c { ci } - { i } " ,
" prompt " : _prompt ( " Bausteine-Seiten-Sichtung " , topic = topic , seiten = seiten , out_path = p ) ,
" role " : " judge " , " capabilities " : " files " ,
" payload " : ( lambda result , p = p : _sichtung_schema ( _json_datei ( p ) ) ) ,
} for i , p in enumerate ( paths , 1 ) ]
ergebnisse = await _race ( topic , f " Seiten-Sichtung { ci } " , slots , 1 ,
_timeout ( " recherche_mapping " , len ( chunk ) ) , provider ,
cancelled = is_cancelled , grace = KONSENS_GRACE )
cs = set ( chunk )
stimmen = [ set ( r ) & cs for r in ( ergebnisse or [ ] ) ] # Noise-Menge je Reader (nur Chunk-Seiten)
R = len ( stimmen )
if not R :
return set ( ) # fail-open: kein Urteil → alles Content
# Seite = Noise nur bei strenger Mehrheit; Gleichstand/Minderheit → Content.
return { page for page in chunk if 2 * sum ( page in s for s in stimmen ) > R }
async def melde ( d , t ) :
set_p ( f " Sichte Seiten { d } / { t } … " )
set_p ( " Sichte Seiten… " , step = _step_idx ( topic , " Quelle aufbereiten " ) )
teile = await _gather_fortschritt ( [ _chunk ( ci , c ) for ci , c in enumerate ( chunks , 1 ) ] , len ( chunks ) , melde )
if is_cancelled ( ) :
return set ( pages )
noise = set ( ) . union ( * ( t for t in teile if isinstance ( t , set ) ) ) if teile else set ( )
content = set ( pages ) - noise
_log ( topic , f " Sichtung: { len ( content ) } Content / { len ( noise ) } Noise von { len ( pages ) } ( { len ( chunks ) } Chunks × { SICHTUNG_READERS } Sichter) " )
return content or set ( pages ) # alles Noise? → fail-open, alle behalten
body = " \n " . join ( zeilen [ 1 : ] ) . strip ( )
hay = f " { url } \n { fn } " . lower ( )
if any ( p in hay for p in CRAWL_KEEP_PATTERNS ) :
content . app end ( fn )
elif any ( p in hay for p in CRAWL_NOISE_PATTERNS ) :
noise . append ( fn )
elif len ( body ) < CRAWL_MIN_CHARS :
noise . append ( fn )
else :
content . append ( fn ) # Default: behalten — alles mit Inhalt bleibt
return content , noise
async def _quelle_aufbereiten ( ctx : GenContext , set_p , files : dict , q : dict , ordner ) - > bool :
@@ -1307,12 +1302,11 @@ async def _quelle_aufbereiten(ctx: GenContext, set_p, files: dict, q: dict, ordn
await asyncio . to_thread ( _pdfs_konvertieren , ordner )
pages = sorted ( set ( _crawl_index ( ordner ) . values ( ) ) )
if pages :
set_p ( " Sichte Seiten… " , step = _step_idx ( topic , " Quelle aufbereiten " ) )
await db . delete_coverage ( topic )
content = await _sichte_inhalt ( ctx , set_p , files , ordner , pages )
if is_cancelled ( ) :
return False
noise = [ p for p in pages if p not in content ]
await db . mark_inhalt ( topic , sorted ( content ) , noise )
content , noise = _sichte_regeln ( ordner , pages ) # deterministischer Regel-Filter
await db . mark_inhalt ( topic , sorted ( content ) , sorted ( noise ) )
_log ( topic , f " Sichtung: { len ( content ) } Content / { len ( noise ) } Noise von { len ( pages ) } (Regeln) " )
await db . set_step_status ( topic , " Quelle aufbereiten " , " fertig " )
return True