Start autonomous Belgium and North Sea RC
This commit is contained in:
@@ -1,12 +1,23 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from importlib import import_module
|
||||
from sqlalchemy import text
|
||||
from fastapi import APIRouter
|
||||
from pathlib import Path
|
||||
from tempfile import NamedTemporaryFile
|
||||
|
||||
from app.schemas.health import HealthResponse, SystemCapabilities
|
||||
from app.providers.registry import list_provider_capabilities
|
||||
from alembic.config import Config
|
||||
from alembic.script import ScriptDirectory
|
||||
from fastapi import APIRouter, Response, status
|
||||
from sqlalchemy import text
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.db.session import get_engine
|
||||
from app.providers.registry import list_provider_capabilities
|
||||
from app.schemas.health import (
|
||||
HealthResponse,
|
||||
SystemCapabilities,
|
||||
SystemCapabilitiesEnvelope,
|
||||
)
|
||||
from app.services.model_registry_service import ModelRegistryService
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
@@ -19,27 +30,136 @@ def _dependency_enabled(module: str) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
@router.get("/health")
|
||||
def readiness() -> HealthResponse:
|
||||
db_status = "ok"
|
||||
def _expected_migration_heads() -> list[str]:
|
||||
backend_root = Path(__file__).resolve().parents[3]
|
||||
config = Config(str(backend_root / "alembic.ini"))
|
||||
config.set_main_option("script_location", str(backend_root / "alembic"))
|
||||
return list(ScriptDirectory.from_config(config).get_heads())
|
||||
|
||||
|
||||
def _database_checks() -> dict[str, str]:
|
||||
checks = {
|
||||
"database": "degraded",
|
||||
"postgis": "degraded",
|
||||
"migration": "degraded",
|
||||
}
|
||||
try:
|
||||
with get_engine().connect() as connection:
|
||||
connection.execute(text("SELECT 1"))
|
||||
checks["database"] = "ok"
|
||||
postgis_version = connection.execute(
|
||||
text("SELECT PostGIS_Version()")
|
||||
).scalar_one()
|
||||
checks["postgis"] = f"ok:{postgis_version}"
|
||||
database_head = connection.execute(
|
||||
text("SELECT version_num FROM alembic_version")
|
||||
).scalar_one()
|
||||
expected_heads = _expected_migration_heads()
|
||||
if len(expected_heads) == 1 and database_head == expected_heads[0]:
|
||||
checks["migration"] = f"ok:{database_head}"
|
||||
else:
|
||||
checks["migration"] = (
|
||||
f"degraded:database={database_head};"
|
||||
f"expected={','.join(expected_heads) or 'none'}"
|
||||
)
|
||||
except Exception:
|
||||
db_status = "degraded"
|
||||
return HealthResponse(status="ok", service="geointel-backend", version="0.1.0", database=db_status)
|
||||
return checks
|
||||
return checks
|
||||
|
||||
|
||||
@router.get("/api/v1/system/capabilities")
|
||||
def capabilities() -> dict:
|
||||
def _storage_check(storage_root: str) -> str:
|
||||
root = Path(storage_root).expanduser()
|
||||
try:
|
||||
root.mkdir(parents=True, exist_ok=True)
|
||||
with NamedTemporaryFile(
|
||||
prefix=".geointel-readiness-",
|
||||
dir=root,
|
||||
delete=True,
|
||||
) as handle:
|
||||
handle.write(b"ok")
|
||||
handle.flush()
|
||||
return "ok"
|
||||
except OSError:
|
||||
return "degraded"
|
||||
|
||||
|
||||
def _readiness_payload() -> HealthResponse:
|
||||
settings = get_settings()
|
||||
checks = _database_checks()
|
||||
checks["storage"] = _storage_check(settings.storage_root)
|
||||
ready = all(
|
||||
value == "ok" or value.startswith("ok:")
|
||||
for value in checks.values()
|
||||
)
|
||||
return HealthResponse(
|
||||
status="ok" if ready else "degraded",
|
||||
service="geointel-backend",
|
||||
version=settings.app_version,
|
||||
build_sha=settings.build_sha,
|
||||
build_time=settings.build_time,
|
||||
database=checks["database"],
|
||||
postgis=checks["postgis"],
|
||||
migration=checks["migration"],
|
||||
storage=checks["storage"],
|
||||
checks=checks,
|
||||
)
|
||||
|
||||
|
||||
@router.get("/health/live", response_model=HealthResponse)
|
||||
def liveness() -> HealthResponse:
|
||||
settings = get_settings()
|
||||
return HealthResponse(
|
||||
status="ok",
|
||||
service="geointel-backend",
|
||||
version=settings.app_version,
|
||||
build_sha=settings.build_sha,
|
||||
build_time=settings.build_time,
|
||||
)
|
||||
|
||||
|
||||
def _readiness_response(response: Response) -> HealthResponse:
|
||||
payload = _readiness_payload()
|
||||
if payload.status != "ok":
|
||||
response.status_code = status.HTTP_503_SERVICE_UNAVAILABLE
|
||||
return payload
|
||||
|
||||
|
||||
@router.get("/health", response_model=HealthResponse)
|
||||
def readiness(response: Response) -> HealthResponse:
|
||||
return _readiness_response(response)
|
||||
|
||||
|
||||
@router.get("/health/ready", response_model=HealthResponse)
|
||||
def readiness_explicit(response: Response) -> HealthResponse:
|
||||
return _readiness_response(response)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/api/v1/system/capabilities",
|
||||
response_model=SystemCapabilitiesEnvelope,
|
||||
)
|
||||
def capabilities() -> SystemCapabilitiesEnvelope:
|
||||
settings = get_settings()
|
||||
providers = [item.to_dict() for item in list_provider_capabilities()]
|
||||
return {"data": SystemCapabilities(
|
||||
postgis=True,
|
||||
rasterio=_dependency_enabled("rasterio"),
|
||||
geopandas=_dependency_enabled("geopandas"),
|
||||
yolo=False,
|
||||
sam=False,
|
||||
grb="bounded",
|
||||
sentinel="planned",
|
||||
providers=providers,
|
||||
).model_dump()}
|
||||
configured_yolo = ModelRegistryService.get_model_capability(
|
||||
settings.yolo_model_id,
|
||||
settings=settings,
|
||||
)
|
||||
yolo_configured = bool(configured_yolo and configured_yolo.configured)
|
||||
yolo_status = configured_yolo.status if configured_yolo else "not_configured"
|
||||
postgis_ready = _database_checks()["postgis"].startswith("ok:")
|
||||
return SystemCapabilitiesEnvelope(
|
||||
data=SystemCapabilities(
|
||||
postgis=postgis_ready,
|
||||
rasterio=_dependency_enabled("rasterio"),
|
||||
geopandas=_dependency_enabled("geopandas"),
|
||||
yolo=yolo_configured,
|
||||
yolo_status=yolo_status,
|
||||
sam=False,
|
||||
grb="bounded",
|
||||
sentinel="planned",
|
||||
version=settings.app_version,
|
||||
build_sha=settings.build_sha,
|
||||
providers=providers,
|
||||
)
|
||||
)
|
||||
|
||||
@@ -12,6 +12,8 @@ class Settings(BaseSettings):
|
||||
|
||||
app_env: str = Field(default="development", validation_alias="GEOINTEL_ENV")
|
||||
app_version: str = Field(default="0.1.0")
|
||||
build_sha: str | None = Field(default=None, validation_alias="GEOINTEL_BUILD_SHA")
|
||||
build_time: str | None = Field(default=None, validation_alias="GEOINTEL_BUILD_TIME")
|
||||
api_prefix: str = Field(default="/api/v1", validation_alias="GEOINTEL_API_PREFIX")
|
||||
database_url: str = Field(
|
||||
default="postgresql+psycopg://geointel:geointel@localhost:5432/geointel?connect_timeout=1",
|
||||
@@ -231,6 +233,11 @@ class Settings(BaseSettings):
|
||||
thematic_raster_max_response_mb: int = Field(default=160, ge=1, validation_alias="THEMATIC_RASTER_MAX_RESPONSE_MB")
|
||||
redis_url: str | None = Field(default=None, validation_alias="REDIS_URL")
|
||||
log_level: str = Field(default="INFO", validation_alias="GEOINTEL_LOG_LEVEL")
|
||||
sql_log_level: str = Field(default="WARNING", validation_alias="GEOINTEL_SQL_LOG_LEVEL")
|
||||
reconcile_interrupted_runs_on_startup: bool = Field(
|
||||
default=False,
|
||||
validation_alias="GEOINTEL_RECONCILE_INTERRUPTED_RUNS_ON_STARTUP",
|
||||
)
|
||||
database_statement_timeout_ms: int = Field(default=5_000, validation_alias="DATABASE_STATEMENT_TIMEOUT_MS")
|
||||
yolo_enabled: bool = Field(default=False, validation_alias="YOLO_ENABLED")
|
||||
yolo_models_dir: str = Field(default="/app/models", validation_alias="YOLO_MODELS_DIR")
|
||||
|
||||
@@ -2,11 +2,13 @@ import logging
|
||||
import sys
|
||||
|
||||
|
||||
def configure_logging(level: str = "INFO") -> None:
|
||||
def configure_logging(level: str = "INFO", sql_level: str = "WARNING") -> None:
|
||||
logging.basicConfig(
|
||||
level=level,
|
||||
format="%(asctime)s | %(levelname)s | %(name)s | %(message)s",
|
||||
stream=sys.stdout,
|
||||
force=True,
|
||||
)
|
||||
for name in ["uvicorn", "uvicorn.error", "uvicorn.access", "sqlalchemy.engine"]:
|
||||
for name in ["uvicorn", "uvicorn.error", "uvicorn.access"]:
|
||||
logging.getLogger(name).setLevel(level)
|
||||
logging.getLogger("sqlalchemy.engine").setLevel(sql_level)
|
||||
|
||||
+50
-7
@@ -1,5 +1,9 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import uuid
|
||||
from contextlib import asynccontextmanager
|
||||
|
||||
from fastapi import FastAPI, HTTPException, Request
|
||||
from fastapi.exceptions import RequestValidationError
|
||||
from fastapi.middleware.cors import CORSMiddleware
|
||||
@@ -9,6 +13,11 @@ from app.api.routes import analysis, areas, assistant, datasets, demo, detection
|
||||
from app.core.config import get_settings
|
||||
from app.core.errors import AppError
|
||||
from app.core.logging import configure_logging
|
||||
from app.db.session import SessionLocal
|
||||
from app.services.runtime_reconciliation_service import RuntimeReconciliationService
|
||||
|
||||
|
||||
logger = logging.getLogger("geointel")
|
||||
|
||||
|
||||
def _to_error_payload(
|
||||
@@ -27,13 +36,33 @@ def _to_error_payload(
|
||||
|
||||
def create_app() -> FastAPI:
|
||||
settings = get_settings()
|
||||
configure_logging(settings.log_level)
|
||||
configure_logging(settings.log_level, settings.sql_log_level)
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(_: FastAPI):
|
||||
if settings.reconcile_interrupted_runs_on_startup:
|
||||
db = SessionLocal()
|
||||
try:
|
||||
result = RuntimeReconciliationService.reconcile(db)
|
||||
logger.info(
|
||||
"Runtime reconciliation completed: jobs=%s analysis_runs=%s",
|
||||
result.interrupted_jobs,
|
||||
result.interrupted_analysis_runs,
|
||||
)
|
||||
except Exception:
|
||||
db.rollback()
|
||||
logger.exception("Runtime reconciliation failed")
|
||||
raise
|
||||
finally:
|
||||
db.close()
|
||||
yield
|
||||
|
||||
app = FastAPI(
|
||||
title="GeoIntel Kempen",
|
||||
title="GeoIntel",
|
||||
version=settings.app_version,
|
||||
docs_url="/docs",
|
||||
redoc_url="/redoc",
|
||||
lifespan=lifespan,
|
||||
)
|
||||
|
||||
app.add_middleware(
|
||||
@@ -60,6 +89,14 @@ def create_app() -> FastAPI:
|
||||
app.include_router(temporal.router, prefix=settings.api_prefix)
|
||||
app.include_router(assistant.router, prefix=settings.api_prefix)
|
||||
|
||||
@app.middleware("http")
|
||||
async def request_identity(request: Request, call_next):
|
||||
request_id = request.headers.get("x-request-id") or str(uuid.uuid4())
|
||||
request.state.request_id = request_id
|
||||
response = await call_next(request)
|
||||
response.headers["x-request-id"] = request_id
|
||||
return response
|
||||
|
||||
@app.exception_handler(AppError)
|
||||
async def app_error(request: Request, exc: AppError): # noqa: ARG001
|
||||
return JSONResponse(
|
||||
@@ -68,7 +105,7 @@ def create_app() -> FastAPI:
|
||||
exc.code,
|
||||
exc.message,
|
||||
exc.details,
|
||||
request_id=request.headers.get("x-request-id"),
|
||||
request_id=request.state.request_id,
|
||||
),
|
||||
)
|
||||
|
||||
@@ -88,7 +125,7 @@ def create_app() -> FastAPI:
|
||||
code,
|
||||
message,
|
||||
details,
|
||||
request_id=request.headers.get("x-request-id"),
|
||||
request_id=request.state.request_id,
|
||||
),
|
||||
)
|
||||
|
||||
@@ -100,19 +137,25 @@ def create_app() -> FastAPI:
|
||||
"VALIDATION_ERROR",
|
||||
"Validation failed",
|
||||
exc.errors(),
|
||||
request_id=request.headers.get("x-request-id"),
|
||||
request_id=request.state.request_id,
|
||||
),
|
||||
)
|
||||
|
||||
@app.exception_handler(Exception)
|
||||
async def unexpected_error(request: Request, exc: Exception): # noqa: ARG001
|
||||
async def unexpected_error(request: Request, exc: Exception):
|
||||
logger.exception(
|
||||
"Unhandled request error request_id=%s method=%s path=%s",
|
||||
request.state.request_id,
|
||||
request.method,
|
||||
request.url.path,
|
||||
)
|
||||
return JSONResponse(
|
||||
status_code=500,
|
||||
content=_to_error_payload(
|
||||
"INTERNAL_ERROR",
|
||||
"Unexpected server error",
|
||||
{"type": exc.__class__.__name__},
|
||||
request_id=request.headers.get("x-request-id"),
|
||||
request_id=request.state.request_id,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@@ -23,7 +23,13 @@ class HealthResponse(BaseModel):
|
||||
status: str
|
||||
service: str
|
||||
version: str
|
||||
build_sha: str | None = None
|
||||
build_time: str | None = None
|
||||
database: str | None = None
|
||||
postgis: str | None = None
|
||||
migration: str | None = None
|
||||
storage: str | None = None
|
||||
checks: dict[str, str] = Field(default_factory=dict)
|
||||
|
||||
|
||||
class SystemCapabilities(BaseModel):
|
||||
@@ -31,7 +37,14 @@ class SystemCapabilities(BaseModel):
|
||||
rasterio: bool
|
||||
geopandas: bool
|
||||
yolo: bool | str
|
||||
yolo_status: str
|
||||
sam: bool | str
|
||||
grb: str
|
||||
sentinel: str
|
||||
version: str
|
||||
build_sha: str | None = None
|
||||
providers: list[ProviderCapability] = Field(default_factory=list)
|
||||
|
||||
|
||||
class SystemCapabilitiesEnvelope(BaseModel):
|
||||
data: SystemCapabilities
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.models import AnalysisRun, Job
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ReconciliationResult:
|
||||
interrupted_jobs: int
|
||||
interrupted_analysis_runs: int
|
||||
|
||||
|
||||
class RuntimeReconciliationService:
|
||||
ERROR_MESSAGE = (
|
||||
"PROCESS_INTERRUPTED: the GeoIntel process restarted before this work "
|
||||
"reached a terminal state"
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def reconcile(
|
||||
db: Session,
|
||||
*,
|
||||
finished_at: datetime | None = None,
|
||||
) -> ReconciliationResult:
|
||||
resolved_finished_at = finished_at or datetime.now(timezone.utc)
|
||||
interrupted_jobs = (
|
||||
db.query(Job)
|
||||
.filter(Job.status == "running")
|
||||
.update(
|
||||
{
|
||||
Job.status: "failed",
|
||||
Job.finished_at: resolved_finished_at,
|
||||
Job.error_message: RuntimeReconciliationService.ERROR_MESSAGE,
|
||||
},
|
||||
synchronize_session=False,
|
||||
)
|
||||
)
|
||||
interrupted_analysis_runs = (
|
||||
db.query(AnalysisRun)
|
||||
.filter(AnalysisRun.status == "running")
|
||||
.update(
|
||||
{
|
||||
AnalysisRun.status: "failed",
|
||||
AnalysisRun.finished_at: resolved_finished_at,
|
||||
AnalysisRun.error_message: RuntimeReconciliationService.ERROR_MESSAGE,
|
||||
},
|
||||
synchronize_session=False,
|
||||
)
|
||||
)
|
||||
db.commit()
|
||||
return ReconciliationResult(
|
||||
interrupted_jobs=interrupted_jobs,
|
||||
interrupted_analysis_runs=interrupted_analysis_runs,
|
||||
)
|
||||
Reference in New Issue
Block a user