Heal UME dock-to-Fabric gaps: partial jobs, startup apply, and one-click UI.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-07 10:33:25 +08:00
parent 74672b733f
commit 8444d85005
11 changed files with 280 additions and 5 deletions

View file

@ -326,6 +326,22 @@ def start_api_sideband_threads() -> None:
try:
fail_stale_topology_running_jobs(db)
client = ume_support._ume_client()
# Empty inventory makes Fabric naming/IP weak — pull inventory first.
try:
from sqlalchemy import func
from .models import UmeInventoryNE
inv_n = int(db.query(func.count(UmeInventoryNE.ne_id)).scalar() or 0)
if inv_n <= 0:
_schedule_log.info(
"topology_auto_sync: inventory empty — syncing inventory first"
)
sync_inventory_full(db, client, trigger_mode="schedule")
except Exception:
_schedule_log.exception(
"topology_auto_sync: pre-inventory sync failed (continuing topology)"
)
sync_topology_full(db, client, trigger_mode="schedule")
_schedule_log.info("topology_auto_sync: sync finished ok")
ume_support._set_runtime_task(
@ -445,5 +461,37 @@ def start_api_sideband_threads() -> None:
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
# Dock already has topology but Fabric/World empty (or last apply partial): heal without waiting 24h.
try:
def _startup_ume_fabric_apply() -> None:
time.sleep(3)
db = SessionLocal()
try:
from .ume_topology_apply import apply_ume_dock_to_fabric_if_needed
out = apply_ume_dock_to_fabric_if_needed(db, reason="startup")
if out:
_schedule_log.info(
"startup: ume dock→fabric apply ok nodes_ensured=%s",
(out.get("fabric_apply") or {}).get("nodes_ensured"),
)
else:
_schedule_log.info("startup: ume dock→fabric apply skipped (no gap)")
except Exception as exc:
_schedule_log.exception("startup: ume dock→fabric apply failed: %s", exc)
finally:
db.close()
t_apply = threading.Thread(
target=_startup_ume_fabric_apply,
name="ume-fabric-apply-startup",
daemon=True,
)
t_apply.start()
_schedule_log.info("started thread %s alive=%s", t_apply.name, t_apply.is_alive())
except Exception as exc:
_schedule_log.exception("startup: ume fabric apply thread init failed: %s", exc)

View file

@ -394,7 +394,7 @@ def _last_finished_job_ended_at(db: Session, domain: str) -> datetime | None:
.filter(
UmeSyncJob.domain == domain,
UmeSyncJob.ended_at.isnot(None),
UmeSyncJob.status.in_(("done", "failed")),
UmeSyncJob.status.in_(("done", "failed", "partial")),
)
.order_by(UmeSyncJob.ended_at.desc())
.limit(50)

View file

@ -204,6 +204,16 @@ def ume_sync_status(
items.append(item)
if item["domain"] and item["domain"] not in latest_by_domain:
latest_by_domain[item["domain"]] = item
topology_fabric: dict[str, Any] = {}
try:
from .ume_topology_apply import ume_topology_apply_gap
topology_fabric = ume_topology_apply_gap(db)
except Exception:
_log.exception("ume sync status: topology_fabric gap failed")
topology_fabric = {"needs_apply": False, "error": "gap_unavailable"}
return {
"total": total,
"page": page,
@ -212,6 +222,7 @@ def ume_sync_status(
"latest_by_domain": latest_by_domain,
"runtime_tasks": _list_runtime_tasks(),
"alarm_subscription": get_subscription_status(),
"topology_fabric": topology_fabric,
}

View file

@ -530,6 +530,8 @@ def sync_topology_full(db: Session, client: UMEClient, *, trigger_mode: str = "m
work.commit()
apply_stats: dict[str, Any] = {}
apply_ok = False
apply_err = ""
try:
from .ume_topology_apply import apply_ume_topology_to_fabric
from .ume_topology_world import ensure_ume_world_and_sbn_folders
@ -540,17 +542,24 @@ def sync_topology_full(db: Session, client: UMEClient, *, trigger_mode: str = "m
apply_stats = apply_ume_topology_to_fabric(apply_db)
world_stats = ensure_ume_world_and_sbn_folders(apply_db)
apply_stats = {**apply_stats, "world": world_stats}
apply_ok = True
finally:
try:
apply_db.close()
except Exception:
pass
except Exception:
except Exception as apply_exc:
apply_err = str(apply_exc)[:500]
apply_stats = {"error": apply_err, "ok": False}
_sync_log.exception(
"topology sync job %s: fabric apply failed (dock sync kept)",
job_id,
)
if not apply_ok:
status = "partial"
err = f"dock_ok_fabric_apply_failed:{apply_err}"[:1024]
pulled = nodes_pulled + links_pulled
inserted = nodes_ins + links_ins
updated = nodes_upd + links_upd
@ -577,12 +586,14 @@ def sync_topology_full(db: Session, client: UMEClient, *, trigger_mode: str = "m
"deleted_topo_nodes": nodes_del,
"deleted_topo_links": links_del,
"fabric_apply": apply_stats,
"fabric_apply_ok": apply_ok,
},
ensure_ascii=False,
)
_sync_log.info(
"topology sync done job=%s nodes=%s/%s/%s links=%s/%s/%s deleted_n=%s deleted_l=%s",
"topology sync done job=%s status=%s nodes=%s/%s/%s links=%s/%s/%s deleted_n=%s deleted_l=%s apply_ok=%s",
job_id,
status,
nodes_pulled,
nodes_ins,
nodes_upd,
@ -591,6 +602,7 @@ def sync_topology_full(db: Session, client: UMEClient, *, trigger_mode: str = "m
links_upd,
nodes_del,
links_del,
apply_ok,
)
finally:
try:

View file

@ -248,6 +248,81 @@ def apply_ume_topology_to_fabric(db: Session) -> dict[str, Any]:
return stats
def ume_topology_apply_gap(db: Session) -> dict[str, Any]:
"""Detect dock→Fabric/World gap for startup heal and UI hints."""
from sqlalchemy import func
from .models import TopoFabricNode, TopoFolder, UmeSyncJob, UmeTopoNode
from .ume_topology_world import get_world_container_folder
dock_me = int(
db.query(func.count(UmeTopoNode.node_id))
.filter(UmeTopoNode.node_type == "TOPO_NODE_ME")
.scalar()
or 0
)
fabric_ume = int(
db.query(func.count(TopoFabricNode.id))
.filter(TopoFabricNode.ume_ne_id.isnot(None), TopoFabricNode.ume_ne_id != "")
.scalar()
or 0
)
world = get_world_container_folder(db)
world_exists = world is not None
latest = (
db.query(UmeSyncJob)
.filter(UmeSyncJob.domain == "topology", UmeSyncJob.ended_at.isnot(None))
.order_by(UmeSyncJob.ended_at.desc())
.first()
)
latest_status = str(getattr(latest, "status", "") or "")
latest_error = str(getattr(latest, "error_message", "") or "")
partial = latest_status == "partial" or latest_error.startswith("dock_ok_fabric_apply_failed")
needs_apply = bool(dock_me > 0 and (fabric_ume == 0 or not world_exists or partial))
return {
"dock_me_count": dock_me,
"fabric_ume_count": fabric_ume,
"world_exists": world_exists,
"world_folder_id": str(getattr(world, "id", "") or ""),
"latest_topology_status": latest_status,
"latest_topology_error": latest_error[:240],
"partial_apply": partial,
"needs_apply": needs_apply,
}
def apply_ume_dock_to_fabric_if_needed(
db: Session, *, reason: str = "manual", force: bool = False
) -> dict[str, Any] | None:
"""Apply dock tables into Fabric/World when gap detected (or force)."""
from .ume_topology_world import ensure_ume_world_and_sbn_folders
gap = ume_topology_apply_gap(db)
if not force and not gap.get("needs_apply"):
return None
_log.info(
"ume dock→fabric apply start reason=%s force=%s gap=%s",
reason,
force,
{k: gap.get(k) for k in ("dock_me_count", "fabric_ume_count", "world_exists", "partial_apply")},
)
apply_stats = apply_ume_topology_to_fabric(db)
world_stats = ensure_ume_world_and_sbn_folders(db)
out = {
"ok": True,
"reason": reason,
"before": gap,
"fabric_apply": apply_stats,
"world": world_stats,
}
_log.info(
"ume dock→fabric apply done reason=%s nodes_ensured=%s",
reason,
(apply_stats or {}).get("nodes_ensured"),
)
return out
def _clear_and_keep(attrs: dict[str, Any]) -> dict[str, Any]:
from .topology_common import _clear_miss_attrs

View file

@ -10,14 +10,25 @@ from sqlalchemy.orm import sessionmaker
import netx_api.models # noqa: F401
from netx_api.db import Base
from netx_api.models import TopoFabricEdge, TopoFabricNode, UmeInventoryNE, UmeTopoLink, UmeTopoNode
from netx_api.models import (
TopoFabricEdge,
TopoFabricNode,
UmeInventoryNE,
UmeSyncJob,
UmeTopoLink,
UmeTopoNode,
)
from netx_api.ume_port_normalize import (
extract_ifnames_from_user_label,
port_keys_compatible,
port_suffix_from_tp_ref,
resolve_link_ifnames,
)
from netx_api.ume_topology_apply import apply_ume_topology_to_fabric
from netx_api.ume_topology_apply import (
apply_ume_dock_to_fabric_if_needed,
apply_ume_topology_to_fabric,
ume_topology_apply_gap,
)
from netx_api.topology_fabric_links import upsert_fabric_edge
@ -230,6 +241,53 @@ class UmeFabricApplyTests(unittest.TestCase):
self.assertEqual(len(manuals), 1)
self.assertEqual(manuals[0].status, "active")
def test_gap_needs_apply_when_dock_only(self):
self._seed_ume()
self.db.commit()
gap = ume_topology_apply_gap(self.db)
self.assertGreater(gap["dock_me_count"], 0)
self.assertEqual(gap["fabric_ume_count"], 0)
self.assertTrue(gap["needs_apply"])
def test_gap_needs_apply_on_partial_job(self):
self._seed_ume()
apply_ume_topology_to_fabric(self.db)
self.db.add(
UmeSyncJob(
domain="topology",
status="partial",
trigger_mode="manual",
started_at=self.now,
ended_at=self.now,
error_message="dock_ok_fabric_apply_failed:boom",
)
)
self.db.commit()
gap = ume_topology_apply_gap(self.db)
self.assertTrue(gap["partial_apply"])
self.assertTrue(gap["needs_apply"])
def test_apply_if_needed_skips_when_healthy(self):
self._seed_ume()
apply_ume_topology_to_fabric(self.db)
# World folder may still be missing → force a clean gap check after ensure.
from netx_api.ume_topology_world import ensure_ume_world_and_sbn_folders
ensure_ume_world_and_sbn_folders(self.db)
self.db.add(
UmeSyncJob(
domain="topology",
status="done",
trigger_mode="manual",
started_at=self.now,
ended_at=self.now,
)
)
self.db.commit()
gap = ume_topology_apply_gap(self.db)
self.assertFalse(gap["needs_apply"])
self.assertIsNone(apply_ume_dock_to_fabric_if_needed(self.db, reason="test"))
if __name__ == "__main__":
unittest.main()

View file

@ -1103,6 +1103,12 @@ const en = {
hidePanel: "Close",
empty: "No sync records",
summary: "{{total}} record(s), {{running}} running",
fabricGapTitle: "Topology dock is newer than Fabric / UME World",
fabricGapHint:
"Dock ME {{dock}}, Fabric UME {{fabric}}. Apply without pulling UME. Last status: {{status}}",
applyFabric: "Apply to topology",
applyingFabric: "Applying…",
applyFabricOk: "Applied to topology",
},
ne: {
title: "Network elements",

View file

@ -1096,6 +1096,12 @@ const zh = {
hidePanel: "关闭",
empty: "暂无同步记录",
summary: "共 {{total}} 条记录,进行中 {{running}}",
fabricGapTitle: "拓扑码头已更新,但尚未灌入拓扑管理(Fabric / UME World)",
fabricGapHint:
"码头 ME {{dock}},Fabric UME {{fabric}};可一键应用(不拉 UME)。上次状态:{{status}}",
applyFabric: "应用到拓扑管理",
applyingFabric: "正在应用…",
applyFabricOk: "已应用到拓扑管理",
},
ne: {
title: "网元清单",

View file

@ -2,6 +2,7 @@ import { useMemo, useState } from "react";
import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query";
import {
apiPost,
applyUmeTopologyToFabric,
disconnectUmeToken,
fetchUmeNe,
cancelUmeAlarmSubscription,
@ -260,6 +261,17 @@ export function UmePage() {
const runningTasks = (syncStatusQuery.data?.items || []).filter((x) => String(x.status || "").toLowerCase() === "running");
const runtimeTasks = syncStatusQuery.data?.runtime_tasks || [];
const topologyFabric = syncStatusQuery.data?.topology_fabric;
const needsFabricApply = Boolean(topologyFabric?.needs_apply);
const applyFabricMut = useMutation({
mutationFn: () => applyUmeTopologyToFabric(),
onSuccess: async () => {
showOk(t("ume.syncStatus.applyFabricOk"));
await queryClient.invalidateQueries({ queryKey: queryKeys.umeSyncStatusAll });
await queryClient.invalidateQueries({ queryKey: queryKeys.topologyTree });
},
onError: (err) => showError(String(err)),
});
const alarmSub: UmeAlarmSubscriptionStatus =
subscriptionStatusQuery.data ??
syncStatusQuery.data?.alarm_subscription ??
@ -532,6 +544,32 @@ export function UmePage() {
</button>
</div>
</div>
{needsFabricApply ? (
<div className="pill pill--high" style={{ marginBottom: 10 }} role="status">
<div style={{ fontWeight: 600 }}>{t("ume.syncStatus.fabricGapTitle")}</div>
<div className="muted" style={{ marginTop: 4 }}>
{t("ume.syncStatus.fabricGapHint")
.replace("{{dock}}", String(topologyFabric?.dock_me_count ?? 0))
.replace("{{fabric}}", String(topologyFabric?.fabric_ume_count ?? 0))
.replace(
"{{status}}",
String(topologyFabric?.latest_topology_status || topologyFabric?.latest_topology_error || "—"),
)}
</div>
<div className="btn-row" style={{ marginTop: 8 }}>
<button
type="button"
className="btn btn--sm"
disabled={applyFabricMut.isPending}
onClick={() => applyFabricMut.mutate()}
>
{applyFabricMut.isPending
? t("ume.syncStatus.applyingFabric")
: t("ume.syncStatus.applyFabric")}
</button>
</div>
</div>
) : null}
<div className="pt-list">
<div className="pt-list-table-wrap">
<table className="data-table pt-list-table">

View file

@ -449,6 +449,14 @@ export const fetchUmeSyncStatus = (params: { page: number; pageSize: number }) =
return apiGet<UmeSyncStatusResponse>(`/v1/ume/sync/status?${p.toString()}`);
};
/** Dock tables → Fabric + UME World (does not pull from UME). */
export const applyUmeTopologyToFabric = () =>
apiPost<{
ok: boolean;
fabric_apply?: Record<string, unknown>;
world?: Record<string, unknown>;
}>("/v1/topology/world/apply-ume", {});
export const fetchUmeTokenStatus = () => apiGet<UmeTokenStatus>("/v1/ume/token/status");
export const refreshUmeToken = () => apiPost<UmeTokenStatus>("/v1/ume/token/refresh", {});
export const disconnectUmeToken = () => apiPost<UmeTokenStatus>("/v1/ume/token/disconnect", {});

View file

@ -112,6 +112,18 @@ export type UmeAlarmSubscriptionStatus = {
ws_logs?: UmeWsLogEntry[];
};
export type UmeTopologyFabricGap = {
dock_me_count?: number;
fabric_ume_count?: number;
world_exists?: boolean;
world_folder_id?: string;
latest_topology_status?: string;
latest_topology_error?: string;
partial_apply?: boolean;
needs_apply?: boolean;
error?: string;
};
export type UmeSyncStatusResponse = {
total?: number;
page?: number;
@ -119,6 +131,7 @@ export type UmeSyncStatusResponse = {
items: UmeSyncJobItem[];
latest_by_domain?: Record<string, UmeSyncJobItem>;
alarm_subscription?: UmeAlarmSubscriptionStatus;
topology_fabric?: UmeTopologyFabricGap;
runtime_tasks?: Array<{
task: string;
status: string;