Spaces:
Sleeping
Sleeping
File size: 15,549 Bytes
3164f14 9d9f0e5 3164f14 9d9f0e5 3164f14 9d9f0e5 3164f14 9d9f0e5 3164f14 9d9f0e5 3164f14 9d9f0e5 3164f14 9d9f0e5 3164f14 9d9f0e5 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 | """Journal d'annotations collaboratif (append-only JSONL, partagé entre sessions).
Source de vérité des avis et, en aval, des datasets SFT/DPO. Délibérément
découplé du gr.State par session : les handlers lisent l'état pipeline pour
construire les événements qu'ils y ajoutent."""
from __future__ import annotations
import json
from copy import deepcopy
from datetime import datetime, timezone
from typing import Any
import gradio as gr
from config import ANNOTATIONS_PATH, METADATA_FIELDS, _ANNOTATIONS_LOCK
from metadata import corrected_from_fields, extract_people, normalize_metadata
from render import pipeline_summary, render_idref_table, status_card
from state import add_event, as_json, empty_state
# --------------------------------------------------------------------------
# Collaborative annotation store (append-only JSONL, shared across sessions).
# This is the source of truth for notation verdicts and, downstream, the
# SFT/DPO training datasets. It is intentionally decoupled from the per-session
# gr.State pipeline: the only coupling is that step handlers *read* pipeline
# state to build the payloads they append here.
# --------------------------------------------------------------------------
# Event schema (one JSON object per line in ANNOTATIONS_PATH):
# ts ISO-8601 UTC timestamp (stamped by append_event)
# annotator free-text reviewer name
# image_url MinIO title-page URL — the primary key everything folds onto
# action "extract" | "validate" | "note"
# step pipeline step the action concerns: vlm|sudoc|idref|dewey|draft
# verdict "ok" | "corrected" | "ko" (validate ; les notes des étapes 3-6
# n'utilisent que ok/ko ; "todo" subsiste dans les événements
# historiques uniquement)
# remark free-text reviewer comment (validate / note)
# doc_type "these" | "memoire" (extract)
# provider,model,prompt (extract)
# raw_json the model's raw extraction dict (extract)
# validated_json the human-accepted dict (validate on step "vlm")
# correction_provider, correction_model (validate on step "vlm", optional:
# set when the corrected JSON came from the tab-2 LLM-assisted
# correction — synthetic gold — rather than manual edits only)
def utc_now_iso() -> str:
return datetime.now(timezone.utc).isoformat()
def append_event(event: dict[str, Any]) -> dict[str, Any]:
"""Stamp and durably append one annotation event. Concurrency-safe: a single
process-wide lock serialises the append, and a lone `\\n`-terminated JSON line
is small enough to write atomically. Returns the stamped event."""
event = {"ts": utc_now_iso(), **event}
line = json.dumps(event, ensure_ascii=False)
with _ANNOTATIONS_LOCK:
with ANNOTATIONS_PATH.open("a", encoding="utf-8") as fh:
fh.write(line + "\n")
return event
def read_events() -> list[dict[str, Any]]:
"""Read the full event log, tolerating a missing file or partial trailing line."""
if not ANNOTATIONS_PATH.exists():
return []
events: list[dict[str, Any]] = []
with _ANNOTATIONS_LOCK:
raw = ANNOTATIONS_PATH.read_text(encoding="utf-8")
for line in raw.splitlines():
line = line.strip()
if not line:
continue
try:
events.append(json.loads(line))
except json.JSONDecodeError:
# Skip a torn line (e.g. a crash mid-write) rather than fail the fold.
continue
return events
def parse_minio_url(url: str) -> tuple[str, str]:
"""Return (collection, doc_type) parsed from a MinIO title-page URL.
e.g. .../images/theses/CESOR/doc_05/... -> ("CESOR", "these").
Falls back to ("", "") when the path does not match the expected layout."""
parts = [p for p in str(url or "").split("/") if p]
for kind, doc_type in (("theses", "these"), ("memoires", "memoire")):
if kind in parts:
idx = parts.index(kind)
collection = parts[idx + 1] if idx + 1 < len(parts) else ""
return collection, doc_type
return "", ""
def _blank_record(url: str) -> dict[str, Any]:
collection, doc_type = parse_minio_url(url)
return {
"image_url": url,
"collection": collection,
"doc_type": doc_type,
"last_extraction": None, # {provider, model, prompt, raw_json, doc_type, ts}
"steps": {}, # step -> {verdict, remark, annotator, ts}
"validated_json": None, # latest human-accepted VLM extraction
"raw_json": None, # raw extraction the validation refers to
"annotators": [],
"last_ts": None,
}
def fold_log(events: list[dict[str, Any]] | None = None) -> dict[str, dict[str, Any]]:
"""Fold the append-only log into current per-image state keyed by image_url.
Later events win for a given (image_url, step); extractions keep only the most
recent raw output. Derived flags (sft_ready, dpo) are computed by consumers."""
events = events if events is not None else read_events()
records: dict[str, dict[str, Any]] = {}
for ev in events:
url = ev.get("image_url")
if not url:
continue
rec = records.setdefault(url, _blank_record(url))
annotator = ev.get("annotator")
if annotator and annotator not in rec["annotators"]:
rec["annotators"].append(annotator)
rec["last_ts"] = ev.get("ts") or rec["last_ts"]
action = ev.get("action")
if action == "extract":
rec["last_extraction"] = {
"provider": ev.get("provider"),
"model": ev.get("model"),
"prompt": ev.get("prompt"),
"raw_json": ev.get("raw_json"),
"doc_type": ev.get("doc_type"),
"ts": ev.get("ts"),
}
if ev.get("doc_type"):
rec["doc_type"] = ev["doc_type"]
rec["raw_json"] = ev.get("raw_json")
elif action in ("validate", "note"):
step = ev.get("step") or "vlm"
rec["steps"][step] = {
"verdict": ev.get("verdict"),
"remark": ev.get("remark") or "",
"annotator": annotator or "",
"ts": ev.get("ts"),
}
if action == "validate" and step == "vlm":
rec["steps"][step]["correction_provider"] = ev.get("correction_provider")
rec["steps"][step]["correction_model"] = ev.get("correction_model")
if ev.get("validated_json") is not None:
rec["validated_json"] = ev.get("validated_json")
if ev.get("raw_json") is not None:
rec["raw_json"] = ev.get("raw_json")
return records
# Verdicts de validation (onglet 2, bouton unique) — deux issues seulement :
# ok extraction correcte telle quelle -> SFT, jamais de paire DPO
# corrected extraction corrigée à la main -> SFT (gold = champs corrigés)
# + paire DPO (rejected = extraction brute)
# (ko subsiste pour les notes qualité des onglets 3-6 ; ko/todo subsistent dans
# le journal pour les événements historiques, mais ne sont plus proposés à la
# validation ni — pour todo — à la notation.)
VALIDATED_VERDICTS = {"ok", "corrected"}
def _comparable(field: str, value: Any) -> Any:
"""Canonicalise une valeur pour la comparaison brut/corrigé : "" vaut null,
l'année est comparée en chaîne (les modèles renvoient tantôt 2015 tantôt
"2015", le formulaire re-parse en entier)."""
if value in ("", None):
return None
if field == "defense_year":
return str(value)
return value
def extraction_differs(extracted: dict[str, Any] | None, corrected: dict[str, Any] | None) -> bool:
"""Compare l'extraction brute et le JSON corrigé dans l'espace du schéma
normalisé, hors `confidence` (non-vérité-terrain) : les différences dues à la
seule normalisation (pipe -> tableau, int -> str, "" -> null) ne comptent pas."""
left = normalize_metadata(extracted or {})
right = normalize_metadata(corrected or {})
return any(
_comparable(f, left.get(f)) != _comparable(f, right.get(f))
for f in METADATA_FIELDS
if f != "confidence"
)
def validate_extraction(
state: dict[str, Any] | None,
annotator: str,
verdict: str,
remark: str,
*field_values: str,
) -> tuple[dict[str, Any], str, str, list[list[str]], str]:
"""Bouton unique « Enregistrer la validation » (onglet 2) : reflète le
formulaire dans state["vlm"]["corrected"] (transmis aux onglets suivants),
journalise le verdict, et alimente les datasets — "ok" (extraction correcte)
→ SFT sans paire DPO ; "corrected" → SFT (gold = champs corrigés) + paire DPO
(rejected = extraction brute). Retourne (state, JSON corrigé, carte de statut,
lignes personnes pour l'onglet 4, résumé pipeline)."""
state = deepcopy(state or empty_state())
vlm = state.get("vlm") or {}
extracted = vlm.get("extracted")
if extracted is None:
raise gr.Error("Lancez d'abord une extraction VLM (onglet 1).")
inp = state.get("input") or {}
url = (inp.get("image_url") or "").strip()
if not url:
raise gr.Error("Une URL d'image est requise pour enregistrer une validation exploitable dans le dataset.")
corrected = corrected_from_fields(*field_values)
state.setdefault("vlm", {})["corrected"] = corrected
verdict = (verdict or "ok").strip().lower()
differs = extraction_differs(extracted, corrected)
# Trace synthetic gold: if the corrected JSON was seeded by the tab-2 LLM-assisted
# correction, record which provider/model produced it (human-reviewed afterwards).
correction = vlm.get("correction") or {}
event = {
"annotator": (annotator or "").strip(),
"image_url": url,
"action": "validate",
"step": "vlm",
"verdict": verdict,
"remark": (remark or "").strip(),
"validated_json": corrected,
"raw_json": extracted,
"doc_type": inp.get("doc_type"),
"provider": inp.get("provider"),
"model": inp.get("model"),
}
if correction.get("provider") and not correction.get("error"):
event["correction_provider"] = correction.get("provider")
event["correction_model"] = correction.get("model")
append_event(event)
if verdict == "corrected":
message = "Extraction corrigée validée → dataset SFT (gold = champs corrigés)"
tone = "ok"
if differs:
message += " + paire DPO (rejected = extraction brute)."
else:
message += (
". ⚠ Aucune différence détectée avec l'extraction brute : "
"aucune paire DPO ne sera générée."
)
tone = "warn"
else:
message = "Extraction validée telle quelle → dataset SFT (pas de paire DPO)."
tone = "ok"
if differs:
message += (
" ⚠ Les champs diffèrent pourtant de l'extraction brute : si vous "
"avez corrigé, ré-enregistrez avec « Extraction corrigée » pour "
"générer la paire DPO."
)
tone = "warn"
message += " Métadonnées transmises aux onglets suivants (Sudoc, IdRef, Dewey)."
state.setdefault("notes", {})["vlm"] = verdict
state = add_event(state, "metadata", "ok", f"Validation enregistrée ({verdict}).")
return (
state,
as_json(corrected),
status_card("Validation", tone, message),
render_idref_table(people=extract_people(corrected)),
pipeline_summary(state),
)
STEP_NOTE_LABELS = {"sudoc": "Sudoc", "idref": "IdRef", "dewey": "Dewey", "draft": "Brouillon"}
# Onglet où se note chaque étape (rappels « avis non enregistré » au changement
# d'onglet ; "vlm" = la validation de l'onglet 2, qui alimente les datasets).
STEP_TAB_NUMBERS = {"vlm": 2, "sudoc": 3, "idref": 4, "dewey": 5, "draft": 6}
def note_button_update(step: str, mode: str):
"""gr.update du bouton de notation d'une étape : "pending" = résultat non
noté (bouton primaire, libellé alerte) ; "saved" = avis enregistré
(ré-enregistrer) ; "idle" = libellé de base."""
label = STEP_NOTE_LABELS.get(step, step)
if mode == "pending":
return gr.update(value=f"⚠ Enregistrer l'avis {label} (non enregistré)", variant="primary")
if mode == "saved":
return gr.update(value=f"Ré-enregistrer l'avis {label}", variant="secondary")
return gr.update(value=f"Enregistrer l'avis {label}", variant="secondary")
def note_pending_updates(step: str):
"""Chaîné (.success, sans état) après l'exécution d'une étape : passe le
bouton de notation en alerte et la carte en orange tant que l'avis n'est pas
enregistré. Le drapeau state["notes"][step] = "pending" est posé par le
handler de l'étape lui-même (pas ici — pas d'aller-retour d'état)."""
label = STEP_NOTE_LABELS.get(step, step)
return (
note_button_update(step, "pending"),
status_card(
f"Notation {label}", "warn",
"Résultat non noté — choisissez OK/KO puis enregistrez votre avis.",
),
)
def warn_pending_notes(state: dict[str, Any] | None, upto_tab: int) -> None:
"""À la sélection d'un onglet : toast non bloquant listant les avis encore
non enregistrés sur les étapes des onglets précédents (state["notes"])."""
notes = (state or {}).get("notes") or {}
pending = [
step for step, tab in STEP_TAB_NUMBERS.items()
if tab < upto_tab and notes.get(step) == "pending"
]
if not pending:
return
parts = [
"validation de l'extraction (onglet 2)" if step == "vlm"
else f"avis {STEP_NOTE_LABELS[step]} (onglet {STEP_TAB_NUMBERS[step]})"
for step in pending
]
gr.Warning("⚠ Non enregistré : " + " ; ".join(parts) + ".")
def save_step_note(state: dict[str, Any] | None, annotator: str, verdict: str, remark: str, step: str):
"""Append a collaborative note on a pipeline step (Sudoc/IdRef/Dewey/Brouillon).
Unlike the VLM verdict, these are pure quality signal and do not feed the datasets.
Retourne (state — drapeau de notation levé, bouton repassé en mode « saved »,
carte de statut)."""
state = deepcopy(state or empty_state())
url = ((state.get("input") or {}).get("image_url") or "").strip()
if not url:
raise gr.Error("Une URL d'image est requise pour enregistrer un avis (lancez d'abord l'extraction de l'onglet 1).")
verdict = (verdict or "ok").strip().lower()
append_event({
"annotator": (annotator or "").strip(),
"image_url": url,
"action": "note",
"step": step,
"verdict": verdict,
"remark": (remark or "").strip(),
})
state.setdefault("notes", {})[step] = verdict
labels = {"ok": "OK", "ko": "KO"}
tone = {"ok": "ok", "ko": "error"}
step_label = STEP_NOTE_LABELS.get(step, step)
return (
state,
note_button_update(step, "saved"),
status_card(f"Notation {step_label}", tone.get(verdict, "idle"), f"Avis « {labels.get(verdict, verdict)} » enregistré."),
pipeline_summary(state),
)
|