@@ -21,7 +21,8 @@ from agents import kill_process
from config import KONSENS_GRACE , KONSENS_MAX_RUNDEN , DEFAULT_PROVIDER
from fsutil import atomic_write_text , atomic_write_json
from jsonio import read_json_file as _json_datei
from paths import arbeit_dir , bausteine_path , project_dir , subbausteine_path , quelle_path , quelle_crawl_dir , safe_ordner
from lernen import FRAGETYPEN
from paths import arbeit_dir , bausteine_path , frage_muster_path , project_dir , subbausteine_path , quelle_path , quelle_crawl_dir , safe_ordner
from crawl import crawl
from pipeline import (
CANCELLED , FAILED , GenContext , _extra , _log , _prompt , _race , _relevanz_schema , _rest_schema ,
@@ -51,6 +52,7 @@ BAUSTEINE_STEPS = (
" Subbausteine finden " , " Subbausteine wählen " , " Subbausteine klären " ,
" Stufen finden " , " Stufen wählen " , " Stufen klären " ,
" Relevanz finden " , " Relevanz wählen " , " Relevanz klären " ,
" Fragen finden " , " Fragen wählen " , " Fragen klären " ,
)
@@ -93,6 +95,20 @@ def subbausteine_titel(topic: str, baustein: str) -> list[str]:
]
def lade_frage_muster ( topic : str , baustein : str ) - > list [ dict ] :
""" Vordefinierte Frage-Muster eines Bausteins aus dem Sidecar (leer = Fallback auf Live). """
fm = _json_datei ( frage_muster_path ( topic ) )
if not isinstance ( fm , dict ) :
return [ ]
return [
{ " subbaustein " : str ( e . get ( " subbaustein " , " " ) ) . strip ( ) ,
" typ " : str ( e . get ( " typ " , " " ) ) . strip ( ) ,
" frage " : frage }
for e in ( fm . get ( baustein ) or [ ] )
if isinstance ( e , dict ) and ( frage := str ( e . get ( " frage " , " " ) ) . strip ( ) )
]
def lade_uebersicht ( topic : str ) - > list [ dict ] :
""" Strukturierte Baustein-Liste für die Übersicht: Titel + Beschreibung + Subbausteine/Stufen.
@@ -132,6 +148,7 @@ def _bausteine_steps(topic: str) -> tuple:
" Subbausteine finden " , " Subbausteine wählen " , " Subbausteine klären " ,
" Stufen finden " , " Stufen wählen " , " Stufen klären " ,
" Relevanz finden " , " Relevanz wählen " , " Relevanz klären " ,
" Fragen finden " , " Fragen wählen " , " Fragen klären " ,
)
mitte = base + ( ( " Ergänzung " , ) if q [ " type " ] == " projekt " else ( ) ) + rest
return ( ( " Quelle laden " , ) if q [ " type " ] == " link " else ( ) ) + mitte
@@ -141,6 +158,40 @@ def _step_idx(topic: str, name: str) -> int:
return _bausteine_steps ( topic ) . index ( name )
# Grobe Anzeige-Phasen: bündeln die Feinschritte (intern bleibt alles feingranular).
# Sonderschritte (Quelle laden, Ergänzung) gehören zur Phase „Inventar".
PHASEN = (
( " Inventar " , ( " Quelle laden " , " Recherche " , " Konsolidierung " , " Klärung " , " Ergänzung " ) ) ,
( " Subbausteine " , ( " Subbausteine finden " , " Subbausteine wählen " , " Subbausteine klären " ) ) ,
( " Stufen " , ( " Stufen finden " , " Stufen wählen " , " Stufen klären " ) ) ,
( " Relevanz " , ( " Relevanz finden " , " Relevanz wählen " , " Relevanz klären " ) ) ,
( " Fragen " , ( " Fragen finden " , " Fragen wählen " , " Fragen klären " ) ) ,
)
def _phasen ( topic : str ) - > list [ tuple [ str , int ] ] :
""" [(grob_label, Anzahl vorhandener Feinschritte)] für die aktuelle Quelle. """
feine = _bausteine_steps ( topic )
return [ ( label , n ) for label , members in PHASEN if ( n := sum ( f in members for f in feine ) ) ]
def _phasen_status ( topic : str , current : int | None ) - > list [ dict ] :
""" Grobe Phasen-Zustände aus dem feinen Fortschritt `current` (None = alles pending,
len(feine) = alles done). → [ { label, state}] mit state done/active/pending. """
out , start = [ ] , 0
for label , n in _phasen ( topic ) :
end = start + n
if current is None or current < start :
state = " pending "
elif current > = end :
state = " done "
else :
state = " active "
out . append ( { " label " : label , " state " : state } )
start = end
return out
def _bausteine_files ( topic : str ) - > dict :
arbeit = arbeit_dir ( topic )
runden = range ( 1 , KONSENS_MAX_RUNDEN + 1 )
@@ -154,18 +205,19 @@ def _bausteine_files(topic: str) -> dict:
" ergaenzung " : arbeit / " ergaenzung.json " ,
" sub_roh " : arbeit / " subbausteine-roh.json " ,
" sidecar " : subbausteine_path ( topic ) ,
" frage_muster " : frage_muster_path ( topic ) ,
}
def _alle_slot_dateien ( files : dict ) - > list [ Path ] :
arbeit = files [ " arbeit " ]
# Subbaustein-/Stufen-Slots sind pro Chunk dynamisch — per Glob einsammeln.
dyn = ( list ( arbeit . glob ( " subbaustein-* " ) ) + list ( arbeit . glob ( " stufe-* " ) ) + list ( arbeit . glob ( " relevanz-* " ) ) ) if arbeit . is_dir ( ) else [ ]
dyn = ( list ( arbeit . glob ( " subbaustein-* " ) ) + list ( arbeit . glob ( " stufe-* " ) ) + list ( arbeit . glob ( " relevanz-* " ) ) + list ( arbeit . glob ( " frage-muster-* " ) ) ) if arbeit . is_dir ( ) else [ ]
return [
* files [ " recherche " ] , files [ " recherche_mapping " ] ,
* ( p for slots in files [ " auswahl " ] . values ( ) for p in slots ) ,
* files [ " mapping " ] . values ( ) , files [ " ergaenzung " ] ,
files [ " sub_roh " ] , files [ " sidecar " ] , * dyn ,
files [ " sub_roh " ] , files [ " sidecar " ] , files [ " frage_muster " ] , * dyn ,
]
@@ -199,38 +251,37 @@ def _resume_step(topic: str) -> int:
sidecar = _json_datei ( files [ " sidecar " ] )
if _sidecar_schema ( sidecar ) is not None :
# Stufen fertig; nur noch Relevanz offen?
if _relevanz_komplett ( sidecar ) :
return len ( _bausteine_steps ( topic ) )
return _step_idx ( topic , " Relevanz finden " )
if not _relevanz_komplett ( sidecar ) :
return _step_idx ( topic , " Relevanz finden " )
# Relevanz fertig; nur noch Frage-Muster offen?
if not _frage_muster_komplett ( topic ) :
return _step_idx ( topic , " Fragen finden " )
return len ( _bausteine_steps ( topic ) )
if _sub_roh_schema ( _json_datei ( files [ " sub_roh " ] ) ) is None :
return _step_idx ( topic , " Subbausteine finden " )
return _step_idx ( topic , " Stufen finden " )
def bausteine_status ( topic : str ) - > dict :
steps = _bausteine_steps ( topic )
# Intern feingranular (Resume/Progress); für die Anzeige zu 5 groben Phasen gebündelt.
feine = _bausteine_steps ( topic )
ready = bausteine_path ( topic ) . exists ( )
generating = topic in _bausteine_progress
partial = False
if generating :
current = _bausteine_step . get ( topic )
states = [
" pending " if current is None else " done " if i < current else " active " if i == current else " pending "
for i in range ( len ( steps ) )
]
elif ready :
states = [ " done " ] * len ( steps )
current = len ( feine ) # alles fertig → alle Phasen done
else :
nx t = _resume_step ( topic )
partial = nx t > 0
states = [ " done " if i < nxt else " pending " for i in range ( len ( steps ) ) ]
curre nt = _resume_step ( topic )
partial = curre nt > 0
return {
" ready " : ready ,
" generating " : generating ,
" progress " : _bausteine_progress . get ( topic ) ,
" error " : _bausteine_errors . get ( topic ) ,
" partial " : partial ,
" steps " : [ { " label " : label , " state " : s } for label , s in zip ( steps , states ) ] ,
" steps " : _phasen_status ( topic , current ) ,
}
@@ -242,12 +293,49 @@ def reset_bausteine(topic: str) -> None:
files = _bausteine_files ( topic )
files [ " final " ] . unlink ( missing_ok = True )
files [ " sidecar " ] . unlink ( missing_ok = True ) # liegt im Themen-Root, nicht in arbeit/
files [ " frage_muster " ] . unlink ( missing_ok = True ) # ebenfalls im Themen-Root
quelle_path ( topic ) . unlink ( missing_ok = True )
shutil . rmtree ( quelle_crawl_dir ( topic ) , ignore_errors = True ) # gecrawlte Link-Quelle
shutil . rmtree ( files [ " arbeit " ] , ignore_errors = True )
_bausteine_errors . pop ( topic , None )
def _reset_ab_phase ( topic : str , phase : int ) - > None :
""" Artefakte AB der groben Phase löschen (1 Inventar … 5 Fragen), frühere behalten.
Danach baut der normale Resume-Flow die fehlenden Phasen neu. Kumulativ: Re-Run ab P
löscht P…5. quelle.json + Crawl bleiben immer (Quellen-Wahl erhalten). """
files = _bausteine_files ( topic )
arbeit = files [ " arbeit " ]
def glob_del ( pat : str ) - > None :
if arbeit . is_dir ( ) :
for p in arbeit . glob ( pat ) :
p . unlink ( missing_ok = True )
if phase < = 5 : # Fragen
files [ " frage_muster " ] . unlink ( missing_ok = True )
glob_del ( " frage-muster-* " )
if phase < = 4 : # Relevanz
glob_del ( " relevanz-* " )
if phase < = 3 : # Stufen + Relevanz teilen die Sidecar → ab Stufen ganz neu
files [ " sidecar " ] . unlink ( missing_ok = True )
glob_del ( " stufe-* " )
else : # ab Relevanz: Stufen behalten, nur Relevanz-Felder strippen
sc = _json_datei ( files [ " sidecar " ] )
if isinstance ( sc , dict ) :
for subs in sc . values ( ) :
for s in ( subs if isinstance ( subs , list ) else [ ] ) :
if isinstance ( s , dict ) :
s . pop ( " relevanz " , None )
atomic_write_json ( files [ " sidecar " ] , sc , indent = 1 )
if phase < = 2 : # Subbausteine
files [ " sub_roh " ] . unlink ( missing_ok = True )
glob_del ( " subbaustein-* " )
if phase < = 1 : # Inventar = kompletter Frischstart (alle Zwischendateien)
for p_alt in _alle_slot_dateien ( files ) :
p_alt . unlink ( missing_ok = True )
def _ergaenzung_schema ( data ) :
""" { " bausteine " : [ { " titel " , " beschreibung " }]} → Liste (leer erlaubt) · sonst None. """
if not isinstance ( data , dict ) or not isinstance ( data . get ( " bausteine " ) , list ) :
@@ -357,6 +445,32 @@ def _relevanz_komplett(data) -> bool:
)
def _frage_muster_schema ( data ) - > list [ dict ] | None :
""" { " muster " : [ { subbaustein, typ, frage}, …]} → Liste valider Einträge · sonst None.
typ muss ein bekannter Fragetyp sein; subbaustein + frage nicht leer. Leere Liste → None.
"""
if not isinstance ( data , dict ) or not isinstance ( data . get ( " muster " ) , list ) :
return None
out = [ ]
for e in data [ " muster " ] :
if not isinstance ( e , dict ) :
return None
sub = str ( e . get ( " subbaustein " , " " ) ) . strip ( )
typ = str ( e . get ( " typ " , " " ) ) . strip ( ) . casefold ( )
frage = str ( e . get ( " frage " , " " ) ) . strip ( )
if not sub or typ not in FRAGETYPEN or not frage :
return None
out . append ( { " subbaustein " : sub , " typ " : typ , " frage " : frage } )
return out or None
def _frage_muster_komplett ( topic : str ) - > bool :
""" Frage-Muster-Sidecar existiert (Build gelaufen)? Einzelne leere Bausteine
fallen zur Prüfungszeit auf Live-Generierung zurück — daher genügt die Datei. """
return isinstance ( _json_datei ( frage_muster_path ( topic ) ) , dict )
def _read ( p : Path ) - > str :
return p . read_text ( encoding = " utf-8 " ) if p . exists ( ) else " "
@@ -741,7 +855,106 @@ async def _relevanz_block(ctx: GenContext, set_p, files: dict, sidecar: dict, in
return relevanz_by_id
async def generate_baustein e( topic : str , instructions : str = " " , provider : str = DEFAULT_PROVIDER ) - > None :
def _norm_frag e( t : str ) - > str :
return " " . join ( str ( t or " " ) . lower ( ) . split ( ) )
async def _frage_muster_block ( ctx : GenContext , set_p , files : dict , sidecar : dict , instructions : str ) - > dict | None :
""" Block E: drei Phasen — Finden (1 Generator je Baustein, parallel), Wählen (Code:
relevante Subs + Dedup), Klären (Kritiker je Baustein bereinigt die Tabelle).
Generativ statt Vote: Muster sind Text, kein 3-Rater-Konsens sinnvoll.
→ { Baustein-Titel: [ { subbaustein, typ, frage}, …]} oder None bei Abbruch. """
topic , is_cancelled = ctx . topic , ctx . is_cancelled
arbeit = files [ " arbeit " ]
# Je Baustein nur RELEVANTE Subbausteine (relevanz != 'rand'; None zählt als relevant).
bausteine = [ ]
for titel , subs in sidecar . items ( ) :
rel = [ s [ " titel " ] for s in subs if isinstance ( s , dict ) and s . get ( " relevanz " ) != " rand " and str ( s . get ( " titel " , " " ) ) . strip ( ) ]
if rel :
bausteine . append ( ( titel , rel ) )
if not bausteine :
return { }
typen_block = " \n " . join ( f " - { k } : { v } " for k , v in FRAGETYPEN . items ( ) )
def roh_path ( c ) :
return arbeit / f " frage-muster-c { c } .json "
def final_path ( c ) :
return arbeit / f " frage-muster-final-c { c } .json "
# Phase „Fragen finden": pro Baustein 1 Generator, alle parallel.
async def _finde ( c , titel , rel ) :
fp = roh_path ( c )
if _frage_muster_schema ( _json_datei ( fp ) ) :
return # Resume
sub_block = " \n " . join ( f " - { s } " for s in rel )
status , _ = await run_single_slot (
ctx , f " Frage-Muster { c } " ,
key = f " bausteine- { topic } -frage-muster-c { c } " ,
prompt = _prompt ( " Frage-Muster-Recherche " , topic = topic , baustein = titel ,
subbausteine = sub_block , typen = typen_block , out_path = fp , extra = _extra ( instructions ) ) ,
role = " fast " , capabilities = " files " ,
payload = lambda result , p = fp : _frage_muster_schema ( _json_datei ( p ) ) ,
timeout = _timeout ( " frage_muster " , len ( rel ) ) ,
)
if status == FAILED :
_log ( topic , f " Frage-Muster Baustein { c } fehlgeschlagen — kein Muster (Fallback Live) " )
set_p ( f " Fragen finden ( { len ( bausteine ) } Bausteine)… " , step = _step_idx ( topic , " Fragen finden " ) )
await asyncio . gather ( * [ _finde ( c , t , r ) for c , ( t , r ) in enumerate ( bausteine , 1 ) ] , return_exceptions = True )
if is_cancelled ( ) :
return None
# Phase „Fragen wählen": Code — nur relevante Subbausteine, Dubletten je Baustein raus.
set_p ( " Fragen wählen… " , step = _step_idx ( topic , " Fragen wählen " ) )
roh_by_c = { }
for c , ( titel , rel ) in enumerate ( bausteine , 1 ) :
eintraege = _frage_muster_schema ( _json_datei ( roh_path ( c ) ) ) or [ ]
rel_set = set ( rel )
gesehen , sauber = set ( ) , [ ]
for e in eintraege :
norm = _norm_frage ( e [ " frage " ] )
if e [ " subbaustein " ] not in rel_set or norm in gesehen :
continue
gesehen . add ( norm )
sauber . append ( e )
roh_by_c [ c ] = sauber
# Phase „Fragen klären": Kritiker je Baustein bereinigt die Tabelle (eindeutig, distinkt).
async def _klaere ( c , titel ) :
roh = roh_by_c [ c ]
fp = final_path ( c )
if _frage_muster_schema ( _json_datei ( fp ) ) or not roh :
return
tabelle = " \n " . join ( f " { i } . [ { e [ ' typ ' ] } ] ( { e [ ' subbaustein ' ] } ) { e [ ' frage ' ] } " for i , e in enumerate ( roh , 1 ) )
status , _ = await run_single_slot (
ctx , f " Frage-Muster-Klärung { c } " ,
key = f " bausteine- { topic } -frage-muster-final-c { c } " ,
prompt = _prompt ( " Frage-Muster-Kritik " , topic = topic , baustein = titel , tabelle = tabelle , out_path = fp , extra = _extra ( instructions ) ) ,
role = " judge " , capabilities = " files " ,
payload = lambda result , p = fp : _frage_muster_schema ( _json_datei ( p ) ) ,
timeout = _timeout ( " frage_muster_check " , len ( roh ) ) ,
)
if status == FAILED :
_log ( topic , f " Frage-Muster-Klärung Baustein { c } fehlgeschlagen — Roh-Muster übernommen " )
set_p ( f " Fragen klären ( { len ( bausteine ) } Bausteine)… " , step = _step_idx ( topic , " Fragen klären " ) )
await asyncio . gather ( * [ _klaere ( c , t ) for c , ( t , _ ) in enumerate ( bausteine , 1 ) ] , return_exceptions = True )
if is_cancelled ( ) :
return None
# Final: geklärte Tabelle je Baustein, Fallback auf Roh-Muster.
ergebnis = { }
for c , ( titel , rel ) in enumerate ( bausteine , 1 ) :
rel_set = set ( rel )
final = _frage_muster_schema ( _json_datei ( final_path ( c ) ) )
if not final :
final = roh_by_c [ c ]
ergebnis [ titel ] = [ e for e in final if e [ " subbaustein " ] in rel_set ]
return ergebnis
async def generate_bausteine ( topic : str , instructions : str = " " , provider : str = DEFAULT_PROVIDER , ab_phase : int | None = None ) - > None :
if topic in _bausteine_progress :
return
_bausteine_progress [ topic ] = " Wartend… "
@@ -769,6 +982,10 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str =
try :
async with _semaphore :
files [ " arbeit " ] . mkdir ( parents = True , exist_ok = True )
# Re-Run ab gewählter Phase: Artefakte ab dort löschen; der Frischstart-Block
# unten wird übersprungen (er würde bei erhaltener Sidecar sonst alles wischen).
if ab_phase is not None :
_reset_ab_phase ( topic , ab_phase )
# Link-Quelle: erst crawlen (gleiche Domain, begrenzt) → wird zur Ordner-Quelle.
if q [ " type " ] == " link " and not _crawl_fertig ( topic ) :
set_p ( " Quelle laden (Crawl)… " , step = _step_idx ( topic , " Quelle laden " ) )
@@ -784,7 +1001,8 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str =
# „Neu erstellen": NUR wenn wirklich alles fertig ist (bausteine.md UND
# Sidecar) → kompletter Frischstart. Liegt bausteine.md ohne Sidecar vor,
# ist das ein Teil-Stand (Block B/C offen) → Resume, nicht wischen.
fertig = final_path . exists ( ) and _sidecar_schema ( _json_datei ( files [ " sidecar " ] ) ) is not None
# Bei explizitem Re-Run (ab_phase) hat _reset_ab_phase das schon erledigt.
fertig = ab_phase is None and final_path . exists ( ) and _sidecar_schema ( _json_datei ( files [ " sidecar " ] ) ) is not None
if fertig :
for p_alt in _alle_slot_dateien ( files ) :
p_alt . unlink ( missing_ok = True )
@@ -1019,6 +1237,19 @@ async def generate_bausteine(topic: str, instructions: str = "", provider: str =
gid + = 1
sub [ " relevanz " ] = relevanz_by_id . get ( gid , " relevant " )
atomic_write_json ( files [ " sidecar " ] , sidecar , indent = 1 )
# Block E: Frage-Muster je relevantem Subbaustein × Typ → eigenes Sidecar.
# Zur Prüfungszeit zieht jeder Agent ein Muster ohne Zurücklegen und formuliert
# daraus eine Frage — distinkte Saat verhindert die Doppelfragen der Live-Generierung.
sidecar = _json_datei ( files [ " sidecar " ] )
if _sidecar_schema ( sidecar ) is not None and _relevanz_komplett ( sidecar ) and not _frage_muster_komplett ( topic ) :
muster = await _frage_muster_block ( ctx , set_p , files , sidecar , instructions )
if is_cancelled ( ) :
abgebrochen ( )
return
if muster is None :
return # Abbruch
atomic_write_json ( files [ " frage_muster " ] , muster , indent = 1 )
except Exception as e :
log . exception ( " [ %s ] Bausteine-Generierung fehlgeschlagen " , topic )
_bausteine_errors [ topic ] = str ( e ) [ : 2000 ]