From b17aef0eae0d2a2de5d676cf24f77cfadf3bfe35 Mon Sep 17 00:00:00 2001 From: oliver Date: Tue, 12 May 2026 10:51:47 +0800 Subject: [PATCH] UME schedule: interval from any finished job; drop trailing sleep; stale running cleanup; paging logs - Use last ended_at (done or failed) for debounce; remove duplicate time.sleep after sync\n- On startup mark running+no ended_at jobs failed (crash/interrupt) so interval logic is not fooled\n- Log marker page progress every 25 pages\n- Log when alarm/inventory scheduler threads start Co-authored-by: Cursor --- netx_api/main.py | 50 ++++++++++++++++++++++++++---------- netx_api/ume_sync_service.py | 9 +++++++ 2 files changed, 46 insertions(+), 13 deletions(-) diff --git a/netx_api/main.py b/netx_api/main.py index 30622cc..b00a01d 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -186,6 +186,32 @@ def _reset_runtime_pause_flags() -> None: _UME_RUNTIME_PAUSED[tid] = False +def _fail_stale_running_sync_jobs_on_startup() -> None: + """Orphan running rows (crashed mid-sync) confuse scheduling; close them so interval uses real ended_at.""" + db = SessionLocal() + try: + rows = ( + db.query(UmeSyncJob) + .filter(UmeSyncJob.status == "running", UmeSyncJob.ended_at.is_(None)) + .all() + ) + if not rows: + return + now_naive = datetime.utcnow() + for row in rows: + row.status = "failed" + row.ended_at = now_naive + msg = str(row.error_message or "").strip() + suffix = "stale_running_reset_on_startup" + row.error_message = (msg + ("; " if msg else "") + suffix)[:1024] + db.commit() + _schedule_log.warning("startup: closed %s orphaned running ume_sync_jobs", len(rows)) + except Exception: + _schedule_log.exception("startup: stale sync job cleanup failed") + finally: + db.close() + + def _sleep_or_until_paused(task_id: str, total_s: float) -> None: """Sleep up to total_s wall seconds, but wake every ~2s to honor pause.""" deadline = time.time() + max(0.0, float(total_s)) @@ -196,13 +222,12 @@ def _sleep_or_until_paused(task_id: str, total_s: float) -> None: time.sleep(min(2.0, max(0.0, deadline - time.time()))) -def _seconds_since_last_done_sync(db: Session, domain: str) -> float | None: - """Seconds since latest completed job (status=done, ended_at set) for domain. None if none.""" +def _seconds_since_last_finished_job(db: Session, domain: str) -> float | None: + """Seconds since latest job with ended_at for domain (done or failed). None if none.""" row = ( db.query(UmeSyncJob) .filter( UmeSyncJob.domain == domain, - UmeSyncJob.status == "done", UmeSyncJob.ended_at.isnot(None), ) .order_by(UmeSyncJob.ended_at.desc()) @@ -224,25 +249,25 @@ def _maybe_wait_for_sync_interval( interval_s: int, label: str, ) -> None: - """If a recent successful sync exists in DB, sleep until interval elapsed since its ended_at.""" + """Sleep until interval elapsed since last finished job (ended_at), if any.""" db = SessionLocal() try: - elapsed = _seconds_since_last_done_sync(db, domain) + elapsed = _seconds_since_last_finished_job(db, domain) finally: db.close() if elapsed is None: - _schedule_log.info("%s: no prior done job for %s, sync now", label, domain) + _schedule_log.info("%s: no prior finished job for %s, sync now", label, domain) return if elapsed >= float(interval_s): - _schedule_log.info("%s: last done %.0fs ago (>= %ss), sync now", label, elapsed, interval_s) + _schedule_log.info("%s: last finished %.0fs ago (>= %ss), sync now", label, elapsed, interval_s) return wait_s = float(interval_s) - elapsed - _schedule_log.info("%s: last done %.0fs ago, wait %.0fs before sync", label, elapsed, wait_s) + _schedule_log.info("%s: last finished %.0fs ago, wait %.0fs before sync", label, elapsed, wait_s) _set_runtime_task( task_id, status="running", last_run_at=datetime.now(timezone.utc), - last_error=f"距上次成功同步未满周期,等待约 {int(wait_s)}s …", + last_error=f"距上次同步结束未满周期,等待约 {int(wait_s)}s …", ) _sleep_or_until_paused(task_id, wait_s) @@ -409,6 +434,7 @@ def sql_ume_query(payload: dict[str, Any] | None = None, db: Session = Depends(g def on_startup() -> None: Base.metadata.create_all(bind=engine) _reset_runtime_pause_flags() + _fail_stale_running_sync_jobs_on_startup() # Best-effort schema evolution for new columns (no migrations framework). # Safe for Postgres (IF NOT EXISTS); ignored on failure. try: @@ -555,11 +581,10 @@ def on_startup() -> None: 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) t2.start() + _schedule_log.info("started thread %s", t2.name) except Exception: pass try: @@ -611,11 +636,10 @@ def on_startup() -> None: 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) t3.start() + _schedule_log.info("started thread %s", t3.name) except Exception: pass diff --git a/netx_api/ume_sync_service.py b/netx_api/ume_sync_service.py index f5f612b..db4d184 100644 --- a/netx_api/ume_sync_service.py +++ b/netx_api/ume_sync_service.py @@ -132,6 +132,15 @@ def _collect_marker_pages( rows = [x for x in rows if isinstance(x, dict)] pages.append(rows) + if page_no == 1 or page_no % 25 == 0: + _sync_log.info( + "ume marker page=%s rows=%s is_end=%s marker_len=%s", + page_no, + len(rows), + getattr(diag, "is_end_of_reply", None), + len(str(getattr(diag, "marker", "") or "")), + ) + # Protection against repeated pages causing infinite loops. cur_sig = "|".join(sorted(_alarm_key(x) for x in rows)) if cur_sig and cur_sig == last_page_signature: