mirror of
https://github.com/hansjone/netx.git
synced 2026-10-11 06:50:49 +08:00
fix(ume): block WSS until startup REST alarm sync completes
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
7a729413ab
commit
e62f4f2c74
3 changed files with 39 additions and 47 deletions
|
|
@ -43,12 +43,10 @@ class Settings(BaseSettings):
|
||||||
ume_sync_alarms_current_enabled: bool = True
|
ume_sync_alarms_current_enabled: bool = True
|
||||||
ume_sync_alarms_current_interval_s: int = 18000
|
ume_sync_alarms_current_interval_s: int = 18000
|
||||||
ume_sync_alarms_current_skip_when_ws: bool = True
|
ume_sync_alarms_current_skip_when_ws: bool = True
|
||||||
# When False (default), WSS connects after delay without a blocking full REST snapshot.
|
# Block WSS until initial REST current-alarm snapshot finishes (see startup gate).
|
||||||
ume_startup_sync_alarms_before_ws: bool = False
|
ume_startup_sync_alarms_before_ws: bool = True
|
||||||
# Defer first REST alarm pull after process start (WSS may connect earlier).
|
# Wait this long after process start before the initial REST snapshot (WSS still blocked).
|
||||||
ume_startup_alarm_sync_delay_s: int = 60
|
ume_startup_alarm_sync_delay_s: int = 60
|
||||||
# After delay, wait for WSS before REST fallback (avoids REST+WSS fighting at boot).
|
|
||||||
ume_startup_wss_grace_s: int = 180
|
|
||||||
ume_alarm_cleared_tombstone_s: int = 300
|
ume_alarm_cleared_tombstone_s: int = 300
|
||||||
ume_alarm_ws_enabled: bool = True
|
ume_alarm_ws_enabled: bool = True
|
||||||
ume_notification_establish_path: str = "/restconf/operations/zte-notifications:establish-subscription"
|
ume_notification_establish_path: str = "/restconf/operations/zte-notifications:establish-subscription"
|
||||||
|
|
|
||||||
|
|
@ -50,6 +50,7 @@ from .ume_alarm_ws import (
|
||||||
get_subscription_status,
|
get_subscription_status,
|
||||||
get_ws_connection_status,
|
get_ws_connection_status,
|
||||||
get_ws_logs,
|
get_ws_logs,
|
||||||
|
is_startup_alarm_sync_pending,
|
||||||
is_wss_active_for_current_alarms,
|
is_wss_active_for_current_alarms,
|
||||||
load_persisted_subscription,
|
load_persisted_subscription,
|
||||||
request_ws_reconnect,
|
request_ws_reconnect,
|
||||||
|
|
@ -286,34 +287,20 @@ def _fail_stale_running_sync_jobs_on_startup() -> None:
|
||||||
db.close()
|
db.close()
|
||||||
|
|
||||||
|
|
||||||
|
def _needs_startup_alarm_sync_before_ws() -> bool:
|
||||||
|
ume_url = str(getattr(settings, "ume_base_url", "") or "").strip()
|
||||||
|
return bool(
|
||||||
|
getattr(settings, "ume_startup_sync_alarms_before_ws", True)
|
||||||
|
and getattr(settings, "ume_alarm_ws_enabled", True)
|
||||||
|
and getattr(settings, "ume_sync_alarms_current_enabled", True)
|
||||||
|
and ume_url
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def _startup_alarm_pull_delay_s() -> int:
|
def _startup_alarm_pull_delay_s() -> int:
|
||||||
return max(0, min(3600, int(getattr(settings, "ume_startup_alarm_sync_delay_s", 60) or 60)))
|
return max(0, min(3600, int(getattr(settings, "ume_startup_alarm_sync_delay_s", 60) or 60)))
|
||||||
|
|
||||||
|
|
||||||
def _startup_wss_grace_s() -> int:
|
|
||||||
return max(0, min(3600, int(getattr(settings, "ume_startup_wss_grace_s", 180) or 180)))
|
|
||||||
|
|
||||||
|
|
||||||
def _ume_alarm_ws_enabled() -> bool:
|
|
||||||
return bool(getattr(settings, "ume_alarm_ws_enabled", True))
|
|
||||||
|
|
||||||
|
|
||||||
def _should_defer_rest_current_sync() -> tuple[bool, str]:
|
|
||||||
"""Skip scheduled REST when WSS owns current alarms or grace period after boot."""
|
|
||||||
if (
|
|
||||||
bool(getattr(settings, "ume_sync_alarms_current_skip_when_ws", True))
|
|
||||||
and is_wss_active_for_current_alarms()
|
|
||||||
):
|
|
||||||
return True, "WSS 实时接收中,已跳过 REST 同步"
|
|
||||||
if not _ume_alarm_ws_enabled():
|
|
||||||
return False, ""
|
|
||||||
grace_end = float(_startup_alarm_pull_delay_s() + _startup_wss_grace_s())
|
|
||||||
if time.monotonic() - _BOOT_MONO < grace_end:
|
|
||||||
left = grace_end - (time.monotonic() - _BOOT_MONO)
|
|
||||||
return True, f"等待 WSS 连接(REST 兜底约 {left:.0f}s 后)"
|
|
||||||
return False, ""
|
|
||||||
|
|
||||||
|
|
||||||
def _wait_until_startup_alarm_pull_allowed(label: str) -> None:
|
def _wait_until_startup_alarm_pull_allowed(label: str) -> None:
|
||||||
delay_s = _startup_alarm_pull_delay_s()
|
delay_s = _startup_alarm_pull_delay_s()
|
||||||
if delay_s <= 0:
|
if delay_s <= 0:
|
||||||
|
|
@ -326,25 +313,14 @@ def _wait_until_startup_alarm_pull_allowed(label: str) -> None:
|
||||||
|
|
||||||
|
|
||||||
def _run_startup_alarm_sync_before_ws() -> None:
|
def _run_startup_alarm_sync_before_ws() -> None:
|
||||||
"""REST-sync current alarms once on boot before WSS connects (avoids stale/reconcile races)."""
|
"""REST-sync current alarms once on boot; WSS gate must already be closed in on_startup."""
|
||||||
ume_url = str(getattr(settings, "ume_base_url", "") or "").strip()
|
if not _needs_startup_alarm_sync_before_ws():
|
||||||
alarms_enabled = bool(getattr(settings, "ume_sync_alarms_current_enabled", True))
|
|
||||||
ws_enabled = bool(getattr(settings, "ume_alarm_ws_enabled", True))
|
|
||||||
before_ws = bool(getattr(settings, "ume_startup_sync_alarms_before_ws", True))
|
|
||||||
|
|
||||||
if not (before_ws and ws_enabled and alarms_enabled and ume_url):
|
|
||||||
complete_startup_alarm_sync_gate()
|
complete_startup_alarm_sync_gate()
|
||||||
return
|
return
|
||||||
|
|
||||||
_wait_until_startup_alarm_pull_allowed("startup_alarm_sync")
|
_wait_until_startup_alarm_pull_allowed("startup_alarm_sync")
|
||||||
if not before_ws:
|
|
||||||
_schedule_log.info("startup: skip REST snapshot before WSS (ume_startup_sync_alarms_before_ws=false)")
|
|
||||||
complete_startup_alarm_sync_gate()
|
|
||||||
return
|
|
||||||
|
|
||||||
begin_startup_alarm_sync_gate()
|
|
||||||
try:
|
try:
|
||||||
_schedule_log.info("startup: syncing current alarms before WSS (legacy mode)")
|
_schedule_log.info("startup: REST current-alarm snapshot (WSS blocked until finished)")
|
||||||
_set_runtime_task(
|
_set_runtime_task(
|
||||||
"alarms_current_auto_sync",
|
"alarms_current_auto_sync",
|
||||||
status="running",
|
status="running",
|
||||||
|
|
@ -707,6 +683,14 @@ def on_startup() -> None:
|
||||||
Base.metadata.create_all(bind=engine)
|
Base.metadata.create_all(bind=engine)
|
||||||
_reset_runtime_pause_flags()
|
_reset_runtime_pause_flags()
|
||||||
_fail_stale_running_sync_jobs_on_startup()
|
_fail_stale_running_sync_jobs_on_startup()
|
||||||
|
if _needs_startup_alarm_sync_before_ws():
|
||||||
|
begin_startup_alarm_sync_gate()
|
||||||
|
_schedule_log.info(
|
||||||
|
"startup: WSS blocked until initial REST current-alarm sync completes (delay=%ss)",
|
||||||
|
_startup_alarm_pull_delay_s(),
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
complete_startup_alarm_sync_gate()
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
try:
|
try:
|
||||||
from .collection_recovery import recover_collection_jobs_on_startup
|
from .collection_recovery import recover_collection_jobs_on_startup
|
||||||
|
|
@ -900,12 +884,22 @@ def on_startup() -> None:
|
||||||
if _runtime_is_paused("alarms_current_auto_sync"):
|
if _runtime_is_paused("alarms_current_auto_sync"):
|
||||||
time.sleep(1)
|
time.sleep(1)
|
||||||
continue
|
continue
|
||||||
defer, defer_reason = _should_defer_rest_current_sync()
|
if is_startup_alarm_sync_pending():
|
||||||
if defer:
|
|
||||||
_refresh_runtime_task_idle(
|
_refresh_runtime_task_idle(
|
||||||
"alarms_current_auto_sync",
|
"alarms_current_auto_sync",
|
||||||
"alarms_current",
|
"alarms_current",
|
||||||
last_error=defer_reason,
|
last_error="启动 REST 全量同步未完成,定时 REST 与 WSS 均待命",
|
||||||
|
)
|
||||||
|
time.sleep(10)
|
||||||
|
continue
|
||||||
|
if (
|
||||||
|
bool(getattr(settings, "ume_sync_alarms_current_skip_when_ws", True))
|
||||||
|
and is_wss_active_for_current_alarms()
|
||||||
|
):
|
||||||
|
_refresh_runtime_task_idle(
|
||||||
|
"alarms_current_auto_sync",
|
||||||
|
"alarms_current",
|
||||||
|
last_error="WSS 实时接收中,已跳过 REST 同步",
|
||||||
)
|
)
|
||||||
time.sleep(max(30, min(alarms_interval_s, 300)))
|
time.sleep(max(30, min(alarms_interval_s, 300)))
|
||||||
continue
|
continue
|
||||||
|
|
|
||||||
|
|
@ -53,8 +53,8 @@ _LOOP_STATUS_COOLDOWN_S = 60.0
|
||||||
_WS_IGNORE_COOLDOWN_S = 60.0
|
_WS_IGNORE_COOLDOWN_S = 60.0
|
||||||
|
|
||||||
# Blocks WSS connect until startup REST sync of current alarms completes (see main.on_startup).
|
# Blocks WSS connect until startup REST sync of current alarms completes (see main.on_startup).
|
||||||
|
# Default closed (clear): WSS must not write until main opens the gate after REST snapshot.
|
||||||
_STARTUP_ALARM_SYNC_GATE = threading.Event()
|
_STARTUP_ALARM_SYNC_GATE = threading.Event()
|
||||||
_STARTUP_ALARM_SYNC_GATE.set()
|
|
||||||
|
|
||||||
|
|
||||||
def begin_startup_alarm_sync_gate() -> None:
|
def begin_startup_alarm_sync_gate() -> None:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue