Spaces:
Build error
fix(backend): Sprint Fix 2 — data integrity, error handling, robustness
Browse files11 issues fixed:
- #9: Validate editorial_status against EditorialStatus enum (pages.py)
- #10: Distinguish "not_configured" vs "error" in provider status, add
error_detail field (model_registry.py)
- #11: Add ondelete="CASCADE"/"SET NULL" to all ForeignKey declarations +
enable PRAGMA foreign_keys=ON for SQLite (corpus.py, job.py,
model_config_db.py, database.py)
- #12: Wrap _write_master in try/except OSError → HTTP 500 (pages.py)
- #19: Catch IntegrityError on ingest commit → HTTP 409 Conflict (ingest.py)
- #20: Fix TOCTOU in corpus_runner: claim pending jobs atomically before
executing, preventing double-execution under concurrency
- #21: Reduce httpx image fetch timeout 60s→30s, add connect timeout 10s
(iiif_fetcher.py)
- #22: Validate master.json structure in search before use (search.py)
- #23: Document return type pattern in jobs.py (ORM→response_model)
- #26: Remove stale provider cache — providers rebuilt on each call
(model_registry.py)
- #34: Validate profile_id exists in profiles_dir at corpus creation
(corpora.py)
Updated tests for timeout change, corpus_runner mock, and profile validation.
563 tests pass, 0 regressions.
https://claude.ai/code/session_01UB4he7RdRPHLvNjky4X8Sw
- backend/app/api/v1/corpora.py +10 -0
- backend/app/api/v1/ingest.py +17 -2
- backend/app/api/v1/jobs.py +1 -1
- backend/app/api/v1/pages.py +13 -3
- backend/app/api/v1/search.py +5 -0
- backend/app/models/corpus.py +2 -2
- backend/app/models/database.py +8 -0
- backend/app/models/job.py +2 -2
- backend/app/models/model_config_db.py +1 -1
- backend/app/services/ai/model_registry.py +19 -13
- backend/app/services/corpus_runner.py +9 -3
- backend/app/services/ingest/iiif_fetcher.py +7 -2
- backend/tests/test_image_pipeline.py +8 -13
- backend/tests/test_job_runner.py +7 -3
- backend/tests/test_security.py +4 -4
|
@@ -74,6 +74,16 @@ async def create_corpus(
|
|
| 74 |
body: CorpusCreate, db: AsyncSession = Depends(get_db)
|
| 75 |
) -> CorpusModel:
|
| 76 |
"""Crée un nouveau corpus. Le slug doit être unique."""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 77 |
# Vérifier unicité du slug
|
| 78 |
existing = await db.execute(
|
| 79 |
select(CorpusModel).where(CorpusModel.slug == body.slug)
|
|
|
|
| 74 |
body: CorpusCreate, db: AsyncSession = Depends(get_db)
|
| 75 |
) -> CorpusModel:
|
| 76 |
"""Crée un nouveau corpus. Le slug doit être unique."""
|
| 77 |
+
# Vérifier que le profil existe dans profiles_dir
|
| 78 |
+
from app.config import settings as _settings
|
| 79 |
+
|
| 80 |
+
profile_path = _settings.profiles_dir / f"{body.profile_id}.json"
|
| 81 |
+
if not profile_path.is_file():
|
| 82 |
+
raise HTTPException(
|
| 83 |
+
status_code=422,
|
| 84 |
+
detail=f"Profil «{body.profile_id}» introuvable dans {_settings.profiles_dir}",
|
| 85 |
+
)
|
| 86 |
+
|
| 87 |
# Vérifier unicité du slug
|
| 88 |
existing = await db.execute(
|
| 89 |
select(CorpusModel).where(CorpusModel.slug == body.slug)
|
|
@@ -20,6 +20,7 @@ import httpx
|
|
| 20 |
from fastapi import APIRouter, Depends, File, HTTPException, UploadFile
|
| 21 |
from pydantic import BaseModel, Field
|
| 22 |
from sqlalchemy import func, select
|
|
|
|
| 23 |
from sqlalchemy.ext.asyncio import AsyncSession
|
| 24 |
|
| 25 |
# 3. local
|
|
@@ -389,7 +390,14 @@ async def ingest_iiif_manifest(
|
|
| 389 |
created.append(page)
|
| 390 |
|
| 391 |
ms.total_pages = (ms.total_pages or 0) + len(created)
|
| 392 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 393 |
|
| 394 |
logger.info(
|
| 395 |
"Manifest IIIF ingéré",
|
|
@@ -443,7 +451,14 @@ async def ingest_iiif_images(
|
|
| 443 |
created.append(page)
|
| 444 |
|
| 445 |
ms.total_pages = (ms.total_pages or 0) + len(created)
|
| 446 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 447 |
|
| 448 |
logger.info(
|
| 449 |
"Images IIIF ingérées",
|
|
|
|
| 20 |
from fastapi import APIRouter, Depends, File, HTTPException, UploadFile
|
| 21 |
from pydantic import BaseModel, Field
|
| 22 |
from sqlalchemy import func, select
|
| 23 |
+
from sqlalchemy.exc import IntegrityError
|
| 24 |
from sqlalchemy.ext.asyncio import AsyncSession
|
| 25 |
|
| 26 |
# 3. local
|
|
|
|
| 390 |
created.append(page)
|
| 391 |
|
| 392 |
ms.total_pages = (ms.total_pages or 0) + len(created)
|
| 393 |
+
try:
|
| 394 |
+
await db.commit()
|
| 395 |
+
except IntegrityError:
|
| 396 |
+
await db.rollback()
|
| 397 |
+
raise HTTPException(
|
| 398 |
+
status_code=409,
|
| 399 |
+
detail="Conflit : certaines pages existent déjà (ingestion concurrente probable)",
|
| 400 |
+
)
|
| 401 |
|
| 402 |
logger.info(
|
| 403 |
"Manifest IIIF ingéré",
|
|
|
|
| 451 |
created.append(page)
|
| 452 |
|
| 453 |
ms.total_pages = (ms.total_pages or 0) + len(created)
|
| 454 |
+
try:
|
| 455 |
+
await db.commit()
|
| 456 |
+
except IntegrityError:
|
| 457 |
+
await db.rollback()
|
| 458 |
+
raise HTTPException(
|
| 459 |
+
status_code=409,
|
| 460 |
+
detail="Conflit : certaines pages existent déjà (ingestion concurrente probable)",
|
| 461 |
+
)
|
| 462 |
|
| 463 |
logger.info(
|
| 464 |
"Images IIIF ingérées",
|
|
@@ -115,7 +115,7 @@ async def run_page(
|
|
| 115 |
page_id: str,
|
| 116 |
background_tasks: BackgroundTasks,
|
| 117 |
db: AsyncSession = Depends(get_db),
|
| 118 |
-
) -> JobModel:
|
| 119 |
"""Lance le pipeline sur une seule page.
|
| 120 |
|
| 121 |
Crée un JobModel (status=pending) et délègue l'exécution à
|
|
|
|
| 115 |
page_id: str,
|
| 116 |
background_tasks: BackgroundTasks,
|
| 117 |
db: AsyncSession = Depends(get_db),
|
| 118 |
+
) -> JobModel: # FastAPI sérialise via response_model=JobResponse
|
| 119 |
"""Lance le pipeline sur une seule page.
|
| 120 |
|
| 121 |
Crée un JobModel (status=pending) et délègue l'exécution à
|
|
@@ -27,7 +27,7 @@ from app.models.corpus import CorpusModel, ManuscriptModel, PageModel
|
|
| 27 |
from app.models.database import get_db
|
| 28 |
from app.schemas.annotation import LayerStatus
|
| 29 |
from app.schemas.corpus_profile import LayerType
|
| 30 |
-
from app.schemas.page_master import PageMaster
|
| 31 |
|
| 32 |
logger = logging.getLogger(__name__)
|
| 33 |
router = APIRouter(prefix="/pages", tags=["pages"])
|
|
@@ -43,7 +43,7 @@ class CorrectionsRequest(BaseModel):
|
|
| 43 |
"""
|
| 44 |
|
| 45 |
ocr_diplomatic_text: str | None = Field(None, max_length=500_000)
|
| 46 |
-
editorial_status:
|
| 47 |
commentary_public: str | None = Field(None, max_length=100_000)
|
| 48 |
commentary_scholarly: str | None = Field(None, max_length=100_000)
|
| 49 |
region_validations: dict[str, str] | None = None
|
|
@@ -356,7 +356,17 @@ async def apply_corrections(
|
|
| 356 |
except ValidationError as exc:
|
| 357 |
raise HTTPException(status_code=422, detail=str(exc)) from exc
|
| 358 |
|
| 359 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 360 |
logger.info(
|
| 361 |
"Corrections appliquées",
|
| 362 |
extra={"page_id": page_id, "version": new_master.editorial.version},
|
|
|
|
| 27 |
from app.models.database import get_db
|
| 28 |
from app.schemas.annotation import LayerStatus
|
| 29 |
from app.schemas.corpus_profile import LayerType
|
| 30 |
+
from app.schemas.page_master import EditorialStatus, PageMaster
|
| 31 |
|
| 32 |
logger = logging.getLogger(__name__)
|
| 33 |
router = APIRouter(prefix="/pages", tags=["pages"])
|
|
|
|
| 43 |
"""
|
| 44 |
|
| 45 |
ocr_diplomatic_text: str | None = Field(None, max_length=500_000)
|
| 46 |
+
editorial_status: EditorialStatus | None = None
|
| 47 |
commentary_public: str | None = Field(None, max_length=100_000)
|
| 48 |
commentary_scholarly: str | None = Field(None, max_length=100_000)
|
| 49 |
region_validations: dict[str, str] | None = None
|
|
|
|
| 356 |
except ValidationError as exc:
|
| 357 |
raise HTTPException(status_code=422, detail=str(exc)) from exc
|
| 358 |
|
| 359 |
+
try:
|
| 360 |
+
_write_master(page_dir, new_master)
|
| 361 |
+
except OSError as exc:
|
| 362 |
+
logger.error(
|
| 363 |
+
"Échec écriture master.json",
|
| 364 |
+
extra={"page_id": page_id, "error": str(exc)},
|
| 365 |
+
)
|
| 366 |
+
raise HTTPException(
|
| 367 |
+
status_code=500,
|
| 368 |
+
detail=f"Impossible d'écrire master.json : {exc}",
|
| 369 |
+
) from exc
|
| 370 |
logger.info(
|
| 371 |
"Corrections appliquées",
|
| 372 |
extra={"page_id": page_id, "version": new_master.editorial.version},
|
|
@@ -117,6 +117,11 @@ async def search_pages(
|
|
| 117 |
except (json.JSONDecodeError, OSError):
|
| 118 |
continue
|
| 119 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 120 |
score, excerpt = _score_master(raw, query_normalized)
|
| 121 |
if score == 0:
|
| 122 |
continue
|
|
|
|
| 117 |
except (json.JSONDecodeError, OSError):
|
| 118 |
continue
|
| 119 |
|
| 120 |
+
# Vérification minimale de la structure attendue
|
| 121 |
+
if not isinstance(raw.get("page_id"), str):
|
| 122 |
+
logger.warning("master.json invalide ignoré : %s", master_path)
|
| 123 |
+
continue
|
| 124 |
+
|
| 125 |
score, excerpt = _score_master(raw, query_normalized)
|
| 126 |
if score == 0:
|
| 127 |
continue
|
|
@@ -44,7 +44,7 @@ class ManuscriptModel(Base):
|
|
| 44 |
|
| 45 |
id: Mapped[str] = mapped_column(String, primary_key=True)
|
| 46 |
corpus_id: Mapped[str] = mapped_column(
|
| 47 |
-
String, ForeignKey("corpora.id"), nullable=False, index=True
|
| 48 |
)
|
| 49 |
shelfmark: Mapped[str | None] = mapped_column(String, nullable=True)
|
| 50 |
title: Mapped[str] = mapped_column(String, nullable=False)
|
|
@@ -69,7 +69,7 @@ class PageModel(Base):
|
|
| 69 |
|
| 70 |
id: Mapped[str] = mapped_column(String, primary_key=True)
|
| 71 |
manuscript_id: Mapped[str] = mapped_column(
|
| 72 |
-
String, ForeignKey("manuscripts.id"), nullable=False, index=True
|
| 73 |
)
|
| 74 |
folio_label: Mapped[str] = mapped_column(String, nullable=False)
|
| 75 |
sequence: Mapped[int] = mapped_column(Integer, nullable=False)
|
|
|
|
| 44 |
|
| 45 |
id: Mapped[str] = mapped_column(String, primary_key=True)
|
| 46 |
corpus_id: Mapped[str] = mapped_column(
|
| 47 |
+
String, ForeignKey("corpora.id", ondelete="CASCADE"), nullable=False, index=True
|
| 48 |
)
|
| 49 |
shelfmark: Mapped[str | None] = mapped_column(String, nullable=True)
|
| 50 |
title: Mapped[str] = mapped_column(String, nullable=False)
|
|
|
|
| 69 |
|
| 70 |
id: Mapped[str] = mapped_column(String, primary_key=True)
|
| 71 |
manuscript_id: Mapped[str] = mapped_column(
|
| 72 |
+
String, ForeignKey("manuscripts.id", ondelete="CASCADE"), nullable=False, index=True
|
| 73 |
)
|
| 74 |
folio_label: Mapped[str] = mapped_column(String, nullable=False)
|
| 75 |
sequence: Mapped[int] = mapped_column(Integer, nullable=False)
|
|
@@ -11,6 +11,7 @@ Les tables sont créées au démarrage de l'application (voir main.py lifespan).
|
|
| 11 |
import logging
|
| 12 |
|
| 13 |
# 2. third-party
|
|
|
|
| 14 |
from sqlalchemy.ext.asyncio import (
|
| 15 |
AsyncSession,
|
| 16 |
async_sessionmaker,
|
|
@@ -29,6 +30,13 @@ engine = create_async_engine(
|
|
| 29 |
connect_args={"check_same_thread": False},
|
| 30 |
)
|
| 31 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 32 |
async_session_factory = async_sessionmaker(
|
| 33 |
engine,
|
| 34 |
expire_on_commit=False,
|
|
|
|
| 11 |
import logging
|
| 12 |
|
| 13 |
# 2. third-party
|
| 14 |
+
from sqlalchemy import event
|
| 15 |
from sqlalchemy.ext.asyncio import (
|
| 16 |
AsyncSession,
|
| 17 |
async_sessionmaker,
|
|
|
|
| 30 |
connect_args={"check_same_thread": False},
|
| 31 |
)
|
| 32 |
|
| 33 |
+
# Activer les clés étrangères SQLite (désactivées par défaut).
|
| 34 |
+
# Nécessaire pour que ondelete="CASCADE" / "SET NULL" fonctionne.
|
| 35 |
+
@event.listens_for(engine.sync_engine, "connect")
|
| 36 |
+
def _set_sqlite_pragma(dbapi_conn, _connection_record):
|
| 37 |
+
cursor = dbapi_conn.execute("PRAGMA foreign_keys=ON")
|
| 38 |
+
cursor.close()
|
| 39 |
+
|
| 40 |
async_session_factory = async_sessionmaker(
|
| 41 |
engine,
|
| 42 |
expire_on_commit=False,
|
|
@@ -28,10 +28,10 @@ class JobModel(Base):
|
|
| 28 |
|
| 29 |
id: Mapped[str] = mapped_column(String, primary_key=True)
|
| 30 |
corpus_id: Mapped[str] = mapped_column(
|
| 31 |
-
String, ForeignKey("corpora.id"), nullable=False, index=True
|
| 32 |
)
|
| 33 |
page_id: Mapped[str | None] = mapped_column(
|
| 34 |
-
String, ForeignKey("pages.id"), nullable=True, index=True
|
| 35 |
)
|
| 36 |
# pending / running / done / failed
|
| 37 |
status: Mapped[str] = mapped_column(String, nullable=False, default="pending")
|
|
|
|
| 28 |
|
| 29 |
id: Mapped[str] = mapped_column(String, primary_key=True)
|
| 30 |
corpus_id: Mapped[str] = mapped_column(
|
| 31 |
+
String, ForeignKey("corpora.id", ondelete="CASCADE"), nullable=False, index=True
|
| 32 |
)
|
| 33 |
page_id: Mapped[str | None] = mapped_column(
|
| 34 |
+
String, ForeignKey("pages.id", ondelete="SET NULL"), nullable=True, index=True
|
| 35 |
)
|
| 36 |
# pending / running / done / failed
|
| 37 |
status: Mapped[str] = mapped_column(String, nullable=False, default="pending")
|
|
@@ -21,7 +21,7 @@ class ModelConfigDB(Base):
|
|
| 21 |
__tablename__ = "model_configs"
|
| 22 |
|
| 23 |
corpus_id: Mapped[str] = mapped_column(
|
| 24 |
-
String, ForeignKey("corpora.id"), primary_key=True
|
| 25 |
)
|
| 26 |
provider_type: Mapped[str] = mapped_column(String, nullable=False)
|
| 27 |
selected_model_id: Mapped[str] = mapped_column(String, nullable=False)
|
|
|
|
| 21 |
__tablename__ = "model_configs"
|
| 22 |
|
| 23 |
corpus_id: Mapped[str] = mapped_column(
|
| 24 |
+
String, ForeignKey("corpora.id", ondelete="CASCADE"), primary_key=True
|
| 25 |
)
|
| 26 |
provider_type: Mapped[str] = mapped_column(String, nullable=False)
|
| 27 |
selected_model_id: Mapped[str] = mapped_column(String, nullable=False)
|
|
@@ -23,27 +23,24 @@ _PROVIDER_DISPLAY_NAMES: dict[ProviderType, str] = {
|
|
| 23 |
}
|
| 24 |
|
| 25 |
|
| 26 |
-
_cached_providers: list[AIProvider] | None = None
|
| 27 |
-
|
| 28 |
-
|
| 29 |
def _build_providers() -> list[AIProvider]:
|
| 30 |
-
"""Construit la liste des providers — imports différés
|
| 31 |
-
global _cached_providers
|
| 32 |
-
if _cached_providers is not None:
|
| 33 |
-
return _cached_providers
|
| 34 |
|
|
|
|
|
|
|
|
|
|
|
|
|
| 35 |
from app.services.ai.provider_google_ai import GoogleAIProvider
|
| 36 |
from app.services.ai.provider_mistral import MistralProvider
|
| 37 |
from app.services.ai.provider_vertex_key import VertexAPIKeyProvider
|
| 38 |
from app.services.ai.provider_vertex_sa import VertexServiceAccountProvider
|
| 39 |
|
| 40 |
-
|
| 41 |
GoogleAIProvider(),
|
| 42 |
VertexAPIKeyProvider(),
|
| 43 |
VertexServiceAccountProvider(),
|
| 44 |
MistralProvider(),
|
| 45 |
]
|
| 46 |
-
return _cached_providers
|
| 47 |
|
| 48 |
|
| 49 |
def get_available_providers() -> list[dict]:
|
|
@@ -59,25 +56,34 @@ def get_available_providers() -> list[dict]:
|
|
| 59 |
"""
|
| 60 |
result: list[dict] = []
|
| 61 |
for provider in _build_providers():
|
| 62 |
-
|
|
|
|
| 63 |
model_count = 0
|
| 64 |
-
|
|
|
|
|
|
|
|
|
|
| 65 |
try:
|
| 66 |
models = provider.list_models()
|
| 67 |
model_count = len(models)
|
|
|
|
|
|
|
| 68 |
except Exception as exc:
|
| 69 |
logger.warning(
|
| 70 |
-
"Provider %s inaccessible : %s",
|
| 71 |
provider.provider_type.value,
|
| 72 |
exc,
|
| 73 |
)
|
| 74 |
-
|
|
|
|
| 75 |
|
| 76 |
result.append({
|
| 77 |
"provider_type": provider.provider_type.value,
|
| 78 |
"display_name": _PROVIDER_DISPLAY_NAMES.get(provider.provider_type, provider.provider_type.value),
|
| 79 |
"available": available,
|
| 80 |
"model_count": model_count,
|
|
|
|
|
|
|
| 81 |
})
|
| 82 |
return result
|
| 83 |
|
|
|
|
| 23 |
}
|
| 24 |
|
| 25 |
|
|
|
|
|
|
|
|
|
|
| 26 |
def _build_providers() -> list[AIProvider]:
|
| 27 |
+
"""Construit la liste des providers — imports différés.
|
|
|
|
|
|
|
|
|
|
| 28 |
|
| 29 |
+
Pas de cache global : la construction est triviale (4 objets légers)
|
| 30 |
+
et l'absence de cache permet de détecter immédiatement les changements
|
| 31 |
+
de variables d'environnement sans redémarrage.
|
| 32 |
+
"""
|
| 33 |
from app.services.ai.provider_google_ai import GoogleAIProvider
|
| 34 |
from app.services.ai.provider_mistral import MistralProvider
|
| 35 |
from app.services.ai.provider_vertex_key import VertexAPIKeyProvider
|
| 36 |
from app.services.ai.provider_vertex_sa import VertexServiceAccountProvider
|
| 37 |
|
| 38 |
+
return [
|
| 39 |
GoogleAIProvider(),
|
| 40 |
VertexAPIKeyProvider(),
|
| 41 |
VertexServiceAccountProvider(),
|
| 42 |
MistralProvider(),
|
| 43 |
]
|
|
|
|
| 44 |
|
| 45 |
|
| 46 |
def get_available_providers() -> list[dict]:
|
|
|
|
| 56 |
"""
|
| 57 |
result: list[dict] = []
|
| 58 |
for provider in _build_providers():
|
| 59 |
+
configured = provider.is_configured()
|
| 60 |
+
available = False
|
| 61 |
model_count = 0
|
| 62 |
+
status = "not_configured"
|
| 63 |
+
error_detail: str | None = None
|
| 64 |
+
|
| 65 |
+
if configured:
|
| 66 |
try:
|
| 67 |
models = provider.list_models()
|
| 68 |
model_count = len(models)
|
| 69 |
+
available = True
|
| 70 |
+
status = "available"
|
| 71 |
except Exception as exc:
|
| 72 |
logger.warning(
|
| 73 |
+
"Provider %s configuré mais inaccessible : %s",
|
| 74 |
provider.provider_type.value,
|
| 75 |
exc,
|
| 76 |
)
|
| 77 |
+
status = "error"
|
| 78 |
+
error_detail = str(exc)
|
| 79 |
|
| 80 |
result.append({
|
| 81 |
"provider_type": provider.provider_type.value,
|
| 82 |
"display_name": _PROVIDER_DISPLAY_NAMES.get(provider.provider_type, provider.provider_type.value),
|
| 83 |
"available": available,
|
| 84 |
"model_count": model_count,
|
| 85 |
+
"status": status,
|
| 86 |
+
"error_detail": error_detail,
|
| 87 |
})
|
| 88 |
return result
|
| 89 |
|
|
@@ -30,15 +30,21 @@ async def execute_corpus_job(corpus_id: str) -> dict:
|
|
| 30 |
Returns:
|
| 31 |
{"total": int, "done": int, "failed": int}
|
| 32 |
"""
|
| 33 |
-
# Collecte
|
|
|
|
|
|
|
| 34 |
async with async_session_factory() as db:
|
| 35 |
result = await db.execute(
|
| 36 |
-
select(JobModel
|
| 37 |
JobModel.corpus_id == corpus_id,
|
| 38 |
JobModel.status == "pending",
|
| 39 |
)
|
| 40 |
)
|
| 41 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 42 |
|
| 43 |
if not job_ids:
|
| 44 |
logger.info(
|
|
|
|
| 30 |
Returns:
|
| 31 |
{"total": int, "done": int, "failed": int}
|
| 32 |
"""
|
| 33 |
+
# Collecte et verrouillage des jobs PENDING : on les passe immédiatement
|
| 34 |
+
# en "claimed" pour éviter qu'un second appel concurrent ne les reprenne
|
| 35 |
+
# (TOCTOU). Chaque job sera ensuite passé en "running" par job_runner.
|
| 36 |
async with async_session_factory() as db:
|
| 37 |
result = await db.execute(
|
| 38 |
+
select(JobModel).where(
|
| 39 |
JobModel.corpus_id == corpus_id,
|
| 40 |
JobModel.status == "pending",
|
| 41 |
)
|
| 42 |
)
|
| 43 |
+
pending_jobs = list(result.scalars().all())
|
| 44 |
+
job_ids: list[str] = [j.id for j in pending_jobs]
|
| 45 |
+
for j in pending_jobs:
|
| 46 |
+
j.status = "claimed"
|
| 47 |
+
await db.commit()
|
| 48 |
|
| 49 |
if not job_ids:
|
| 50 |
logger.info(
|
|
@@ -9,7 +9,7 @@ import httpx
|
|
| 9 |
|
| 10 |
logger = logging.getLogger(__name__)
|
| 11 |
|
| 12 |
-
_DEFAULT_TIMEOUT =
|
| 13 |
|
| 14 |
_HEADERS = {
|
| 15 |
"User-Agent": (
|
|
@@ -36,7 +36,12 @@ def fetch_iiif_image(url: str, timeout: float = _DEFAULT_TIMEOUT) -> bytes:
|
|
| 36 |
httpx.RequestError: pour toute autre erreur réseau.
|
| 37 |
"""
|
| 38 |
logger.info("Fetching IIIF image", extra={"url": url})
|
| 39 |
-
response = httpx.get(
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 40 |
response.raise_for_status()
|
| 41 |
logger.info(
|
| 42 |
"IIIF image fetched",
|
|
|
|
| 9 |
|
| 10 |
logger = logging.getLogger(__name__)
|
| 11 |
|
| 12 |
+
_DEFAULT_TIMEOUT = 30.0 # secondes (connect 10s + read 30s)
|
| 13 |
|
| 14 |
_HEADERS = {
|
| 15 |
"User-Agent": (
|
|
|
|
| 36 |
httpx.RequestError: pour toute autre erreur réseau.
|
| 37 |
"""
|
| 38 |
logger.info("Fetching IIIF image", extra={"url": url})
|
| 39 |
+
response = httpx.get(
|
| 40 |
+
url,
|
| 41 |
+
headers=_HEADERS,
|
| 42 |
+
follow_redirects=True,
|
| 43 |
+
timeout=httpx.Timeout(timeout, connect=10.0),
|
| 44 |
+
)
|
| 45 |
response.raise_for_status()
|
| 46 |
logger.info(
|
| 47 |
"IIIF image fetched",
|
|
@@ -270,18 +270,11 @@ def test_fetch_iiif_image_success():
|
|
| 270 |
result = fetch_iiif_image("https://example.com/image.jpg")
|
| 271 |
|
| 272 |
assert result == fake_bytes
|
| 273 |
-
mock_get.
|
| 274 |
-
|
| 275 |
-
|
| 276 |
-
|
| 277 |
-
|
| 278 |
-
"+https://huggingface.co/spaces/Ma-Ri-Ba-Ku/iiif-studio)"
|
| 279 |
-
),
|
| 280 |
-
"Accept": "image/jpeg,image/png,image/*,*/*",
|
| 281 |
-
},
|
| 282 |
-
follow_redirects=True,
|
| 283 |
-
timeout=60.0,
|
| 284 |
-
)
|
| 285 |
|
| 286 |
|
| 287 |
def test_fetch_iiif_image_http_error():
|
|
@@ -321,7 +314,9 @@ def test_fetch_iiif_image_custom_timeout():
|
|
| 321 |
fetch_iiif_image("https://example.com/img.jpg", timeout=120.0)
|
| 322 |
|
| 323 |
_, kwargs = mock_get.call_args
|
| 324 |
-
|
|
|
|
|
|
|
| 325 |
|
| 326 |
|
| 327 |
# ---------------------------------------------------------------------------
|
|
|
|
| 270 |
result = fetch_iiif_image("https://example.com/image.jpg")
|
| 271 |
|
| 272 |
assert result == fake_bytes
|
| 273 |
+
_, kwargs = mock_get.call_args
|
| 274 |
+
assert kwargs["follow_redirects"] is True
|
| 275 |
+
# Timeout is now an httpx.Timeout object (connect=10s, read=30s)
|
| 276 |
+
assert kwargs["timeout"].connect == 10.0
|
| 277 |
+
assert kwargs["timeout"].read == 30.0
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 278 |
|
| 279 |
|
| 280 |
def test_fetch_iiif_image_http_error():
|
|
|
|
| 314 |
fetch_iiif_image("https://example.com/img.jpg", timeout=120.0)
|
| 315 |
|
| 316 |
_, kwargs = mock_get.call_args
|
| 317 |
+
# Custom timeout wraps in httpx.Timeout(120.0, connect=10.0)
|
| 318 |
+
assert kwargs["timeout"].read == 120.0
|
| 319 |
+
assert kwargs["timeout"].connect == 10.0
|
| 320 |
|
| 321 |
|
| 322 |
# ---------------------------------------------------------------------------
|
|
@@ -528,10 +528,11 @@ async def test_corpus_runner_calls_execute_per_job(monkeypatch):
|
|
| 528 |
async def execute(self, stmt):
|
| 529 |
_call_count[0] += 1
|
| 530 |
if _call_count[0] == 1:
|
| 531 |
-
# Premier appel : retourne les
|
| 532 |
-
|
|
|
|
| 533 |
else:
|
| 534 |
-
# Second appel : retourne les objets JobModel avec statut
|
| 535 |
rows = [_FakeJob("job-alpha", "done"), _FakeJob("job-beta", "done")]
|
| 536 |
|
| 537 |
class _Result:
|
|
@@ -542,6 +543,9 @@ async def test_corpus_runner_calls_execute_per_job(monkeypatch):
|
|
| 542 |
return _Scalars()
|
| 543 |
return _Result()
|
| 544 |
|
|
|
|
|
|
|
|
|
|
| 545 |
def _mock_factory():
|
| 546 |
return _FakeSession()
|
| 547 |
|
|
|
|
| 528 |
async def execute(self, stmt):
|
| 529 |
_call_count[0] += 1
|
| 530 |
if _call_count[0] == 1:
|
| 531 |
+
# Premier appel : retourne les objets JobModel pending
|
| 532 |
+
# (corpus_runner itère maintenant les objets, pas les IDs)
|
| 533 |
+
rows = [_FakeJob("job-alpha", "pending"), _FakeJob("job-beta", "pending")]
|
| 534 |
else:
|
| 535 |
+
# Second appel : retourne les objets JobModel avec statut final
|
| 536 |
rows = [_FakeJob("job-alpha", "done"), _FakeJob("job-beta", "done")]
|
| 537 |
|
| 538 |
class _Result:
|
|
|
|
| 543 |
return _Scalars()
|
| 544 |
return _Result()
|
| 545 |
|
| 546 |
+
async def commit(self):
|
| 547 |
+
pass
|
| 548 |
+
|
| 549 |
def _mock_factory():
|
| 550 |
return _FakeSession()
|
| 551 |
|
|
@@ -142,7 +142,7 @@ async def test_ssrf_localhost(async_client):
|
|
| 142 |
"""Un manifest_url pointant vers localhost doit être rejeté."""
|
| 143 |
# Créer un corpus d'abord
|
| 144 |
create = await async_client.post("/api/v1/corpora", json={
|
| 145 |
-
"slug": "ssrf-test", "title": "SSRF", "profile_id": "
|
| 146 |
})
|
| 147 |
cid = create.json()["id"]
|
| 148 |
|
|
@@ -157,7 +157,7 @@ async def test_ssrf_localhost(async_client):
|
|
| 157 |
async def test_ssrf_metadata_ip(async_client):
|
| 158 |
"""Un manifest_url vers 169.254.x.x (cloud metadata) doit être rejeté."""
|
| 159 |
create = await async_client.post("/api/v1/corpora", json={
|
| 160 |
-
"slug": "ssrf-meta", "title": "SSRF", "profile_id": "
|
| 161 |
})
|
| 162 |
cid = create.json()["id"]
|
| 163 |
|
|
@@ -171,7 +171,7 @@ async def test_ssrf_metadata_ip(async_client):
|
|
| 171 |
async def test_ssrf_file_scheme(async_client):
|
| 172 |
"""Un manifest_url avec file:// doit être rejeté."""
|
| 173 |
create = await async_client.post("/api/v1/corpora", json={
|
| 174 |
-
"slug": "ssrf-file", "title": "SSRF", "profile_id": "
|
| 175 |
})
|
| 176 |
cid = create.json()["id"]
|
| 177 |
|
|
@@ -207,7 +207,7 @@ async def test_search_query_max_length_ok(async_client):
|
|
| 207 |
async def test_model_id_too_long(async_client):
|
| 208 |
"""Un model_id >256 chars doit être rejeté."""
|
| 209 |
create = await async_client.post("/api/v1/corpora", json={
|
| 210 |
-
"slug": "model-test", "title": "T", "profile_id": "
|
| 211 |
})
|
| 212 |
cid = create.json()["id"]
|
| 213 |
|
|
|
|
| 142 |
"""Un manifest_url pointant vers localhost doit être rejeté."""
|
| 143 |
# Créer un corpus d'abord
|
| 144 |
create = await async_client.post("/api/v1/corpora", json={
|
| 145 |
+
"slug": "ssrf-test", "title": "SSRF", "profile_id": "medieval-illuminated",
|
| 146 |
})
|
| 147 |
cid = create.json()["id"]
|
| 148 |
|
|
|
|
| 157 |
async def test_ssrf_metadata_ip(async_client):
|
| 158 |
"""Un manifest_url vers 169.254.x.x (cloud metadata) doit être rejeté."""
|
| 159 |
create = await async_client.post("/api/v1/corpora", json={
|
| 160 |
+
"slug": "ssrf-meta", "title": "SSRF", "profile_id": "medieval-illuminated",
|
| 161 |
})
|
| 162 |
cid = create.json()["id"]
|
| 163 |
|
|
|
|
| 171 |
async def test_ssrf_file_scheme(async_client):
|
| 172 |
"""Un manifest_url avec file:// doit être rejeté."""
|
| 173 |
create = await async_client.post("/api/v1/corpora", json={
|
| 174 |
+
"slug": "ssrf-file", "title": "SSRF", "profile_id": "medieval-illuminated",
|
| 175 |
})
|
| 176 |
cid = create.json()["id"]
|
| 177 |
|
|
|
|
| 207 |
async def test_model_id_too_long(async_client):
|
| 208 |
"""Un model_id >256 chars doit être rejeté."""
|
| 209 |
create = await async_client.post("/api/v1/corpora", json={
|
| 210 |
+
"slug": "model-test", "title": "T", "profile_id": "medieval-illuminated",
|
| 211 |
})
|
| 212 |
cid = create.json()["id"]
|
| 213 |
|