from __future__ import annotations from dataclasses import dataclass from datetime import datetime, timezone from sqlalchemy.orm import Session from app.models import AoiOperation, AoiOperationPartition, AnalysisRun, Job @dataclass(frozen=True) class ReconciliationResult: interrupted_jobs: int interrupted_analysis_runs: int resumed_aoi_partitions: int exhausted_aoi_partitions: 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, ) ) resumed_aoi_partitions = ( db.query(AoiOperationPartition) .filter( AoiOperationPartition.status == "running", AoiOperationPartition.attempt_count < AoiOperationPartition.max_attempts, ) .update( { AoiOperationPartition.status: "queued", AoiOperationPartition.error_message: RuntimeReconciliationService.ERROR_MESSAGE, AoiOperationPartition.started_at: None, }, synchronize_session=False, ) ) exhausted_aoi_partitions = ( db.query(AoiOperationPartition) .filter( AoiOperationPartition.status == "running", AoiOperationPartition.attempt_count >= AoiOperationPartition.max_attempts, ) .update( { AoiOperationPartition.status: "failed", AoiOperationPartition.finished_at: resolved_finished_at, AoiOperationPartition.error_message: RuntimeReconciliationService.ERROR_MESSAGE, }, synchronize_session=False, ) ) db.query(AoiOperation).filter(AoiOperation.status == "running").update( {AoiOperation.status: "queued"}, synchronize_session=False ) db.commit() return ReconciliationResult( interrupted_jobs=interrupted_jobs, interrupted_analysis_runs=interrupted_analysis_runs, resumed_aoi_partitions=resumed_aoi_partitions, exhausted_aoi_partitions=exhausted_aoi_partitions, )