mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 03:10:46 +08:00
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 <cursoragent@cursor.com>
This commit is contained in:
parent
61f6c99f93
commit
34ee677433
2 changed files with 29 additions and 0 deletions
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue