generated from john/python-template
gpt-5.3 codex review: Phase 7 and the addition of the new test-effectiveness-auditor skill.
Quality Gate / gate (push) Failing after 12s
Quality Gate / gate (push) Failing after 12s
This commit is contained in:
@@ -1,16 +1,30 @@
|
||||
"""Health endpoint routes."""
|
||||
|
||||
from fastapi import APIRouter
|
||||
from fastapi import Request
|
||||
|
||||
from transcription.worker import resolve_worker_health
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
def healthz() -> dict[str, str]:
|
||||
"""Return a simple health status payload."""
|
||||
return {"status": "ok"}
|
||||
def healthz(request: Request) -> dict[str, object]:
|
||||
"""Return health status with worker-liveness signal."""
|
||||
worker = resolve_worker_health(request.app.state)
|
||||
payload: dict[str, object] = {
|
||||
"status": "ok",
|
||||
"worker": {
|
||||
"state": worker.state,
|
||||
},
|
||||
}
|
||||
if worker.error_id is not None:
|
||||
payload["worker"]["error_id"] = worker.error_id
|
||||
if worker.error_category is not None:
|
||||
payload["worker"]["error_category"] = worker.error_category
|
||||
return payload
|
||||
|
||||
|
||||
@router.get("/healthz")
|
||||
def healthz_route() -> dict[str, str]:
|
||||
def healthz_route(request: Request) -> dict[str, object]:
|
||||
"""Route wrapper for health status payload."""
|
||||
return healthz()
|
||||
return healthz(request)
|
||||
|
||||
@@ -45,12 +45,14 @@ async def _lifespan(app: FastAPI):
|
||||
|
||||
settings.upload_dir.mkdir(parents=True, exist_ok=True)
|
||||
settings.prompt_dir.mkdir(parents=True, exist_ok=True)
|
||||
settings.log_dir.mkdir(parents=True, exist_ok=True)
|
||||
settings.database_backup_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
await _recover_stale_processing_jobs(app)
|
||||
|
||||
async with AsyncExitStack() as stack:
|
||||
stack.push_async_callback(dispose_database_runtime)
|
||||
stop_event, worker_notifier = await stack.enter_async_context(
|
||||
stop_event, worker_notifier, worker_health = await stack.enter_async_context(
|
||||
worker_consumer_lifespan(
|
||||
session_factory=app.state.runtime.session_factory,
|
||||
poll_interval_seconds=1.0,
|
||||
@@ -58,6 +60,7 @@ async def _lifespan(app: FastAPI):
|
||||
)
|
||||
app.state.worker_stop_event = stop_event
|
||||
app.state.worker_notifier = worker_notifier
|
||||
app.state.worker_health = worker_health
|
||||
yield
|
||||
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ once at startup. Provider-specific defaults (model names, base URLs)
|
||||
are resolved by the provider adapters, not here.
|
||||
"""
|
||||
|
||||
import copy
|
||||
import logging.config
|
||||
from collections.abc import Sequence
|
||||
from enum import StrEnum
|
||||
@@ -42,7 +43,7 @@ class SqliteSettings(BaseModel):
|
||||
model_config = ConfigDict(extra="forbid", frozen=True)
|
||||
|
||||
driver: Literal["sqlite"] = "sqlite"
|
||||
path: NonEmptyStr = "app.db"
|
||||
path: NonEmptyStr = "./data/transcription.db"
|
||||
|
||||
|
||||
class PostgresSettings(BaseModel):
|
||||
@@ -78,6 +79,10 @@ class Settings(BaseSettings):
|
||||
port: int = 8000
|
||||
log_level: Literal["critical", "error", "warning", "info", "debug", "trace"] = "info"
|
||||
reload: bool = False
|
||||
log_dir: Path = Path("./data/logs")
|
||||
log_file_name: NonEmptyStr = "transcription.log"
|
||||
log_file_max_bytes: int = Field(default=10 * 1024 * 1024, gt=0)
|
||||
log_file_backup_count: int = Field(default=5, ge=1)
|
||||
|
||||
# --- AI provider ---
|
||||
provider: Provider = Provider.OPENROUTER
|
||||
@@ -99,15 +104,16 @@ class Settings(BaseSettings):
|
||||
sqlite_check_same_thread: bool = False
|
||||
|
||||
# --- filesystem paths ---
|
||||
upload_dir: Path = Path("./uploads")
|
||||
upload_dir: Path = Path("./data")
|
||||
prompt_dir: Path = Path("./prompts")
|
||||
homepage_dir: Path = Path("./data/homepage")
|
||||
database_backup_dir: Path = Path("./data/backups")
|
||||
|
||||
# --- worker reliability ---
|
||||
worker_max_retries: int = Field(default=0, ge=0)
|
||||
# Bounded only from below. Vision transcription of a dense page routinely runs
|
||||
# well past twenty seconds, so an upper cap here would silently fail real work.
|
||||
worker_provider_timeout_seconds: float = Field(default=180.0, gt=0.0)
|
||||
worker_provider_timeout_seconds: float = Field(default=30.0, gt=0.0)
|
||||
worker_min_transcription_chars: int = Field(default=0, ge=0)
|
||||
worker_min_transcription_lines: int = Field(default=0, ge=0)
|
||||
worker_fail_on_finish_reason_length: bool = False
|
||||
@@ -193,16 +199,24 @@ LOGGING_CONFIG: dict[str, Any] = {
|
||||
"class": "logging.StreamHandler",
|
||||
"formatter": "standard",
|
||||
"stream": "ext://sys.stdout",
|
||||
},
|
||||
"file": {
|
||||
"class": "logging.handlers.RotatingFileHandler",
|
||||
"formatter": "standard",
|
||||
"filename": str((Path("./data/logs") / "transcription.log")),
|
||||
"maxBytes": 10 * 1024 * 1024,
|
||||
"backupCount": 5,
|
||||
"encoding": "utf-8",
|
||||
}
|
||||
},
|
||||
"root": {
|
||||
"level": "INFO",
|
||||
"handlers": ["console"],
|
||||
"handlers": ["console", "file"],
|
||||
},
|
||||
"loggers": {
|
||||
"transcription": {
|
||||
"level": "DEBUG",
|
||||
"handlers": ["console"],
|
||||
"handlers": ["console", "file"],
|
||||
"propagate": False,
|
||||
}
|
||||
},
|
||||
@@ -211,8 +225,13 @@ LOGGING_CONFIG: dict[str, Any] = {
|
||||
|
||||
def configure_logging(settings: Settings | None = None) -> None:
|
||||
"""Configure root logging once at startup."""
|
||||
cfg = LOGGING_CONFIG.copy()
|
||||
cfg = copy.deepcopy(LOGGING_CONFIG)
|
||||
active_settings = settings or get_settings()
|
||||
active_settings.log_dir.mkdir(parents=True, exist_ok=True)
|
||||
file_handler = cfg["handlers"]["file"]
|
||||
file_handler["filename"] = str(active_settings.log_dir / active_settings.log_file_name)
|
||||
file_handler["maxBytes"] = active_settings.log_file_max_bytes
|
||||
file_handler["backupCount"] = active_settings.log_file_backup_count
|
||||
cfg["loggers"]["transcription"]["level"] = active_settings.log_level.upper()
|
||||
logging.config.dictConfig(cfg)
|
||||
logger.debug("Logging configured")
|
||||
|
||||
@@ -8,6 +8,8 @@ from collections.abc import AsyncGenerator
|
||||
from contextlib import asynccontextmanager
|
||||
from contextlib import contextmanager
|
||||
from contextlib import suppress
|
||||
from dataclasses import dataclass
|
||||
from typing import Literal
|
||||
from typing import Protocol
|
||||
from typing import runtime_checkable
|
||||
|
||||
@@ -24,6 +26,48 @@ from .services.workflows import process_next_queued_job as process_next_queued_j
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
WorkerHealthState = Literal["starting", "running", "stopped", "failed", "unknown"]
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class WorkerHealthSnapshot:
|
||||
"""Structured worker-loop health status for API visibility."""
|
||||
|
||||
state: WorkerHealthState
|
||||
error_id: str | None = None
|
||||
error_category: str | None = None
|
||||
|
||||
|
||||
class WorkerHealth:
|
||||
"""Mutable worker-loop health state owned by app lifespan."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._state: WorkerHealthState = "starting"
|
||||
self._error_id: str | None = None
|
||||
self._error_category: str | None = None
|
||||
|
||||
def mark_running(self) -> None:
|
||||
self._state = "running"
|
||||
self._error_id = None
|
||||
self._error_category = None
|
||||
|
||||
def mark_stopped(self) -> None:
|
||||
if self._state != "failed":
|
||||
self._state = "stopped"
|
||||
|
||||
def mark_failed(self, error: AppError) -> None:
|
||||
self._state = "failed"
|
||||
self._error_id = error.error_id
|
||||
self._error_category = error.category.value
|
||||
|
||||
def snapshot(self) -> WorkerHealthSnapshot:
|
||||
return WorkerHealthSnapshot(
|
||||
state=self._state,
|
||||
error_id=self._error_id,
|
||||
error_category=self._error_category,
|
||||
)
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class WorkerNotifier(Protocol):
|
||||
"""Abstraction for signaling the worker loop about new work."""
|
||||
@@ -59,28 +103,42 @@ def resolve_worker_notifier(state: object) -> WorkerNotifier:
|
||||
return NoopWorkerNotifier()
|
||||
|
||||
|
||||
def resolve_worker_health(state: object) -> WorkerHealthSnapshot:
|
||||
"""Resolve worker health from app-like state objects."""
|
||||
health = getattr(state, "worker_health", None)
|
||||
if isinstance(health, WorkerHealth):
|
||||
return health.snapshot()
|
||||
if isinstance(health, WorkerHealthSnapshot):
|
||||
return health
|
||||
if health is not None:
|
||||
logger.warning("Ignoring worker_health of unsupported type %r; reporting unknown.", type(health))
|
||||
return WorkerHealthSnapshot(state="unknown")
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def worker_consumer_lifespan(
|
||||
*,
|
||||
session_factory: async_sessionmaker[AsyncSession] | None = None,
|
||||
poll_interval_seconds: float = 1.0,
|
||||
) -> AsyncGenerator[tuple[asyncio.Event, WorkerNotifier]]:
|
||||
) -> AsyncGenerator[tuple[asyncio.Event, WorkerNotifier, WorkerHealth]]:
|
||||
"""Start and stop the worker consumer loop for app lifespan."""
|
||||
stop_event = asyncio.Event()
|
||||
wake_event = asyncio.Event()
|
||||
worker_notifier: WorkerNotifier = EventWorkerNotifier(wake_event)
|
||||
worker_health = WorkerHealth()
|
||||
worker_task = asyncio.create_task(
|
||||
run_worker_loop(
|
||||
session_factory=session_factory,
|
||||
stop_event=stop_event,
|
||||
wake_event=wake_event,
|
||||
poll_interval_seconds=poll_interval_seconds,
|
||||
worker_health=worker_health,
|
||||
)
|
||||
)
|
||||
worker_notifier.notify()
|
||||
|
||||
try:
|
||||
yield stop_event, worker_notifier
|
||||
yield stop_event, worker_notifier, worker_health
|
||||
finally:
|
||||
stop_event.set()
|
||||
worker_notifier.notify()
|
||||
@@ -126,6 +184,7 @@ async def run_worker_loop(
|
||||
stop_event: asyncio.Event | None = None,
|
||||
wake_event: asyncio.Event | None = None,
|
||||
poll_interval_seconds: float = 1.0,
|
||||
worker_health: WorkerHealth | None = None,
|
||||
) -> None:
|
||||
"""Run worker loop until stop_event is set.
|
||||
|
||||
@@ -142,9 +201,13 @@ async def run_worker_loop(
|
||||
"""
|
||||
services = ServiceBundle.from_session_factory(session_factory)
|
||||
try:
|
||||
if worker_health is not None:
|
||||
worker_health.mark_running()
|
||||
while True:
|
||||
if stop_event is not None and stop_event.is_set():
|
||||
logger.info("Worker stop event received")
|
||||
if worker_health is not None:
|
||||
worker_health.mark_stopped()
|
||||
return
|
||||
|
||||
if wake_event is not None:
|
||||
@@ -177,7 +240,11 @@ async def run_worker_loop(
|
||||
error.error_id,
|
||||
error.category.value,
|
||||
)
|
||||
if worker_health is not None:
|
||||
worker_health.mark_failed(error)
|
||||
finally:
|
||||
if worker_health is not None:
|
||||
worker_health.mark_stopped()
|
||||
await services.aclose()
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user