#!/usr/bin/env python3 """Watch the release index build and reclaim known stuck job-876 tails.""" from __future__ import annotations import os import sys import time from pathlib import Path ROOT = Path(__file__).resolve().parents[1] DOMAINCHECK_ROOT = ROOT / "domainCheck" if str(DOMAINCHECK_ROOT) not in sys.path: sys.path.insert(0, str(DOMAINCHECK_ROOT)) from app.utils.database import Database # noqa: E402 ENV_FILE = Path("/etc/default/domaincheck-worker") INDEX_NAME = "idx_detect_job_items_release_node_job" JOB_ID = 876 TARGET_NODES = [ "mainland-controller-01-ba", "mainland-controller-01-bd", "mainland-controller-01-aw", "mainland-controller-01-ae", ] POLL_SECONDS = 15 def load_env(path: Path) -> None: if not path.exists(): return for raw in path.read_text(encoding="utf-8").splitlines(): line = raw.strip() if not line or line.startswith("#") or "=" not in line: continue key, value = line.split("=", 1) os.environ[key] = value def log(message: str) -> None: stamp = time.strftime("%Y-%m-%d %H:%M:%S") print(f"[{stamp}] {message}", flush=True) def fetch_index_state(db: Database) -> tuple[bool, bool, str | None, int | None, int | None]: conn = db.get_connection() try: cur = conn.cursor() try: cur.execute( """ SELECT i.indisvalid, i.indisready, p.phase, p.blocks_done, p.blocks_total FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid LEFT JOIN pg_stat_progress_create_index p ON p.index_relid = i.indexrelid WHERE c.relname = %s """, (INDEX_NAME,), ) row = cur.fetchone() finally: cur.close() finally: db.close(conn=conn) if not row: return False, False, None, None, None return bool(row[0]), bool(row[1]), row[2], row[3], row[4] def main() -> int: load_env(ENV_FILE) db = Database() log(f"watch start index={INDEX_NAME} job_id={JOB_ID} targets={','.join(TARGET_NODES)}") while True: try: valid, ready, phase, blocks_done, blocks_total = fetch_index_state(db) log( "index_state " f"valid={int(valid)} ready={int(ready)} " f"phase={phase or 'none'} blocks_done={blocks_done} blocks_total={blocks_total}" ) if valid: break except Exception as exc: # pragma: no cover - operational script log(f"index_state_error error={exc!r}") db = Database() time.sleep(POLL_SECONDS) released_total = 0 for node in TARGET_NODES: try: released = db.release_detect_job_items_for_node_job(node, JOB_ID) released_total += int(released or 0) log(f"release node={node} released={released}") except Exception as exc: # pragma: no cover - operational script log(f"release_error node={node} error={exc!r}") try: status = db.refresh_detect_job_status(JOB_ID) log(f"refresh_status job_id={JOB_ID} status={status}") except Exception as exc: # pragma: no cover - operational script log(f"refresh_status_error job_id={JOB_ID} error={exc!r}") log(f"done released_total={released_total}") return 0 if __name__ == "__main__": raise SystemExit(main())