Files
getDomain/tools/watch_job876_tail_release.py

119 lines
3.5 KiB
Python

#!/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())