Spaces:
Build error
Build error
| """ | |
| ExΓ©cution rΓ©elle du pipeline sur un job page (Sprint 4 β Session C). | |
| Point d'entrΓ©e principal : execute_page_job(job_id) | |
| SΓ©quence stricte (CLAUDE.md Β§8 pipeline) : | |
| 1. job β RUNNING | |
| 2. Charger page / manuscrit / corpus depuis BDD | |
| 3. Charger CorpusProfile depuis profiles/{profile_id}.json | |
| 4. Charger ModelConfig depuis BDD (erreur explicite si absent) | |
| 5. fetch_and_normalize() β ImageDerivativeInfo | |
| 6. run_primary_analysis() β PageMaster (+ double stockage R05) | |
| 7. generate_alto() + write_alto() | |
| 8. page.processing_status β ANALYZED | |
| 9. job β DONE | |
| Sur toute exception : job β FAILED + error_message, page β ERROR. | |
| Aucun Γ©chec silencieux (CLAUDE.md Β§7). | |
| """ | |
| # 1. stdlib | |
| import asyncio | |
| import json | |
| import logging | |
| from datetime import datetime, timezone | |
| from pathlib import Path | |
| # 2. third-party | |
| from sqlalchemy.ext.asyncio import AsyncSession | |
| # 3. local | |
| from app import config as _config_module | |
| from app.models.corpus import CorpusModel, ManuscriptModel, PageModel | |
| from app.models.database import async_session_factory | |
| from app.models.job import JobModel | |
| from app.models.model_config_db import ModelConfigDB | |
| from app.schemas.corpus_profile import CorpusProfile | |
| from app.schemas.image import ImageSourceInfo | |
| from app.schemas.model_config import ModelConfig, ProviderType | |
| from app.services.image.normalizer import ( | |
| create_derivatives, | |
| fetch_ai_derivative_bytes, | |
| fetch_and_normalize, | |
| ) | |
| logger = logging.getLogger(__name__) | |
| # Racine du projet β rΓ©solue via .resolve() pour garantir un chemin absolu | |
| # mΓͺme si __file__ est relatif au CWD (comportement variable selon l'environnement | |
| # d'exΓ©cution : Docker, HuggingFace Spaces, tests...). | |
| # job_runner.py est Γ backend/app/services/job_runner.py β 4 parents β racine. | |
| _PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent.parent | |
| # ββ Point d'entrΓ©e public ββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def execute_page_job(job_id: str, db: AsyncSession | None = None) -> None: | |
| """BackgroundTask : exΓ©cute le pipeline complet sur une page. | |
| Args: | |
| job_id: identifiant du JobModel en BDD. | |
| db: session optionnelle (fournie dans les tests ; None = nouvelle session). | |
| """ | |
| if db is None: | |
| async with async_session_factory() as session: | |
| await _run_job_impl(job_id, session) | |
| else: | |
| await _run_job_impl(job_id, db) | |
| # ββ ImplΓ©mentation interne (testable directement) ββββββββββββββββββββββββββ | |
| async def _run_job_impl(job_id: str, db: AsyncSession) -> None: | |
| """ExΓ©cution du pipeline sur un job, avec la session fournie. | |
| ExposΓ© (prΓ©fixe _ conservΓ©) pour les tests unitaires. | |
| """ | |
| # ββ 1. Charger le job, passer status β RUNNING ββββββββββββββββββββββββββ | |
| job = await db.get(JobModel, job_id) | |
| if job is None: | |
| logger.error("Job introuvable β exΓ©cution abandonnΓ©e", extra={"job_id": job_id}) | |
| return | |
| job.status = "running" | |
| job.started_at = datetime.now(timezone.utc) | |
| await db.commit() | |
| page: PageModel | None = None | |
| try: | |
| # ββ 2. Charger page / manuscrit / corpus βββββββββββββββββββββββββββββ | |
| if job.page_id is None: | |
| raise ValueError("Ce job n'a pas de page_id β impossible d'exΓ©cuter le pipeline") | |
| page = await db.get(PageModel, job.page_id) | |
| if page is None: | |
| raise ValueError(f"Page introuvable en BDD : {job.page_id}") | |
| manuscript = await db.get(ManuscriptModel, page.manuscript_id) | |
| if manuscript is None: | |
| raise ValueError(f"Manuscrit introuvable en BDD : {page.manuscript_id}") | |
| corpus = await db.get(CorpusModel, manuscript.corpus_id) | |
| if corpus is None: | |
| raise ValueError(f"Corpus introuvable en BDD : {manuscript.corpus_id}") | |
| # ββ 3. Charger le CorpusProfile ββββββββββββββββββββββββββββββββββββββ | |
| # settings.profiles_dir est la source canonique du chemin (config.py). | |
| # RΓ©solu depuis PROFILES_DIR en Docker, ou _REPO_ROOT/profiles en local. | |
| profile_path = _config_module.settings.profiles_dir / f"{corpus.profile_id}.json" | |
| if not profile_path.exists(): | |
| raise FileNotFoundError( | |
| f"Fichier de profil introuvable : {profile_path}. " | |
| f"Profil attendu : Β«{corpus.profile_id}Β»" | |
| ) | |
| profile_data = json.loads(profile_path.read_text(encoding="utf-8")) | |
| corpus_profile = CorpusProfile.model_validate(profile_data) | |
| # ββ 4. Charger le ModelConfig (erreur explicite si absent) βββββββββββ | |
| model_db = await db.get(ModelConfigDB, corpus.id) | |
| if model_db is None: | |
| raise ValueError( | |
| f"Aucun modèle IA configuré pour le corpus «{corpus.id}». " | |
| "Sélectionnez un modèle via PUT /api/v1/corpora/{id}/model avant " | |
| "de lancer le pipeline." | |
| ) | |
| model_config = ModelConfig( | |
| corpus_id=corpus.id, | |
| selected_model_id=model_db.selected_model_id, | |
| selected_model_display_name=model_db.selected_model_display_name, | |
| provider=ProviderType(model_db.provider_type), | |
| supports_vision=model_db.supports_vision, | |
| last_fetched_at=model_db.updated_at, | |
| ) | |
| # ββ 5. Obtenir l'image pour l'IA βββββββββββββββββββββββββββββββββββββ | |
| data_dir = _config_module.settings.data_dir | |
| image_source = page.image_master_path or "" | |
| from app.services.ai.analyzer import run_primary_analysis | |
| if page.iiif_service_url: | |
| # ββ Mode IIIF natif : fetch en mΓ©moire, zΓ©ro stockage ββββββββββββ | |
| deriv_bytes, deriv_w, deriv_h = await asyncio.to_thread( | |
| fetch_ai_derivative_bytes, | |
| iiif_service_url=page.iiif_service_url, | |
| fallback_url=None, | |
| ) | |
| image_source_info = ImageSourceInfo( | |
| original_url=image_source or page.iiif_service_url, | |
| iiif_service_url=page.iiif_service_url, | |
| manifest_url=page.manifest_url, | |
| is_iiif=True, | |
| original_width=page.canvas_width or deriv_w, | |
| original_height=page.canvas_height or deriv_h, | |
| ) | |
| # ββ 6. Analyse primaire IA (R05 : double stockage) βββββββββββββββ | |
| page_master = await asyncio.to_thread( | |
| run_primary_analysis, | |
| derivative_image_bytes=deriv_bytes, | |
| derivative_width=deriv_w, | |
| derivative_height=deriv_h, | |
| corpus_profile=corpus_profile, | |
| model_config=model_config, | |
| page_id=page.id, | |
| manuscript_id=manuscript.id, | |
| corpus_slug=corpus.slug, | |
| folio_label=page.folio_label, | |
| sequence=page.sequence, | |
| image_info=image_source_info, | |
| base_data_dir=data_dir, | |
| project_root=_PROJECT_ROOT, | |
| ) | |
| elif image_source.startswith(("http://", "https://")): | |
| # ββ Mode fallback URL : tΓ©lΓ©charge + stocke sur disque (legacy) ββ | |
| image_info = await asyncio.to_thread( | |
| fetch_and_normalize, | |
| image_source, corpus.slug, page.folio_label, data_dir | |
| ) | |
| page_master = await asyncio.to_thread( | |
| run_primary_analysis, | |
| derivative_image_path=Path(image_info.derivative_path), | |
| corpus_profile=corpus_profile, | |
| model_config=model_config, | |
| page_id=page.id, | |
| manuscript_id=manuscript.id, | |
| corpus_slug=corpus.slug, | |
| folio_label=page.folio_label, | |
| sequence=page.sequence, | |
| image_info=image_info, | |
| base_data_dir=data_dir, | |
| project_root=_PROJECT_ROOT, | |
| ) | |
| elif image_source: | |
| # ββ Mode fichier local (upload) ββββββββββββββββββββββββββββββββββ | |
| source_path = Path(image_source).resolve() | |
| data_dir_resolved = data_dir.resolve() | |
| if not source_path.is_relative_to(data_dir_resolved): | |
| raise ValueError( | |
| f"Chemin image hors du rΓ©pertoire de donnΓ©es interdit : " | |
| f"{image_source!r} (rΓ©solu : {source_path})" | |
| ) | |
| source_bytes = source_path.read_bytes() | |
| image_info = await asyncio.to_thread( | |
| create_derivatives, | |
| source_bytes, image_source, corpus.slug, page.folio_label, data_dir | |
| ) | |
| page_master = await asyncio.to_thread( | |
| run_primary_analysis, | |
| derivative_image_path=Path(image_info.derivative_path), | |
| corpus_profile=corpus_profile, | |
| model_config=model_config, | |
| page_id=page.id, | |
| manuscript_id=manuscript.id, | |
| corpus_slug=corpus.slug, | |
| folio_label=page.folio_label, | |
| sequence=page.sequence, | |
| image_info=image_info, | |
| base_data_dir=data_dir, | |
| project_root=_PROJECT_ROOT, | |
| ) | |
| else: | |
| raise ValueError( | |
| f"La page {page.id} n'a pas d'image source " | |
| "(ni iiif_service_url, ni image_master_path)" | |
| ) | |
| # ββ 6b. Index pour la recherche βββββββββββββββββββββββββββββββββββββ | |
| from app.services.search.indexer import index_page | |
| await index_page(db, page_master) | |
| # ββ 7. GΓ©nΓ©rer et Γ©crire l'ALTO XML ββββββββββββββββββββββββββββββββββ | |
| from app.services.export.alto import generate_alto, write_alto | |
| alto_xml = generate_alto(page_master) | |
| alto_path = ( | |
| data_dir | |
| / "corpora" | |
| / corpus.slug | |
| / "pages" | |
| / page.folio_label | |
| / "alto.xml" | |
| ) | |
| write_alto(alto_xml, alto_path) | |
| # ββ 8. Page β ANALYZED βββββββββββββββββββββββββββββββββββββββββββββββ | |
| page.processing_status = "ANALYZED" | |
| if page_master.ocr is not None: | |
| page.confidence_summary = page_master.ocr.confidence | |
| # ββ 9. Job β DONE ββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| job.status = "done" | |
| job.finished_at = datetime.now(timezone.utc) | |
| await db.commit() | |
| logger.info( | |
| "Job terminé avec succès", | |
| extra={ | |
| "job_id": job_id, | |
| "page_id": page.id, | |
| "folio": page.folio_label, | |
| "corpus": corpus.slug, | |
| }, | |
| ) | |
| except Exception as exc: | |
| logger.error( | |
| "Γchec du job", | |
| extra={"job_id": job_id, "error": str(exc)}, | |
| exc_info=True, | |
| ) | |
| job.status = "failed" | |
| job.error_message = str(exc) | |
| job.finished_at = datetime.now(timezone.utc) | |
| if page is not None: | |
| page.processing_status = "ERROR" | |
| try: | |
| await db.commit() | |
| except Exception as commit_exc: | |
| logger.error( | |
| "Impossible de persister l'Γ©tat d'Γ©chec du job", | |
| extra={"job_id": job_id, "commit_error": str(commit_exc)}, | |
| ) | |