From 34ee677433ca006ddc20487edc37187cf2737e1c Mon Sep 17 00:00:00 2001 From: oliver Date: Tue, 12 May 2026 10:21:01 +0800 Subject: [PATCH] UME scheduled sync: commit running job before long UME pull; schedule logs Makes schedule attempts visible in ume_sync_jobs while HTTP pagination runs; log iteration/sleep/finish for alarms_current and inventory auto loops. Co-authored-by: Cursor --- netx_api/main.py | 17 +++++++++++++++++ netx_api/ume_sync_service.py | 12 ++++++++++++ 2 files changed, 29 insertions(+) diff --git a/netx_api/main.py b/netx_api/main.py index 6d86ad7..c7fdde7 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -2,11 +2,14 @@ from __future__ import annotations import csv import json +import logging from datetime import datetime, timezone from io import StringIO import time import re import threading + +_schedule_log = logging.getLogger("netx.ume.schedule") from fastapi import Depends, FastAPI, File, HTTPException, Query, UploadFile from fastapi.responses import Response from sqlalchemy import text as sql_text @@ -446,6 +449,10 @@ def on_startup() -> None: if _runtime_is_paused("alarms_current_auto_sync"): time.sleep(1) continue + _schedule_log.info( + "alarms_current_auto_sync: iteration start (sleep_after=%ss)", + interval_s, + ) _set_runtime_task( "alarms_current_auto_sync", status="running", @@ -456,6 +463,7 @@ def on_startup() -> None: try: client = _ume_client() sync_alarms_current(db, client, trigger_mode="schedule") + _schedule_log.info("alarms_current_auto_sync: sync finished ok") _set_runtime_task( "alarms_current_auto_sync", status="running", @@ -465,12 +473,14 @@ def on_startup() -> None: finally: db.close() except Exception as exc: + _schedule_log.exception("alarms_current_auto_sync: sync failed: %s", exc) _set_runtime_task( "alarms_current_auto_sync", status="error", last_run_at=datetime.now(timezone.utc), last_error=str(exc)[:240], ) + _schedule_log.info("alarms_current_auto_sync: sleeping %ss", interval_s) time.sleep(interval_s) t2 = threading.Thread(target=_alarms_current_sync_loop, name="ume-alarms-current-sync", daemon=True) @@ -489,6 +499,10 @@ def on_startup() -> None: if _runtime_is_paused("inventory_auto_sync"): time.sleep(1) continue + _schedule_log.info( + "inventory_auto_sync: iteration start (sleep_after=%ss)", + interval_s, + ) _set_runtime_task( "inventory_auto_sync", status="running", @@ -499,6 +513,7 @@ def on_startup() -> None: try: client = _ume_client() sync_inventory_full(db, client, trigger_mode="schedule") + _schedule_log.info("inventory_auto_sync: sync finished ok") _set_runtime_task( "inventory_auto_sync", status="running", @@ -508,12 +523,14 @@ def on_startup() -> None: finally: db.close() except Exception as exc: + _schedule_log.exception("inventory_auto_sync: sync failed: %s", exc) _set_runtime_task( "inventory_auto_sync", status="error", last_run_at=datetime.now(timezone.utc), last_error=str(exc)[:240], ) + _schedule_log.info("inventory_auto_sync: sleeping %ss", interval_s) time.sleep(interval_s) t3 = threading.Thread(target=_inventory_auto_sync_loop, name="ume-inventory-auto-sync", daemon=True) diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index de1065a..f5f612b 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -1,9 +1,12 @@ from __future__ import annotations import json +import logging from datetime import datetime, timezone from typing import Any +_sync_log = logging.getLogger("netx.ume.sync") + from sqlalchemy.orm import Session from .config import settings @@ -168,6 +171,8 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = " job = _build_sync_job("inventory", trigger_mode) db.add(job) db.flush() + db.commit() + _sync_log.info("inventory sync job %s committed as running (trigger=%s)", getattr(job, "id", "?"), trigger_mode) pulled = inserted = updated = 0 try: limit_max = int(getattr(settings, "ume_limit_max", 5000) or 5000) @@ -280,6 +285,13 @@ def _sync_alarms_common( ) db.add(batch) db.flush() + db.commit() + _sync_log.info( + "alarms sync job domain=%s id=%s committed as running (trigger=%s)", + domain, + getattr(job, "id", "?"), + trigger_mode, + ) pulled = inserted = updated = 0 deleted_stale_current = 0 paging_mode = "marker"