"""Requests to ows-abacus-account.""" import datetime as dt import json import math from owsresponse import response from abacus_contract.connectors.airflow import run_airflow_cli_command from abacus_contract.constants import error def _iso(s: str) -> dt.datetime: return dt.datetime.fromisoformat(s.replace('Z', '+00:00')) def runtimes_from_stdout(stdout_text: str) -> dict: """Parse the stdout from the airflow CLI to get runtimes.""" runs = json.loads(stdout_text) durs = [] for r in runs: s, e = r.get('start_date'), r.get('end_date') if s and e: durs.append((_iso(e) - _iso(s)).total_seconds()) if not durs: return { 'count': 0, 'avg_seconds': None, 'median_seconds': None, 'p95_seconds': None, } durs.sort() p95 = durs[math.ceil(0.95 * len(durs)) - 1] return { 'count': len(durs), 'average_run_time_seconds': sum(durs) / len(durs), 'median_run_time_seconds': ( durs[len(durs) // 2] if len(durs) % 2 else (durs[len(durs) // 2 - 1] + durs[len(durs) // 2]) / 2 ), 'p95_run_time_seconds': p95, 'min_run_time_seconds': durs[0], 'max_run_time_seconds': durs[-1], } def get_average_dag_run_time(dag_id: str, days_back: int = 30) -> response.Response: """Get average dag run time for the last 30 days.""" until = dt.datetime.now().isoformat() + 'Z' since = (dt.datetime.now() - dt.timedelta(days=days_back)).isoformat() + 'Z' request_body = f'dags list-runs -d {dag_id} -s {since} -e {until} --state success --output json' result = run_airflow_cli_command(request_body) errors = result['errors'] output = result['output'] if errors: return response.create_error_response( code=error.ERROR_DAG_RUN_TIMES, message=f'Error getting DAG run times: {errors}', status=500, ) runtime = runtimes_from_stdout(output) reponse = {'dag_id': dag_id, **runtime} return response.Response(message=reponse, status=200)