diff --git a/netx_api/config.py b/netx_api/config.py index 948ec32..ee00f2e 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -75,6 +75,8 @@ class Settings(BaseSettings): # Config sync (periodic running-config backup into DB) config_sync_scheduler_enabled: bool = True config_sync_scheduler_tick_sec: int = 60 + # After process start / unexpected restart, wait before any scheduled sync. + config_sync_startup_grace_sec: int = 3600 # Managed NE exec: max CLI commands per request (lab can raise; hard-capped in ne_exec). ne_exec_max_commands: int = 5 # WebCRT interactive terminal sessions diff --git a/netx_api/config_sync_recovery.py b/netx_api/config_sync_recovery.py index 8323b18..f39f563 100644 --- a/netx_api/config_sync_recovery.py +++ b/netx_api/config_sync_recovery.py @@ -3,66 +3,112 @@ from __future__ import annotations import logging +from datetime import datetime from sqlalchemy.orm import Session from .config_sync_runner import dispatch_cycle from .config_sync_service import finalize_cycle, sync_cycle_progress from .models import ConfigSyncCycle, ConfigSyncTask -from datetime import datetime _log = logging.getLogger("netx.config_sync.recovery") +_ACTIVE = ("running", "paused", "pending") + + +def _close_cycle(db: Session, cycle: ConfigSyncCycle, *, reason: str) -> None: + """Fail in-flight work and mark cycle terminal so it cannot block forever.""" + cycle_id = str(cycle.id) + now = datetime.utcnow() + for task in ( + db.query(ConfigSyncTask) + .filter( + ConfigSyncTask.cycle_id == cycle_id, + ConfigSyncTask.status.in_(("running", "pending")), + ) + .all() + ): + was_running = str(task.status) == "running" + task.status = "fail" if was_running else "cancelled" + task.message = reason + task.ended_at = now + db.commit() + sync_cycle_progress(db, cycle_id) + db.refresh(cycle) + cycle.status = "fail" + cycle.error_message = reason + cycle.ended_at = now + db.commit() + def recover_config_sync_on_startup(db: Session) -> int: """ - Mark orphaned running tasks as fail(orphan_recovered), then resume pending - tasks for cycles still marked running/paused. + Resume at most one interrupted cycle after process restart. + + Rules: + - Only one active cycle (running/pending/paused) may exist; older actives are closed. + - Orphan ``running`` tasks are re-queued to ``pending`` and continued (crash 续跑). + - ``paused`` cycles stay paused (no auto dispatch) but still occupy the single-flight slot. + - New scheduled cycles remain blocked while this active cycle exists. """ cycles = ( db.query(ConfigSyncCycle) - .filter(ConfigSyncCycle.status.in_(("running", "paused", "pending"))) + .filter(ConfigSyncCycle.status.in_(_ACTIVE)) .all() ) - resumed = 0 - for cycle in cycles: - cycle_id = str(cycle.id) - orphans = ( - db.query(ConfigSyncTask) - .filter(ConfigSyncTask.cycle_id == cycle_id, ConfigSyncTask.status == "running") - .all() + cycles = sorted(cycles, key=lambda c: c.created_at or datetime.min) + if not cycles: + return 0 + + # Single-flight hygiene: keep newest, close older interrupted cycles. + primary = cycles[-1] + for stale in cycles[:-1]: + _log.warning( + "config_sync recovery closing older active cycle=%s (keep=%s)", + stale.id, + primary.id, ) - for task in orphans: - task.status = "fail" - task.message = "orphan_recovered" - task.ended_at = datetime.utcnow() - if orphans: - db.commit() - _log.info("config_sync recovery cycle=%s orphaned_tasks=%s", cycle_id, len(orphans)) + _close_cycle(db, stale, reason="superseded_active_cycle") - sync_cycle_progress(db, cycle_id) - db.refresh(cycle) + cycle_id = str(primary.id) + orphans = ( + db.query(ConfigSyncTask) + .filter(ConfigSyncTask.cycle_id == cycle_id, ConfigSyncTask.status == "running") + .all() + ) + for task in orphans: + task.status = "pending" + task.message = "requeued_after_restart" + task.started_at = None + task.ended_at = None + if orphans: + db.commit() + _log.info("config_sync recovery cycle=%s requeued_orphans=%s", cycle_id, len(orphans)) - if str(cycle.status) == "paused": - continue + sync_cycle_progress(db, cycle_id) + db.refresh(primary) - pending = ( - db.query(ConfigSyncTask) - .filter(ConfigSyncTask.cycle_id == cycle_id, ConfigSyncTask.status == "pending") - .count() - ) - if pending <= 0: - if str(cycle.status) in ("running", "pending"): - finalize_cycle(db, cycle_id) - continue + if str(primary.status) == "paused": + _log.info("config_sync recovery cycle=%s stays paused (blocks new cycles)", cycle_id) + return 0 - if str(cycle.status) == "pending": - cycle.status = "running" - if not cycle.started_at: - cycle.started_at = datetime.utcnow() - db.commit() + pending = ( + db.query(ConfigSyncTask) + .filter(ConfigSyncTask.cycle_id == cycle_id, ConfigSyncTask.status == "pending") + .all() + ) + if not pending: + finalize_cycle(db, cycle_id) + _log.info("config_sync recovery cycle=%s finalized (no pending)", cycle_id) + return 0 - n = dispatch_cycle(cycle_id) - resumed += n - _log.info("config_sync recovery resumed cycle=%s pending=%s", cycle_id, n) - return resumed + if str(primary.status) == "pending": + primary.status = "running" + if not primary.started_at: + primary.started_at = datetime.utcnow() + db.commit() + + task_ids = [str(t.id) for t in pending] + n = dispatch_cycle(cycle_id) + _log.info("config_sync recovery resumed cycle=%s pending=%s dispatched=%s", cycle_id, len(task_ids), n) + return n diff --git a/netx_api/config_sync_runner.py b/netx_api/config_sync_runner.py index 4167ad1..5a0077a 100644 --- a/netx_api/config_sync_runner.py +++ b/netx_api/config_sync_runner.py @@ -70,7 +70,7 @@ def _update_task(task_id: str, **fields: Any) -> None: db.close() -def _update_task(task_id: str, **fields: Any) -> None: +def _claim_task(cycle_id: str, task_id: str) -> bool: for attempt in range(10): db = SessionLocal() try: diff --git a/netx_api/config_sync_scheduler.py b/netx_api/config_sync_scheduler.py index 633ba0d..4f052b6 100644 --- a/netx_api/config_sync_scheduler.py +++ b/netx_api/config_sync_scheduler.py @@ -5,7 +5,7 @@ from __future__ import annotations import logging import threading import time -from datetime import datetime +from datetime import datetime, timedelta from uuid import uuid4 from .config import settings @@ -22,20 +22,50 @@ from .models import ConfigSyncCycle, ConfigSyncTask _log = logging.getLogger("netx.config_sync.scheduler") _stop = threading.Event() _thread: threading.Thread | None = None +_BOOT_MONO = time.monotonic() def _utcnow() -> datetime: return datetime.utcnow() +def startup_grace_remaining_sec() -> float: + grace = max(0, int(settings.config_sync_startup_grace_sec or 0)) + elapsed = time.monotonic() - _BOOT_MONO + return max(0.0, float(grace) - elapsed) + + +def in_startup_grace() -> bool: + return startup_grace_remaining_sec() > 0 + + +def startup_grace_until() -> datetime | None: + rem = startup_grace_remaining_sec() + if rem <= 0: + return None + return _utcnow() + timedelta(seconds=rem) + + def try_start_scheduled_cycle() -> str | None: - """Create and dispatch a scheduled cycle if policy is due. Returns cycle id or None.""" + """Create and dispatch a scheduled cycle if policy is due. Returns cycle id or None. + + Never starts while another cycle is active (running/pending/paused), including + a cycle being resumed after crash. Startup grace only delays *new* scheduled runs. + """ + if in_startup_grace(): + return None db = SessionLocal() try: policy = ensure_policy(db) if not policy.enabled: return None - if has_running_cycle(db): + active = has_running_cycle(db) + if active: + _log.debug( + "config_sync schedule skip: active cycle=%s status=%s", + active.id, + active.status, + ) return None due = next_due_at(db, policy) if due is not None and due > _utcnow(): @@ -85,7 +115,8 @@ def try_start_scheduled_cycle() -> str | None: def _loop() -> None: tick = max(15, int(settings.config_sync_scheduler_tick_sec or 60)) - _log.info("config_sync scheduler started tick=%ss", tick) + grace = max(0, int(settings.config_sync_startup_grace_sec or 0)) + _log.info("config_sync scheduler started tick=%ss startup_grace=%ss", tick, grace) while not _stop.is_set(): try: if bool(settings.config_sync_scheduler_enabled): diff --git a/netx_api/config_sync_service.py b/netx_api/config_sync_service.py index 42e4863..5c1fbb3 100644 --- a/netx_api/config_sync_service.py +++ b/netx_api/config_sync_service.py @@ -47,7 +47,7 @@ def _utcnow() -> datetime: def ensure_policy(db: Session) -> ConfigSyncPolicy: row = db.get(ConfigSyncPolicy, POLICY_ID) if row is None: - row = ConfigSyncPolicy(id=POLICY_ID) + row = ConfigSyncPolicy(id=POLICY_ID, enabled=False) db.add(row) db.commit() db.refresh(row) @@ -202,15 +202,21 @@ def expand_targets(db: Session, policy: ConfigSyncPolicy) -> list[dict[str, str] return out -def has_running_cycle(db: Session) -> ConfigSyncCycle | None: +def has_active_cycle(db: Session) -> ConfigSyncCycle | None: + """Any non-terminal cycle occupies the single-flight slot (incl. paused).""" return ( db.query(ConfigSyncCycle) - .filter(ConfigSyncCycle.status.in_(("running", "pending"))) + .filter(ConfigSyncCycle.status.in_(("running", "pending", "paused"))) .order_by(ConfigSyncCycle.created_at.desc()) .first() ) +def has_running_cycle(db: Session) -> ConfigSyncCycle | None: + """Backward-compatible alias: treat paused as active so a new cycle cannot start.""" + return has_active_cycle(db) + + def last_finished_cycle(db: Session) -> ConfigSyncCycle | None: return ( db.query(ConfigSyncCycle) @@ -224,6 +230,8 @@ def next_due_at(db: Session, policy: ConfigSyncPolicy | None = None) -> datetime pol = policy or ensure_policy(db) if not pol.enabled: return None + from .config_sync_scheduler import startup_grace_until + last = ( db.query(ConfigSyncCycle) .filter(ConfigSyncCycle.status == "success", ConfigSyncCycle.ended_at.isnot(None)) @@ -232,8 +240,14 @@ def next_due_at(db: Session, policy: ConfigSyncPolicy | None = None) -> datetime ) days = max(1, int(pol.interval_days or 3)) if last and last.ended_at: - return last.ended_at + timedelta(days=days) - return _utcnow() + due = last.ended_at + timedelta(days=days) + else: + # Never synced successfully: do not fire immediately on enable / first boot. + due = _utcnow() + timedelta(days=days) + grace_until = startup_grace_until() + if grace_until is not None and due < grace_until: + return grace_until + return due def create_cycle(db: Session, body: ConfigSyncCycleCreate) -> ConfigSyncCycleOut: diff --git a/netx_api/main.py b/netx_api/main.py index 11f46e1..5659920 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -840,7 +840,7 @@ def on_startup() -> None: ensure_policy(db) cfg_resumed = recover_config_sync_on_startup(db) if cfg_resumed: - _schedule_log.info("startup: resumed %s pending config_sync tasks", cfg_resumed) + _schedule_log.info("startup: resumed %s config_sync task(s) from interrupted cycle", cfg_resumed) except Exception: _schedule_log.exception("startup: ne collection / config_sync recovery failed") finally: diff --git a/netx_api/models.py b/netx_api/models.py index d573a89..7d01fd8 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -492,7 +492,7 @@ class ConfigSyncPolicy(Base): __tablename__ = "config_sync_policy" id: Mapped[int] = mapped_column(Integer, primary_key=True, default=1) - enabled: Mapped[bool] = mapped_column(Boolean, default=True) + enabled: Mapped[bool] = mapped_column(Boolean, default=False) interval_days: Mapped[int] = mapped_column(Integer, default=3) concurrency: Mapped[int] = mapped_column(Integer, default=5) scope_mode: Mapped[str] = mapped_column(String(32), default="all") # all | selected diff --git a/tests/test_config_sync.py b/tests/test_config_sync.py index 1ffa259..52e2ff3 100644 --- a/tests/test_config_sync.py +++ b/tests/test_config_sync.py @@ -211,33 +211,81 @@ class ConfigSyncRecoveryTests(unittest.TestCase): @patch("netx_api.config_sync_recovery.dispatch_cycle", return_value=2) @patch("netx_api.config_sync_recovery.finalize_cycle") @patch("netx_api.config_sync_recovery.sync_cycle_progress") - def test_orphan_running_tasks_marked_fail(self, _sync, _fin, dispatch): + def test_requeues_orphans_and_resumes(self, _sync, _fin, dispatch): cycle = MagicMock() cycle.id = "c1" cycle.status = "running" + cycle.created_at = datetime.now(timezone.utc) cycle.started_at = datetime.now(timezone.utc) + cycle.error_message = "" + cycle.ended_at = None orphan = MagicMock() orphan.status = "running" orphan.message = "" + orphan.started_at = datetime.now(timezone.utc) orphan.ended_at = None + pending = MagicMock() + pending.status = "pending" + pending.id = "t2" + db = MagicMock() db.query.side_effect = [ MagicMock(filter=MagicMock(return_value=MagicMock(all=MagicMock(return_value=[cycle])))), MagicMock(filter=MagicMock(return_value=MagicMock(all=MagicMock(return_value=[orphan])))), - MagicMock(filter=MagicMock(return_value=MagicMock(count=MagicMock(return_value=2)))), + MagicMock( + filter=MagicMock( + return_value=MagicMock(all=MagicMock(return_value=[orphan, pending])) + ) + ), ] db.refresh = MagicMock() resumed = recover_config_sync_on_startup(db) - self.assertEqual(orphan.status, "fail") - self.assertEqual(orphan.message, "orphan_recovered") - self.assertIsNotNone(orphan.ended_at) + self.assertEqual(orphan.status, "pending") + self.assertEqual(orphan.message, "requeued_after_restart") + self.assertIsNone(orphan.started_at) self.assertEqual(resumed, 2) dispatch.assert_called_once_with("c1") + @patch("netx_api.config_sync_recovery.dispatch_cycle") + @patch("netx_api.config_sync_recovery.finalize_cycle") + @patch("netx_api.config_sync_recovery.sync_cycle_progress") + def test_closes_older_active_keeps_newest(self, _sync, _fin, dispatch): + old = MagicMock() + old.id = "old" + old.status = "running" + old.created_at = datetime(2026, 1, 1, tzinfo=timezone.utc) + old.error_message = "" + old.ended_at = None + + new = MagicMock() + new.id = "new" + new.status = "running" + new.created_at = datetime(2026, 1, 2, tzinfo=timezone.utc) + new.started_at = datetime.now(timezone.utc) + new.error_message = "" + new.ended_at = None + + db = MagicMock() + db.query.side_effect = [ + MagicMock(filter=MagicMock(return_value=MagicMock(all=MagicMock(return_value=[old, new])))), + MagicMock(filter=MagicMock(return_value=MagicMock(all=MagicMock(return_value=[])))), + MagicMock(filter=MagicMock(return_value=MagicMock(all=MagicMock(return_value=[])))), + MagicMock(filter=MagicMock(return_value=MagicMock(all=MagicMock(return_value=[])))), + ] + db.refresh = MagicMock() + dispatch.return_value = 0 + + recover_config_sync_on_startup(db) + + self.assertEqual(old.status, "fail") + self.assertEqual(old.error_message, "superseded_active_cycle") + _fin.assert_called() + dispatch.assert_not_called() + if __name__ == "__main__": unittest.main() diff --git a/web/WEB.md b/web/WEB.md index 3cbcb55..bc4f742 100644 --- a/web/WEB.md +++ b/web/WEB.md @@ -98,6 +98,10 @@ src/ - 存储:PostgreSQL `ne_config_snapshot` / `ne_config_history`(zlib BYTEA) - 范围:ManagedNE + UME CLI 目标;厂商固定只读命令矩阵 - 调度:`NETX_CONFIG_SYNC_SCHEDULER_ENABLED`(默认开),周期天数策略可配(默认 3 天) +- 默认策略:`enabled=false`(首次无自动任务,需在页面手动开启周期调度或点「立即同步」) +- 单飞:同一时刻只允许一个 `running|pending|paused` 周期;上轮未结束时不会开启新周期 +- 崩溃续跑:启动时把中断的 `running` 任务重新入队并继续,占用单飞槽位,避免与新周期重叠 +- 进程启动宽限:`NETX_CONFIG_SYNC_STARTUP_GRACE_SEC`(默认 3600)仅约束**新建**自动周期,不影响续跑 - 前端:`/network/tasks/config-sync`(看板)+ `/network/configs`(查看) ## WebCRT diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 984887d..a342363 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -55,6 +55,8 @@ const en = { configSync: "Config sync", portTraffic: "Port traffic", }, + collapseNav: "Collapse sidebar", + expandNav: "Expand sidebar", placeholder: { title: "Coming soon", body: "This capability is not available yet.", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index c096464..0c33358 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -55,6 +55,8 @@ const zh = { configSync: "配置同步", portTraffic: "端口流量监控", }, + collapseNav: "折叠侧栏", + expandNav: "展开侧栏", placeholder: { title: "功能规划中", body: "该能力尚未开放。", diff --git a/web/src/index.css b/web/src/index.css index e54efb3..b87cc07 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -2670,6 +2670,89 @@ pre { background: #f7f9fc; border-right: 1px solid #d0d7e2; min-height: 0; + transition: width 0.18s ease; +} + +.network-shell.is-nav-collapsed .network-nav { + width: 48px; +} + +.network-nav__head { + display: flex; + align-items: center; + gap: 8px; + padding: 10px 8px 8px 12px; + border-bottom: 1px solid #e2e8f0; + min-height: 44px; +} + +.network-shell.is-nav-collapsed .network-nav__head { + justify-content: center; + padding: 10px 4px; +} + +.network-nav__title { + flex: 1; + min-width: 0; + font-size: 13px; + font-weight: 700; + color: #0f172a; + white-space: nowrap; + overflow: hidden; + text-overflow: ellipsis; +} + +.app-main .network-nav__collapse-btn, +.network-nav__collapse-btn { + flex-shrink: 0; + width: 28px; + height: 28px; + padding: 0; + border: 1px solid #cbd5e1; + border-radius: 6px; + background: #fff; + color: #475569; + cursor: pointer; + line-height: 1; +} + +.app-main .network-nav__collapse-btn:hover, +.network-nav__collapse-btn:hover { + background: #eef2f7; + color: #0f172a; +} + +.network-nav__rail { + display: flex; + flex-direction: column; + align-items: center; + gap: 6px; + padding: 10px 4px; + overflow: auto; +} + +.network-nav__rail-link { + display: flex; + align-items: center; + justify-content: center; + width: 32px; + height: 32px; + border-radius: 8px; + text-decoration: none; + color: #475569; + font-size: 12px; + font-weight: 700; + background: transparent; +} + +.network-nav__rail-link:hover { + background: #eef2f7; + color: #0f172a; +} + +.network-nav__rail-link.is-active { + background: #e3effc; + color: #1565c0; } .network-nav__scroll { diff --git a/web/src/pages/ConfigSyncPage.tsx b/web/src/pages/ConfigSyncPage.tsx index 0f63af9..920b031 100644 --- a/web/src/pages/ConfigSyncPage.tsx +++ b/web/src/pages/ConfigSyncPage.tsx @@ -31,7 +31,7 @@ export function ConfigSyncPage() { const [taskStatus, setTaskStatus] = useState(""); const [taskKeyword, setTaskKeyword] = useState(""); - const [enabled, setEnabled] = useState(true); + const [enabled, setEnabled] = useState(false); const [intervalDays, setIntervalDays] = useState(3); const [concurrency, setConcurrency] = useState(5); const [scopeMode, setScopeMode] = useState<"all" | "selected">("all"); diff --git a/web/src/pages/network/NetworkLayout.tsx b/web/src/pages/network/NetworkLayout.tsx index 216da16..4d2c089 100644 --- a/web/src/pages/network/NetworkLayout.tsx +++ b/web/src/pages/network/NetworkLayout.tsx @@ -3,6 +3,8 @@ import { NavLink, Outlet, useLocation } from "react-router-dom"; import { NETWORK_NAV, type NetworkNavGroupId } from "../../config/networkNav"; import { useI18n } from "../../i18n"; +const SIDEBAR_KEY = "netx.network.sidebarCollapsed"; + function groupContainsPath(groupId: NetworkNavGroupId, pathname: string): boolean { const group = NETWORK_NAV.find((g) => g.id === groupId); if (!group) return false; @@ -11,9 +13,18 @@ function groupContainsPath(groupId: NetworkNavGroupId, pathname: string): boolea ); } +function readSidebarCollapsed(): boolean { + try { + return localStorage.getItem(SIDEBAR_KEY) === "1"; + } catch { + return false; + } +} + export function NetworkLayout() { const { t } = useI18n(); const { pathname } = useLocation(); + const [sidebarCollapsed, setSidebarCollapsed] = useState(readSidebarCollapsed); const [openGroups, setOpenGroups] = useState>(() => { const init = {} as Record; for (const g of NETWORK_NAV) { @@ -35,50 +46,95 @@ export function NetworkLayout() { setOpenGroups((prev) => ({ ...prev, [id]: !prev[id] })); }; + const toggleSidebar = () => { + setSidebarCollapsed((prev) => { + const next = !prev; + try { + localStorage.setItem(SIDEBAR_KEY, next ? "1" : "0"); + } catch { + // ignore + } + return next; + }); + }; + return ( -
+
+ ); + })} + + ) : ( + + )}