From 4bb1df8299e56487d440d18e9b284e1b486c1408 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 6 Aug 2026 18:51:44 +0800 Subject: [PATCH] Make topology resume kick sync immediately and keep schedule across restarts. Wire topology_auto_sync into resume/pause force-sync hints, ignore stale_running cleanup rows when computing the interval clock, and wait a full interval instead of pulling topology on every cold start. Co-authored-by: Cursor --- netx_api/ume_support.py | 35 ++++++++++++++++++++++++++++------- netx_api/ume_sync_router.py | 4 ++-- 2 files changed, 30 insertions(+), 9 deletions(-) diff --git a/netx_api/ume_support.py b/netx_api/ume_support.py index 5e5f71c..aa2114f 100644 --- a/netx_api/ume_support.py +++ b/netx_api/ume_support.py @@ -358,20 +358,29 @@ def _sleep_or_until_paused(task_id: str, total_s: float) -> None: def _last_finished_job_ended_at(db: Session, domain: str) -> datetime | None: - """Latest finished sync job end time for domain (success or failed).""" - row = ( + """Latest finished sync job end time for domain (success or failed). + + Ignores startup/reaper cleanup rows (``stale_running_*``) so process restart + does not reset the schedule clock to "just now". + """ + rows = ( db.query(UmeSyncJob) .filter( UmeSyncJob.domain == domain, UmeSyncJob.ended_at.isnot(None), ) .order_by(UmeSyncJob.ended_at.desc()) - .limit(1) - .first() + .limit(50) + .all() ) - if not row or row.ended_at is None: - return None - return _ensure_utc(row.ended_at) + for row in rows: + err = str(getattr(row, "error_message", "") or "") + if "stale_running" in err: + continue + if row.ended_at is None: + continue + return _ensure_utc(row.ended_at) + return None def _seconds_since_last_finished_job(db: Session, domain: str) -> float | None: @@ -419,6 +428,18 @@ def _maybe_wait_for_sync_interval( db.close() _refresh_runtime_task_idle(task_id, domain) if elapsed is None: + # Topology dumps are heavy: do not pull on every process start when there is + # no prior real finished job (e.g. only orphan cleanup rows). Wait one full + # interval unless resume/kick skipped debounce above. + if domain == "topology": + _schedule_log.info( + "%s: no prior finished job for %s, wait full %ss (no startup pull)", + label, + domain, + interval_s, + ) + _sleep_or_until_paused(task_id, float(interval_s)) + return _schedule_log.info("%s: no prior finished job for %s, sync now", label, domain) return if elapsed >= float(interval_s): diff --git a/netx_api/ume_sync_router.py b/netx_api/ume_sync_router.py index a5d5120..5c74904 100644 --- a/netx_api/ume_sync_router.py +++ b/netx_api/ume_sync_router.py @@ -221,7 +221,7 @@ def ume_runtime_task_pause(task: str) -> dict[str, Any]: if tid not in UME_KNOWN_RUNTIME_TASKS: raise HTTPException(status_code=404, detail="unknown_runtime_task") _runtime_pause_task(tid) - if tid in ("alarms_current_auto_sync", "inventory_auto_sync"): + if tid in ("alarms_current_auto_sync", "inventory_auto_sync", "topology_auto_sync"): _clear_force_resume_hints(tid) if tid == "alarms_current_ws_consumer": request_ws_reconnect() @@ -237,7 +237,7 @@ def ume_runtime_task_resume(task: str) -> dict[str, Any]: if tid not in UME_KNOWN_RUNTIME_TASKS: raise HTTPException(status_code=404, detail="unknown_runtime_task") _runtime_resume_task(tid) - if tid in ("alarms_current_auto_sync", "inventory_auto_sync"): + if tid in ("alarms_current_auto_sync", "inventory_auto_sync", "topology_auto_sync"): _request_force_sync_after_resume(tid) resume_hint = RT_RESUMED_SYNC_SOON elif tid == "alarms_current_ws_consumer":