Split UME and alarms routes out of main; align UI scopes.

Extract ume_support/ume_router and alarms_router so main stays startup-focused, add topology fabric/views import surfaces, and gate managed-NE writes by ne:write.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-02 16:30:29 +08:00
parent 633a9d55bd
commit 8180f069ea
11 changed files with 2365 additions and 2090 deletions

View file

@ -11,7 +11,7 @@
## Runtime
- Ensure PostgreSQL backup policy exists (daily logical backup + retention).
- Prefer Alembic: `alembic upgrade head` and `NETX_SKIP_LEGACY_STARTUP_DDL=true`.
- Prefer Alembic: see [docs/ALEMBIC.md](docs/ALEMBIC.md). After `alembic upgrade head`, set `NETX_SKIP_LEGACY_STARTUP_DDL=true`.
- Optional: `NETX_RUN_INLINE_SCHEDULERS=false` and run `python -m netx_api.worker` for collectors.
- Run `oclaw` and `netx` under process managers (systemd/Windows service/pm2 equivalent).
- Enable auto-restart and startup-at-boot for both services.

22
docs/ALEMBIC.md Normal file
View file

@ -0,0 +1,22 @@
# Alembic (schema migrations)
netx historically evolved the schema with startup `ALTER TABLE … IF NOT EXISTS`.
Alembic is the preferred path going forward.
## Commands
```powershell
cd netx
.\.venv\Scripts\alembic.exe upgrade head
```
## Env
| Variable | Meaning |
|----------|---------|
| `NETX_SKIP_LEGACY_STARTUP_DDL=true` | Skip the large ad-hoc ALTER block in API startup (keep auth `scopes` column ensures). Use after `alembic upgrade head`. |
| `NETX_DATABASE_URL` | Same URL Alembic reads via `netx_api.config.settings`. |
Fresh lab installs can keep the legacy startup DDL (`false`, default) until you adopt Alembic in your deploy checklist.
Revision for capability scopes: `alembic/versions/20260802_scopes.py`.

502
netx_api/alarms_router.py Normal file
View file

@ -0,0 +1,502 @@
"""Legacy Excel alarm import, batches, diagnostics, and AP analyze routes."""
from __future__ import annotations
import csv
import json
from datetime import datetime, timezone
from io import StringIO
from typing import Any
from fastapi import APIRouter, Depends, File, HTTPException, Query, UploadFile
from fastapi.responses import Response
from sqlalchemy.orm import Session
from .ap_client import analyze_with_oclaw
from .config import settings
from .db import get_db
from .importer import aggregate_alarms, import_alarm_excel, query_alarms
from .models import AiAnalyzeHistory, AlarmBatch, AlarmNorm, ImportErrorRow, ImportJob
from .parser_config import load_parser_config
from .schemas import (
AiAnalyzeHistoryItem,
AiAnalyzeHistoryResponse,
AlarmAggregateBucket,
AlarmAggregateResponse,
AlarmItem,
AlarmQueryResponse,
BatchSummary,
ImportJobItem,
ImportJobListResponse,
)
router = APIRouter(tags=["alarms-import"])
parser_cfg = load_parser_config()
def _ensure_utc(dt: datetime | None) -> datetime | None:
if dt is None:
return None
if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
@router.post("/v1/alarms/import", response_model=BatchSummary)
async def import_alarms(file: UploadFile = File(...), db: Session = Depends(get_db)) -> BatchSummary:
filename = str(file.filename or "alarm.xlsx")
if not filename.lower().endswith((".xlsx", ".xls")):
raise HTTPException(status_code=400, detail="only_excel_supported_in_phase1")
content = await file.read()
if not content:
raise HTTPException(status_code=400, detail="empty_file")
batch = import_alarm_excel(db, filename=filename, content=content, parser=parser_cfg)
try:
job = ImportJob(
kind="alarms",
file_name=filename,
batch_id=str(batch.batch_id),
ok=1,
summary=f"success={int(batch.success_rows)} failed={int(batch.failed_rows)}",
)
db.add(job)
db.commit()
except Exception:
db.rollback()
return BatchSummary(
batch_id=str(batch.batch_id),
total_rows=int(batch.total_rows or 0),
success_rows=int(batch.success_rows or 0),
failed_rows=int(batch.failed_rows or 0),
status=str(batch.status or ""),
created_at=_ensure_utc(batch.created_at) or datetime.now(timezone.utc),
)
@router.post("/v1/logs/import")
async def import_logs(file: UploadFile = File(...)) -> dict:
# Placeholder for Phase 2: logs parsing + storage + query.
filename = str(file.filename or "logs.zip")
if not filename:
raise HTTPException(status_code=400, detail="filename_required")
raise HTTPException(status_code=501, detail="logs_import_not_implemented")
@router.get("/v1/jobs", response_model=ImportJobListResponse)
def list_jobs(limit: int = Query(default=20, ge=1, le=100), db: Session = Depends(get_db)) -> ImportJobListResponse:
rows = db.query(ImportJob).order_by(ImportJob.created_at.desc()).limit(limit).all()
items = [
ImportJobItem(
id=int(x.id),
kind=str(x.kind),
file_name=str(x.file_name or ""),
batch_id=str(x.batch_id) if x.batch_id else None,
ok=bool(int(x.ok or 0)),
summary=str(x.summary or ""),
created_at=_ensure_utc(x.created_at) or datetime.now(timezone.utc),
)
for x in rows
]
return ImportJobListResponse(items=items)
@router.get("/v1/batches")
def list_batches(limit: int = Query(default=20, ge=1, le=100), db: Session = Depends(get_db)) -> dict:
rows = db.query(AlarmBatch).order_by(AlarmBatch.created_at.desc()).limit(limit).all()
return {
"items": [
{
"batch_id": x.batch_id,
"source_file": x.source_file,
"status": x.status,
"total_rows": x.total_rows,
"success_rows": x.success_rows,
"failed_rows": x.failed_rows,
"created_at": (_ensure_utc(x.created_at) or datetime.now(timezone.utc)).isoformat(),
}
for x in rows
]
}
@router.get("/v1/batches/{batch_id}/errors.csv")
def download_batch_errors(batch_id: str, db: Session = Depends(get_db)):
rows = (
db.query(ImportErrorRow)
.filter(ImportErrorRow.batch_id == batch_id)
.order_by(ImportErrorRow.id.asc())
.all()
)
if not rows:
raise HTTPException(status_code=404, detail="batch_or_errors_not_found")
buf = StringIO()
writer = csv.writer(buf)
writer.writerow(["row_no", "reason", "raw_json"])
for r in rows:
writer.writerow([r.row_no, r.reason, r.raw_json])
return Response(
content=buf.getvalue(),
media_type="text/csv",
headers={"content-disposition": f'attachment; filename="batch_{batch_id}_errors.csv"'},
)
@router.delete("/v1/batches/{batch_id}")
def delete_batch(batch_id: str, db: Session = Depends(get_db)) -> dict:
batch = db.get(AlarmBatch, batch_id)
if not batch:
raise HTTPException(status_code=404, detail="batch_not_found")
try:
alarms_deleted = int(
db.query(AlarmNorm).filter(AlarmNorm.batch_id == batch_id).delete(synchronize_session=False)
)
errors_deleted = int(
db.query(ImportErrorRow).filter(ImportErrorRow.batch_id == batch_id).delete(synchronize_session=False)
)
jobs_deleted = int(
db.query(ImportJob).filter(ImportJob.batch_id == batch_id).delete(synchronize_session=False)
)
db.delete(batch)
db.commit()
return {
"ok": True,
"batch_id": batch_id,
"deleted": {
"batch": 1,
"alarms": alarms_deleted,
"errors": errors_deleted,
"jobs": jobs_deleted,
},
}
except Exception as exc:
db.rollback()
raise HTTPException(status_code=500, detail=f"delete_batch_failed: {exc}") from exc
@router.delete("/v1/batches")
def delete_all_batches(db: Session = Depends(get_db)) -> dict:
try:
alarms_deleted = int(db.query(AlarmNorm).delete(synchronize_session=False))
errors_deleted = int(db.query(ImportErrorRow).delete(synchronize_session=False))
jobs_deleted = int(db.query(ImportJob).delete(synchronize_session=False))
batches_deleted = int(db.query(AlarmBatch).delete(synchronize_session=False))
db.commit()
return {
"ok": True,
"deleted": {
"batches": batches_deleted,
"alarms": alarms_deleted,
"errors": errors_deleted,
"jobs": jobs_deleted,
},
}
except Exception as exc:
db.rollback()
raise HTTPException(status_code=500, detail=f"delete_all_batches_failed: {exc}") from exc
@router.get("/v1/diagnostics")
def diagnostics(
batch_id: str = Query(...),
lang: str | None = Query(default=None),
db: Session = Depends(get_db),
) -> dict:
sev_rows = aggregate_alarms(db, group_by="severity_norm", batch_id=batch_id)
code_rows = aggregate_alarms(db, group_by="alarm_code", batch_id=batch_id)[:10]
ne_rows = aggregate_alarms(db, group_by="ne_name", batch_id=batch_id)[:10]
total = sum(count for _, count in sev_rows)
lang_norm = _normalize_netx_lang(lang)
proto_counts: dict[str, int] = {}
for name, desc, code, raw in (
db.query(AlarmNorm.alarm_name, AlarmNorm.description, AlarmNorm.alarm_code, AlarmNorm.raw_json)
.filter(AlarmNorm.batch_id == batch_id)
.all()
):
blob = " | ".join([str(code or ""), str(name or ""), str(desc or ""), str(raw or "")])
k = _protocol_bucket_label(blob, lang=lang_norm)
proto_counts[k] = int(proto_counts.get(k, 0)) + 1
protocol_summary = sorted(proto_counts.items(), key=lambda x: x[1], reverse=True)[:10]
return {
"batch_id": batch_id,
"total_alarms": int(total),
"severity_summary": [{"key": k, "count": v} for k, v in sev_rows],
"top_alarm_codes": [{"key": k, "count": v} for k, v in code_rows],
"top_ne": [{"key": k, "count": v} for k, v in ne_rows],
"protocol_summary": [{"key": k, "count": v} for k, v in protocol_summary],
}
@router.post("/v1/ap/analyze")
def ap_analyze(payload: dict, db: Session = Depends(get_db)) -> dict:
batch_id = str(payload.get("batch_id") or "").strip()
question = str(payload.get("question") or "").strip()
if not batch_id or not question:
raise HTTPException(status_code=400, detail="batch_id_and_question_required")
diag = diagnostics(batch_id=batch_id, db=db)
analysis_request_id = str(payload.get("analysis_request_id") or "").strip()
filters_obj = payload.get("filters") if isinstance(payload.get("filters"), dict) else {}
req = {
"analysis_request_id": analysis_request_id,
"question": question,
"dataset_ref": {
"batch_id": batch_id,
"filters": filters_obj or {},
},
"context": {
"severity_summary": diag["severity_summary"],
"top_alarm_codes": diag["top_alarm_codes"],
"top_ne": diag["top_ne"],
"protocol_summary": diag.get("protocol_summary", []),
"findings": diag.get("findings", []),
},
"constraints": payload.get("constraints") or {"language": "zh-CN", "max_points": 6},
"interaction_mode": "expert",
"specialist": "ops",
}
ok = False
err = ""
oclaw_resp: dict[str, Any] | None = None
try:
oclaw_resp = analyze_with_oclaw(req)
ok = bool(oclaw_resp.get("ok")) if isinstance(oclaw_resp, dict) else False
except Exception as exc:
err = str(exc)
# Persist Q&A history (best-effort; never block response).
try:
answer = ""
if isinstance(oclaw_resp, dict):
answer = str(oclaw_resp.get("answer") or "").strip()
row = AiAnalyzeHistory(
analysis_request_id=analysis_request_id,
batch_id=batch_id,
question=question,
filters_json=json.dumps(filters_obj or {}, ensure_ascii=False),
ok=1 if ok else 0,
answer=answer,
error=err,
evidence_json=json.dumps(diag or {}, ensure_ascii=False),
created_at=datetime.utcnow(),
)
db.add(row)
db.commit()
except Exception:
db.rollback()
if not ok:
return {
"ok": False,
"error": err or "oclaw_bridge_unavailable",
"fallback_diagnostics": diag,
"batch_id": batch_id,
"question": question,
}
return {"ok": True, "batch_id": batch_id, "question": question, "diagnostics": diag, "oclaw": oclaw_resp}
@router.get("/v1/ap/history", response_model=AiAnalyzeHistoryResponse)
def ap_history(
batch_id: str | None = Query(default=None),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=20, ge=1, le=100),
db: Session = Depends(get_db),
) -> AiAnalyzeHistoryResponse:
q = db.query(AiAnalyzeHistory)
if batch_id and str(batch_id).strip():
q = q.filter(AiAnalyzeHistory.batch_id == str(batch_id).strip())
total = int(q.count())
rows = (
q.order_by(AiAnalyzeHistory.id.desc())
.offset((int(page) - 1) * int(page_size))
.limit(int(page_size))
.all()
)
items: list[AiAnalyzeHistoryItem] = []
for r in rows:
try:
filters = json.loads(str(r.filters_json or "{}"))
except Exception:
filters = {}
items.append(
AiAnalyzeHistoryItem(
id=int(r.id),
analysis_request_id=str(r.analysis_request_id or ""),
batch_id=str(r.batch_id or ""),
question=str(r.question or ""),
filters=filters if isinstance(filters, dict) else {},
ok=bool(int(r.ok or 0) == 1),
answer=str(r.answer or ""),
error=str(r.error or ""),
created_at=_ensure_utc(r.created_at) or datetime.now(timezone.utc),
)
)
return AiAnalyzeHistoryResponse(total=total, page=page, page_size=page_size, items=items)
@router.get("/v1/alarms", response_model=AlarmQueryResponse)
def list_alarms(
batch_id: str | None = Query(default=None),
alarm_code: str | None = Query(default=None),
severity: str | None = Query(default=None),
ne_name: str | None = Query(default=None),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=50, ge=1, le=200),
db: Session = Depends(get_db),
) -> AlarmQueryResponse:
total, rows = query_alarms(
db,
batch_id=batch_id,
alarm_code=alarm_code,
severity=severity,
ne_name=ne_name,
page=page,
page_size=page_size,
)
items = [
AlarmItem(
id=x.id,
batch_id=x.batch_id,
row_no=x.row_no,
alarm_time=_ensure_utc(x.alarm_time) or datetime.now(timezone.utc),
severity_norm=x.severity_norm,
severity_raw=x.severity_raw,
ne_name=x.ne_name,
alarm_code=x.alarm_code,
description=x.description,
ack_state=x.ack_state,
)
for x in rows
]
return AlarmQueryResponse(total=total, page=page, page_size=page_size, items=items)
@router.get("/v1/alarms/fields")
def alarms_fields() -> dict:
"""List all columns in alarms_norm for power querying."""
cols = []
try:
cols = [str(c.name) for c in AlarmNorm.__table__.columns] # type: ignore[attr-defined]
except Exception:
cols = []
return {"items": cols}
def _serialize_alarm_row(row: AlarmNorm) -> dict[str, Any]:
out: dict[str, Any] = {}
for c in AlarmNorm.__table__.columns: # type: ignore[attr-defined]
name = str(c.name)
v = getattr(row, name, None)
if hasattr(v, "isoformat"):
try:
if isinstance(v, datetime):
out[name] = (_ensure_utc(v) or v).isoformat()
else:
out[name] = v.isoformat() # datetime/date
continue
except Exception:
pass
out[name] = v
return out
@router.get("/v1/alarms/raw")
def alarms_raw(
batch_id: str | None = Query(default=None),
severity: str | None = Query(default=None),
alarm_code: str | None = Query(default=None),
ne_name: str | None = Query(default=None),
q: str | None = Query(default=None, description="free text contains on alarm_code/ne_name/description/service"),
order_by: str = Query(default="alarm_time"),
order: str = Query(default="desc"),
page: int = Query(default=1, ge=1),
page_size: int = Query(default=50, ge=1, le=200),
db: Session = Depends(get_db),
) -> dict:
"""
Power query: return **all columns** for alarms_norm rows.
Safety constraints:
- batch_id is required (avoid unbounded scans)
- order_by is whitelisted
- page_size capped
"""
bid = str(batch_id or "").strip()
if not bid:
raise HTTPException(status_code=400, detail="batch_id_required")
stmt = db.query(AlarmNorm).filter(AlarmNorm.batch_id == bid)
if severity and str(severity).strip():
stmt = stmt.filter(AlarmNorm.severity_norm == str(severity).strip())
if alarm_code and str(alarm_code).strip():
stmt = stmt.filter(AlarmNorm.alarm_code.contains(str(alarm_code).strip()))
if ne_name and str(ne_name).strip():
stmt = stmt.filter(AlarmNorm.ne_name.contains(str(ne_name).strip()))
if q and str(q).strip():
qw = str(q).strip()
stmt = stmt.filter(
(AlarmNorm.alarm_code.contains(qw))
| (AlarmNorm.ne_name.contains(qw))
| (AlarmNorm.description.contains(qw))
| (AlarmNorm.service.contains(qw))
)
allowed_order_by = {
"id": AlarmNorm.id,
"alarm_time": AlarmNorm.alarm_time,
"severity_norm": AlarmNorm.severity_norm,
"ne_name": AlarmNorm.ne_name,
"alarm_code": AlarmNorm.alarm_code,
}
col = allowed_order_by.get(str(order_by or "").strip(), AlarmNorm.alarm_time)
if str(order or "").strip().lower() == "asc":
stmt = stmt.order_by(col.asc())
else:
stmt = stmt.order_by(col.desc())
total = int(stmt.count())
rows = (
stmt.offset((int(page) - 1) * int(page_size))
.limit(int(page_size))
.all()
)
return {
"total": total,
"page": int(page),
"page_size": int(page_size),
"items": [_serialize_alarm_row(r) for r in rows],
}
@router.get("/v1/alarms/aggregate", response_model=AlarmAggregateResponse)
def alarms_aggregate(
group_by: str = Query(default="severity_norm"),
batch_id: str | None = Query(default=None),
db: Session = Depends(get_db),
) -> AlarmAggregateResponse:
try:
rows = aggregate_alarms(db, group_by=group_by, batch_id=batch_id)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
return AlarmAggregateResponse(
group_by=group_by,
buckets=[AlarmAggregateBucket(key=k, count=v) for k, v in rows],
)
@router.get("/v1/batches/{batch_id}")
def get_batch(batch_id: str, db: Session = Depends(get_db)) -> dict:
batch = db.get(AlarmBatch, batch_id)
if not batch:
raise HTTPException(status_code=404, detail="batch_not_found")
errors = (
db.query(ImportErrorRow)
.filter(ImportErrorRow.batch_id == batch_id)
.order_by(ImportErrorRow.id.asc())
.limit(20)
.all()
)
return {
"batch": BatchSummary.model_validate(batch, from_attributes=True).model_dump(),
"errors_preview": [
{"row_no": e.row_no, "reason": e.reason, "raw_json": e.raw_json}
for e in errors
],
}

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,27 @@
"""Topology fabric node/edge operations (narrow import surface for routers)."""
from __future__ import annotations
from .topology_service import (
get_discover_job,
get_fabric_neighborhood,
get_fabric_summary,
list_fabric_edges,
list_fabric_nodes,
merge_duplicate_fabric_nodes,
refresh_fabric_stats,
start_discover_job,
upsert_fabric_edge,
)
__all__ = [
"get_discover_job",
"get_fabric_neighborhood",
"get_fabric_summary",
"list_fabric_edges",
"list_fabric_nodes",
"merge_duplicate_fabric_nodes",
"refresh_fabric_stats",
"start_discover_job",
"upsert_fabric_edge",
]

View file

@ -41,32 +41,34 @@ from .topology_schemas import (
ViewPopulateRequest,
ViewPositionsPatch,
)
from .topology_service import (
from .topology_fabric import (
get_discover_job,
get_fabric_neighborhood,
get_fabric_summary,
list_fabric_edges,
list_fabric_nodes,
merge_duplicate_fabric_nodes,
refresh_fabric_stats,
start_discover_job,
upsert_fabric_edge,
)
from .topology_views import (
add_nodes_to_view,
bootstrap_topology_tree,
create_folder,
create_view,
delete_folder,
delete_view,
get_discover_job,
get_fabric_neighborhood,
get_fabric_summary,
get_topology_tree,
get_view_graph,
list_fabric_edges,
list_fabric_nodes,
list_views,
merge_duplicate_fabric_nodes,
patch_view_edge_style,
patch_view_positions,
populate_view,
project_fabric_neighbors_to_view,
refresh_fabric_stats,
remove_view_nodes,
start_discover_job,
update_folder,
update_view,
upsert_fabric_edge,
)
router = APIRouter(prefix="/v1/topology", tags=["topology"])

View file

@ -0,0 +1,45 @@
"""Topology folder tree + leaf view operations.
Re-exports view/tree APIs from ``topology_service`` so routers can depend on a
narrower module boundary while the monolith file is gradually split.
"""
from __future__ import annotations
from .topology_service import (
add_nodes_to_view,
bootstrap_topology_tree,
create_folder,
create_view,
delete_folder,
delete_view,
get_topology_tree,
get_view_graph,
list_views,
patch_view_edge_style,
patch_view_positions,
populate_view,
project_fabric_neighbors_to_view,
remove_view_nodes,
update_folder,
update_view,
)
__all__ = [
"add_nodes_to_view",
"bootstrap_topology_tree",
"create_folder",
"create_view",
"delete_folder",
"delete_view",
"get_topology_tree",
"get_view_graph",
"list_views",
"patch_view_edge_style",
"patch_view_positions",
"populate_view",
"project_fabric_neighbors_to_view",
"remove_view_nodes",
"update_folder",
"update_view",
]

1156
netx_api/ume_router.py Normal file

File diff suppressed because it is too large Load diff

523
netx_api/ume_support.py Normal file
View file

@ -0,0 +1,523 @@
"""UME shared client, runtime task state, and sync helpers (used by router + startup)."""
from __future__ import annotations
import logging
import re
import threading
import time
from datetime import datetime, timezone
from typing import Any
from fastapi import HTTPException
from sqlalchemy.orm import Session
from .config import settings
from .db import SessionLocal
from .models import UmeAlarmCurrent, UmeInventoryNE, UmeSyncJob
from .runtime_task_messages import (
RT_ALARMS_SYNC_IN_PROGRESS_SKIP,
RT_KEEPALIVE_FAILED,
RT_OCLAW_FWD_DISABLED,
RT_PULLING_ALARMS_CURRENT,
RT_PULLING_INVENTORY,
RT_RESUMED,
RT_RESUMED_OCLAW_WSS_RECONNECT,
RT_RESUMED_SYNC_SOON,
RT_RESUMED_WSS_RECONNECT,
RT_STARTUP_ALARM_SYNC_BEFORE_WS,
RT_STARTUP_GATE_WAITING,
RT_UME_WS_DISABLED_NO_BASE_URL,
RT_WSS_ACTIVE_SKIP_REST,
)
from .ume_alarm_ws import (
begin_startup_alarm_sync_gate,
complete_startup_alarm_sync_gate,
is_startup_alarm_sync_pending,
is_wss_active_for_current_alarms,
)
from .ume_client import UMEClient
from .ume_sync_service import sync_alarms_current, sync_inventory_full
from .ume_token_store import (
clear_shared_token,
load_shared_token,
release_refresh_lock,
save_shared_token,
try_acquire_refresh_lock,
wait_for_token_update,
)
_schedule_log = logging.getLogger("netx.ume.schedule")
_UME_CLIENT_SINGLETON = UMEClient(
token_loader=lambda: load_shared_token(),
token_saver=lambda token, exp: save_shared_token(token, exp),
token_clearer=lambda: clear_shared_token(),
lock_acquirer=lambda: try_acquire_refresh_lock(),
lock_releaser=lambda: release_refresh_lock(),
token_waiter=lambda min_exp: wait_for_token_update(min_expires_at_epoch_s=float(min_exp)),
)
_UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = {
"token_keepalive": {"task": "token_keepalive", "status": "init", "last_run_at": None, "last_error": ""},
"alarms_current_auto_sync": {"task": "alarms_current_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
"alarms_current_ws_consumer": {"task": "alarms_current_ws_consumer", "status": "init", "last_run_at": None, "last_error": ""},
"oclaw_alarm_forwarder": {"task": "oclaw_alarm_forwarder", "status": "init", "last_run_at": None, "last_error": ""},
"inventory_auto_sync": {"task": "inventory_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
}
_UME_WS_STOP_EVENT: threading.Event | None = None
_UME_RUNTIME_PAUSED: dict[str, bool] = {}
UME_KNOWN_RUNTIME_TASKS: tuple[str, ...] = tuple(_UME_RUNTIME_TASKS.keys())
_UME_RUNTIME_LOCK = threading.Lock()
# Debounce skip / wake for scheduled sync threads (resume should not wait full interval).
_UME_DEBOUNCE_MUTEX = threading.Lock()
_UME_SYNC_SKIP_DEBOUNCE: set[str] = set()
_UME_DEBOUNCE_WAKE: dict[str, threading.Event] = {}
def _debounce_wake_event(task_id: str) -> threading.Event:
with _UME_DEBOUNCE_MUTEX:
ev = _UME_DEBOUNCE_WAKE.get(task_id)
if ev is None:
ev = threading.Event()
_UME_DEBOUNCE_WAKE[task_id] = ev
return ev
def _request_force_sync_after_resume(task_id: str) -> None:
"""Skip next debounce wait and interrupt an in-progress debounce sleep (UI 开始)."""
with _UME_DEBOUNCE_MUTEX:
_UME_SYNC_SKIP_DEBOUNCE.add(task_id)
try:
_debounce_wake_event(task_id).set()
except Exception:
pass
def _clear_force_resume_hints(task_id: str) -> None:
"""Pause: drop pending skip/wake so state is predictable."""
with _UME_DEBOUNCE_MUTEX:
_UME_SYNC_SKIP_DEBOUNCE.discard(task_id)
try:
_debounce_wake_event(task_id).clear()
except Exception:
pass
def _reset_debounce_wakeup() -> None:
with _UME_DEBOUNCE_MUTEX:
_UME_SYNC_SKIP_DEBOUNCE.clear()
for ev in _UME_DEBOUNCE_WAKE.values():
try:
ev.clear()
except Exception:
pass
def _set_runtime_task(task: str, *, status: str, last_run_at: datetime | None = None, last_error: str = "") -> None:
with _UME_RUNTIME_LOCK:
item = _UME_RUNTIME_TASKS.get(task, {"task": task, "status": "init", "last_run_at": None, "last_error": ""})
item["status"] = str(status or "unknown")
if last_run_at is not None:
item["last_run_at"] = last_run_at
item["last_error"] = str(last_error or "")
_UME_RUNTIME_TASKS[task] = item
def _runtime_is_paused(task: str) -> bool:
with _UME_RUNTIME_LOCK:
return bool(_UME_RUNTIME_PAUSED.get(str(task or "").strip()))
def _runtime_pause_task(task: str) -> None:
tid = str(task or "").strip()
with _UME_RUNTIME_LOCK:
if tid not in _UME_RUNTIME_TASKS:
raise KeyError(tid)
_UME_RUNTIME_PAUSED[tid] = True
def _runtime_resume_task(task: str) -> None:
tid = str(task or "").strip()
with _UME_RUNTIME_LOCK:
_UME_RUNTIME_PAUSED[tid] = False
def _format_runtime_interval_label(seconds: int) -> str:
s = max(1, int(seconds))
if s >= 3600 and s % 3600 == 0:
h = s // 3600
return f"{h} h"
if s >= 60 and s % 60 == 0:
m = s // 60
return f"{m} min"
return f"{s}s"
def _runtime_task_interval_fields(task_id: str) -> tuple[int | None, str]:
"""Effective loop interval as configured at process start (matches startup clamps)."""
if task_id == "token_keepalive":
if not bool(getattr(settings, "ume_keepalive_enabled", True)):
return None, "disabled"
interval_s = int(getattr(settings, "ume_keepalive_interval_s", 600) or 600)
eff = max(30, min(interval_s, 3600))
return eff, _format_runtime_interval_label(eff)
if task_id == "alarms_current_auto_sync":
if not bool(getattr(settings, "ume_sync_alarms_current_enabled", True)):
return None, "disabled"
interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000)
eff = max(30, min(interval_s, 86400))
return eff, _format_runtime_interval_label(eff)
if task_id == "alarms_current_ws_consumer":
if not bool(getattr(settings, "ume_alarm_ws_enabled", True)):
return None, "disabled"
return None, "realtime"
if task_id == "oclaw_alarm_forwarder":
if not is_forwarder_enabled():
return None, "disabled"
return None, "realtime"
if task_id == "inventory_auto_sync":
if not bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)):
return None, "disabled"
hours = int(getattr(settings, "ume_sync_inventory_every_hours", 48) or 48)
hours = max(1, min(hours, 168))
eff = int(hours * 3600)
return eff, _format_runtime_interval_label(eff)
return None, "—"
def _list_runtime_tasks() -> list[dict[str, Any]]:
with _UME_RUNTIME_LOCK:
out: list[dict[str, Any]] = []
for v in _UME_RUNTIME_TASKS.values():
task_id = str(v.get("task") or "")
paused = bool(_UME_RUNTIME_PAUSED.get(task_id))
eff_status = "paused" if paused else str(v.get("status") or "unknown")
ts = _ensure_utc(v.get("last_run_at")) if isinstance(v.get("last_run_at"), datetime) else None
interval_s, interval_label = _runtime_task_interval_fields(task_id)
out.append(
{
"task": task_id,
"status": eff_status,
"paused": paused,
"last_run_at": ts.isoformat() if ts else None,
"last_error": str(v.get("last_error") or ""),
"interval_s": interval_s,
"interval_label": interval_label,
}
)
return out
def _ensure_utc(dt: datetime | None) -> datetime | None:
if dt is None:
return None
# All timestamps are stored as UTC in DB (naive). Treat naive as UTC.
if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
try:
return dt.astimezone(timezone.utc)
except Exception:
return dt
def _reset_runtime_pause_flags() -> None:
"""Ensure no task is stuck paused in memory after process boot (pause is not persisted)."""
with _UME_RUNTIME_LOCK:
for tid in UME_KNOWN_RUNTIME_TASKS:
_UME_RUNTIME_PAUSED[tid] = False
_reset_debounce_wakeup()
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 _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:
return max(0, min(3600, int(getattr(settings, "ume_startup_alarm_sync_delay_s", 60) or 60)))
def _wait_until_startup_alarm_pull_allowed(label: str) -> None:
delay_s = _startup_alarm_pull_delay_s()
if delay_s <= 0:
return
remaining = float(delay_s) - (time.monotonic() - _BOOT_MONO)
if remaining <= 0:
return
_schedule_log.info("%s: defer alarm pull %.0fs after process start", label, remaining)
time.sleep(remaining)
def _run_startup_alarm_sync_before_ws() -> None:
"""REST-sync current alarms once on boot; WSS gate must already be closed in on_startup."""
if not _needs_startup_alarm_sync_before_ws():
complete_startup_alarm_sync_gate()
return
_wait_until_startup_alarm_pull_allowed("startup_alarm_sync")
try:
_schedule_log.info("startup: REST current-alarm snapshot (WSS blocked until finished)")
_set_runtime_task(
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=RT_STARTUP_ALARM_SYNC_BEFORE_WS,
)
db = SessionLocal()
try:
client = _ume_client()
sync_alarms_current(db, client, trigger_mode="schedule", wss_active=False)
_schedule_log.info("startup: current alarms sync completed, WSS may connect")
_set_runtime_task(
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="",
)
finally:
db.close()
except RuntimeError as exc:
if str(exc) != "alarms_current_sync_busy":
raise
_schedule_log.warning("startup: skip REST before WSS — sync already in progress")
except Exception as exc:
_schedule_log.exception("startup: current alarms sync before WSS failed: %s", exc)
_set_runtime_task(
"alarms_current_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=str(exc)[:240],
)
finally:
complete_startup_alarm_sync_gate()
def _sleep_or_until_paused(task_id: str, total_s: float) -> None:
"""Sleep up to total_s wall seconds; honor pause; wake early on resume (debounce interrupt)."""
deadline = time.time() + max(0.0, float(total_s))
ev = _debounce_wake_event(task_id)
ev.clear()
while time.time() < deadline:
if _runtime_is_paused(task_id):
time.sleep(1)
continue
remaining = deadline - time.time()
if remaining <= 0:
break
timeout = min(2.0, remaining)
if ev.wait(timeout=timeout):
ev.clear()
_schedule_log.info("%s: debounce wait interrupted (resume)", task_id)
with _UME_DEBOUNCE_MUTEX:
_UME_SYNC_SKIP_DEBOUNCE.discard(task_id)
return
if ev.is_set():
ev.clear()
def _last_finished_job_ended_at(db: Session, domain: str) -> datetime | None:
"""Latest finished sync job end time for domain (success or failed)."""
row = (
db.query(UmeSyncJob)
.filter(
UmeSyncJob.domain == domain,
UmeSyncJob.ended_at.isnot(None),
)
.order_by(UmeSyncJob.ended_at.desc())
.limit(1)
.first()
)
if not row or row.ended_at is None:
return None
return _ensure_utc(row.ended_at)
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."""
end = _last_finished_job_ended_at(db, domain)
if end is None:
return None
return max(0.0, (datetime.now(timezone.utc) - end).total_seconds())
def _refresh_runtime_task_idle(task_id: str, domain: str, *, last_error: str | None = None) -> None:
"""Mark scheduled sync task running; last_run_at = last finished job time (idle / debounce)."""
with _UME_RUNTIME_LOCK:
prev_error = str((_UME_RUNTIME_TASKS.get(task_id) or {}).get("last_error") or "")
db = SessionLocal()
try:
ended = _last_finished_job_ended_at(db, domain)
finally:
db.close()
_set_runtime_task(
task_id,
status="idle",
last_run_at=ended,
last_error=prev_error if last_error is None else last_error,
)
def _maybe_wait_for_sync_interval(
*,
task_id: str,
domain: str,
interval_s: int,
label: str,
) -> None:
"""Sleep until interval elapsed since last finished job (ended_at), if any."""
with _UME_DEBOUNCE_MUTEX:
if task_id in _UME_SYNC_SKIP_DEBOUNCE:
_UME_SYNC_SKIP_DEBOUNCE.discard(task_id)
_schedule_log.info("%s: debounce skipped (resume/kick)", label)
return
db = SessionLocal()
try:
elapsed = _seconds_since_last_finished_job(db, domain)
finally:
db.close()
_refresh_runtime_task_idle(task_id, domain)
if elapsed is None:
_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 finished %.0fs ago (>= %ss), sync now", label, elapsed, interval_s)
return
wait_s = float(interval_s) - elapsed
_schedule_log.info("%s: last finished %.0fs ago, wait %.0fs before sync", label, elapsed, wait_s)
_sleep_or_until_paused(task_id, wait_s)
def _parse_time(text: str | None) -> datetime | None:
s = str(text or "").strip()
if not s:
return None
s2 = s.replace("Z", "+00:00")
try:
dt = datetime.fromisoformat(s2)
return _ensure_utc(dt)
except Exception:
return None
def _aggregate_rows(items: list[Any], key_fn) -> list[dict[str, Any]]:
bucket: dict[str, int] = {}
for item in items:
key = str(key_fn(item) or "").strip()
if not key:
key = "unknown"
bucket[key] = int(bucket.get(key, 0)) + 1
return [{"key": k, "count": v} for k, v in sorted(bucket.items(), key=lambda kv: kv[1], reverse=True)]
def _ume_alarm_host_name(
alarm: UmeAlarmCurrent | UmeAlarmHistory,
ne: UmeInventoryNE | None = None,
) -> str:
hn = str(getattr(alarm, "host_name", "") or "").strip()
if hn:
return hn
if ne is not None:
return str(getattr(ne, "host_name", "") or "").strip()
return ""
def _ume_alarm_ne_group_key(
alarm: UmeAlarmCurrent | UmeAlarmHistory,
ne: UmeInventoryNE | None,
) -> str:
return (
_ume_alarm_host_name(alarm, ne)
or (str(ne.user_label if ne else "") or "").strip()
or (str(ne.ne_name if ne else "") or "").strip()
or str(alarm.ne_id or "").strip()
or "unknown"
)
_PROTOCOL_BUCKET_ZH: dict[str, str] = {
"IP/MPLS": "IP/MPLS",
"ETH": "ETH",
"OTN/Optical": "OTN/光",
"Clock": "时钟",
"Power": "电源",
"Other": "其他",
}
def _classify_protocol_bucket(text: str) -> str:
"""Canonical English protocol/technology bucket id."""
t = (text or "").upper()
if any(x in t for x in ("BGP", "OSPF", "ISIS", "LDP", "MPLS", "L3VPN", "VPN")):
return "IP/MPLS"
if any(x in t for x in ("ETH", "GE", "10GE", "25GE", "40GE", "100GE", "XGE")):
return "ETH"
if any(x in t for x in ("OTN", "ODU", "OCH", "OMS", "OSC", "DWDM", "WDM", "ROADM")):
return "OTN/Optical"
if any(x in t for x in ("CLOCK", "SYNC", "PTP", "1588", "BITS", "TOD")):
return "Clock"
if any(x in t for x in ("PWR", "POWER", "PSU", "BAT", "BATT")):
return "Power"
return "Other"
def _protocol_bucket_label(text: str, *, lang: str = "zh") -> str:
key = _classify_protocol_bucket(text)
if str(lang or "").strip().lower().startswith("en"):
return key
return _PROTOCOL_BUCKET_ZH.get(key, key)
def _normalize_netx_lang(lang: str | None) -> str:
return "en" if str(lang or "").strip().lower().startswith("en") else "zh"
def _ume_client() -> UMEClient:
return _UME_CLIENT_SINGLETON
def _ume_error_kind(err: str) -> str:
low = str(err or "").lower()
if "401" in low or "403" in low or "password" in low or "auth" in low:
return "auth_failed"
if "timeout" in low:
return "timeout"
if "tls" in low or "certificate" in low or "ssl" in low:
return "tls_failed"
if "connect" in low or "name or service not known" in low:
return "connect_failed"
if "handshake" in low:
return "handshake_failed"
return "other"

View file

@ -29,6 +29,7 @@ export const MODULES: readonly ModuleDefinition[] = [
descKey: "workbench.cards.umeSyncDesc",
iconTone: "blue",
titleKey: "layout.titleUme",
requiredScope: "alarms:read",
},
{
moduleId: "ne",
@ -38,6 +39,7 @@ export const MODULES: readonly ModuleDefinition[] = [
descKey: "workbench.cards.managedNeDesc",
iconTone: "green",
titleKey: "layout.titleManagedNe",
requiredScope: "ne:read",
},
{
moduleId: "network",
@ -47,6 +49,7 @@ export const MODULES: readonly ModuleDefinition[] = [
descKey: "workbench.cards.networkDesc",
iconTone: "slate",
titleKey: "layout.titleNetwork",
requiredScope: "ne:read",
},
{
moduleId: "topology",
@ -56,6 +59,7 @@ export const MODULES: readonly ModuleDefinition[] = [
descKey: "workbench.cards.topologyDesc",
iconTone: "amber",
titleKey: "layout.titleTopology",
requiredScope: "ne:read",
},
{
moduleId: "webcrt",
@ -76,6 +80,7 @@ export const MODULES: readonly ModuleDefinition[] = [
iconTone: "amber",
titleKey: "layout.titlePortTrafficWall",
workbenchHidden: true,
requiredScope: "ne:read",
},
{
moduleId: "users",

View file

@ -23,6 +23,7 @@ import { HopProxyFields, emptyHopProxyFields, type HopProxyFieldsState } from ".
import { queryKeys } from "../constants/queryKeys";
import { useI18n } from "../i18n";
import { useToast } from "../hooks/useToast";
import { useAuth } from "../auth/AuthContext";
import type { ManagedNeItem } from "../types";
import { pageCount } from "../utils/display";
import { formatSystemTime } from "../utils/time";
@ -131,6 +132,8 @@ function connectStatusClass(status: string): string {
export function NePage() {
const { t } = useI18n();
const { showOk, showError } = useToast();
const { hasScope, isAdmin } = useAuth();
const canWriteNe = isAdmin || hasScope("ne:write");
const queryClient = useQueryClient();
const importRef = useRef<HTMLInputElement>(null);
@ -578,7 +581,7 @@ export function NePage() {
<h2>{t("managedNe.title")}</h2>
<div className="panel__toolbar-end">
<div className="panel__actions">
<button type="button" onClick={openCreate} disabled={!credsOk}>
<button type="button" onClick={openCreate} disabled={!credsOk || !canWriteNe}>
{t("managedNe.add")}
</button>
{SHOW_UME_MANAGED_SYNC ? (
@ -618,7 +621,7 @@ export function NePage() {
<button
type="button"
onClick={() => importRef.current?.click()}
disabled={!credsOk || importMutation.isPending}
disabled={!credsOk || importMutation.isPending || !canWriteNe}
>
{importMutation.isPending ? t("managedNe.importing") : t("managedNe.importBtn")}
</button>
@ -642,7 +645,7 @@ export function NePage() {
</button>
<button
type="button"
disabled={selected.length === 0 || batchHopMutation.isPending}
disabled={selected.length === 0 || batchHopMutation.isPending || !canWriteNe}
onClick={() => {
if (selected.length === 0) {
showError(t("managedNe.hop.selectRequired"));
@ -656,7 +659,7 @@ export function NePage() {
</button>
<button
type="button"
disabled={selected.length === 0 || batchAccountMutation.isPending}
disabled={selected.length === 0 || batchAccountMutation.isPending || !canWriteNe}
onClick={() => {
if (selected.length === 0) {
showError(t("managedNe.account.selectRequired"));
@ -671,7 +674,7 @@ export function NePage() {
<button
type="button"
className="btn btn--danger"
disabled={selected.length === 0 || batchDeleteMutation.isPending}
disabled={selected.length === 0 || batchDeleteMutation.isPending || !canWriteNe}
onClick={() => {
if (selected.length === 0) {
showError(t("managedNe.batchDeleteSelectRequired"));
@ -851,12 +854,13 @@ export function NePage() {
>
{t("managedNe.connectDetail")}
</button>
<button type="button" onClick={() => openEdit(row)}>
<button type="button" onClick={() => openEdit(row)} disabled={!canWriteNe}>
{t("managedNe.edit")}
</button>
<button
type="button"
className="btn--danger"
disabled={!canWriteNe}
onClick={() => {
if (window.confirm(t("managedNe.confirmDelete"))) deleteMutation.mutate(row.id);
}}