diff --git a/netx_api/config_sync_recovery.py b/netx_api/config_sync_recovery.py index f39f563..16862e0 100644 --- a/netx_api/config_sync_recovery.py +++ b/netx_api/config_sync_recovery.py @@ -8,7 +8,14 @@ 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 .config_sync_service import ( + DEFAULT_CYCLE_KEEP, + _cycle_keep_value, + ensure_policy, + finalize_cycle, + prune_config_sync_cycles, + sync_cycle_progress, +) from .models import ConfigSyncCycle, ConfigSyncTask _log = logging.getLogger("netx.config_sync.recovery") @@ -57,7 +64,18 @@ def recover_config_sync_on_startup(db: Session) -> int: .all() ) cycles = sorted(cycles, key=lambda c: c.created_at or datetime.min) + def _prune() -> None: + try: + keep = _cycle_keep_value(ensure_policy(db)) + except Exception: + keep = DEFAULT_CYCLE_KEEP + try: + prune_config_sync_cycles(db, keep=keep) + except Exception: + _log.exception("config_sync recovery prune failed") + if not cycles: + _prune() return 0 # Single-flight hygiene: keep newest, close older interrupted cycles. @@ -69,6 +87,7 @@ def recover_config_sync_on_startup(db: Session) -> int: primary.id, ) _close_cycle(db, stale, reason="superseded_active_cycle") + _prune() cycle_id = str(primary.id) orphans = ( diff --git a/netx_api/config_sync_schemas.py b/netx_api/config_sync_schemas.py index 1ab214c..67418dc 100644 --- a/netx_api/config_sync_schemas.py +++ b/netx_api/config_sync_schemas.py @@ -20,6 +20,7 @@ class ConfigSyncPolicyOut(BaseModel): scope_mode: str selected_targets: list[ConfigSyncTargetRef] = Field(default_factory=list) history_keep: int + cycle_keep: int = 30 updated_at: datetime | None = None @@ -30,6 +31,7 @@ class ConfigSyncPolicyUpdate(BaseModel): scope_mode: Literal["all", "selected"] | None = None selected_targets: list[ConfigSyncTargetRef] | None = None history_keep: int | None = Field(default=None, ge=0, le=30) + cycle_keep: int | None = Field(default=None, ge=0, le=200) class ConfigSyncCycleCreate(BaseModel): diff --git a/netx_api/config_sync_service.py b/netx_api/config_sync_service.py index 80a40e7..8db370d 100644 --- a/netx_api/config_sync_service.py +++ b/netx_api/config_sync_service.py @@ -47,6 +47,9 @@ def _utcnow() -> datetime: return datetime.utcnow() +DEFAULT_CYCLE_KEEP = 30 + + def ensure_policy(db: Session) -> ConfigSyncPolicy: row = db.get(ConfigSyncPolicy, POLICY_ID) if row is None: @@ -57,6 +60,35 @@ def ensure_policy(db: Session) -> ConfigSyncPolicy: return row +def prune_config_sync_cycles(db: Session, *, keep: int = DEFAULT_CYCLE_KEEP) -> int: + """Delete finished cycles beyond ``keep`` (newest kept). Active cycles always retained.""" + keep = max(0, min(200, int(keep))) + finished = ( + db.query(ConfigSyncCycle) + .filter(ConfigSyncCycle.status.in_(("success", "fail", "cancelled"))) + .order_by(ConfigSyncCycle.created_at.desc()) + .all() + ) + to_drop = finished if keep == 0 else finished[keep:] + if not to_drop: + return 0 + dropped = 0 + for cycle in to_drop: + cid = str(cycle.id) + db.query(ConfigSyncTask).filter(ConfigSyncTask.cycle_id == cid).delete( + synchronize_session=False + ) + db.delete(cycle) + dropped += 1 + if dropped: + db.commit() + return dropped + + +def _cycle_keep_value(row: ConfigSyncPolicy) -> int: + return max(0, min(200, int(getattr(row, "cycle_keep", None) or DEFAULT_CYCLE_KEEP))) + + def _targets_from_json(raw: Any) -> list[ConfigSyncTargetRef]: items: list[ConfigSyncTargetRef] = [] if not isinstance(raw, list): @@ -80,6 +112,7 @@ def policy_to_out(row: ConfigSyncPolicy) -> ConfigSyncPolicyOut: scope_mode=str(row.scope_mode or "all"), selected_targets=_targets_from_json(row.selected_targets), history_keep=max(0, min(30, int(row.history_keep if row.history_keep is not None else 3))), + cycle_keep=_cycle_keep_value(row), updated_at=row.updated_at, ) @@ -107,9 +140,12 @@ def update_policy(db: Session, body: ConfigSyncPolicyUpdate) -> ConfigSyncPolicy ] if "history_keep" in data and data["history_keep"] is not None: row.history_keep = max(0, min(30, int(data["history_keep"]))) + if "cycle_keep" in data and data["cycle_keep"] is not None: + row.cycle_keep = max(0, min(200, int(data["cycle_keep"]))) row.updated_at = _utcnow() db.commit() db.refresh(row) + prune_config_sync_cycles(db, keep=_cycle_keep_value(row)) return policy_to_out(row) @@ -443,7 +479,12 @@ def stop_cycle(db: Session, cycle_id: str) -> ConfigSyncCycleOut: _release_pool(cycle_id) except Exception: pass - return cycle_to_out(row) + out = cycle_to_out(row) + try: + prune_config_sync_cycles(db, keep=_cycle_keep_value(ensure_policy(db))) + except Exception: + _log.exception("prune_config_sync_cycles after stop failed") + return out def dashboard(db: Session) -> ConfigSyncDashboardOut: @@ -690,3 +731,7 @@ def finalize_cycle(db: Session, cycle_id: str) -> None: cycle.error_message = "" cycle.ended_at = _utcnow() db.commit() + try: + prune_config_sync_cycles(db, keep=_cycle_keep_value(ensure_policy(db))) + except Exception: + _log.exception("prune_config_sync_cycles after finish failed") diff --git a/netx_api/main.py b/netx_api/main.py index c65e8ce..1ee43cf 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -829,6 +829,9 @@ def on_startup() -> None: ) except Exception: pass + conn.exec_driver_sql( + "ALTER TABLE config_sync_policy ADD COLUMN IF NOT EXISTS cycle_keep INTEGER DEFAULT 30" + ) except Exception: _schedule_log.exception("startup: auth/port_traffic/topology schema migration failed") _reset_runtime_pause_flags() diff --git a/netx_api/models.py b/netx_api/models.py index a22a33e..37d8f4a 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -693,6 +693,8 @@ class ConfigSyncPolicy(Base): scope_mode: Mapped[str] = mapped_column(String(32), default="all") # all | selected selected_targets: Mapped[list] = mapped_column(_JsonType, default=list) history_keep: Mapped[int] = mapped_column(Integer, default=3) + # Finished sync cycles to retain (newest kept); active cycles always kept. + cycle_keep: Mapped[int] = mapped_column(Integer, default=30) updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) diff --git a/tests/test_config_sync.py b/tests/test_config_sync.py index 52e2ff3..f92c7ee 100644 --- a/tests/test_config_sync.py +++ b/tests/test_config_sync.py @@ -3,13 +3,17 @@ from __future__ import annotations import unittest -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from unittest.mock import MagicMock, patch +from uuid import uuid4 from netx_api.config_sync_codec import compress_text, decompress_text from netx_api.config_sync_commands import command_list, commands_for_vendor, normalize_vendor_key from netx_api.config_sync_recovery import recover_config_sync_on_startup from netx_api.config_sync_runner import _claim_task, _save_success_snapshot +from netx_api.config_sync_service import prune_config_sync_cycles +from netx_api.db import Base, SessionLocal, engine +from netx_api.models import ConfigSyncCycle, ConfigSyncTask class ConfigSyncCommandsTests(unittest.TestCase): @@ -208,10 +212,11 @@ class ConfigSyncClaimTests(unittest.TestCase): class ConfigSyncRecoveryTests(unittest.TestCase): + @patch("netx_api.config_sync_recovery.prune_config_sync_cycles") @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_requeues_orphans_and_resumes(self, _sync, _fin, dispatch): + def test_requeues_orphans_and_resumes(self, _sync, _fin, dispatch, _prune): cycle = MagicMock() cycle.id = "c1" cycle.status = "running" @@ -250,10 +255,11 @@ class ConfigSyncRecoveryTests(unittest.TestCase): self.assertEqual(resumed, 2) dispatch.assert_called_once_with("c1") + @patch("netx_api.config_sync_recovery.prune_config_sync_cycles") @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): + def test_closes_older_active_keeps_newest(self, _sync, _fin, dispatch, _prune): old = MagicMock() old.id = "old" old.status = "running" @@ -287,5 +293,66 @@ class ConfigSyncRecoveryTests(unittest.TestCase): dispatch.assert_not_called() +class ConfigSyncCyclePruneTests(unittest.TestCase): + @classmethod + def setUpClass(cls) -> None: + Base.metadata.create_all(bind=engine) + with engine.begin() as conn: + try: + conn.exec_driver_sql( + "ALTER TABLE config_sync_policy ADD COLUMN IF NOT EXISTS cycle_keep INTEGER DEFAULT 30" + ) + except Exception: + pass + + def setUp(self) -> None: + self.db = SessionLocal() + self.db.query(ConfigSyncTask).delete() + self.db.query(ConfigSyncCycle).delete() + self.db.commit() + + def tearDown(self) -> None: + self.db.close() + + def test_prune_cycles_keeps_newest_and_active(self) -> None: + now = datetime.utcnow() + for i in range(5): + cid = uuid4().hex + self.db.add( + ConfigSyncCycle( + id=cid, + trigger_mode="manual", + status="success", + created_at=now - timedelta(minutes=5 - i), + ended_at=now - timedelta(minutes=5 - i), + ) + ) + self.db.add( + ConfigSyncTask( + id=uuid4().hex, + cycle_id=cid, + source="managed", + target_id=f"ne-{i}", + status="success", + ) + ) + open_id = uuid4().hex + self.db.add( + ConfigSyncCycle( + id=open_id, + trigger_mode="manual", + status="running", + created_at=now, + ) + ) + self.db.commit() + dropped = prune_config_sync_cycles(self.db, keep=2) + self.assertEqual(dropped, 3) + left = {r.id for r in self.db.query(ConfigSyncCycle).all()} + self.assertIn(open_id, left) + self.assertEqual(len(left), 3) + self.assertEqual(self.db.query(ConfigSyncTask).count(), 2) + + if __name__ == "__main__": unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 604a2f4..6f14e95 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -304,7 +304,8 @@ const en = { enabled: "Enable scheduled sync", intervalDays: "Interval (days)", concurrency: "Concurrency", - historyKeep: "History keep", + historyKeep: "Config versions to keep", + cycleKeep: "Jobs to keep", scope: "Scope", scopeAll: "All NEs", scopeSelected: "Selected NEs", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index b7716ba..1330e27 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -304,7 +304,8 @@ const zh = { enabled: "启用周期调度", intervalDays: "周期(天)", concurrency: "并发", - historyKeep: "历史保留", + historyKeep: "配置历史保留", + cycleKeep: "任务保留数", scope: "范围", scopeAll: "全部网元", scopeSelected: "指定网元", diff --git a/web/src/pages/ConfigSyncPage.tsx b/web/src/pages/ConfigSyncPage.tsx index cdf79e9..9b473f7 100644 --- a/web/src/pages/ConfigSyncPage.tsx +++ b/web/src/pages/ConfigSyncPage.tsx @@ -37,6 +37,7 @@ export function ConfigSyncPage() { const [concurrency, setConcurrency] = useState(5); const [scopeMode, setScopeMode] = useState<"all" | "selected">("all"); const [historyKeep, setHistoryKeep] = useState(3); + const [cycleKeep, setCycleKeep] = useState(30); const [selectedMap, setSelectedMap] = useState>({}); const [policyHydrated, setPolicyHydrated] = useState(false); @@ -63,6 +64,7 @@ export function ConfigSyncPage() { setConcurrency(Number(p.concurrency || 5)); setScopeMode(p.scope_mode === "selected" ? "selected" : "all"); setHistoryKeep(Number(p.history_keep ?? 3)); + setCycleKeep(Math.max(0, Math.min(200, Number(p.cycle_keep ?? 30)))); const map: Record = {}; for (const ref of p.selected_targets || []) { map[`${ref.source}:${ref.id}`] = { source: ref.source, id: ref.id }; @@ -117,6 +119,7 @@ export function ConfigSyncPage() { concurrency, scope_mode: scopeMode, history_keep: historyKeep, + cycle_keep: cycleKeep, selected_targets: Object.values(selectedMap), }), onSuccess: async (saved) => { @@ -127,6 +130,7 @@ export function ConfigSyncPage() { setConcurrency(Number(saved.concurrency || 5)); setScopeMode(saved.scope_mode === "selected" ? "selected" : "all"); setHistoryKeep(Number(saved.history_keep ?? 3)); + setCycleKeep(Math.max(0, Math.min(200, Number(saved.cycle_keep ?? 30)))); const map: Record = {}; for (const ref of saved.selected_targets || []) { map[`${ref.source}:${ref.id}`] = { source: ref.source, id: ref.id }; @@ -313,6 +317,16 @@ export function ConfigSyncPage() { onChange={(e) => setHistoryKeep(Math.max(0, Math.min(30, Number(e.target.value) || 0)))} /> +